diff --git a/.changeset/fn-8426-question-tool-wait.md b/.changeset/fn-8426-question-tool-wait.md new file mode 100644 index 0000000000..823b4349a5 --- /dev/null +++ b/.changeset/fn-8426-question-tool-wait.md @@ -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. diff --git a/packages/engine/src/__tests__/await-input-sentinel.test.ts b/packages/engine/src/__tests__/await-input-sentinel.test.ts index 037c7730dc..159197b12b 100644 --- a/packages/engine/src/__tests__/await-input-sentinel.test.ts +++ b/packages/engine/src/__tests__/await-input-sentinel.test.ts @@ -1,5 +1,5 @@ import { describe, expect, it } from "vitest"; -import { parseAwaitInputSentinel } from "../executor.js"; +import { parseAwaitInputQuestionToolCall, parseAwaitInputSentinel } from "../executor.js"; describe("parseAwaitInputSentinel (U6)", () => { 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(); }); }); + +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(); + }); +}); diff --git a/packages/engine/src/__tests__/ce-workflow-step-executor.test.ts b/packages/engine/src/__tests__/ce-workflow-step-executor.test.ts index d013f8faeb..288ceea239 100644 --- a/packages/engine/src/__tests__/ce-workflow-step-executor.test.ts +++ b/packages/engine/src/__tests__/ce-workflow-step-executor.test.ts @@ -48,7 +48,10 @@ type CapturedSession = { * Make createFnAgent capture its session-creation args and return a mock session * 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 }, +): { last?: CapturedSession; all: CapturedSession[] } { const holder: { last?: CapturedSession; all: CapturedSession[] } = { all: [] }; mockedCreateFnAgent.mockImplementation(async (opts: any) => { const captured: CapturedSession = { @@ -69,6 +72,12 @@ function captureSession(output = '{"verdict":"APPROVE","notes":""}'): { last?: C return () => {}; }, 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(() => {}); + } for (const fn of listeners) { fn({ type: "message_update", @@ -857,6 +866,44 @@ describe("CE workflow-step executor integration", () => { // ── Item 5: FUSION_HEADLESS gating on stepEnv ─────────────────────────────── 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([ ["code-review group", { id: "custom-check", name: "Implementation Check", optionalGroupId: "code-review" }], ["browser-verification group", { id: "custom-check", name: "Implementation Check", optionalGroupId: "browser-verification" }], diff --git a/packages/engine/src/executor.ts b/packages/engine/src/executor.ts index 11a260d376..28c52f3b8e 100644 --- a/packages/engine/src/executor.ts +++ b/packages/engine/src/executor.ts @@ -1054,6 +1054,42 @@ export function parseAwaitInputSentinel(output: string | undefined): string | nu 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 | 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; + 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 * 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 = ""; 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) => { if (event.type === "message_update") { 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") { agentLogger.onToolStart(event.toolName, event.args as Record | undefined); + if (!unattended && detectedQuestion === null) { + const question = parseAwaitInputQuestionToolCall( + event.toolName, + event.args as Record | undefined, + ); + if (question) { + detectedQuestion = question; + resolveQuestion?.("await-input"); + } + } } if (event.type === "tool_execution_end") { 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([ promptPromise.then(() => "completed" as const), 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") { executorLog.warn(`${task.id}: workflow step '${workflowStep.name}' (${attemptLabel}) timed out after ${timeoutMs}ms — disposing session`); await this.store.logEntry(