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 { 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();
|
||||
});
|
||||
});
|
||||
|
||||
@@ -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" }],
|
||||
|
||||
@@ -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(
|
||||
|
||||
Reference in New Issue
Block a user