fix(FN-8426): wait for answers to agent questions
Convert supported runtime question-tool calls into Fusion's durable awaiting-user-input contract so workflow execution cannot continue while the operator question is unanswered. Fusion-Task-Id: FN-8426
This commit is contained in:
7
.changeset/fn-8426-question-tool-wait.md
Normal file
7
.changeset/fn-8426-question-tool-wait.md
Normal file
@@ -0,0 +1,7 @@
|
|||||||
|
---
|
||||||
|
"@runfusion/fusion": patch
|
||||||
|
---
|
||||||
|
|
||||||
|
summary: Keep workflow tasks paused while an agent question is awaiting an operator response.
|
||||||
|
category: fix
|
||||||
|
dev: Converts supported runtime question-tool calls into the durable workflow await-input contract.
|
||||||
@@ -1,5 +1,5 @@
|
|||||||
import { describe, expect, it } from "vitest";
|
import { describe, expect, it } from "vitest";
|
||||||
import { parseAwaitInputSentinel } from "../executor.js";
|
import { parseAwaitInputQuestionToolCall, parseAwaitInputSentinel } from "../executor.js";
|
||||||
|
|
||||||
describe("parseAwaitInputSentinel (U6)", () => {
|
describe("parseAwaitInputSentinel (U6)", () => {
|
||||||
it("extracts the question from a well-formed sentinel block", () => {
|
it("extracts the question from a well-formed sentinel block", () => {
|
||||||
@@ -35,3 +35,20 @@ describe("parseAwaitInputSentinel (U6)", () => {
|
|||||||
expect(parseAwaitInputSentinel("===FUSION_AWAIT_INPUT===\nno closing marker")).toBeNull();
|
expect(parseAwaitInputSentinel("===FUSION_AWAIT_INPUT===\nno closing marker")).toBeNull();
|
||||||
});
|
});
|
||||||
});
|
});
|
||||||
|
|
||||||
|
describe("parseAwaitInputQuestionToolCall", () => {
|
||||||
|
it("normalizes supported single and multi-question tool calls", () => {
|
||||||
|
expect(parseAwaitInputQuestionToolCall("request_user_input", {
|
||||||
|
questions: [
|
||||||
|
{ question: "Which layout should compact tablets use?" },
|
||||||
|
{ prompt: "Should desktop keep three panes?" },
|
||||||
|
],
|
||||||
|
})).toBe("Which layout should compact tablets use?\n\nShould desktop keep three panes?");
|
||||||
|
expect(parseAwaitInputQuestionToolCall("AskUserQuestion", { message: "Proceed?" })).toBe("Proceed?");
|
||||||
|
});
|
||||||
|
|
||||||
|
it("does not classify ordinary tools or malformed question calls", () => {
|
||||||
|
expect(parseAwaitInputQuestionToolCall("bash", { question: "ignored" })).toBeNull();
|
||||||
|
expect(parseAwaitInputQuestionToolCall("ask_user", { question: " " })).toBeNull();
|
||||||
|
});
|
||||||
|
});
|
||||||
|
|||||||
@@ -48,7 +48,10 @@ type CapturedSession = {
|
|||||||
* Make createFnAgent capture its session-creation args and return a mock session
|
* Make createFnAgent capture its session-creation args and return a mock session
|
||||||
* that emits the given output line, then resolves. Returns the capture holder.
|
* that emits the given output line, then resolves. Returns the capture holder.
|
||||||
*/
|
*/
|
||||||
function captureSession(output = '{"verdict":"APPROVE","notes":""}'): { last?: CapturedSession; all: CapturedSession[] } {
|
function captureSession(
|
||||||
|
output = '{"verdict":"APPROVE","notes":""}',
|
||||||
|
questionTool?: { name: string; args: Record<string, unknown> },
|
||||||
|
): { last?: CapturedSession; all: CapturedSession[] } {
|
||||||
const holder: { last?: CapturedSession; all: CapturedSession[] } = { all: [] };
|
const holder: { last?: CapturedSession; all: CapturedSession[] } = { all: [] };
|
||||||
mockedCreateFnAgent.mockImplementation(async (opts: any) => {
|
mockedCreateFnAgent.mockImplementation(async (opts: any) => {
|
||||||
const captured: CapturedSession = {
|
const captured: CapturedSession = {
|
||||||
@@ -69,6 +72,12 @@ function captureSession(output = '{"verdict":"APPROVE","notes":""}'): { last?: C
|
|||||||
return () => {};
|
return () => {};
|
||||||
},
|
},
|
||||||
prompt: vi.fn(async () => {
|
prompt: vi.fn(async () => {
|
||||||
|
if (questionTool) {
|
||||||
|
for (const fn of listeners) {
|
||||||
|
fn({ type: "tool_execution_start", toolName: questionTool.name, args: questionTool.args });
|
||||||
|
}
|
||||||
|
await new Promise<void>(() => {});
|
||||||
|
}
|
||||||
for (const fn of listeners) {
|
for (const fn of listeners) {
|
||||||
fn({
|
fn({
|
||||||
type: "message_update",
|
type: "message_update",
|
||||||
@@ -857,6 +866,44 @@ describe("CE workflow-step executor integration", () => {
|
|||||||
|
|
||||||
// ── Item 5: FUSION_HEADLESS gating on stepEnv ───────────────────────────────
|
// ── Item 5: FUSION_HEADLESS gating on stepEnv ───────────────────────────────
|
||||||
describe("executeWorkflowStep FUSION_HEADLESS (U3)", () => {
|
describe("executeWorkflowStep FUSION_HEADLESS (U3)", () => {
|
||||||
|
it("parks a board workflow step when its runtime calls a user-question tool", async () => {
|
||||||
|
const store = createMockStore();
|
||||||
|
const live = baseStepTask();
|
||||||
|
store.getTask.mockResolvedValue(live as any);
|
||||||
|
const { executor } = makeExecutor(store);
|
||||||
|
captureSession("", {
|
||||||
|
name: "request_user_input",
|
||||||
|
args: { questions: [{ question: "Should compact tablets use one pane or two?" }] },
|
||||||
|
});
|
||||||
|
|
||||||
|
const result = await (executor as any).runGraphCustomNode(
|
||||||
|
{
|
||||||
|
id: "plan",
|
||||||
|
kind: "prompt",
|
||||||
|
column: "in-progress",
|
||||||
|
config: {
|
||||||
|
executor: "skill",
|
||||||
|
skillName: "compound-engineering:ce-plan",
|
||||||
|
prompt: "Plan the responsive layout.",
|
||||||
|
},
|
||||||
|
},
|
||||||
|
live,
|
||||||
|
{},
|
||||||
|
undefined,
|
||||||
|
);
|
||||||
|
|
||||||
|
expect(result).toEqual({ outcome: "failure", value: "awaiting-user-input" });
|
||||||
|
expect(store.updateTask).toHaveBeenCalledWith(
|
||||||
|
"FN-CE-1",
|
||||||
|
expect.objectContaining({
|
||||||
|
status: "awaiting-user-input",
|
||||||
|
paused: true,
|
||||||
|
pausedReason: expect.stringContaining("Should compact tablets use one pane or two?"),
|
||||||
|
}),
|
||||||
|
undefined,
|
||||||
|
);
|
||||||
|
});
|
||||||
|
|
||||||
it.each([
|
it.each([
|
||||||
["code-review group", { id: "custom-check", name: "Implementation Check", optionalGroupId: "code-review" }],
|
["code-review group", { id: "custom-check", name: "Implementation Check", optionalGroupId: "code-review" }],
|
||||||
["browser-verification group", { id: "custom-check", name: "Implementation Check", optionalGroupId: "browser-verification" }],
|
["browser-verification group", { id: "custom-check", name: "Implementation Check", optionalGroupId: "browser-verification" }],
|
||||||
|
|||||||
@@ -1054,6 +1054,42 @@ export function parseAwaitInputSentinel(output: string | undefined): string | nu
|
|||||||
return question ? question : null;
|
return question ? question : null;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
const USER_QUESTION_TOOL_NAMES = new Set([
|
||||||
|
"askuserquestion",
|
||||||
|
"ask_user",
|
||||||
|
"ask_followup_question",
|
||||||
|
"request_user_input",
|
||||||
|
"elicit",
|
||||||
|
"ask_question",
|
||||||
|
"fn_ask_question",
|
||||||
|
]);
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Normalize a question-tool invocation into the same durable await-input
|
||||||
|
* contract used by skill sentinels. Some runtimes expose an interactive
|
||||||
|
* question tool even though Fusion workflow-step sessions have no synchronous
|
||||||
|
* listener; detecting the call at the session event boundary prevents the
|
||||||
|
* task from continuing after the unanswered question is rendered.
|
||||||
|
*/
|
||||||
|
export function parseAwaitInputQuestionToolCall(
|
||||||
|
toolName: string,
|
||||||
|
args: Record<string, unknown> | undefined,
|
||||||
|
): string | null {
|
||||||
|
if (!USER_QUESTION_TOOL_NAMES.has(toolName.trim().toLowerCase()) || !args) return null;
|
||||||
|
|
||||||
|
const records = Array.isArray(args.questions) ? args.questions : [args];
|
||||||
|
const questions = records.flatMap((value) => {
|
||||||
|
if (!value || typeof value !== "object" || Array.isArray(value)) return [];
|
||||||
|
const record = value as Record<string, unknown>;
|
||||||
|
const question = [record.question, record.prompt, record.message, record.text, record.title]
|
||||||
|
.find((candidate): candidate is string => typeof candidate === "string" && candidate.trim().length > 0)
|
||||||
|
?.trim();
|
||||||
|
return question ? [question] : [];
|
||||||
|
});
|
||||||
|
|
||||||
|
return questions.length > 0 ? questions.join("\n\n") : null;
|
||||||
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* (U2 / KTD-2) Fusion workflow-step conventions preamble, prepended to a skill
|
* (U2 / KTD-2) Fusion workflow-step conventions preamble, prepended to a skill
|
||||||
* step's prompt at the skill-prompt build path (runGraphCustomNode). It teaches
|
* step's prompt at the skill-prompt build path (runGraphCustomNode). It teaches
|
||||||
@@ -16602,6 +16638,11 @@ You have access to the file system to review changes.${inlineFixBlock}${verdictB
|
|||||||
|
|
||||||
let output = "";
|
let output = "";
|
||||||
const deltaNormalizer = createStreamingDeltaNormalizer();
|
const deltaNormalizer = createStreamingDeltaNormalizer();
|
||||||
|
let detectedQuestion: string | null = null;
|
||||||
|
let resolveQuestion: ((value: "await-input") => void) | undefined;
|
||||||
|
const questionPromise = new Promise<"await-input">((resolve) => {
|
||||||
|
resolveQuestion = resolve;
|
||||||
|
});
|
||||||
session.subscribe((event) => {
|
session.subscribe((event) => {
|
||||||
if (event.type === "message_update") {
|
if (event.type === "message_update") {
|
||||||
const msgEvent = event.assistantMessageEvent;
|
const msgEvent = event.assistantMessageEvent;
|
||||||
@@ -16620,6 +16661,16 @@ You have access to the file system to review changes.${inlineFixBlock}${verdictB
|
|||||||
}
|
}
|
||||||
if (event.type === "tool_execution_start") {
|
if (event.type === "tool_execution_start") {
|
||||||
agentLogger.onToolStart(event.toolName, event.args as Record<string, unknown> | undefined);
|
agentLogger.onToolStart(event.toolName, event.args as Record<string, unknown> | undefined);
|
||||||
|
if (!unattended && detectedQuestion === null) {
|
||||||
|
const question = parseAwaitInputQuestionToolCall(
|
||||||
|
event.toolName,
|
||||||
|
event.args as Record<string, unknown> | undefined,
|
||||||
|
);
|
||||||
|
if (question) {
|
||||||
|
detectedQuestion = question;
|
||||||
|
resolveQuestion?.("await-input");
|
||||||
|
}
|
||||||
|
}
|
||||||
}
|
}
|
||||||
if (event.type === "tool_execution_end") {
|
if (event.type === "tool_execution_end") {
|
||||||
agentLogger.onToolEnd(event.toolName, event.isError, event.result);
|
agentLogger.onToolEnd(event.toolName, event.isError, event.result);
|
||||||
@@ -16645,8 +16696,18 @@ You have access to the file system to review changes.${inlineFixBlock}${verdictB
|
|||||||
const outcome = await Promise.race([
|
const outcome = await Promise.race([
|
||||||
promptPromise.then(() => "completed" as const),
|
promptPromise.then(() => "completed" as const),
|
||||||
timeoutPromise,
|
timeoutPromise,
|
||||||
|
questionPromise,
|
||||||
]);
|
]);
|
||||||
|
|
||||||
|
if (outcome === "await-input" && detectedQuestion) {
|
||||||
|
try { session.dispose(); } catch { /* best-effort */ }
|
||||||
|
await agentLogger.flush();
|
||||||
|
return {
|
||||||
|
success: true,
|
||||||
|
output: `===FUSION_AWAIT_INPUT===\n${detectedQuestion}\n===END_FUSION_AWAIT_INPUT===`,
|
||||||
|
};
|
||||||
|
}
|
||||||
|
|
||||||
if (outcome === "timeout") {
|
if (outcome === "timeout") {
|
||||||
executorLog.warn(`${task.id}: workflow step '${workflowStep.name}' (${attemptLabel}) timed out after ${timeoutMs}ms — disposing session`);
|
executorLog.warn(`${task.id}: workflow step '${workflowStep.name}' (${attemptLabel}) timed out after ${timeoutMs}ms — disposing session`);
|
||||||
await this.store.logEntry(
|
await this.store.logEntry(
|
||||||
|
|||||||
Reference in New Issue
Block a user