From 514ccd304cec1feed373bba58593ccfc872d7fd6 Mon Sep 17 00:00:00 2001 From: Phil Larson Date: Wed, 15 Jul 2026 18:07:37 -0700 Subject: [PATCH] Recover malformed agent interview responses (#2146) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit ## Summary - preserve valid onboarding JSON returned in Pi thinking-only assistant blocks - retry one bounded JSON-only reformat turn when the model returns prose or malformed output - keep streamed output as a final extraction fallback instead of overwriting it with an empty content array ## Verification - `pnpm --filter @fusion/dashboard exec vitest run src/__tests__/agent-onboarding.test.ts` — 20 passed - `pnpm --filter @fusion/dashboard typecheck` - `pnpm lint` - `pnpm check:changesets --strict` - live local-runtime AI Interview produced a structured Hermes/computer-use onboarding question after restart Follow-up to #2142, which fixed the missing planning-model fallback and runtime-hint prompt. ## Summary by CodeRabbit * **Bug Fixes** * Improved agent onboarding recovery when assistant replies include thinking-only content or malformed JSON. * Added a single automatic retry that re-formats invalid output into valid onboarding JSON. * Preserved structured “thinking” content as part of valid onboarding responses. * Normalized optional onboarding fields so null/empty/whitespace-only values are treated as missing. * Tightened Hermes automation so the runtime hint is set exactly to `hermes`. * **Tests** * Added onboarding event synchronization and expanded coverage for recovery and field normalization. --- .changeset/calm-interview-recovery.md | 7 + packages/core/src/prompt-overrides.ts | 5 + .../src/__tests__/agent-onboarding.test.ts | 269 +++++++++++++++++- packages/dashboard/src/__tests__/chat.test.ts | 15 +- packages/dashboard/src/agent-onboarding.ts | 195 ++++++++++--- 5 files changed, 449 insertions(+), 42 deletions(-) create mode 100644 .changeset/calm-interview-recovery.md diff --git a/.changeset/calm-interview-recovery.md b/.changeset/calm-interview-recovery.md new file mode 100644 index 0000000000..6139c9de97 --- /dev/null +++ b/.changeset/calm-interview-recovery.md @@ -0,0 +1,7 @@ +--- +"@runfusion/fusion": patch +--- + +summary: Recover agent interviews when models return thinking-only or malformed JSON responses. +category: fix +dev: Preserves structured thinking output and retries one JSON-only reformat turn before surfacing an error. diff --git a/packages/core/src/prompt-overrides.ts b/packages/core/src/prompt-overrides.ts index 18848a9bb0..ab695d0c97 100644 --- a/packages/core/src/prompt-overrides.ts +++ b/packages/core/src/prompt-overrides.ts @@ -301,6 +301,10 @@ For questions: For completion: {\n "type": "complete",\n "data": {\n "title": "Task title",\n "description": "Detailed description",\n "suggestedSize": "S|M|L",\n "suggestedDependencies": [],\n "keyDeliverables": ["Item 1", "Item 2"]\n }\n}`, }, + /** + * FNXC:AgentOnboardingRuntime 2026-07-15-15:25: + * Agent onboarding must emit the exact `hermes` runtime hint for computer-use, desktop-automation, and UI-testing agents so runtime resolution selects the Hermes tool surface instead of treating a descriptive label as an unknown runtime. + */ "agent-onboarding-system": { key: "agent-onboarding-system", name: "Agent Onboarding System", @@ -325,6 +329,7 @@ Rules: - Prefer structuring instructionsText with these markdown sections when drafting: ## Description, ## Expertise, ## Priorities, ## Boundaries, ## Communication, ## Collaboration & Escalation - Freeform instructionsText is still acceptable for compatibility; sectioned structure is preferred for new agents - modelHint and runtimeHint are optional draft suggestions only (not final runtime selection) +- When the user requests Hermes, computer use, desktop automation, or UI testing, use the exact runtimeHint "hermes"; never invent a descriptive runtime name - heartbeatProcedurePath, heartbeatIntervalMs, and heartbeatEnabled are optional draft hints only.`, }, "subtask-breakdown-system": { diff --git a/packages/dashboard/src/__tests__/agent-onboarding.test.ts b/packages/dashboard/src/__tests__/agent-onboarding.test.ts index 572819e566..c994c9f1fa 100644 --- a/packages/dashboard/src/__tests__/agent-onboarding.test.ts +++ b/packages/dashboard/src/__tests__/agent-onboarding.test.ts @@ -28,6 +28,7 @@ vi.mock("@fusion/engine", () => ({ import { __resetAgentOnboardingState, + agentOnboardingStreamManager, cancelAgentOnboardingSession, createAgentOnboardingSessionPrompt, getAgentOnboardingSession, @@ -38,6 +39,7 @@ import { retryAgentOnboardingSession, SessionNotFoundError, startAgentOnboardingSession, + stopAgentOnboardingGeneration, } from "../agent-onboarding.js"; function createMockAgent(responses: string[]) { @@ -65,6 +67,26 @@ async function waitFor(check: () => boolean, timeoutMs = 2000): Promise { } } +/* +FNXC:AgentOnboarding 2026-07-15-14:32: +Recovery regressions must await the onboarding event seam rather than poll wall-clock session state. Buffered-event replay covers generations that complete before the test subscribes. + +FNXC:AgentOnboarding 2026-07-15-16:42: +Live subscribers receive AgentOnboardingStreamEvent objects keyed by `type`, while SessionEventBuffer replay records are keyed by `event` with JSON-serialized `data`. The test seam must honor both shapes. +*/ +function waitForOnboardingEvent(sessionId: string, eventTypes: string[]): Promise { + if (agentOnboardingStreamManager.getBufferedEvents(sessionId, 0).some((event) => eventTypes.includes(event.event))) { + return Promise.resolve(); + } + return new Promise((resolve) => { + const unsubscribe = agentOnboardingStreamManager.subscribe(sessionId, (event) => { + if (!eventTypes.includes(event.type)) return; + unsubscribe(); + resolve(); + }); + }); +} + function createSkillPluginRunner(skills: Array<{ name: string; enabled?: boolean }>) { return { getPluginSkills: () => skills.map((skill) => ({ pluginId: "fusion-plugin-compound-engineering", skill })), @@ -80,6 +102,7 @@ describe("agent-onboarding", () => { afterEach(() => { __resetAgentOnboardingState(); + vi.useRealTimers(); }); it("parses question responses", () => { @@ -173,6 +196,31 @@ describe("agent-onboarding", () => { ).toThrow(/Invalid summary/); }); + it("normalizes absent-like optional string hints", () => { + const parsed = parseAgentOnboardingResponse( + JSON.stringify({ + type: "complete", + data: { + name: "Hermes Desktop Tester", + role: "executor", + instructionsText: "Exercise desktop workflows with Hermes computer-use tools.", + thinkingLevel: "medium", + maxTurns: 25, + heartbeatProcedurePath: " ", + modelHint: "", + runtimeHint: null, + }, + }), + ); + + expect(parsed.type).toBe("complete"); + if (parsed.type === "complete") { + expect(parsed.data.heartbeatProcedurePath).toBeUndefined(); + expect(parsed.data.modelHint).toBeUndefined(); + expect(parsed.data.runtimeHint).toBeUndefined(); + } + }); + it("rejects malformed rich draft fields", () => { expect(() => parseAgentOnboardingResponse( @@ -184,7 +232,7 @@ describe("agent-onboarding", () => { instructionsText: "Valid instructions", thinkingLevel: "medium", maxTurns: 20, - heartbeatProcedurePath: "", + heartbeatProcedurePath: 42, }, }), ), @@ -290,6 +338,7 @@ describe("agent-onboarding", () => { const options = mockCreateFnAgent.mock.calls.at(-1)?.[0] as { systemPrompt?: string }; expect(options.systemPrompt).toContain('"runtimeHint"'); expect(options.systemPrompt).toContain("optional draft suggestions"); + expect(options.systemPrompt).toContain('use the exact runtimeHint "hermes"'); expect(options.systemPrompt).not.toContain("Do not include runtimeMode/model/runtimeHint"); }); @@ -406,6 +455,224 @@ describe("agent-onboarding", () => { expect(options.skillSelection?.requestedSkillNames).toEqual(["fusion"]); }); + it("preserves valid JSON from thinking-only assistant responses", async () => { + const response = JSON.stringify({ + type: "question", + data: { id: "goal", type: "text", question: "What is the primary goal?" }, + }); + const messages: Array<{ + role: string; + content: Array<{ type: "thinking"; thinking: string }>; + }> = []; + mockCreateFnAgent.mockImplementationOnce(async (options: unknown) => { + const callbacks = options as { onThinking?: (delta: string) => void }; + return { + session: { + state: { messages }, + prompt: vi.fn(async () => { + callbacks.onThinking?.(response); + messages.push({ role: "assistant", content: [{ type: "thinking", thinking: response }] }); + }), + dispose: vi.fn(), + }, + }; + }); + + const sessionId = await startAgentOnboardingSession( + "127.0.0.1", + { intent: "thinking-only response", existingAgents: [], templates: [] }, + process.cwd(), + ); + + await waitForOnboardingEvent(sessionId, ["question", "error"]); + + const session = getAgentOnboardingSession(sessionId); + expect(session?.error).toBeUndefined(); + expect(session?.currentQuestion?.id).toBe("goal"); + }); + + it("prefers parseable thinking JSON over non-JSON text in the same assistant response", async () => { + const response = JSON.stringify({ + type: "question", + data: { id: "goal", type: "text", question: "What is the primary goal?" }, + }); + const messages: Array<{ + role: string; + content: Array<{ type: "text"; text: string } | { type: "thinking"; thinking: string }>; + }> = []; + const prompt = vi.fn(async () => { + messages.push({ + role: "assistant", + content: [ + { type: "thinking", thinking: response }, + { type: "text", text: "I worked through the onboarding request." }, + ], + }); + }); + mockCreateFnAgent.mockResolvedValueOnce({ + session: { state: { messages }, prompt, dispose: vi.fn() }, + }); + + const sessionId = await startAgentOnboardingSession( + "127.0.0.1", + { intent: "mixed thinking and text response", existingAgents: [], templates: [] }, + process.cwd(), + ); + + await waitForOnboardingEvent(sessionId, ["question", "error"]); + + const session = getAgentOnboardingSession(sessionId); + expect(session?.error).toBeUndefined(); + expect(session?.currentQuestion?.id).toBe("goal"); + expect(prompt).toHaveBeenCalledTimes(1); + }); + + it("uses streamed output when assistant string content is blank", async () => { + const response = JSON.stringify({ + type: "question", + data: { id: "goal", type: "text", question: "What is the primary goal?" }, + }); + const messages: Array<{ role: string; content: string }> = []; + let prompt: ReturnType; + mockCreateFnAgent.mockImplementationOnce(async (options: unknown) => { + const callbacks = options as { onText?: (delta: string) => void }; + prompt = vi.fn(async () => { + callbacks.onText?.(response); + messages.push({ role: "assistant", content: " " }); + }); + return { + session: { state: { messages }, prompt, dispose: vi.fn() }, + }; + }); + + const sessionId = await startAgentOnboardingSession( + "127.0.0.1", + { intent: "blank assistant content", existingAgents: [], templates: [] }, + process.cwd(), + ); + + await waitForOnboardingEvent(sessionId, ["question", "error"]); + + const session = getAgentOnboardingSession(sessionId); + expect(session?.error).toBeUndefined(); + expect(session?.currentQuestion?.id).toBe("goal"); + expect(prompt!).toHaveBeenCalledTimes(1); + }); + + it("retries once when the model response is not valid JSON", async () => { + const agent = createMockAgent([ + "I can help design that agent.", + JSON.stringify({ + type: "question", + data: { id: "goal", type: "text", question: "What is the primary goal?" }, + }), + ]); + mockCreateFnAgent.mockResolvedValueOnce(agent); + + const sessionId = await startAgentOnboardingSession( + "127.0.0.1", + { intent: "recover malformed output", existingAgents: [], templates: [] }, + process.cwd(), + ); + + await waitForOnboardingEvent(sessionId, ["question", "error"]); + + const session = getAgentOnboardingSession(sessionId); + expect(session?.error).toBeUndefined(); + expect(session?.currentQuestion?.id).toBe("goal"); + expect(agent.session.prompt).toHaveBeenCalledTimes(2); + }); + + it("settles a stalled reformat turn when its timeout expires", async () => { + vi.useFakeTimers(); + const messages: Array<{ role: string; content: string }> = []; + let callCount = 0; + const prompt = vi.fn(() => { + callCount += 1; + if (callCount === 1) { + messages.push({ role: "assistant", content: "not valid JSON" }); + return Promise.resolve(); + } + return new Promise(() => {}); + }); + const dispose = vi.fn(); + let expiredOnText: ((delta: string) => void) | undefined; + mockCreateFnAgent.mockImplementationOnce(async (options: unknown) => { + expiredOnText = (options as { onText?: (delta: string) => void }).onText; + return { + session: { state: { messages }, prompt, dispose }, + }; + }); + + const sessionId = await startAgentOnboardingSession( + "127.0.0.1", + { intent: "stalled reformat response", existingAgents: [], templates: [] }, + process.cwd(), + ); + await vi.advanceTimersByTimeAsync(0); + expect(prompt).toHaveBeenCalledTimes(2); + + await vi.advanceTimersByTimeAsync(120_000); + + const session = getAgentOnboardingSession(sessionId) as { error?: string; agent?: unknown; thinkingOutput: string } | undefined; + expect(session?.error).toBe("AI generation timed out. You can retry."); + expect(session?.agent).toBeUndefined(); + expect(dispose).toHaveBeenCalledTimes(1); + expiredOnText?.("stale expired output"); + expect(session?.thinkingOutput).toBe(""); + expect(agentOnboardingStreamManager.getBufferedEvents(sessionId, 0)).toContainEqual( + expect.objectContaining({ event: "error", data: JSON.stringify("AI generation timed out. You can retry.") }), + ); + }); + + it("settles a stopped generation without letting late cleanup detach a newer retry", async () => { + const response = JSON.stringify({ + type: "question", + data: { id: "stale", type: "text", question: "Stale question?" }, + }); + const expiredMessages: Array<{ role: string; content: string }> = []; + let resolveExpiredPrompt!: () => void; + const expiredPrompt = vi.fn(() => new Promise((resolve) => { + resolveExpiredPrompt = () => { + expiredMessages.push({ role: "assistant", content: response }); + resolve(); + }; + })); + let signalRetryPromptStarted!: () => void; + const retryPromptStarted = new Promise((resolve) => { + signalRetryPromptStarted = resolve; + }); + const retryPrompt = vi.fn(() => { + signalRetryPromptStarted(); + return new Promise(() => {}); + }); + mockCreateFnAgent + .mockResolvedValueOnce({ session: { state: { messages: expiredMessages }, prompt: expiredPrompt, dispose: vi.fn() } }) + .mockResolvedValueOnce({ session: { state: { messages: [] }, prompt: retryPrompt, dispose: vi.fn() } }); + + const sessionId = await startAgentOnboardingSession( + "127.0.0.1", + { intent: "stop and retry", existingAgents: [], templates: [] }, + process.cwd(), + ); + expect(expiredPrompt).toHaveBeenCalledTimes(1); + expect(stopAgentOnboardingGeneration(sessionId)).toBe(true); + + const retry = retryAgentOnboardingSession(sessionId); + await retryPromptStarted; + expect(retryPrompt).toHaveBeenCalledTimes(1); + + resolveExpiredPrompt(); + await Promise.resolve(); + await Promise.resolve(); + await Promise.resolve(); + + expect(getAgentOnboardingSession(sessionId)?.currentQuestion).toBeUndefined(); + expect(stopAgentOnboardingGeneration(sessionId)).toBe(true); + await retry; + expect(getAgentOnboardingSession(sessionId)?.error).toBe("Generation stopped by user. You can retry."); + }); + it("progresses through start -> question -> response -> final summary", async () => { mockCreateFnAgent.mockResolvedValueOnce( createMockAgent([ diff --git a/packages/dashboard/src/__tests__/chat.test.ts b/packages/dashboard/src/__tests__/chat.test.ts index 9fcecfc818..54b6fb5484 100644 --- a/packages/dashboard/src/__tests__/chat.test.ts +++ b/packages/dashboard/src/__tests__/chat.test.ts @@ -39,11 +39,13 @@ vi.mock("@fusion/core", () => ({ })); /* FNXC:DashboardChatTests 2026-07-12-08:15: -chat.ts has 14 named runtime imports + `import * as engineModule` from @fusion/engine (lines 43-59). The engine module transitively imports many @fusion/core exports (AWAITING_APPROVAL_PAUSE_REASON, THINKING_LEVELS, etc.), so loading the real engine against a partial @fusion/core mock throws. Since this test only exercises resolveFileReferences (no real AI calls), stub every engine export chat.ts references so the real engine module never loads. + +FNXC:DashboardChatTests 2026-07-15-16:15: +chat.ts now has 25 named runtime imports plus `import * as engineModule` from @fusion/engine. Keep this isolated mock complete as chat gains tool factories so the Gate's mock-completeness invariant does not regress. */ vi.mock("@fusion/engine", () => ({ createFnAgent: vi.fn(), @@ -58,6 +60,17 @@ vi.mock("@fusion/engine", () => ({ createChatArtifactTools: vi.fn(() => []), createChatTaskDocumentTools: vi.fn(() => []), createWorkflowAuthoringTools: vi.fn(() => []), + createTaskCreateTool: vi.fn(), + createTaskListTool: vi.fn(), + createTaskShowTool: vi.fn(), + createTaskSearchTool: vi.fn(), + createListAgentsTool: vi.fn(), + createDelegateTaskTool: vi.fn(), + createGetAgentConfigTool: vi.fn(), + createWebFetchTool: vi.fn(), + createGoalRetrievalTools: vi.fn(() => []), + createMemoryTools: vi.fn(() => []), + createResearchTools: vi.fn(() => []), resolveMcpServersForStore: vi.fn(async () => ({ servers: [], errors: [] })), resolveExecutorThinkingLevel: vi.fn(() => undefined), /* diff --git a/packages/dashboard/src/agent-onboarding.ts b/packages/dashboard/src/agent-onboarding.ts index 2fb36c8631..79b516c8b9 100644 --- a/packages/dashboard/src/agent-onboarding.ts +++ b/packages/dashboard/src/agent-onboarding.ts @@ -65,7 +65,27 @@ type SkillSelectionPluginRunner = Parameters>; @@ -99,6 +120,7 @@ interface Session { error?: string; history: Array<{ question: PlanningQuestion; response: Record }>; thinkingOutput: string; + agentEpoch: number; agent?: OnboardingAgent; rootDir: string; modelProvider?: string; @@ -110,7 +132,12 @@ interface Session { } const sessions = new Map(); -const activeGenerations = new Map(); +type ActiveGeneration = { + abortController: AbortController; + timer: NodeJS.Timeout; + reject: (reason?: unknown) => void; +}; +const activeGenerations = new Map(); export class AgentOnboardingStreamManager extends EventEmitter { private readonly sessions = new Map>(); @@ -171,6 +198,13 @@ function repairJson(text: string): string { return text.replace(/,\s*([}\]])/g, "$1"); } +function normalizeOptionalSummaryString(value: unknown, field: string): string | undefined { + if (value === undefined || value === null) return undefined; + if (typeof value !== "string") throw new Error(`Invalid summary.${field}`); + const trimmed = value.trim(); + return trimmed || undefined; +} + export function parseAgentOnboardingResponse(text: string): { type: "question"; data: PlanningQuestion } | { type: "complete"; data: AgentOnboardingSummary } { const candidate = extractJsonCandidate(text); if (!candidate) throw new Error("AI returned no valid JSON"); @@ -207,12 +241,10 @@ export function parseAgentOnboardingResponse(text: string): { type: "question"; throw new Error("Invalid summary.maxTurns"); } - if (data.heartbeatProcedurePath !== undefined) { - if (typeof data.heartbeatProcedurePath !== "string" || !data.heartbeatProcedurePath.trim()) { - throw new Error("Invalid summary.heartbeatProcedurePath"); - } - data.heartbeatProcedurePath = data.heartbeatProcedurePath.trim(); - } + data.heartbeatProcedurePath = normalizeOptionalSummaryString( + data.heartbeatProcedurePath, + "heartbeatProcedurePath", + ); if (data.heartbeatIntervalMs !== undefined) { if (typeof data.heartbeatIntervalMs !== "number" || !Number.isInteger(data.heartbeatIntervalMs) || data.heartbeatIntervalMs <= 0) { @@ -224,13 +256,8 @@ export function parseAgentOnboardingResponse(text: string): { type: "question"; throw new Error("Invalid summary.heartbeatEnabled"); } - if (data.modelHint !== undefined && typeof data.modelHint !== "string") { - throw new Error("Invalid summary.modelHint"); - } - - if (data.runtimeHint !== undefined && typeof data.runtimeHint !== "string") { - throw new Error("Invalid summary.runtimeHint"); - } + data.modelHint = normalizeOptionalSummaryString(data.modelHint, "modelHint"); + data.runtimeHint = normalizeOptionalSummaryString(data.runtimeHint, "runtimeHint"); return { type: "complete", data: data as AgentOnboardingSummary }; } @@ -307,6 +334,7 @@ export async function startAgentOnboardingSession( }), history: [], thinkingOutput: "", + agentEpoch: 0, rootDir, modelProvider, modelId, @@ -324,6 +352,7 @@ export async function startAgentOnboardingSession( } async function createAgentOnboardingAgent(session: Session, store?: TaskStore): Promise { + const agentEpoch = ++session.agentEpoch; const systemPrompt = resolvePrompt("agent-onboarding-system", session.promptOverrides) || AGENT_ONBOARDING_SYSTEM_PROMPT; const skillContext = buildSessionSkillContextSync(null, "executor", session.rootDir, session.pluginRunner); const mcpServers = (await resolveMcpServersForStore(store ?? {})).servers; @@ -346,10 +375,16 @@ async function createAgentOnboardingAgent(session: Session, store?: TaskStore): */ ...(skillContext.skillSelectionContext ? { skillSelection: skillContext.skillSelectionContext } : {}), onThinking: (delta: string) => { + /* + FNXC:AgentOnboarding 2026-07-15-16:48: + Provider cancellation is best-effort. Ignore callbacks from an invalidated agent epoch so an expired prompt cannot contaminate a fresh retry's streamed output or SSE timeline. + */ + if (session.agentEpoch !== agentEpoch) return; session.thinkingOutput += delta; agentOnboardingStreamManager.broadcast(session.id, { type: "thinking", data: delta }); }, onText: (delta: string) => { + if (session.agentEpoch !== agentEpoch) return; session.thinkingOutput += delta; agentOnboardingStreamManager.broadcast(session.id, { type: "thinking", data: delta }); }, @@ -363,47 +398,122 @@ async function runGenerationWithTimeout(session: Session, operation: () => Pr existing.abortController.abort(); } const abortController = new AbortController(); + /* + FNXC:AgentOnboarding 2026-07-15-16:36: + A generation timeout must settle the wrapper even when the provider ignores cancellation and its prompt promise never resolves. Race the operation against an explicit rejection; abort remains best-effort cleanup, while continueConversation owns the single terminal error event. + */ + let rejectGeneration!: (reason?: unknown) => void; + const interruption = new Promise((_, reject) => { + rejectGeneration = reject; + }); const timer = setTimeout(() => { - session.error = "AI generation timed out. You can retry."; - agentOnboardingStreamManager.broadcast(session.id, { type: "error", data: session.error }); abortController.abort(); + rejectGeneration(new AgentOnboardingGenerationTimeoutError()); }, GENERATION_TIMEOUT_MS); - activeGenerations.set(session.id, { abortController, timer }); + /* + FNXC:AgentOnboarding 2026-07-15-17:32: + User stop must settle the active prompt race even when the provider ignores cancellation. Keep a per-generation reject handle and identity-guard cleanup so a late stopped prompt cannot publish stale state or delete a newer retry's registration. + */ + const activeGeneration: ActiveGeneration = { abortController, timer, reject: rejectGeneration }; + activeGenerations.set(session.id, activeGeneration); try { - return await operation(); + return await Promise.race([operation(), interruption]); } finally { clearTimeout(timer); - activeGenerations.delete(session.id); + if (activeGenerations.get(session.id) === activeGeneration) { + activeGenerations.delete(session.id); + } } } +type OnboardingMessage = { + role: string; + content?: string | Array< + | { type: "text"; text: string } + | { type: "thinking"; thinking: string } + | { type: string } + >; +}; + +function extractLastAssistantResponse(messages: unknown, streamedOutput: string): string { + const assistant = (Array.isArray(messages) ? messages : []) + .filter((message): message is OnboardingMessage => ( + typeof message === "object" + && message !== null + && "role" in message + && (message as { role?: unknown }).role === "assistant" + )) + .pop(); + if (typeof assistant?.content === "string") return assistant.content.trim() || streamedOutput; + if (!Array.isArray(assistant?.content)) return streamedOutput; + + const textContent = assistant.content + .filter((block): block is { type: "text"; text: string } => block.type === "text" && "text" in block && typeof block.text === "string") + .map((block) => block.text) + .join(""); + const thinkingContent = assistant.content + .filter((block): block is { type: "thinking"; thinking: string } => block.type === "thinking" && "thinking" in block && typeof block.thinking === "string") + .map((block) => block.thinking) + .join(""); + /* + FNXC:AgentOnboarding 2026-07-15-16:15: + Pi can emit valid onboarding JSON in a thinking block alongside explanatory text. Select the first parseable text, thinking, or streamed candidate; only preserve the historical text-first fallback when none parses so the bounded reformat turn still receives the model's visible response. + */ + const candidates = [textContent, thinkingContent, streamedOutput] + .map((candidate) => candidate.trim()) + .filter((candidate) => candidate.length > 0); + for (const candidate of candidates) { + try { + parseAgentOnboardingResponse(candidate); + return candidate; + } catch { + // Try the next model-output surface before invoking the bounded recovery turn. + } + } + return candidates[0] ?? streamedOutput; +} + async function continueConversation(session: Session, message: string): Promise { if (!session.agent) throw new Error("Session agent not initialized"); const agent = session.agent; session.thinkingOutput = ""; try { - await runGenerationWithTimeout(session, async () => { - await agent.session.prompt(message); - const assistant = (agent.session.state.messages as Array<{ role: string; content?: string | Array<{ type: string; text: string }> }>).filter((m) => m.role === "assistant").pop(); - let responseText = session.thinkingOutput; - if (assistant?.content) { - if (typeof assistant.content === "string") responseText = assistant.content; - else responseText = assistant.content.filter((c) => c.type === "text").map((c) => c.text).join(""); - } - const parsed = parseAgentOnboardingResponse(responseText); - session.error = undefined; - session.updatedAt = new Date(); - if (parsed.type === "question") { - session.currentQuestion = parsed.data; - agentOnboardingStreamManager.broadcast(session.id, { type: "question", data: parsed.data }); - } else { - session.summary = parsed.data; - session.currentQuestion = undefined; - agentOnboardingStreamManager.broadcast(session.id, { type: "summary", data: parsed.data }); - agentOnboardingStreamManager.broadcast(session.id, { type: "complete" }); - } - }); + await runGenerationWithTimeout(session, () => agent.session.prompt(message)); + let responseText = extractLastAssistantResponse(agent.session.state.messages, session.thinkingOutput); + let parsed: ReturnType; + try { + parsed = parseAgentOnboardingResponse(responseText); + } catch { + /* + FNXC:AgentOnboarding 2026-07-15-14:32: + A malformed first interview response must not consume the recovery turn's generation budget. Run the reformat prompt as a distinct timed generation so it receives the full timeout and remains independently stoppable. + */ + session.thinkingOutput = ""; + await runGenerationWithTimeout(session, () => agent.session.prompt(REFORMAT_PROMPT)); + responseText = extractLastAssistantResponse(agent.session.state.messages, session.thinkingOutput); + parsed = parseAgentOnboardingResponse(responseText); + } + session.error = undefined; + session.updatedAt = new Date(); + if (parsed.type === "question") { + session.currentQuestion = parsed.data; + agentOnboardingStreamManager.broadcast(session.id, { type: "question", data: parsed.data }); + } else { + session.summary = parsed.data; + session.currentQuestion = undefined; + agentOnboardingStreamManager.broadcast(session.id, { type: "summary", data: parsed.data }); + agentOnboardingStreamManager.broadcast(session.id, { type: "complete" }); + } } catch (err) { + if (err instanceof AgentOnboardingGenerationStoppedError) { + return; + } + if (err instanceof AgentOnboardingGenerationTimeoutError) { + const expiredAgent = session.agent; + session.agentEpoch += 1; + session.agent = undefined; + try { expiredAgent?.session.dispose?.(); } catch { /* best-effort provider cancellation */ } + } session.error = err instanceof Error ? err.message : String(err); agentOnboardingStreamManager.broadcast(session.id, { type: "error", data: session.error }); } @@ -481,9 +591,14 @@ export function stopAgentOnboardingGeneration(sessionId: string): boolean { activeGenerations.delete(sessionId); const session = sessions.get(sessionId); if (session) { + const expiredAgent = session.agent; + session.agentEpoch += 1; + session.agent = undefined; + try { expiredAgent?.session.dispose?.(); } catch { /* best-effort provider cancellation */ } session.error = "Generation stopped by user. You can retry."; agentOnboardingStreamManager.broadcast(session.id, { type: "error", data: session.error }); } + active.reject(new AgentOnboardingGenerationStoppedError()); return true; }