diff --git a/apps/mobile/src/lib/threadActivity.test.ts b/apps/mobile/src/lib/threadActivity.test.ts index 8e2019da8..e6ac19090 100644 --- a/apps/mobile/src/lib/threadActivity.test.ts +++ b/apps/mobile/src/lib/threadActivity.test.ts @@ -550,6 +550,109 @@ describe("buildThreadFeed", () => { ]); }); + it.each([ + { + outcome: "completed", + message: "Prime Agent finished without sending a final response.", + tone: "info" as const, + status: null, + }, + { + outcome: "failed", + message: "Prime Agent stopped before sending a final response.", + tone: "error" as const, + status: null, + }, + ])("keeps the $outcome missing-response row outside settled-turn folding", (fixture) => { + const turnId = TurnId.make(`turn-${fixture.outcome}`); + const thread = makeThread({ + id: ThreadId.make(`thread-${fixture.outcome}`), + projectId: ProjectId.make("project-1"), + title: "Missing final response", + latestTurn: { + turnId, + state: fixture.outcome === "completed" ? "completed" : "error", + requestedAt: "2026-04-01T00:00:00.000Z", + startedAt: "2026-04-01T00:00:01.000Z", + completedAt: "2026-04-01T00:00:18.000Z", + assistantMessageId: null, + }, + activities: [ + makeActivity({ + id: EventId.make(`tool-${fixture.outcome}`), + kind: "tool.completed", + tone: "tool", + summary: "Read files", + createdAt: "2026-04-01T00:00:05.000Z", + turnId, + payload: { itemType: "file_read", status: "completed" }, + }), + makeActivity({ + id: EventId.make(`missing-response-${fixture.outcome}`), + kind: "turn.response.missing", + tone: fixture.tone, + summary: fixture.message, + createdAt: "2026-04-01T00:00:18.000Z", + turnId, + payload: { outcome: fixture.outcome }, + }), + ], + }); + + const feed = buildThreadFeed(thread); + const presented = deriveThreadFeedPresentation(feed, thread.latestTurn, new Set()); + expect(presented.map((entry) => entry.id)).toEqual([ + `turn-fold:${turnId}`, + `missing-response-${fixture.outcome}`, + ]); + expect(presented[1]).toMatchObject({ + type: "activity-group", + activities: [ + { + summary: fixture.message, + status: fixture.status, + toolLike: false, + terminalResponseNotice: true, + }, + ], + }); + }); + + it("folds ordinary runtime activity instead of treating it as a terminal response notice", () => { + const turnId = TurnId.make("turn-runtime-warning"); + const thread = makeThread({ + id: ThreadId.make("thread-runtime-warning"), + projectId: ProjectId.make("project-1"), + title: "Runtime warning", + latestTurn: { + turnId, + state: "completed", + requestedAt: "2026-04-01T00:00:00.000Z", + startedAt: "2026-04-01T00:00:01.000Z", + completedAt: "2026-04-01T00:00:18.000Z", + assistantMessageId: null, + }, + activities: [ + makeActivity({ + id: EventId.make("ordinary-runtime-warning"), + kind: "runtime.warning", + summary: "Reconnecting", + createdAt: "2026-04-01T00:00:18.000Z", + turnId, + payload: { message: "Reconnecting" }, + }), + ], + }); + + const feed = buildThreadFeed(thread); + expect(feed[0]).not.toMatchObject({ + activities: [{ terminalResponseNotice: true }], + }); + expect(deriveThreadFeedPresentation(feed, thread.latestTurn, new Set())).toEqual([ + expect.objectContaining({ type: "turn-fold", turnId }), + ]); + }); + it("measures a steer-superseded turn from its user boundary through trailing work", () => { const firstTurnId = TurnId.make("turn-1"); const secondTurnId = TurnId.make("turn-2"); diff --git a/apps/mobile/src/lib/threadActivity.ts b/apps/mobile/src/lib/threadActivity.ts index 1e93225f9..c951e8231 100644 --- a/apps/mobile/src/lib/threadActivity.ts +++ b/apps/mobile/src/lib/threadActivity.ts @@ -58,6 +58,8 @@ export interface ThreadFeedActivity { | "zap"; readonly toolLike: boolean; readonly status: "success" | "failure" | "neutral" | null; + /** Terminal response notice that must remain outside settled-turn work folding. */ + readonly terminalResponseNotice?: boolean; } const MAX_VISIBLE_WORK_LOG_ENTRIES = 1; @@ -561,7 +563,8 @@ function normalizeCompactToolLabel(value: string): string { return value.replace(/\s+(?:complete|completed)\s*$/i, "").trim(); } -function workLogEntryIsToolLike(entry: WorkLogEntry): boolean { +function workLogEntryIsToolLike(entry: WorkLogEntry | DerivedWorkLogEntry): boolean { + if ("activityKind" in entry && entry.activityKind === "turn.response.missing") return false; if (entry.tone === "tool" || entry.tone === "thinking" || entry.tone === "error") { return true; } @@ -1101,6 +1104,18 @@ function groupAdjacentActivities(entries: ReadonlyArray): Th continue; } + if (entry.activity.terminalResponseNotice === true) { + grouped.push({ + type: "activity-group", + id: entry.id, + createdAt: entry.createdAt, + turnId: entry.turnId, + activities: [entry.activity], + }); + openGroupActivities = null; + continue; + } + if (openGroupActivities !== null && openGroupTurnId === entry.turnId) { openGroupActivities.push(entry.activity); continue; @@ -1210,7 +1225,16 @@ function deriveThreadFeedTurnFolds( const terminalAssistantMessageId = terminalAssistantMessageIdByTurn.get(turnId); const hiddenEntryIds = new Set( - entries.filter((entry) => entry.id !== terminalAssistantMessageId).map((entry) => entry.id), + entries + .filter( + (entry) => + entry.id !== terminalAssistantMessageId && + !( + entry.type === "activity-group" && + entry.activities.some((activity) => activity.terminalResponseNotice === true) + ), + ) + .map((entry) => entry.id), ); if (hiddenEntryIds.size === 0) { continue; @@ -1585,6 +1609,9 @@ export function buildThreadFeed( icon: workEntryIcon(entry), toolLike: workLogEntryIsToolLike(entry), status: workEntryStatus(entry), + ...(entry.activityKind === "turn.response.missing" + ? { terminalResponseNotice: true } + : {}), }, }; }), diff --git a/apps/server/scripts/acp-mock-agent.ts b/apps/server/scripts/acp-mock-agent.ts index 8c9cf3b1d..c7c943327 100644 --- a/apps/server/scripts/acp-mock-agent.ts +++ b/apps/server/scripts/acp-mock-agent.ts @@ -16,6 +16,7 @@ const exitLogPath = process.env.T3_ACP_EXIT_LOG_PATH; const emitToolCalls = process.env.T3_ACP_EMIT_TOOL_CALLS === "1"; const emitInterleavedAssistantToolCalls = process.env.T3_ACP_EMIT_INTERLEAVED_ASSISTANT_TOOL_CALLS === "1"; +const omitInterleavedFinalText = process.env.T3_ACP_OMIT_INTERLEAVED_FINAL_TEXT === "1"; const emitGenericToolPlaceholders = process.env.T3_ACP_EMIT_GENERIC_TOOL_PLACEHOLDERS === "1"; const emitAskQuestion = process.env.T3_ACP_EMIT_ASK_QUESTION === "1"; const emitXAiAskUserQuestion = process.env.T3_ACP_EMIT_XAI_ASK_USER_QUESTION === "1"; @@ -620,13 +621,15 @@ const program = Effect.gen(function* () { }, }); - yield* agent.client.sessionUpdate({ - sessionId: requestedSessionId, - update: { - sessionUpdate: "agent_message_chunk", - content: { type: "text", text: "after tool" }, - }, - }); + if (!omitInterleavedFinalText) { + yield* agent.client.sessionUpdate({ + sessionId: requestedSessionId, + update: { + sessionUpdate: "agent_message_chunk", + content: { type: "text", text: "after tool" }, + }, + }); + } return { stopReason: "end_turn" }; } diff --git a/apps/server/src/orchestration/Layers/CheckpointReactor.test.ts b/apps/server/src/orchestration/Layers/CheckpointReactor.test.ts index b732ffbfc..d36b5e95c 100644 --- a/apps/server/src/orchestration/Layers/CheckpointReactor.test.ts +++ b/apps/server/src/orchestration/Layers/CheckpointReactor.test.ts @@ -530,6 +530,10 @@ describe("CheckpointReactor", () => { (entry) => entry.latestTurn?.turnId === "turn-1" && entry.checkpoints.length === 1, ); expect(thread.checkpoints[0]?.checkpointTurnCount).toBe(1); + expect( + (thread.checkpoints[0] as { readonly assistantMessageId: string | null } | undefined) + ?.assistantMessageId, + ).toBeNull(); expect( gitRefExists(harness.cwd, checkpointRefForThreadTurn(ThreadId.make("thread-1"), 0)), ).toBe(true); diff --git a/apps/server/src/orchestration/Layers/CheckpointReactor.ts b/apps/server/src/orchestration/Layers/CheckpointReactor.ts index c5bbb71ae..8c165f023 100644 --- a/apps/server/src/orchestration/Layers/CheckpointReactor.ts +++ b/apps/server/src/orchestration/Layers/CheckpointReactor.ts @@ -311,8 +311,7 @@ const make = Effect.gen(function* () { input.assistantMessageId ?? input.thread.messages .toReversed() - .find((entry) => entry.role === "assistant" && entry.turnId === input.turnId)?.id ?? - MessageId.make(`assistant:${input.turnId}`); + .find((entry) => entry.role === "assistant" && entry.turnId === input.turnId)?.id; yield* orchestrationEngine.dispatch({ type: "thread.turn.diff.complete", @@ -323,7 +322,7 @@ const make = Effect.gen(function* () { checkpointRef: targetCheckpointRef, status: input.status, files, - assistantMessageId, + ...(assistantMessageId === undefined ? {} : { assistantMessageId }), checkpointTurnCount: input.turnCount, createdAt: input.createdAt, }); diff --git a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.activity.test.ts b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.activity.test.ts index cfe2ed65e..1702ff19f 100644 --- a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.activity.test.ts +++ b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.activity.test.ts @@ -17,6 +17,125 @@ const base = { threadId: ThreadId.make("thread-1"), }; +describe("runtimeEventToActivities missing final response", () => { + it.each([ + { + type: "runtime.warning" as const, + provider: ProviderDriverKind.make("primeAgent"), + outcome: "completed" as const, + message: "Prime Agent finished without sending a final response.", + tone: "info" as const, + }, + { + type: "runtime.error" as const, + provider: ProviderDriverKind.make("primeAgent"), + outcome: "failed" as const, + message: "Prime Agent stopped before sending a final response.", + tone: "error" as const, + }, + ])("projects the $outcome marker without provider detail", (fixture) => { + const [activity] = runtimeEventToActivities({ + ...base, + provider: fixture.provider, + type: fixture.type, + eventId: EventId.make(`missing-response-${fixture.outcome}`), + turnId: TurnId.make(`turn-${fixture.outcome}`), + payload: { + message: fixture.message, + detail: { kind: "missing-final-response", outcome: fixture.outcome }, + }, + } satisfies ProviderRuntimeEvent); + + expect(activity).toMatchObject({ + kind: "turn.response.missing", + summary: fixture.message, + tone: fixture.tone, + payload: { outcome: fixture.outcome }, + turnId: `turn-${fixture.outcome}`, + }); + expect(activity?.payload).toEqual({ outcome: fixture.outcome }); + expect(JSON.stringify(activity)).not.toContain("missing-final-response"); + }); + + it.each([ + { + name: "an extra field", + type: "runtime.warning" as const, + outcome: "completed" as const, + detail: { kind: "missing-final-response", outcome: "completed", private: "raw" }, + }, + { + name: "a mismatched outcome", + type: "runtime.warning" as const, + outcome: "failed" as const, + detail: { kind: "missing-final-response", outcome: "failed" }, + }, + { + name: "a non-object detail", + type: "runtime.error" as const, + outcome: "failed" as const, + detail: "missing-final-response", + }, + ])("leaves $type unchanged when its marker has $name", (fixture) => { + const [activity] = runtimeEventToActivities({ + ...base, + type: fixture.type, + eventId: EventId.make(`malformed-${fixture.outcome}`), + payload: { + message: "Ordinary provider event", + detail: fixture.detail, + }, + } satisfies ProviderRuntimeEvent); + + expect(activity?.kind).toBe(fixture.type); + }); + + it("does not reclassify another provider using the same detail shape", () => { + const [activity] = runtimeEventToActivities({ + ...base, + type: "runtime.error", + eventId: EventId.make("other-provider-missing-marker"), + payload: { + message: "Codex provider error", + detail: { kind: "missing-final-response", outcome: "failed" }, + }, + } satisfies ProviderRuntimeEvent); + + expect(activity).toMatchObject({ + kind: "runtime.error", + summary: "Runtime error", + payload: { message: "Codex provider error" }, + }); + }); + + it("preserves ordinary warning and error activity presentations", () => { + const [warning] = runtimeEventToActivities({ + ...base, + type: "runtime.warning", + eventId: EventId.make("ordinary-warning"), + payload: { message: "Reconnecting", detail: { willRetry: true } }, + } satisfies ProviderRuntimeEvent); + const [error] = runtimeEventToActivities({ + ...base, + type: "runtime.error", + eventId: EventId.make("ordinary-error"), + payload: { message: "Provider failed", detail: { private: "unchanged-drop" } }, + } satisfies ProviderRuntimeEvent); + + expect(warning).toMatchObject({ + kind: "runtime.warning", + summary: "Reconnecting", + payload: { message: "Reconnecting", detail: { willRetry: true } }, + }); + expect(error).toMatchObject({ + kind: "runtime.error", + summary: "Runtime error", + payload: { message: "Provider failed" }, + }); + expect(error?.payload).toEqual({ message: "Provider failed" }); + }); +}); + describe("runtimeEventToActivities task progress", () => { it("persists usage independently from replaceable activity", () => { const taskId = RuntimeTaskId.make("agent-1"); diff --git a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts index f3390cee7..b7bb5f5a3 100644 --- a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts +++ b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts @@ -3179,6 +3179,65 @@ describe("ProviderRuntimeIngestion", () => { expect(activityPayload?.message).toBe("runtime activity exploded"); }); + it("persists a textless completion notice without fabricating an assistant message", async () => { + const harness = await createHarness(); + const now = "2026-01-01T00:00:00.000Z"; + const turnId = asTurnId("turn-textless-completion"); + + harness.emit({ + type: "turn.started", + eventId: asEventId("evt-textless-turn-started"), + provider: ProviderDriverKind.make("primeAgent"), + createdAt: now, + threadId: asThreadId("thread-1"), + turnId, + payload: {}, + }); + harness.emit({ + type: "runtime.warning", + eventId: asEventId("evt-textless-completion-notice"), + provider: ProviderDriverKind.make("primeAgent"), + createdAt: "2026-01-01T00:00:01.000Z", + threadId: asThreadId("thread-1"), + turnId, + payload: { + message: "Prime Agent finished without sending a final response.", + detail: { kind: "missing-final-response", outcome: "completed" }, + }, + }); + harness.emit({ + type: "turn.completed", + eventId: asEventId("evt-textless-turn-completed"), + provider: ProviderDriverKind.make("primeAgent"), + createdAt: "2026-01-01T00:00:02.000Z", + threadId: asThreadId("thread-1"), + turnId, + payload: { state: "completed" }, + }); + + const thread = await waitForThread( + harness.readModel, + (entry) => + entry.latestTurn?.turnId === turnId && + entry.latestTurn.state === "completed" && + entry.activities.some((activity) => activity.id === "evt-textless-completion-notice"), + ); + const activity = thread.activities.find( + (entry: ProviderRuntimeTestActivity) => entry.id === "evt-textless-completion-notice", + ); + expect(activity).toMatchObject({ + kind: "turn.response.missing", + summary: "Prime Agent finished without sending a final response.", + payload: { outcome: "completed" }, + turnId, + }); + expect(Object.keys((activity?.payload ?? {}) as Record)).toEqual(["outcome"]); + expect( + thread.messages.some((message) => message.role === "assistant" && message.turnId === turnId), + ).toBe(false); + expect(thread.latestTurn?.assistantMessageId).toBeNull(); + }); + it("keeps the session running when a runtime.warning arrives during an active turn", async () => { const harness = await createHarness(); const now = "2026-01-01T00:00:00.000Z"; diff --git a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts index d7e8db740..5e2af5461 100644 --- a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts +++ b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts @@ -236,6 +236,33 @@ function truncateDetail(value: string, limit = 180): string { return value.length > limit ? `${value.slice(0, limit - 3)}...` : value; } +function missingFinalResponseOutcome( + event: Extract, +): "completed" | "failed" | undefined { + if (event.provider !== "primeAgent") return undefined; + const detail = event.payload.detail; + if (!detail || typeof detail !== "object" || Array.isArray(detail)) { + return undefined; + } + const marker = detail as Record; + const keys = Object.keys(marker); + if ( + keys.length !== 2 || + !keys.includes("kind") || + !keys.includes("outcome") || + marker.kind !== "missing-final-response" + ) { + return undefined; + } + if (event.type === "runtime.warning" && marker.outcome === "completed") { + return "completed"; + } + if (event.type === "runtime.error" && marker.outcome === "failed") { + return "failed"; + } + return undefined; +} + function normalizeProposedPlanMarkdown(planMarkdown: string | undefined): string | undefined { const trimmed = planMarkdown?.trim(); if (!trimmed) { @@ -657,6 +684,21 @@ export function runtimeEventToActivities( } case "runtime.error": { + const missingResponseOutcome = missingFinalResponseOutcome(event); + if (missingResponseOutcome !== undefined) { + return [ + { + id: event.eventId, + createdAt: event.createdAt, + tone: "error", + kind: "turn.response.missing", + summary: truncateDetail(event.payload.message, 120), + payload: { outcome: missingResponseOutcome }, + turnId: toTurnId(event.turnId) ?? null, + ...maybeSequence, + }, + ]; + } return [ { id: event.eventId, @@ -694,6 +736,21 @@ export function runtimeEventToActivities( } case "runtime.warning": { + const missingResponseOutcome = missingFinalResponseOutcome(event); + if (missingResponseOutcome !== undefined) { + return [ + { + id: event.eventId, + createdAt: event.createdAt, + tone: "info", + kind: "turn.response.missing", + summary: truncateDetail(event.payload.message, 120), + payload: { outcome: missingResponseOutcome }, + turnId: toTurnId(event.turnId) ?? null, + ...maybeSequence, + }, + ]; + } return [ { id: event.eventId, diff --git a/apps/server/src/provider/Layers/PrimeAgentAdapter.test.ts b/apps/server/src/provider/Layers/PrimeAgentAdapter.test.ts index df658cc0d..9ef24e847 100644 --- a/apps/server/src/provider/Layers/PrimeAgentAdapter.test.ts +++ b/apps/server/src/provider/Layers/PrimeAgentAdapter.test.ts @@ -20,6 +20,7 @@ import { ProviderInstanceId, type ProviderRuntimeEvent, ThreadId, + TurnId, } from "@t3tools/contracts"; import { ServerConfig } from "../../config.ts"; @@ -437,7 +438,7 @@ exec ${process.execPath} ${mockAgentPath} "$@" const failedPrompt = yield* failingAdapter .sendTurn({ threadId: failingThreadId, input: "fail", attachments: [] }) .pipe(Effect.result); - assert.equal(failedPrompt._tag, "Failure"); + assert.equal(failedPrompt._tag, "Success"); const failedTurnId = failureEvents.find( (event) => event.threadId === failingThreadId && event.type === "turn.started", )?.turnId; @@ -532,3 +533,228 @@ exec ${process.execPath} ${mockAgentPath} "$@" yield* Fiber.interrupt(eventFiber); }).pipe(Effect.scoped, Effect.provide(testLayer)), ); + +it.effect( + "emits safe missing-final-response notices only for authoritative textless terminals", + () => + Effect.gen(function* () { + const tempDir = yield* Effect.promise(() => + NodeFSP.mkdtemp(NodePath.join(NodeOS.tmpdir(), "prime-agent-acp-textless-test-")), + ); + const wrapperPath = NodePath.join(tempDir, "private-native-path", "fake-prime-agent.sh"); + yield* Effect.promise(() => + NodeFSP.mkdir(NodePath.dirname(wrapperPath), { recursive: true }), + ); + yield* Effect.promise(() => + NodeFSP.writeFile( + wrapperPath, + `#!/bin/sh +exec ${process.execPath} ${mockAgentPath} "$@" +`, + "utf8", + ), + ); + yield* Effect.promise(() => NodeFSP.chmod(wrapperPath, 0o755)); + + let caseIndex = 0; + const runTerminalTurn = Effect.fn("PrimeAgentAdapter.test.runTerminalTurn")(function* ( + label: string, + environment: NodeJS.ProcessEnv, + ) { + caseIndex += 1; + const adapter = yield* makePrimeAgentAdapter(decodeSettings({ binaryPath: wrapperPath }), { + instanceId: ProviderInstanceId.make(`primeAgent-textless-${caseIndex}`), + environment: { ...process.env, ...environment }, + }); + const threadId = ThreadId.make(`textless-${label}-${caseIndex}`); + const completed = yield* Deferred.make(); + const events: Array = []; + const eventFiber = yield* adapter.streamEvents.pipe( + Stream.runForEach((event) => + Effect.gen(function* () { + events.push(event); + if (event.type === "turn.completed" && event.threadId === threadId) { + yield* Deferred.succeed(completed, undefined).pipe(Effect.ignore); + } + }), + ), + Effect.forkChild, + ); + yield* Effect.yieldNow; + + yield* adapter.startSession({ + threadId, + provider: ProviderDriverKind.make("primeAgent"), + cwd: process.cwd(), + runtimeMode: "full-access", + }); + const result = yield* adapter + .sendTurn({ threadId, input: label, attachments: [] }) + .pipe(Effect.result); + yield* Deferred.await(completed); + const turnId = events.find( + (event) => event.type === "turn.started" && event.threadId === threadId, + )?.turnId; + assert.isDefined(turnId); + const turnEvents = events.filter((event) => event.turnId === turnId); + const providerThread = yield* adapter.readThread(threadId); + + yield* adapter.stopSession(threadId); + yield* Fiber.interrupt(eventFiber); + return { result, turnEvents, providerThread }; + }); + + const toolOnly = yield* runTerminalTurn("tool-only", { + T3_ACP_EMIT_INTERLEAVED_ASSISTANT_TOOL_CALLS: "1", + T3_ACP_OMIT_INTERLEAVED_FINAL_TEXT: "1", + }); + assert.equal(toolOnly.result._tag, "Success"); + assert.deepEqual( + toolOnly.turnEvents + .filter((event) => event.type === "content.delta") + .map((event) => event.payload.delta), + ["before tool"], + ); + const toolOnlyWarning = toolOnly.turnEvents.find((event) => event.type === "runtime.warning"); + assert.equal(toolOnlyWarning?.type, "runtime.warning"); + if (toolOnlyWarning?.type === "runtime.warning") { + assert.deepEqual(toolOnlyWarning.payload, { + message: "Prime Agent finished without sending a final response.", + detail: { kind: "missing-final-response", outcome: "completed" }, + }); + const serializedNotice = encodeUnknownJsonString(toolOnlyWarning); + assert.notInclude(serializedNotice, "echo"); + assert.notInclude(serializedNotice, "tool-call-1"); + assert.notInclude(serializedNotice, wrapperPath); + } + const toolOnlyTypes = toolOnly.turnEvents.map((event) => event.type); + assert.isBelow( + toolOnlyTypes.lastIndexOf("runtime.warning"), + toolOnlyTypes.lastIndexOf("turn.completed"), + ); + assert.equal(toolOnly.turnEvents.at(-1)?.type, "turn.completed"); + const toolOnlyTerminal = toolOnly.turnEvents.find(isTurnCompletedEvent); + assert.equal(toolOnlyTerminal?.payload.state, "completed"); + assert.isTrue(toolOnly.providerThread.turns.every((turn) => turn.items.length === 0)); + + const normalText = yield* runTerminalTurn("normal-text", { + T3_ACP_PROMPT_RESPONSE_TEXT: "normal final response", + }); + assert.equal(normalText.result._tag, "Success"); + assert.deepEqual( + normalText.turnEvents + .filter((event) => event.type === "content.delta") + .map((event) => event.payload.delta), + ["normal final response"], + ); + assert.isFalse( + normalText.turnEvents.some( + (event) => event.type === "runtime.warning" || event.type === "runtime.error", + ), + ); + + const failed = yield* runTerminalTurn("failed", { T3_ACP_FAIL_PROMPT: "1" }); + assert.equal(failed.result._tag, "Success"); + const serializedFailureResult = encodeUnknownJsonString(failed.result); + assert.notInclude(serializedFailureResult, "Mock prompt failure"); + assert.notInclude(serializedFailureResult, wrapperPath); + const failureNotice = failed.turnEvents.find((event) => event.type === "runtime.error"); + assert.equal(failureNotice?.type, "runtime.error"); + if (failureNotice?.type === "runtime.error") { + assert.deepEqual(failureNotice.payload, { + message: "Prime Agent stopped before sending a final response.", + detail: { kind: "missing-final-response", outcome: "failed" }, + }); + const serializedNotice = encodeUnknownJsonString(failureNotice); + assert.notInclude(serializedNotice, "Mock prompt failure"); + assert.notInclude(serializedNotice, wrapperPath); + } + const serializedFailureEvents = encodeUnknownJsonString(failed.turnEvents); + assert.notInclude(serializedFailureEvents, "Mock prompt failure"); + assert.notInclude(serializedFailureEvents, wrapperPath); + const failureTypes = failed.turnEvents.map((event) => event.type); + assert.isBelow( + failureTypes.lastIndexOf("runtime.error"), + failureTypes.lastIndexOf("turn.completed"), + ); + assert.equal(failed.turnEvents.at(-1)?.type, "turn.completed"); + const failedTerminal = failed.turnEvents.find(isTurnCompletedEvent); + assert.equal(failedTerminal?.payload.state, "failed"); + + const privateThought = yield* runTerminalTurn("private-thought", { + T3_ACP_EMIT_PRIVATE_THOUGHT_CHUNK: "1", + T3_ACP_PROMPT_RESPONSE_TEXT: " ", + }); + assert.equal(privateThought.result._tag, "Success"); + const privateThoughtSerialized = encodeUnknownJsonString(privateThought.turnEvents); + assert.notInclude(privateThoughtSerialized, "private-thought-sentinel"); + assert.notInclude(privateThoughtSerialized, "agent_thought_chunk"); + const privateThoughtWarning = privateThought.turnEvents.find( + (event) => event.type === "runtime.warning", + ); + assert.equal(privateThoughtWarning?.type, "runtime.warning"); + if (privateThoughtWarning?.type === "runtime.warning") { + assert.deepEqual(privateThoughtWarning.payload.detail, { + kind: "missing-final-response", + outcome: "completed", + }); + } + + const cancelledAdapter = yield* makePrimeAgentAdapter( + decodeSettings({ binaryPath: wrapperPath }), + { + instanceId: ProviderInstanceId.make("primeAgent-textless-cancelled"), + environment: { ...process.env, T3_ACP_HANG_PROMPT_FOREVER: "1" }, + }, + ); + const cancelledThreadId = ThreadId.make("textless-cancelled"); + const cancelledStarted = yield* Deferred.make(); + const cancelledCompleted = yield* Deferred.make(); + const cancelledEvents: Array = []; + const cancelledEventFiber = yield* cancelledAdapter.streamEvents.pipe( + Stream.runForEach((event) => + Effect.gen(function* () { + cancelledEvents.push(event); + if ( + event.type === "turn.started" && + event.threadId === cancelledThreadId && + event.turnId !== undefined + ) { + yield* Deferred.succeed(cancelledStarted, event.turnId).pipe(Effect.ignore); + } + if (event.type === "turn.completed" && event.threadId === cancelledThreadId) { + yield* Deferred.succeed(cancelledCompleted, undefined).pipe(Effect.ignore); + } + }), + ), + Effect.forkChild, + ); + yield* Effect.yieldNow; + yield* cancelledAdapter.startSession({ + threadId: cancelledThreadId, + provider: ProviderDriverKind.make("primeAgent"), + cwd: process.cwd(), + runtimeMode: "full-access", + }); + const cancelledTurnFiber = yield* cancelledAdapter + .sendTurn({ threadId: cancelledThreadId, input: "cancel", attachments: [] }) + .pipe(Effect.forkChild); + const cancelledTurnId = yield* Deferred.await(cancelledStarted); + yield* cancelledAdapter.interruptTurn(cancelledThreadId, cancelledTurnId); + yield* Deferred.await(cancelledCompleted); + yield* Fiber.join(cancelledTurnFiber); + const cancelledTurnEvents = cancelledEvents.filter( + (event) => event.turnId === cancelledTurnId, + ); + assert.isFalse( + cancelledTurnEvents.some( + (event) => event.type === "runtime.warning" || event.type === "runtime.error", + ), + ); + const cancelledTerminal = cancelledTurnEvents.find(isTurnCompletedEvent); + assert.equal(cancelledTerminal?.payload.state, "cancelled"); + + yield* cancelledAdapter.stopSession(cancelledThreadId); + yield* Fiber.interrupt(cancelledEventFiber); + }).pipe(Effect.scoped, Effect.provide(testLayer)), +); diff --git a/apps/server/src/provider/Layers/PrimeAgentAdapter.ts b/apps/server/src/provider/Layers/PrimeAgentAdapter.ts index c60ce1b44..0c9eb8132 100644 --- a/apps/server/src/provider/Layers/PrimeAgentAdapter.ts +++ b/apps/server/src/provider/Layers/PrimeAgentAdapter.ts @@ -55,6 +55,12 @@ import { isPrimeAgentCompatibleResumeCursor, PRIME_AGENT_ACP_RESUME_CURSOR, } from "../prime/PrimeAgentResumeCursor.ts"; +import { + PRIME_AGENT_FINISHED_WITHOUT_FINAL_RESPONSE, + PRIME_AGENT_STOPPED_WITHOUT_FINAL_RESPONSE, + PRIME_AGENT_TURN_FAILED, + primeAgentMissingFinalResponseDetail, +} from "../prime/PrimeAgentTerminalResponse.ts"; import { type EventNdjsonLogger, makeEventNdjsonLogger } from "./EventNdjsonLogger.ts"; const PROVIDER = ProviderDriverKind.make("primeAgent"); @@ -72,6 +78,7 @@ interface PrimeAgentActiveTurn { readonly id: TurnId; readonly cancellation: Deferred.Deferred; cancellationRequested: boolean; + hasPublicAssistantTextAfterLatestToolBoundary: boolean; } interface PrimeAgentSessionContext { @@ -233,12 +240,13 @@ export function makePrimeAgentAdapter( turnId: TurnId, outcome: | { readonly state: "completed"; readonly stopReason: EffectAcpSchema.StopReason | null } - | { readonly state: "failed"; readonly errorMessage: string } + | { + readonly state: "failed"; + readonly errorMessage: string; + readonly terminalFailure?: boolean; + } | { readonly state: "cancelled" }, - completedPrompt?: { - readonly prompt: ReadonlyArray; - readonly result: EffectAcpSchema.PromptResponse; - }, + recordCompletedTurn = false, ) => Effect.gen(function* () { if ( @@ -267,8 +275,30 @@ export function makePrimeAgentAdapter( ctx.stopRequested || ctx.activeTurn.cancellationRequested ? ({ state: "cancelled" } as const) : outcome; - if (completedPrompt && effectiveOutcome.state !== "cancelled") { - ctx.turns.push({ id: turnId, items: [completedPrompt] }); + if (recordCompletedTurn && effectiveOutcome.state !== "cancelled") { + ctx.turns.push({ id: turnId, items: [] }); + } + + const missingFinalResponse = !ctx.activeTurn.hasPublicAssistantTextAfterLatestToolBoundary; + const shouldEmitMissingFinalResponseNotice = + missingFinalResponse && + (effectiveOutcome.state === "completed" || + (effectiveOutcome.state === "failed" && effectiveOutcome.terminalFailure === true)); + if (shouldEmitMissingFinalResponseNotice) { + const failed = effectiveOutcome.state === "failed"; + yield* offerRuntimeEvent({ + type: failed ? "runtime.error" : "runtime.warning", + ...(yield* makeEventStamp()), + provider: PROVIDER, + threadId: ctx.threadId, + turnId, + payload: { + message: failed + ? PRIME_AGENT_STOPPED_WITHOUT_FINAL_RESPONSE + : PRIME_AGENT_FINISHED_WITHOUT_FINAL_RESPONSE, + detail: primeAgentMissingFinalResponseDetail(failed ? "failed" : "completed"), + }, + }); } const { activeTurnId: _activeTurnId, ...readySession } = ctx.session; @@ -287,7 +317,13 @@ export function makePrimeAgentAdapter( turnId, payload: effectiveOutcome.state === "failed" - ? { state: "failed", errorMessage: effectiveOutcome.errorMessage } + ? { + state: "failed", + errorMessage: + missingFinalResponse && effectiveOutcome.terminalFailure === true + ? PRIME_AGENT_STOPPED_WITHOUT_FINAL_RESPONSE + : effectiveOutcome.errorMessage, + } : effectiveOutcome.state === "cancelled" ? { state: "cancelled", stopReason: "cancelled" } : { @@ -303,15 +339,19 @@ export function makePrimeAgentAdapter( turnId: TurnId, outcome: | { readonly state: "completed"; readonly stopReason: EffectAcpSchema.StopReason | null } - | { readonly state: "failed"; readonly errorMessage: string } + | { + readonly state: "failed"; + readonly errorMessage: string; + readonly terminalFailure?: boolean; + } | { readonly state: "cancelled" }, - completedPrompt?: { - readonly prompt: ReadonlyArray; - readonly result: EffectAcpSchema.PromptResponse; - }, + recordCompletedTurn = false, ) => Effect.uninterruptible( - withThreadLock(ctx.threadId, settleActiveTurnLocked(ctx, turnId, outcome, completedPrompt)), + withThreadLock( + ctx.threadId, + settleActiveTurnLocked(ctx, turnId, outcome, recordCompletedTurn), + ), ); /** Must be called while holding the thread lock. */ @@ -534,6 +574,9 @@ export function makePrimeAgentAdapter( return; } case "ToolCallUpdated": + if (ctx.activeTurn?.id === notificationTurnId) { + ctx.activeTurn.hasPublicAssistantTextAfterLatestToolBoundary = false; + } yield* logNative(ctx.threadId, "session/update", event.rawPayload); yield* offerRuntimeEvent( makeAcpToolCallEvent({ @@ -547,6 +590,9 @@ export function makePrimeAgentAdapter( ); return; case "ContentDelta": + if (event.text.trim().length > 0 && ctx.activeTurn?.id === notificationTurnId) { + ctx.activeTurn.hasPublicAssistantTextAfterLatestToolBoundary = true; + } yield* logNative(ctx.threadId, "session/update", event.rawPayload); yield* offerRuntimeEvent( makeAcpContentDeltaEvent({ @@ -723,6 +769,7 @@ export function makePrimeAgentAdapter( id: turnId, cancellation: yield* Deferred.make(), cancellationRequested: false, + hasPublicAssistantTextAfterLatestToolBoundary: false, }; ctx.activeTurn = activeTurn; ctx.lastPlanFingerprint = undefined; @@ -762,7 +809,7 @@ export function makePrimeAgentAdapter( result.stopReason === "cancelled" ? { state: "cancelled" } : { state: "completed", stopReason: result.stopReason ?? null }, - { prompt, result }, + true, ); if (!settled && !activeTurn.cancellationRequested && !ctx.stopRequested) { return yield* new ProviderAdapterRequestError({ @@ -775,13 +822,21 @@ export function makePrimeAgentAdapter( }); return yield* restore(promptEffect).pipe( - Effect.catch((error) => + Effect.catch(() => Effect.gen(function* () { yield* settleActiveTurn(ctx, turnId, { state: "failed", - errorMessage: error.message, + errorMessage: PRIME_AGENT_TURN_FAILED, + terminalFailure: true, }); - return yield* error; + // Admission succeeded and the runtime event stream already + // carries the authoritative failed terminal. Returning the + // admitted turn prevents a second turn-start failure activity. + return { + threadId: input.threadId, + turnId, + resumeCursor: ctx.session.resumeCursor, + }; }), ), Effect.ensuring( diff --git a/apps/server/src/provider/prime/PrimeAgentDaemonAdapter.test.ts b/apps/server/src/provider/prime/PrimeAgentDaemonAdapter.test.ts index 990941af2..556e2c94f 100644 --- a/apps/server/src/provider/prime/PrimeAgentDaemonAdapter.test.ts +++ b/apps/server/src/provider/prime/PrimeAgentDaemonAdapter.test.ts @@ -38,6 +38,7 @@ import type { PrimeDaemonEvent, PrimeDaemonMessage } from "./PrimeAgentDaemonEve import type { PrimeAgentDaemonManager } from "./PrimeAgentDaemonManager.ts"; import { makePrimeAgentDaemonAdapter, + PRIME_AGENT_FAILED_RUN_SETTLEMENT_GRACE_MS, PRIME_AGENT_SIDE_QUESTION_TIMEOUT_MS, type PrimeAgentDaemonAdapterLiveOptions, } from "./PrimeAgentDaemonAdapter.ts"; @@ -4473,6 +4474,14 @@ describe("PrimeAgentDaemonAdapter", () => { "session.state.changed", ]); expect(turnEvents.filter((event) => event.type === "turn.completed")).toHaveLength(1); + const assistantStarts = turnEvents.filter( + (event) => + event.type === "item.started" && event.payload.itemType === "assistant_message", + ); + expect(assistantStarts.map((event) => event.itemId)).toEqual([ + `assistant:${result.turnId}:segment:0`, + `assistant:${result.turnId}:segment:1`, + ]); expect(turnEvents.find((event) => event.type === "turn.completed")).toMatchObject({ payload: { state: "completed", @@ -4517,6 +4526,200 @@ describe("PrimeAgentDaemonAdapter", () => { ).pipe(Effect.provide(testLayer)), ); + it.effect("keeps an automatic reconnect continuation attached to the original turn", () => + Effect.scoped( + Effect.gen(function* () { + const captures = makeCaptures(); + const adapter = yield* makePrimeAgentDaemonAdapter(decodeSettings({}), manager, { + instanceId, + runtimeFactory: fakeRuntimeFactory(captures), + }); + const subscription = yield* subscribe(adapter); + yield* adapter.startSession({ threadId, cwd: process.cwd(), runtimeMode: "full-access" }); + yield* awaitObservedType(subscription.observed, "thread.started"); + + const running = yield* adapter + .sendTurn({ threadId, input: "continue after reconnect" }) + .pipe(Effect.forkChild); + const started = yield* awaitObservedType(subscription.observed, "turn.started"); + yield* Queue.take(captures.promptObserved!); + + const failedMessage = { + ...assistantMessage("", "error"), + errorMessage: "PRIVATE reconnect teardown cause", + } satisfies PrimeDaemonMessage; + yield* offer(captures, { _tag: "MessageStarted", message: failedMessage }); + yield* offer(captures, { _tag: "MessageCompleted", message: failedMessage }); + yield* offer(captures, { _tag: "RunCompleted", messages: [failedMessage] }); + yield* offer(captures, { + ...initialSnapshot(), + state: { ...initialSnapshot().state, isStreaming: true, retryAttempt: 1 }, + }); + yield* offer(captures, { _tag: "SessionInfoChanged", name: "active resync barrier" }); + yield* awaitObservedType(subscription.observed, "thread.metadata.updated"); + yield* TestClock.adjust(PRIME_AGENT_FAILED_RUN_SETTLEMENT_GRACE_MS + 1); + expect(running.pollUnsafe()).toBeUndefined(); + + const toolOnlyMessage = assistantMessage("preparing the tool", "toolUse"); + yield* offer(captures, { _tag: "MessageStarted", message: toolOnlyMessage }); + yield* offer(captures, { _tag: "MessageCompleted", message: toolOnlyMessage }); + yield* offer(captures, { _tag: "RunCompleted", messages: [toolOnlyMessage] }); + yield* offer(captures, { _tag: "RunStarted" }); + + const finalMessage = assistantMessage("the recovered final response"); + yield* offer(captures, { + _tag: "AssistantStream", + phase: "delta", + kind: "text", + delta: finalMessage.text, + }); + yield* offer(captures, { _tag: "MessageCompleted", message: finalMessage }); + yield* offer(captures, { _tag: "RunCompleted", messages: [finalMessage] }); + + const result = yield* Fiber.join(running); + expect(result.turnId).toBe(started.turnId); + const turnEvents = subscription.events.filter((event) => event.turnId === result.turnId); + expect(turnEvents.filter((event) => event.type === "turn.completed")).toHaveLength(1); + const recoveredItemId = `assistant:${result.turnId}:segment:2`; + const recoveredStartIndex = turnEvents.findIndex( + (event) => event.type === "item.started" && event.itemId === recoveredItemId, + ); + const recoveredDeltaIndex = turnEvents.findIndex( + (event) => + event.type === "content.delta" && + event.itemId === recoveredItemId && + event.payload.streamKind === "assistant_text" && + event.payload.delta === finalMessage.text, + ); + expect(recoveredStartIndex).toBeGreaterThanOrEqual(0); + expect(recoveredDeltaIndex).toBeGreaterThan(recoveredStartIndex); + expect( + turnEvents.some( + (event) => + (event.type === "runtime.warning" || event.type === "runtime.error") && + typeof event.payload.detail === "object" && + event.payload.detail !== null && + "kind" in event.payload.detail && + event.payload.detail.kind === "missing-final-response", + ), + ).toBe(false); + expect(encodeUnknownJson(subscription.events)).not.toContain( + "PRIVATE reconnect teardown cause", + ); + const providerThread = yield* adapter.readThread(threadId); + expect(providerThread.turns.every((turn) => turn.items.length === 0)).toBe(true); + expect(encodeUnknownJson(providerThread)).not.toContain("PRIVATE reconnect teardown cause"); + yield* Fiber.interrupt(subscription.fiber); + }), + ).pipe(Effect.provide(testLayer)), + ); + + it.effect("settles a genuinely failed daemon run after the bounded retry handoff", () => + Effect.scoped( + Effect.gen(function* () { + const captures = makeCaptures(); + const adapter = yield* makePrimeAgentDaemonAdapter(decodeSettings({}), manager, { + instanceId, + runtimeFactory: fakeRuntimeFactory(captures), + }); + const subscription = yield* subscribe(adapter); + yield* adapter.startSession({ threadId, cwd: process.cwd(), runtimeMode: "full-access" }); + yield* awaitObservedType(subscription.observed, "thread.started"); + + const running = yield* adapter + .sendTurn({ threadId, input: "fail safely" }) + .pipe(Effect.forkChild); + const started = yield* awaitObservedType(subscription.observed, "turn.started"); + yield* Queue.take(captures.promptObserved!); + const failedMessage = { + ...assistantMessage("", "error"), + errorMessage: "PRIVATE provider error", + } satisfies PrimeDaemonMessage; + yield* offer(captures, { _tag: "MessageCompleted", message: failedMessage }); + yield* offer(captures, { _tag: "RunCompleted", messages: [failedMessage] }); + yield* offer(captures, { _tag: "SessionInfoChanged", name: "failure barrier" }); + yield* awaitObservedType(subscription.observed, "thread.metadata.updated"); + expect(running.pollUnsafe()).toBeUndefined(); + + yield* TestClock.adjust(PRIME_AGENT_FAILED_RUN_SETTLEMENT_GRACE_MS); + const result = yield* Fiber.join(running); + expect(result.turnId).toBe(started.turnId); + const turnEvents = subscription.events.filter((event) => event.turnId === result.turnId); + expect(turnEvents.find((event) => event.type === "runtime.error")).toMatchObject({ + payload: { + message: "Prime Agent stopped before sending a final response.", + detail: { kind: "missing-final-response", outcome: "failed" }, + }, + }); + expect(turnEvents.find((event) => event.type === "turn.completed")).toMatchObject({ + payload: { + state: "failed", + errorMessage: "Prime Agent stopped before sending a final response.", + }, + }); + expect(encodeUnknownJson(turnEvents)).not.toContain("PRIVATE provider error"); + yield* Fiber.interrupt(subscription.fiber); + }), + ).pipe(Effect.provide(testLayer)), + ); + + it.effect("keeps input admitted during a failed-run handoff attached to the turn", () => + Effect.scoped( + Effect.gen(function* () { + const captures = makeCaptures(); + const adapter = yield* makePrimeAgentDaemonAdapter(decodeSettings({}), manager, { + instanceId, + runtimeFactory: fakeRuntimeFactory(captures), + }); + const subscription = yield* subscribe(adapter); + yield* adapter.startSession({ threadId, cwd: process.cwd(), runtimeMode: "full-access" }); + yield* awaitObservedType(subscription.observed, "thread.started"); + + const running = yield* adapter + .sendTurn({ threadId, input: "initial" }) + .pipe(Effect.forkChild); + const started = yield* awaitObservedType(subscription.observed, "turn.started"); + yield* Queue.take(captures.promptObserved!); + const failedMessage = { + ...assistantMessage("", "error"), + errorMessage: "PRIVATE transient failure", + } satisfies PrimeDaemonMessage; + yield* offer(captures, { _tag: "RunCompleted", messages: [failedMessage] }); + yield* offer(captures, { _tag: "SessionInfoChanged", name: "handoff input barrier" }); + yield* awaitObservedType(subscription.observed, "thread.metadata.updated"); + + const steered = yield* adapter.sendTurn({ threadId, input: "continue this turn" }); + expect(steered.turnId).toBe(started.turnId); + yield* TestClock.adjust(PRIME_AGENT_FAILED_RUN_SETTLEMENT_GRACE_MS + 1); + expect(running.pollUnsafe()).toBeUndefined(); + expect(subscription.events.filter((event) => event.type === "turn.completed")).toHaveLength( + 0, + ); + + yield* offer(captures, { _tag: "RunStarted" }); + yield* offer(captures, { + _tag: "QueueChanged", + queuedCount: 0, + steeringCount: 0, + followUpCount: 0, + }); + const finalMessage = assistantMessage("continued final response"); + yield* offer(captures, { _tag: "MessageStarted", message: finalMessage }); + yield* offer(captures, { _tag: "MessageCompleted", message: finalMessage }); + yield* offer(captures, { _tag: "RunCompleted", messages: [finalMessage] }); + const result = yield* Fiber.join(running); + expect(result.turnId).toBe(started.turnId); + expect( + subscription.events.filter( + (event) => event.type === "turn.completed" && event.turnId === result.turnId, + ), + ).toHaveLength(1); + expect(encodeUnknownJson(subscription.events)).not.toContain("PRIVATE transient failure"); + yield* Fiber.interrupt(subscription.fiber); + }), + ).pipe(Effect.provide(testLayer)), + ); + it.effect("steers an active daemon run without opening or settling another turn", () => Effect.scoped( Effect.gen(function* () { @@ -5191,12 +5394,24 @@ describe("PrimeAgentDaemonAdapter", () => { .pipe(Effect.forkChild); yield* awaitObservedType(subscription.observed, "turn.started"); yield* Queue.take(captures.promptObserved!); + const preToolMessage = assistantMessage("partial response before bash"); + yield* offer(captures, { _tag: "MessageStarted", message: preToolMessage }); + yield* offer(captures, { _tag: "MessageCompleted", message: preToolMessage }); + yield* offer(captures, { + _tag: "BashStarted", + command: "PRIVATE command", + excludeFromContext: false, + transient: false, + }); yield* offer(captures, { _tag: "ExtensionRequest", request: { id: "native-close-secret", method: "confirm", title: "Still pending?" }, }); const requested = yield* awaitObservedType(subscription.observed, "interaction.requested"); - yield* offer(captures, { _tag: "SessionClosed", error: "daemon closed" }); + yield* offer(captures, { + _tag: "SessionClosed", + error: "PRIVATE daemon closed /tmp/native.sock", + }); const resolved = yield* awaitObservedType(subscription.observed, "interaction.resolved"); yield* awaitObservedType(subscription.observed, "turn.completed"); yield* awaitObservedType(subscription.observed, "session.exited"); @@ -5224,6 +5439,24 @@ describe("PrimeAgentDaemonAdapter", () => { expect(subscription.events.filter((event) => event.type === "session.exited")).toHaveLength( 1, ); + expect( + subscription.events.find( + (event) => + event.type === "runtime.error" && + event.turnId === result.turnId && + typeof event.payload.detail === "object" && + event.payload.detail !== null && + "kind" in event.payload.detail && + event.payload.detail.kind === "missing-final-response", + ), + ).toBeDefined(); + expect(subscription.events.find((event) => event.type === "session.exited")).toMatchObject({ + payload: { + exitKind: "error", + reason: "Prime Agent session closed unexpectedly.", + }, + }); + expect(encodeUnknownJson(subscription.events)).not.toContain("PRIVATE"); const stoppedAgain = yield* adapter.stopSession(threadId).pipe(Effect.result); expect(stoppedAgain._tag).toBe("Failure"); yield* Fiber.interrupt(subscription.fiber); diff --git a/apps/server/src/provider/prime/PrimeAgentDaemonAdapter.ts b/apps/server/src/provider/prime/PrimeAgentDaemonAdapter.ts index 3a8b45465..974e4c6d9 100644 --- a/apps/server/src/provider/prime/PrimeAgentDaemonAdapter.ts +++ b/apps/server/src/provider/prime/PrimeAgentDaemonAdapter.ts @@ -17,6 +17,7 @@ import { type ProviderSession, ProviderDriverKind, ProviderInstanceId, + RuntimeItemId, RuntimeRequestId, SessionInteractionRequest, SessionInteractionRequestId, @@ -82,6 +83,11 @@ import { mapPrimeAgentContextUsageDraft, mapPrimeAgentDaemonRuntimeEventDrafts, } from "./PrimeAgentDaemonRuntimeEvents.ts"; +import { + PRIME_AGENT_STOPPED_WITHOUT_FINAL_RESPONSE, + PRIME_AGENT_TURN_FAILED, + primeAgentMissingFinalResponseDetail, +} from "./PrimeAgentTerminalResponse.ts"; import { PRIME_AGENT_PERMISSION_EXTENSION_FILENAME, PRIME_AGENT_PERMISSION_EXTENSION_MARKER_COMMAND, @@ -109,6 +115,7 @@ import { const PROVIDER = ProviderDriverKind.make("primeAgent"); const SESSION_STATS_TIMEOUT_MS = 1_000; const MODEL_DISCOVERY_TIMEOUT_MS = 15_000; +export const PRIME_AGENT_FAILED_RUN_SETTLEMENT_GRACE_MS = 3_000; export const PRIME_AGENT_SIDE_QUESTION_TIMEOUT_MS = 2 * 60_000; const PRIME_AGENT_SIDE_QUESTION_MAX_ACTIVE = 4; const unavailableSessionGoal: SessionGoalUpdatedPayload = { @@ -144,6 +151,16 @@ interface PrimeAgentDaemonActiveTurn { readonly completed: Deferred.Deferred; cancellationRequested: boolean; assistantTextStreamed: boolean; + nextAssistantMessageSequence: number; + activeAssistantItemId: RuntimeItemId | undefined; + lastAssistantHadRenderableText: boolean; + runCompletionHandoffSequence: number; + pendingRunCompletionHandoff: + | { + readonly sequence: number; + readonly event: Extract; + } + | undefined; queuedInputCount: number; awaitingQueuedRun: boolean; queuedActionObserved: boolean; @@ -400,6 +417,19 @@ type TurnOutcome = | { readonly state: "failed"; readonly errorMessage: string } | { readonly state: "cancelled" }; +type PrimeAgentRunCompletedEvent = Extract; + +function primeAgentRunCompletedNeedsHandoff(event: PrimeAgentRunCompletedEvent): boolean { + const lastAssistant = event.messages.findLast((message) => message.role === "assistant"); + if (lastAssistant?.stopReason === "aborted") return false; + return ( + lastAssistant?.stopReason === "error" || + (lastAssistant?.errorMessage?.trim().length ?? 0) > 0 || + lastAssistant?.stopReason === "toolUse" || + (lastAssistant?.toolCalls.length ?? 0) > 0 + ); +} + function runtimeOperationError( threadId: ThreadId, method: string, @@ -853,25 +883,71 @@ export function makePrimeAgentDaemonAdapter( turn: PrimeAgentDaemonActiveTurn | undefined, ) => Effect.gen(function* () { + let allocatedAssistantItemLazily = false; + const allocateAssistantItemId = () => { + if (turn === undefined) return undefined; + const itemId = RuntimeItemId.make( + `assistant:${turn.id}:segment:${turn.nextAssistantMessageSequence}`, + ); + turn.nextAssistantMessageSequence += 1; + turn.activeAssistantItemId = itemId; + return itemId; + }; if ( turn !== undefined && event._tag === "MessageStarted" && event.message.role === "assistant" ) { + // Prime emits one assistant message per model/tool loop. A turn-scoped + // item id strands later final text at the first message's timestamp. + // Use an opaque subscriber-local sequence instead of native identity. + allocateAssistantItemId(); turn.assistantTextStreamed = false; + } else if ( + turn !== undefined && + turn.activeAssistantItemId === undefined && + ((event._tag === "AssistantStream" && event.kind === "text") || + (event._tag === "MessageCompleted" && event.message.role === "assistant")) + ) { + // The public stream is ordered, but reconnect recovery can omit a + // start. Allocate lazily rather than merging into an earlier item. + allocateAssistantItemId(); + turn.assistantTextStreamed = false; + allocatedAssistantItemLazily = true; } const compactionScope = event._tag === "CompactionStarted" || event._tag === "CompactionCompleted" ? context.activeCompactionScope : undefined; const runtimeTurnId = compactionScope === undefined ? turn?.id : compactionScope.turnId; + if ( + allocatedAssistantItemLazily && + runtimeTurnId !== undefined && + turn?.activeAssistantItemId !== undefined + ) { + yield* offerRuntimeEvent({ + type: "item.started", + ...(yield* makeEventStamp()), + provider: PROVIDER, + providerInstanceId: boundInstanceId, + threadId: context.threadId, + turnId: runtimeTurnId, + itemId: turn.activeAssistantItemId, + payload: { itemType: "assistant_message", status: "inProgress" }, + }); + } const drafts = mapPrimeAgentDaemonRuntimeEventDrafts({ event, provider: PROVIDER, providerInstanceId: boundInstanceId, threadId: context.threadId, ...(runtimeTurnId === undefined ? {} : { turnId: runtimeTurnId }), - ...(turn === undefined ? {} : { assistantTextStreamed: turn.assistantTextStreamed }), + ...(turn === undefined + ? {} + : { + assistantItemId: turn.activeAssistantItemId, + assistantTextStreamed: turn.assistantTextStreamed, + }), }); for (const draft of drafts) { yield* offerRuntimeEvent({ ...draft, ...(yield* makeEventStamp()) }); @@ -881,10 +957,22 @@ export function makePrimeAgentDaemonAdapter( event._tag === "AssistantStream" && event.kind === "text" && event.phase === "delta" && - event.delta !== undefined + event.delta !== undefined && + event.delta.trim().length > 0 ) { turn.assistantTextStreamed = true; } + if ( + turn !== undefined && + event._tag === "MessageCompleted" && + event.message.role === "assistant" + ) { + turn.lastAssistantHadRenderableText = + event.message.text.trim().length > 0 && + event.message.stopReason !== "toolUse" && + event.message.toolCalls.length === 0; + turn.activeAssistantItemId = undefined; + } }); /** Must be called with the thread lock held. Reservations leave only in stream finalizers. */ @@ -1054,9 +1142,25 @@ export function makePrimeAgentDaemonAdapter( const effectiveOutcome: TurnOutcome = context.stopRequested || turn.cancellationRequested ? { state: "cancelled" } : outcome; + turn.pendingRunCompletionHandoff = undefined; + if (effectiveOutcome.state === "failed" && !turn.lastAssistantHadRenderableText) { + yield* offerRuntimeEvent({ + type: "runtime.error", + ...(yield* makeEventStamp()), + provider: PROVIDER, + providerInstanceId: boundInstanceId, + threadId: context.threadId, + turnId: turn.id, + payload: { + message: PRIME_AGENT_STOPPED_WITHOUT_FINAL_RESPONSE, + class: "provider_error", + detail: primeAgentMissingFinalResponseDetail("failed"), + }, + }); + } if (effectiveOutcome.state === "completed") { yield* publishDrafts(context, effectiveOutcome.event, turn); - context.turns.push({ id: turn.id, items: [effectiveOutcome.event] }); + context.turns.push({ id: turn.id, items: [] }); } else { yield* offerRuntimeEvent({ type: "turn.completed", @@ -1094,6 +1198,48 @@ export function makePrimeAgentDaemonAdapter( withThreadLock(context.threadId, settleActiveTurnLocked(context, turn, outcome)), ); + const promotePendingRunCompletionToQueuedRun = (turn: PrimeAgentDaemonActiveTurn) => { + const pending = turn.pendingRunCompletionHandoff; + if (pending === undefined) return; + turn.completedRunMessages.push(...pending.event.messages); + turn.pendingRunCompletionHandoff = undefined; + turn.awaitingQueuedRun = true; + turn.queuedActionObserved = false; + }; + + const schedulePendingRunCompletionHandoff = ( + context: PrimeAgentDaemonSessionContext, + turn: PrimeAgentDaemonActiveTurn, + sequence: number, + ) => + Effect.forkIn( + Effect.sleep(PRIME_AGENT_FAILED_RUN_SETTLEMENT_GRACE_MS).pipe( + Effect.andThen( + withThreadLock( + context.threadId, + Effect.gen(function* () { + const pending = turn.pendingRunCompletionHandoff; + if (pending === undefined || pending.sequence !== sequence) return; + turn.pendingRunCompletionHandoff = undefined; + const completionEvent = + turn.completedRunMessages.length === 0 + ? pending.event + : { + ...pending.event, + messages: [...turn.completedRunMessages, ...pending.event.messages], + }; + const settled = yield* settleActiveTurnLocked(context, turn, { + state: "completed", + event: completionEvent, + }); + if (settled) yield* refreshContextUsage(context).pipe(Effect.forkDetach); + }), + ), + ), + ), + context.scope, + ).pipe(Effect.asVoid); + /** Must be called with the thread lock held. */ const clearPendingInteractionsLocked = ( context: PrimeAgentDaemonSessionContext, @@ -1255,7 +1401,18 @@ export function makePrimeAgentDaemonAdapter( turn.queuedInputCount === 0 && !event.state.inputQueue.activeAction && !event.state.isStreaming; - if (turn.awaitingQueuedRun && authoritativeIdle) { + if (turn.pendingRunCompletionHandoff !== undefined) { + if (event.state.isStreaming) { + // The continuation may have started while disconnected, + // so the resync snapshot is an authoritative RunStarted. + turn.completedRunMessages.push( + ...turn.pendingRunCompletionHandoff.event.messages, + ); + turn.pendingRunCompletionHandoff = undefined; + } + // An idle snapshot can be the gap before RunStarted. Keep + // the bounded handoff alive until the event or timeout. + } else if (turn.awaitingQueuedRun && authoritativeIdle) { const explicitClear = context.inputQueueClearPending; context.inputQueueClearPending = false; if (explicitClear) { @@ -1610,6 +1767,17 @@ export function makePrimeAgentDaemonAdapter( turn.queuedActionObserved = false; return false; } + if (primeAgentRunCompletedNeedsHandoff(event)) { + turn.lastAssistantHadRenderableText = false; + // A daemon/kernel reconnect can emit a non-final agent_end + // and immediately continue the same public run. Keep the Pylon + // turn bound briefly so the following RunStarted is not orphaned. + turn.runCompletionHandoffSequence += 1; + const sequence = turn.runCompletionHandoffSequence; + turn.pendingRunCompletionHandoff = { sequence, event }; + yield* schedulePendingRunCompletionHandoff(context, turn, sequence); + return false; + } const completionEvent = turn.completedRunMessages.length === 0 ? event @@ -1750,10 +1918,23 @@ export function makePrimeAgentDaemonAdapter( let publishEvent = true; if (event._tag === "RunStarted") { context.nativeRunActive = true; - } else if (event._tag === "BashStarted") { + if (turn?.pendingRunCompletionHandoff !== undefined) { + turn.completedRunMessages.push(...turn.pendingRunCompletionHandoff.event.messages); + turn.pendingRunCompletionHandoff = undefined; + } + } else if ( + event._tag === "ToolStarted" || + event._tag === "ToolProgress" || + event._tag === "ToolCompleted" || + (event._tag === "AssistantStream" && event.kind === "toolCall") + ) { + if (turn !== undefined) turn.lastAssistantHadRenderableText = false; + } else if (event._tag === "BashStarted" || event._tag === "BashOutput") { context.nativeBashActive = true; + if (turn !== undefined) turn.lastAssistantHadRenderableText = false; } else if (event._tag === "BashCompleted") { context.nativeBashActive = false; + if (turn !== undefined) turn.lastAssistantHadRenderableText = false; } else if (event._tag === "GoalUpdated") { publishEvent = false; yield* updateGoalProjection( @@ -2586,6 +2767,7 @@ export function makePrimeAgentDaemonAdapter( ), ); activeTurn.queuedInputCount += 1; + promotePendingRunCompletionToQueuedRun(activeTurn); return { _tag: "Steered" as const, result: { @@ -2615,6 +2797,11 @@ export function makePrimeAgentDaemonAdapter( completed: yield* Deferred.make(), cancellationRequested: false, assistantTextStreamed: false, + nextAssistantMessageSequence: 0, + activeAssistantItemId: undefined, + lastAssistantHadRenderableText: false, + runCompletionHandoffSequence: 0, + pendingRunCompletionHandoff: undefined, queuedInputCount: 0, awaitingQueuedRun: false, queuedActionObserved: false, @@ -2675,18 +2862,20 @@ export function makePrimeAgentDaemonAdapter( }); return yield* restore(runPrompt).pipe( - Effect.catch((error) => + Effect.catch(() => Effect.gen(function* () { const cancelled = turn.cancellationRequested || turn.controller.signal.aborted; - const settled = yield* settleActiveTurn( + yield* settleActiveTurn( context, turn, cancelled ? { state: "cancelled" } - : { state: "failed", errorMessage: error.message }, + : { state: "failed", errorMessage: PRIME_AGENT_TURN_FAILED }, ); - if (cancelled || !settled) return result; - return yield* error; + // The prompt was admitted. Its runtime events already carry + // the authoritative failed terminal, so do not reclassify it + // as a second turn-start failure in orchestration. + return result; }), ), Effect.onInterrupt(() => @@ -3629,6 +3818,7 @@ export function makePrimeAgentDaemonAdapter( }; yield* updateInputQueueProjection(context, next); turn.queuedInputCount = Math.max(1, next.steeringCount + next.followUpCount); + promotePendingRunCompletionToQueuedRun(turn); return context.inputQueue; } if (Exit.isSuccess(reconciled)) { diff --git a/apps/server/src/provider/prime/PrimeAgentDaemonRuntimeEvents.test.ts b/apps/server/src/provider/prime/PrimeAgentDaemonRuntimeEvents.test.ts index 20402412f..ebead1fbd 100644 --- a/apps/server/src/provider/prime/PrimeAgentDaemonRuntimeEvents.test.ts +++ b/apps/server/src/provider/prime/PrimeAgentDaemonRuntimeEvents.test.ts @@ -109,6 +109,25 @@ describe("mapPrimeAgentDaemonRuntimeEventDrafts", () => { ]); }); + it("uses distinct caller-assigned ids for native assistant segments", () => { + const first = RuntimeItemId.make("assistant:turn-1:segment:0"); + const second = RuntimeItemId.make("assistant:turn-1:segment:1"); + const firstDraft = mapPrimeAgentDaemonRuntimeEventDrafts({ + ...context, + assistantItemId: first, + event: { _tag: "MessageStarted", message: assistant({ text: "first" }) }, + }); + const secondDraft = mapPrimeAgentDaemonRuntimeEventDrafts({ + ...context, + assistantItemId: second, + event: { _tag: "MessageStarted", message: assistant({ text: "final" }) }, + }); + + expect(firstDraft[0]).toMatchObject({ type: "item.started", itemId: first }); + expect(secondDraft[0]).toMatchObject({ type: "item.started", itemId: second }); + expect(firstDraft[0]?.itemId).not.toBe(secondDraft[0]?.itemId); + }); + it("backfills final assistant text only when streaming was explicitly absent", () => { const event = { _tag: "MessageCompleted", @@ -386,7 +405,7 @@ describe("mapPrimeAgentDaemonRuntimeEventDrafts", () => { }); }); - it("does not complete a replayed run and fails an active run with no assistant", () => { + it("keeps replayed runs detached and marks active textless success without inventing prose", () => { expect( mapPrimeAgentDaemonRuntimeEventDrafts({ provider, @@ -404,12 +423,22 @@ describe("mapPrimeAgentDaemonRuntimeEventDrafts", () => { }, ]); - expect( - mapPrimeAgentDaemonRuntimeEventDrafts({ - ...context, - event: { _tag: "RunCompleted", messages: [] }, - }), - ).toEqual([ + const textless = mapPrimeAgentDaemonRuntimeEventDrafts({ + ...context, + event: { _tag: "RunCompleted", messages: [] }, + }); + expect(textless).toEqual([ + { + provider, + providerInstanceId, + threadId, + turnId, + type: "runtime.warning", + payload: { + message: "Prime Agent finished without sending a final response.", + detail: { kind: "missing-final-response", outcome: "completed" }, + }, + }, { provider, providerInstanceId, @@ -417,8 +446,15 @@ describe("mapPrimeAgentDaemonRuntimeEventDrafts", () => { turnId, type: "turn.completed", payload: { - state: "failed", - errorMessage: "Prime Agent completed the run without an assistant message.", + state: "completed", + usage: { + inputTokens: 0, + outputTokens: 0, + cachedInputTokens: 0, + cacheWriteTokens: 0, + totalTokens: 0, + }, + totalCostUsd: 0, }, }, { @@ -430,6 +466,21 @@ describe("mapPrimeAgentDaemonRuntimeEventDrafts", () => { payload: { state: "ready" }, }, ]); + + const toolOnly = mapPrimeAgentDaemonRuntimeEventDrafts({ + ...context, + event: { + _tag: "RunCompleted", + messages: [assistant({ text: "before tool", stopReason: "toolUse" })], + }, + }); + expect(toolOnly[0]).toMatchObject({ + type: "runtime.warning", + payload: { + message: "Prime Agent finished without sending a final response.", + detail: { kind: "missing-final-response", outcome: "completed" }, + }, + }); }); it("does not turn aborted or error run endings into successful turns", () => { @@ -448,14 +499,30 @@ describe("mapPrimeAgentDaemonRuntimeEventDrafts", () => { _tag: "RunCompleted", messages: [ assistant({ stopReason: "toolUse" }), - assistant({ stopReason: "error", errorMessage: "quota exhausted" }), + assistant({ text: "", stopReason: "error", errorMessage: "quota exhausted" }), ], }, }); - expect(failed[0]).toMatchObject({ + expect(failed[0]).toEqual({ + provider, + providerInstanceId, + threadId, + turnId, + type: "runtime.error", + payload: { + message: "Prime Agent stopped before sending a final response.", + class: "provider_error", + detail: { kind: "missing-final-response", outcome: "failed" }, + }, + }); + expect(failed[1]).toMatchObject({ type: "turn.completed", - payload: { state: "failed", errorMessage: "quota exhausted" }, + payload: { + state: "failed", + errorMessage: "Prime Agent stopped before sending a final response.", + }, }); + expect(JSON.stringify(failed)).not.toContain("quota exhausted"); expect( mapPrimeAgentDaemonRuntimeEventDrafts({ @@ -652,6 +719,43 @@ describe("mapPrimeAgentDaemonRuntimeEventDrafts", () => { ).toEqual([]); }); + it("uses fixed lifecycle copy instead of native connection or close errors", () => { + const connection = mapPrimeAgentDaemonRuntimeEventDrafts({ + ...context, + event: { + _tag: "ConnectionStatus", + status: "reconnecting", + error: "PRIVATE socket path /tmp/native.sock", + }, + }); + const closed = mapPrimeAgentDaemonRuntimeEventDrafts({ + ...context, + event: { _tag: "SessionClosed", error: "PRIVATE daemon teardown cause" }, + }); + + expect(connection).toEqual([ + { + provider, + providerInstanceId, + threadId, + turnId, + type: "session.state.changed", + payload: { state: "starting", reason: "Prime Agent connection is unavailable." }, + }, + ]); + expect(closed).toEqual([ + { + provider, + providerInstanceId, + threadId, + turnId, + type: "session.exited", + payload: { exitKind: "error", reason: "Prime Agent session closed unexpectedly." }, + }, + ]); + expect(JSON.stringify([...connection, ...closed])).not.toContain("PRIVATE"); + }); + it("never includes native raw payloads and ignores replay, presentation, and duplicate streams", () => { const ignoredEvents: ReadonlyArray = [ { _tag: "TurnStarted" }, diff --git a/apps/server/src/provider/prime/PrimeAgentDaemonRuntimeEvents.ts b/apps/server/src/provider/prime/PrimeAgentDaemonRuntimeEvents.ts index 0112f9911..dc48920ca 100644 --- a/apps/server/src/provider/prime/PrimeAgentDaemonRuntimeEvents.ts +++ b/apps/server/src/provider/prime/PrimeAgentDaemonRuntimeEvents.ts @@ -15,6 +15,12 @@ import type { PrimeDaemonUsage, } from "./PrimeAgentDaemonEvents.ts"; import type { PrimeAgentDaemonSessionStats } from "./PrimeAgentDaemonSessionRuntime.ts"; +import { + PRIME_AGENT_FINISHED_WITHOUT_FINAL_RESPONSE, + PRIME_AGENT_STOPPED_WITHOUT_FINAL_RESPONSE, + PRIME_AGENT_TURN_FAILED, + primeAgentMissingFinalResponseDetail, +} from "./PrimeAgentTerminalResponse.ts"; type RuntimeEventDraft = Event extends ProviderRuntimeEvent ? Omit @@ -28,6 +34,7 @@ interface RuntimeEventContext { readonly providerInstanceId?: ProviderInstanceId | undefined; readonly threadId: ThreadId; readonly turnId?: TurnId | undefined; + readonly assistantItemId?: RuntimeItemId | undefined; } const MAX_TEXT_LENGTH = 100_000; @@ -105,8 +112,11 @@ export function mapPrimeAgentContextUsageDraft(input: { }; } -function assistantItemId(turnId: TurnId | undefined): RuntimeItemId | undefined { - return turnId === undefined ? undefined : RuntimeItemId.make(`assistant:${turnId}`); +function assistantItemId(input: RuntimeEventContext): RuntimeItemId | undefined { + return ( + input.assistantItemId ?? + (input.turnId === undefined ? undefined : RuntimeItemId.make(`assistant:${input.turnId}`)) + ); } function reasoningItemId(turnId: TurnId | undefined): RuntimeItemId | undefined { @@ -318,6 +328,7 @@ export function mapPrimeAgentDaemonRuntimeEventDrafts(input: { readonly providerInstanceId?: ProviderInstanceId | undefined; readonly threadId: ThreadId; readonly turnId?: TurnId | undefined; + readonly assistantItemId?: RuntimeItemId | undefined; readonly assistantTextStreamed?: boolean | undefined; }): ReadonlyArray { const context: RuntimeEventContext = input; @@ -337,45 +348,67 @@ export function mapPrimeAgentDaemonRuntimeEventDrafts(input: { const runMessages = assistantMessages(event.messages); const message = runMessages.at(-1); - if (message === undefined) { - return [ - { - ...base, - type: "turn.completed", - payload: { - state: "failed", - errorMessage: "Prime Agent completed the run without an assistant message.", - }, - }, - ready, - ]; - } - - const errorMessage = boundedNonEmpty(message.errorMessage); + const nativeError = boundedNonEmpty(message?.errorMessage); const state = - message.stopReason === "aborted" + message?.stopReason === "aborted" ? "cancelled" - : message.stopReason === "error" || errorMessage !== undefined + : message?.stopReason === "error" || nativeError !== undefined ? "failed" : "completed"; + const hasFinalResponse = + boundedNonEmpty(message?.text) !== undefined && + message?.stopReason !== "toolUse" && + message?.toolCalls.length === 0; + const missingFinalResponse = state !== "cancelled" && !hasFinalResponse; const usage = aggregateAssistantUsage(runMessages); const completed: PrimeAgentRuntimeEventDraft = { ...base, type: "turn.completed", payload: { state, - stopReason: message.stopReason, + ...(message === undefined ? {} : { stopReason: message.stopReason }), usage: turnUsage(usage), totalCostUsd: usage.totalCostUsd, - ...(errorMessage === undefined ? {} : { errorMessage }), + ...(state !== "failed" + ? {} + : { + errorMessage: missingFinalResponse + ? PRIME_AGENT_STOPPED_WITHOUT_FINAL_RESPONSE + : PRIME_AGENT_TURN_FAILED, + }), }, }; - return [completed, ready]; + const terminalNotice: ReadonlyArray = !missingFinalResponse + ? [] + : state === "failed" + ? [ + { + ...base, + type: "runtime.error", + payload: { + message: PRIME_AGENT_STOPPED_WITHOUT_FINAL_RESPONSE, + class: "provider_error", + detail: primeAgentMissingFinalResponseDetail("failed"), + }, + }, + ] + : [ + { + ...base, + type: "runtime.warning", + payload: { + message: PRIME_AGENT_FINISHED_WITHOUT_FINAL_RESPONSE, + detail: primeAgentMissingFinalResponseDetail("completed"), + }, + }, + ]; + return [...terminalNotice, completed, ready]; } + case "TurnStarted": return []; case "MessageStarted": { - const itemId = event.message.role === "assistant" ? assistantItemId(input.turnId) : undefined; + const itemId = event.message.role === "assistant" ? assistantItemId(context) : undefined; return itemId === undefined ? [] : [ @@ -388,10 +421,11 @@ export function mapPrimeAgentDaemonRuntimeEventDrafts(input: { ]; } case "MessageCompleted": { - const itemId = event.message.role === "assistant" ? assistantItemId(input.turnId) : undefined; + const itemId = event.message.role === "assistant" ? assistantItemId(context) : undefined; if (itemId === undefined || event.message.role !== "assistant") return []; - const errorMessage = boundedNonEmpty(event.message.errorMessage); - const failed = event.message.stopReason === "error" || errorMessage !== undefined; + const failed = + event.message.stopReason === "error" || + boundedNonEmpty(event.message.errorMessage) !== undefined; const completed: PrimeAgentRuntimeEventDraft = { ...base, type: "item.completed", @@ -399,7 +433,6 @@ export function mapPrimeAgentDaemonRuntimeEventDrafts(input: { payload: { itemType: "assistant_message", status: failed ? "failed" : "completed", - ...(errorMessage === undefined ? {} : { detail: errorMessage }), }, }; const finalText = @@ -419,7 +452,7 @@ export function mapPrimeAgentDaemonRuntimeEventDrafts(input: { case "AssistantStream": { if (event.kind === "toolCall" || event.phase === "done" || event.phase === "error") return []; if (event.kind === "text") { - const itemId = assistantItemId(input.turnId); + const itemId = assistantItemId(context); return event.phase !== "delta" || itemId === undefined || event.delta === undefined ? [] : [ @@ -600,27 +633,27 @@ export function mapPrimeAgentDaemonRuntimeEventDrafts(input: { }, ]; case "ConnectionStatus": { - const reason = boundedNonEmpty(event.error, MAX_SCALAR_LENGTH); + const hasError = boundedNonEmpty(event.error, MAX_SCALAR_LENGTH) !== undefined; return [ { ...base, type: "session.state.changed", payload: { state: event.status === "connected" ? "ready" : "starting", - ...(reason === undefined ? {} : { reason }), + ...(hasError ? { reason: "Prime Agent connection is unavailable." } : {}), }, }, ]; } case "SessionClosed": { - const reason = boundedNonEmpty(event.error); + const hasError = boundedNonEmpty(event.error) !== undefined; return [ { ...base, type: "session.exited", payload: { - exitKind: reason === undefined ? "graceful" : "error", - ...(reason === undefined ? {} : { reason }), + exitKind: hasError ? "error" : "graceful", + ...(hasError ? { reason: "Prime Agent session closed unexpectedly." } : {}), }, }, ]; diff --git a/apps/server/src/provider/prime/PrimeAgentTerminalResponse.ts b/apps/server/src/provider/prime/PrimeAgentTerminalResponse.ts new file mode 100644 index 000000000..4b7d5b63a --- /dev/null +++ b/apps/server/src/provider/prime/PrimeAgentTerminalResponse.ts @@ -0,0 +1,10 @@ +export const PRIME_AGENT_FINISHED_WITHOUT_FINAL_RESPONSE = + "Prime Agent finished without sending a final response."; +export const PRIME_AGENT_STOPPED_WITHOUT_FINAL_RESPONSE = + "Prime Agent stopped before sending a final response."; +export const PRIME_AGENT_TURN_FAILED = "Prime Agent turn failed."; + +export const primeAgentMissingFinalResponseDetail = (outcome: "completed" | "failed") => ({ + kind: "missing-final-response" as const, + outcome, +}); diff --git a/apps/web/src/components/chat/MessagesTimeline.logic.test.ts b/apps/web/src/components/chat/MessagesTimeline.logic.test.ts index 1905fc87a..f5ec7e83f 100644 --- a/apps/web/src/components/chat/MessagesTimeline.logic.test.ts +++ b/apps/web/src/components/chat/MessagesTimeline.logic.test.ts @@ -541,6 +541,60 @@ describe("deriveMessagesTimelineRows", () => { ).toBeDefined(); }); + it("keeps a missing-response work row visible outside settled-turn folding", () => { + const rows = deriveMessagesTimelineRows({ + timelineEntries: [ + { + id: "tool-entry", + kind: "work", + createdAt: "2026-01-01T00:00:05Z", + entry: { + id: "tool-entry", + createdAt: "2026-01-01T00:00:05Z", + turnId: "turn-1" as never, + label: "Read files", + tone: "tool", + sourceActivityKind: "tool.completed", + }, + }, + { + id: "missing-response-entry", + kind: "work", + createdAt: "2026-01-01T00:00:20Z", + entry: { + id: "missing-response-entry", + createdAt: "2026-01-01T00:00:20Z", + turnId: "turn-1" as never, + label: "Prime Agent finished without sending a final response.", + tone: "info", + sourceActivityKind: "turn.response.missing", + }, + }, + ], + latestTurn: { + turnId: "turn-1" as never, + state: "completed", + startedAt: "2026-01-01T00:00:00Z", + completedAt: "2026-01-01T00:00:20Z", + }, + isWorking: false, + activeTurnStartedAt: null, + turnDiffSummaryByAssistantMessageId: new Map(), + revertTurnCountByUserMessageId: new Map(), + }); + + expect(rows.map((row) => row.id)).toEqual(["turn-fold:turn-1", "missing-response-entry"]); + expect(rows[1]).toMatchObject({ + kind: "work", + groupedEntries: [ + { + sourceActivityKind: "turn.response.missing", + label: "Prime Agent finished without sending a final response.", + }, + ], + }); + }); + it("derives a sane duration for a steer-superseded turn with one instant commentary message", () => { // A steer ends the previous turn early: its only message completes the // instant it is created, and trailing work entries land after it. The diff --git a/apps/web/src/components/chat/MessagesTimeline.logic.ts b/apps/web/src/components/chat/MessagesTimeline.logic.ts index eab2f2a16..c320a0c7b 100644 --- a/apps/web/src/components/chat/MessagesTimeline.logic.ts +++ b/apps/web/src/components/chat/MessagesTimeline.logic.ts @@ -2,6 +2,7 @@ import * as Equal from "effect/Equal"; import { formatDuration, workEntryIndicatesToolNeutralStatus, + workLogEntryIsMissingResponse, workLogEntryIsToolLike, type TimelineEntry, type TurnPlanEntry, @@ -393,7 +394,10 @@ function deriveTurnFolds(input: { // Agent-spawn CTA rows never fold: workflows outlive their launching // turn (dynamic spawns, background execution), and folding the CTA // when the turn settles makes a still-running fleet invisible. - if (entry.kind === "work" && entry.entry.agentSpawn !== undefined) { + if ( + entry.kind === "work" && + (entry.entry.agentSpawn !== undefined || workLogEntryIsMissingResponse(entry.entry)) + ) { continue; } hiddenEntryIds.add(entry.id); @@ -508,8 +512,10 @@ export function deriveMessagesTimelineRows(input: { while (cursor < input.timelineEntries.length) { const nextEntry = input.timelineEntries[cursor]; if ( + workLogEntryIsMissingResponse(timelineEntry.entry) || !nextEntry || nextEntry.kind !== "work" || + workLogEntryIsMissingResponse(nextEntry.entry) || collapsedEntryIds.has(nextEntry.id) || foldsByAnchorEntryId.has(nextEntry.id) ) { diff --git a/apps/web/src/components/chat/MessagesTimeline.tsx b/apps/web/src/components/chat/MessagesTimeline.tsx index 1aba4ad30..4161f1ab3 100644 --- a/apps/web/src/components/chat/MessagesTimeline.tsx +++ b/apps/web/src/components/chat/MessagesTimeline.tsx @@ -36,6 +36,7 @@ import { workEntryIndicatesToolFailure, workEntryIndicatesToolNeutralStatus, workEntryIndicatesToolSuccess, + workLogEntryIsMissingResponse, workLogEntryIsToolLike, } from "../../session-logic"; import { type TurnDiffSummary } from "../../types"; @@ -1375,11 +1376,16 @@ const WorkGroupSection = memo(function WorkGroupSection({ [groupedEntries], ); const onlyToolEntries = nonEmptyEntries.every((entry) => workLogEntryIsToolLike(entry)); - const groupLabel = onlyToolEntries - ? nonEmptyEntries.length === 1 - ? "1 tool call" - : `${nonEmptyEntries.length} tool calls` - : "Work Log"; + const onlyResponseNotices = nonEmptyEntries.every((entry) => + workLogEntryIsMissingResponse(entry), + ); + const groupLabel = onlyResponseNotices + ? "Response status" + : onlyToolEntries + ? nonEmptyEntries.length === 1 + ? "1 tool call" + : `${nonEmptyEntries.length} tool calls` + : "Work Log"; if (nonEmptyEntries.length === 0) return null; @@ -2267,6 +2273,7 @@ const PlainWorkEntryRow = memo(function PlainWorkEntryRow(props: { const displayText = preview ? `${heading} - ${preview}` : heading; const expandedBody = buildToolCallExpandedBody(workEntry, workspaceRoot); const canExpand = expandedBody !== null; + const missingResponse = workLogEntryIsMissingResponse(workEntry); const showFailedIndicator = workEntryIndicatesToolFailure(workEntry); const showDestructiveRowStyle = showFailedIndicator && @@ -2353,13 +2360,13 @@ const PlainWorkEntryRow = memo(function PlainWorkEntryRow(props: { render={ } > - Failed + {missingResponse ? "No final response" : "Failed"} ) : showSuccessIndicator ? ( diff --git a/apps/web/src/session-logic.test.ts b/apps/web/src/session-logic.test.ts index 5670bebee..bd7708350 100644 --- a/apps/web/src/session-logic.test.ts +++ b/apps/web/src/session-logic.test.ts @@ -723,6 +723,57 @@ describe("workEntryIndicatesToolFailure", () => { }); describe("deriveWorkLogEntries", () => { + it.each([ + { + outcome: "completed", + message: "Prime Agent finished without sending a final response.", + tone: "info" as const, + }, + { + outcome: "failed", + message: "Prime Agent stopped before sending a final response.", + tone: "error" as const, + }, + ])("presents a provider-neutral $outcome missing-response row", (fixture) => { + const [entry] = deriveWorkLogEntries([ + makeActivity({ + id: `missing-response-${fixture.outcome}`, + kind: "turn.response.missing", + summary: fixture.message, + tone: fixture.tone, + payload: { outcome: fixture.outcome }, + turnId: "turn-1", + }), + ]); + + expect(entry).toMatchObject({ + id: `missing-response-${fixture.outcome}`, + label: fixture.message, + tone: fixture.tone, + sourceActivityKind: "turn.response.missing", + turnId: "turn-1", + }); + expect(entry?.detail).toBeUndefined(); + expect(entry && workLogEntryIsToolLike(entry)).toBe(false); + }); + + it("does not treat ordinary provider activity as a missing-response row", () => { + const [entry] = deriveWorkLogEntries([ + makeActivity({ + id: "ordinary-runtime-warning", + kind: "runtime.warning", + summary: "Reconnecting", + tone: "info", + payload: { message: "Reconnecting", detail: { willRetry: true } }, + }), + ]); + + expect(entry).toMatchObject({ + label: "Reconnecting", + sourceActivityKind: "runtime.warning", + }); + }); + it("omits tool started entries and keeps completed entries", () => { const activities: OrchestrationThreadActivity[] = [ makeActivity({ diff --git a/apps/web/src/session-logic.ts b/apps/web/src/session-logic.ts index 0f1308e12..14fcc7715 100644 --- a/apps/web/src/session-logic.ts +++ b/apps/web/src/session-logic.ts @@ -174,7 +174,12 @@ export type TimelineEntry = entry: WorkLogEntry; }; +export function workLogEntryIsMissingResponse(entry: WorkLogEntry): boolean { + return entry.sourceActivityKind === "turn.response.missing"; +} + export function workLogEntryIsToolLike(entry: WorkLogEntry): boolean { + if (workLogEntryIsMissingResponse(entry)) return false; if (entry.tone === "tool" || entry.tone === "thinking" || entry.tone === "error") { return true; } diff --git a/docs/internals/prime-agent-daemon-parity.md b/docs/internals/prime-agent-daemon-parity.md index 5eb44508f..209adcf52 100644 --- a/docs/internals/prime-agent-daemon-parity.md +++ b/docs/internals/prime-agent-daemon-parity.md @@ -12,6 +12,15 @@ Daemon mode is currently disabled on Windows. Prime Agent 0.7.2's named-pipe tra Prime event queues are bounded to 256 entries and preserve FIFO delivery with backpressure during normal operation; no sliding or dropping queue is used. Pre-snapshot admission reserves capacity cumulatively and fails closed on overflow. Adapter teardown allows one second for a stalled consumer to drain, then logs a structured forced-shutdown outcome and every rejected terminal event's safe type/thread identity before closing. This explicit exceptional path prevents scope deadlock; it is not silent loss. +Daemon assistant messages receive opaque subscriber-local segment identifiers so separate model/tool +cycles within one Pylon turn retain their chronological positions without persisting Prime's native +identifiers. An error- or tool-terminated native run completion without a final public response is held for a bounded three-second handoff because Prime +may reconnect and automatically continue that same turn; a following native run remains attached to +the original Pylon turn, while an exhausted handoff settles as a real failure. Daemon and ACP modes +both emit a fixed provider-neutral status before an authoritative terminal event when no public final +assistant text follows the latest tool boundary. Reasoning, tool data, native errors, and identifiers +are never used to synthesize assistant prose. + | Public API outcome | Pylon status | Decision | | -------------------------------------------------------------------------------------------------------------------------------------------------------------------------- | --------------------------------------------------------- | ------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------- | | `attach`, root `subscribe`, `getState`, `getInitialSnapshot`, `dispose` | Integrated internally | Own one private daemon session per active Pylon thread, resume exact verified identity, reconnect through public snapshots/events, and close with the thread scope. The root connection's `getMessages` transcript API is never called; Pylon's event-sourced transcript remains authoritative. These are lifecycle primitives, not client RPCs. | diff --git a/docs/user/providers-prime-agent.md b/docs/user/providers-prime-agent.md index d1090d381..17cbfcf82 100644 --- a/docs/user/providers-prime-agent.md +++ b/docs/user/providers-prime-agent.md @@ -43,6 +43,15 @@ ACP compatibility mode instead, because the daemon API cannot safely preserve ar arguments. Pylon shows that fallback in the provider status rather than silently discarding the arguments. +## Turn Completion + +A Prime turn can contain several assistant segments around tool work. Pylon keeps those segments in +native order, so a final response appears after the work that preceded it instead of being appended +to an older message higher in the thread. If Prime authoritatively finishes without public assistant +text after its latest tool activity, Pylon shows **Prime Agent finished without sending a final +response.** as a status row rather than inventing an assistant reply. Failures use the corresponding +stopped status; cancellation remains cancellation. + On Windows, Pylon currently uses ACP compatibility mode because Prime Agent 0.7.2's public named-pipe daemon transport does not expose a verifiable per-user ACL or authenticated handshake. Native daemon mode remains fail-closed there until the transport can prevent another local OS user from impersonating or connecting to the daemon. Pylon uses a short-lived, prompt-free Prime Agent RPC process to bootstrap the configured-model