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:
gsxdsm
2026-07-20 13:12:22 -07:00
parent 625dbc6c5f
commit 1b8b7f617e
4 changed files with 134 additions and 2 deletions

View 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.

View File

@@ -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();
});
});

View File

@@ -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<string, unknown> },
): { 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<void>(() => {});
}
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" }],

View File

@@ -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<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
* 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<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") {
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(