From 46f08ee0c12c9bff805a2e80913fcce184e98b3a Mon Sep 17 00:00:00 2001 From: Mike Olson Date: Mon, 20 Jul 2026 12:58:54 -0400 Subject: [PATCH 1/2] test(orchestrator): Align post-merge CTM fixtures --- apps/server/src/cli/config.test.ts | 4 ++-- .../tool_call_read_only_on_request/claude_transcript.ndjson | 2 +- .../tool_call_restricted_granular/claude_transcript.ndjson | 2 +- 3 files changed, 4 insertions(+), 4 deletions(-) diff --git a/apps/server/src/cli/config.test.ts b/apps/server/src/cli/config.test.ts index 349dd42a61d..754bdd8a16a 100644 --- a/apps/server/src/cli/config.test.ts +++ b/apps/server/src/cli/config.test.ts @@ -125,7 +125,7 @@ it.layer(NodeServices.layer)("cli config resolution", (it) => { tailscaleServeEnabled: false, tailscaleServePort: 443, }); - assert.equal(resolved.stateDir, join(baseDir, "userdata")); + assert.equal(resolved.stateDir, join(baseDir, "userdata-v2")); }), ); @@ -195,7 +195,7 @@ it.layer(NodeServices.layer)("cli config resolution", (it) => { tailscaleServeEnabled: true, tailscaleServePort: 8443, }); - assert.equal(resolved.dbPath, join(baseDir, "userdata", "state.sqlite")); + assert.equal(resolved.dbPath, join(baseDir, "userdata-v2", "state.sqlite")); }), ); diff --git a/apps/server/src/orchestration-v2/testkit/fixtures/tool_call_read_only_on_request/claude_transcript.ndjson b/apps/server/src/orchestration-v2/testkit/fixtures/tool_call_read_only_on_request/claude_transcript.ndjson index e9b0b981eec..033f08ff9ba 100644 --- a/apps/server/src/orchestration-v2/testkit/fixtures/tool_call_read_only_on_request/claude_transcript.ndjson +++ b/apps/server/src/orchestration-v2/testkit/fixtures/tool_call_read_only_on_request/claude_transcript.ndjson @@ -1,5 +1,5 @@ {"type":"transcript_start","provider":"claudeAgent","protocol":"claude-agent-sdk.query","version":"0.2.111","scenario":"tool_call_read_only_on_request","metadata":{"prompts":["Create or overwrite .codex-probe-write-action.txt with exactly this text: codex app-server approval fixture. Use a local shell command or file edit only, then briefly report what happened. Do not read package metadata, use GitHub, use web, or use MCP."],"model":"claude-sonnet-4-6","nativeSessionId":"2b73817d-f7ab-41a7-a533-31b13b698b0d","queryMode":"streaming","tools":"claude_code","permissionMode":"default","enablePermissionCallback":true,"generatedBy":"recordClaudeAgentSdkReplayTranscript"}} -{"type":"expect_outbound","label":"query.open","frame":{"type":"query.open","options":{"model":"claude-sonnet-4-6","tools":{"type":"preset","preset":"claude_code"},"permissionMode":"default","sessionId":"2b73817d-f7ab-41a7-a533-31b13b698b0d"}}} +{"type":"expect_outbound","label":"query.open","frame":{"type":"query.open","options":{"model":"claude-sonnet-4-6","tools":["Read","Glob","Grep"],"permissionMode":"default","allowedTools":["Read","Glob","Grep"],"sessionId":"2b73817d-f7ab-41a7-a533-31b13b698b0d"}}} {"type":"expect_outbound","label":"prompt.offer:1","frame":{"type":"prompt.offer","message":{"type":"user","message":{"role":"user","content":"Create or overwrite .codex-probe-write-action.txt with exactly this text: codex app-server approval fixture. Use a local shell command or file edit only, then briefly report what happened. Do not read package metadata, use GitHub, use web, or use MCP."},"parent_tool_use_id":null}}} {"type":"emit_inbound","label":"system","frame":{"type":"system","subtype":"hook_started","hook_id":"ae3f0d68-becd-4b08-871b-f857a9ca1f90","hook_name":"SessionStart:startup","hook_event":"SessionStart","uuid":"1d299042-40fc-44af-8459-0a80baa7b7de","session_id":"2b73817d-f7ab-41a7-a533-31b13b698b0d"}} {"type":"emit_inbound","label":"system","frame":{"type":"system","subtype":"hook_response","hook_id":"ae3f0d68-becd-4b08-871b-f857a9ca1f90","hook_name":"SessionStart:startup","hook_event":"SessionStart","output":"","stdout":"","stderr":"","exit_code":0,"outcome":"success","uuid":"5c410b64-6b03-49c4-9991-a8e624d0fe12","session_id":"2b73817d-f7ab-41a7-a533-31b13b698b0d"}} diff --git a/apps/server/src/orchestration-v2/testkit/fixtures/tool_call_restricted_granular/claude_transcript.ndjson b/apps/server/src/orchestration-v2/testkit/fixtures/tool_call_restricted_granular/claude_transcript.ndjson index 348a77c2bbd..abd246cea4e 100644 --- a/apps/server/src/orchestration-v2/testkit/fixtures/tool_call_restricted_granular/claude_transcript.ndjson +++ b/apps/server/src/orchestration-v2/testkit/fixtures/tool_call_restricted_granular/claude_transcript.ndjson @@ -1,5 +1,5 @@ {"type":"transcript_start","provider":"claudeAgent","protocol":"claude-agent-sdk.query","version":"0.2.111","scenario":"tool_call_restricted_granular","metadata":{"prompts":["Create or overwrite .codex-probe-write-action.txt with exactly this text: codex app-server approval fixture. Use a local shell command or file edit only, then briefly report what happened. Do not read package metadata, use GitHub, use web, or use MCP."],"model":"claude-sonnet-4-6","nativeSessionId":"ffd3df75-f3cf-433d-bb67-1dd1e9db316b","queryMode":"streaming","tools":"claude_code","permissionMode":"default","enablePermissionCallback":true,"generatedBy":"recordClaudeAgentSdkReplayTranscript"}} -{"type":"expect_outbound","label":"query.open","frame":{"type":"query.open","options":{"model":"claude-sonnet-4-6","tools":{"type":"preset","preset":"claude_code"},"permissionMode":"default","sessionId":"ffd3df75-f3cf-433d-bb67-1dd1e9db316b"}}} +{"type":"expect_outbound","label":"query.open","frame":{"type":"query.open","options":{"model":"claude-sonnet-4-6","tools":["Read","Glob","Grep"],"permissionMode":"default","sessionId":"ffd3df75-f3cf-433d-bb67-1dd1e9db316b"}}} {"type":"expect_outbound","label":"prompt.offer:1","frame":{"type":"prompt.offer","message":{"type":"user","message":{"role":"user","content":"Create or overwrite .codex-probe-write-action.txt with exactly this text: codex app-server approval fixture. Use a local shell command or file edit only, then briefly report what happened. Do not read package metadata, use GitHub, use web, or use MCP."},"parent_tool_use_id":null}}} {"type":"emit_inbound","label":"system","frame":{"type":"system","subtype":"hook_started","hook_id":"ba333d12-6a9a-4a5e-9cb5-2d6ce18cd0f8","hook_name":"SessionStart:startup","hook_event":"SessionStart","uuid":"544c0cee-2195-49a9-ab30-088c3d3d5e83","session_id":"ffd3df75-f3cf-433d-bb67-1dd1e9db316b"}} {"type":"emit_inbound","label":"system","frame":{"type":"system","subtype":"hook_response","hook_id":"ba333d12-6a9a-4a5e-9cb5-2d6ce18cd0f8","hook_name":"SessionStart:startup","hook_event":"SessionStart","output":"","stdout":"","stderr":"","exit_code":0,"outcome":"success","uuid":"09419836-e234-4231-a15d-927acb1fbdb1","session_id":"ffd3df75-f3cf-433d-bb67-1dd1e9db316b"}} From 5e0f0d50bb81c5876d26f1108e08d1f90c1dc2dc Mon Sep 17 00:00:00 2001 From: Mike Olson Date: Tue, 21 Jul 2026 01:10:10 -0400 Subject: [PATCH 2/2] fix(orchestrator): Preserve post-interrupt recovery state --- .../Adapters/ClaudeAdapterV2.test.ts | 163 ++- .../Adapters/ClaudeAdapterV2.ts | 29 +- .../Adapters/CodexAdapterV2.test.ts | 339 ++++++- .../Adapters/CodexAdapterV2.ts | 112 ++- .../ProviderTurnStartService.ts | 9 +- .../RunExecutionService.test.ts | 931 +++++++++++++++++- .../orchestration-v2/RunExecutionService.ts | 345 +++++-- .../turn_interrupt_mid_tool/codex_output.ts | 16 +- 8 files changed, 1866 insertions(+), 78 deletions(-) diff --git a/apps/server/src/orchestration-v2/Adapters/ClaudeAdapterV2.test.ts b/apps/server/src/orchestration-v2/Adapters/ClaudeAdapterV2.test.ts index f27daa9ca9d..8abe1cee74b 100644 --- a/apps/server/src/orchestration-v2/Adapters/ClaudeAdapterV2.test.ts +++ b/apps/server/src/orchestration-v2/Adapters/ClaudeAdapterV2.test.ts @@ -1113,14 +1113,19 @@ describe("ClaudeAdapterV2 background wake turns", () => { uuid: "00000000-0000-4000-8000-000000000101", session_id: WAKE_NATIVE_SESSION, }); - const makeResultFrame = (input: { readonly uuid: string; readonly result: string }) => + const makeResultFrame = (input: { + readonly uuid: string; + readonly result: string; + readonly numTurns?: number; + readonly origin?: { readonly kind: "task-notification" }; + }) => claudeSdkFrame({ type: "result", subtype: "success", duration_ms: 10, duration_api_ms: 10, is_error: false, - num_turns: 1, + num_turns: input.numTurns ?? 1, result: input.result, stop_reason: "end_turn", total_cost_usd: 0, @@ -1134,6 +1139,7 @@ describe("ClaudeAdapterV2 background wake turns", () => { permission_denials: [], uuid: input.uuid, session_id: WAKE_NATIVE_SESSION, + ...(input.origin === undefined ? {} : { origin: input.origin }), }); const turnOneResult = makeResultFrame({ uuid: "00000000-0000-4000-8000-000000000102", @@ -1152,6 +1158,15 @@ describe("ClaudeAdapterV2 background wake turns", () => { const wakeResult = makeResultFrame({ uuid: "00000000-0000-4000-8000-000000000104", result: WAKE_RESULT_TEXT, + origin: { kind: "task-notification" }, + }); + const STALE_TASK_NOTIFICATION_RESULT_TEXT = + "Stale task-notification origin text that must not appear."; + const staleTaskNotificationResult = makeResultFrame({ + uuid: "00000000-0000-4000-8000-000000000106", + result: STALE_TASK_NOTIFICATION_RESULT_TEXT, + numTurns: 0, + origin: { kind: "task-notification" }, }); const awaitUntil = (predicate: () => boolean, label: string): Effect.Effect => @@ -1432,6 +1447,150 @@ describe("ClaudeAdapterV2 background wake turns", () => { ), ); + it.effect("ignores a live task-notification origin result during a normal user turn", () => + Effect.scoped( + Effect.gen(function* () { + const harness = yield* makeWakeHarness; + const now = yield* DateTime.now; + const probeAssistantText = "Probe after stale task-notification result."; + const recoveryAssistantText = "Recovered after the interrupt; continuing."; + const staleResultText = STALE_TASK_NOTIFICATION_RESULT_TEXT; + const hasMessageText = (text: string) => + harness.events.some( + (event) => event.type === "message.updated" && event.message.text === text, + ); + + yield* harness.runtime.startTurn( + makeClaudeTestTurnInput({ + threadId: harness.threadId, + providerThread: harness.providerThread, + now, + attemptId: RunAttemptId.make("attempt-claude-stale-notif-1"), + text: "Continue after interrupt.", + attachments: [], + }), + ); + yield* awaitUntil(() => harness.offeredMessages.length === 1, "recovery prompt offered"); + + // Live interleaving seen after interrupt recovery: a stale stopped + // task_notification and its task-notification-origin result arrive + // before the real root assistant stream. + yield* Queue.offer( + harness.sdkMessages, + claudeSdkFrame({ + type: "system", + subtype: "task_notification", + task_id: "task-stale-stopped", + status: "stopped", + output_file: "/tmp/task-stale-stopped.log", + summary: "", + uuid: "00000000-0000-4000-8000-000000000107", + session_id: WAKE_NATIVE_SESSION, + }), + ); + yield* Queue.offer(harness.sdkMessages, staleTaskNotificationResult); + // Queue-ordered probe: once this assistant text is emitted, the stale + // origin result ahead of it has been consumed. + yield* Queue.offer( + harness.sdkMessages, + claudeSdkFrame({ + type: "assistant", + message: { + role: "assistant", + content: [{ type: "text", text: probeAssistantText }], + }, + parent_tool_use_id: null, + uuid: "00000000-0000-4000-8000-00000000010a", + session_id: WAKE_NATIVE_SESSION, + }), + ); + + yield* awaitUntil( + () => hasMessageText(probeAssistantText), + "probe assistant after stale task-notification result", + ); + assert.lengthOf(harness.terminalEvents(), 0); + assert.isFalse(hasMessageText(staleResultText)); + + yield* Queue.offer( + harness.sdkMessages, + claudeSdkFrame({ + type: "assistant", + message: { + role: "assistant", + content: [{ type: "text", text: recoveryAssistantText }], + }, + parent_tool_use_id: null, + uuid: "00000000-0000-4000-8000-000000000108", + session_id: WAKE_NATIVE_SESSION, + }), + ); + yield* Queue.offer( + harness.sdkMessages, + makeResultFrame({ + uuid: "00000000-0000-4000-8000-000000000109", + result: recoveryAssistantText, + }), + ); + + yield* awaitUntil(() => harness.terminalEvents().length === 1, "user turn terminal"); + assert.equal(harness.terminalEvents()[0]?.status, "completed"); + assert.isTrue(hasMessageText(recoveryAssistantText)); + assert.isFalse(hasMessageText(staleResultText)); + }).pipe(Effect.provide(Layer.merge(idAllocatorLayer, NodeServices.layer))), + ), + ); + + it.effect("terminalizes a continuation turn from a task-notification origin wake result", () => + Effect.scoped( + Effect.gen(function* () { + const harness = yield* makeWakeHarness; + const now = yield* DateTime.now; + + yield* harness.runtime.startTurn( + makeClaudeTestTurnInput({ + threadId: harness.threadId, + providerThread: harness.providerThread, + now, + attemptId: RunAttemptId.make("attempt-claude-notif-origin-2a"), + text: "Run the build in the background.", + attachments: [], + }), + ); + yield* Queue.offer(harness.sdkMessages, wakeTaskStarted); + yield* Queue.offer(harness.sdkMessages, turnOneResult); + yield* awaitUntil(() => harness.terminalEvents().length === 1, "first turn terminal"); + yield* Queue.offer(harness.sdkMessages, wakeNotification); + yield* Queue.offer(harness.sdkMessages, wakeResult); + yield* awaitUntil(() => harness.continuationRequests.length === 1, "continuation request"); + + yield* harness.runtime.startTurn( + makeClaudeTestTurnInput({ + threadId: harness.threadId, + providerThread: harness.providerThread, + now, + attemptId: RunAttemptId.make("attempt-claude-notif-origin-2b"), + text: "Background task completed.", + attachments: [], + providerTurnOrdinal: 2, + messageCreatedBy: "agent", + messageCreationSource: "provider", + }), + ); + + yield* awaitUntil(() => harness.terminalEvents().length === 2, "continuation terminal"); + assert.equal(harness.terminalEvents()[1]?.status, "completed"); + assert.lengthOf(harness.offeredMessages, 1); + assert.isTrue( + harness.events.some( + (event) => event.type === "message.updated" && event.message.text === WAKE_RESULT_TEXT, + ), + ); + assert.isFalse(yield* harness.hasPendingBackgroundWork); + }).pipe(Effect.provide(Layer.merge(idAllocatorLayer, NodeServices.layer))), + ), + ); + it.effect("clears the pending task when the wake notification carries no summary", () => Effect.scoped( Effect.gen(function* () { diff --git a/apps/server/src/orchestration-v2/Adapters/ClaudeAdapterV2.ts b/apps/server/src/orchestration-v2/Adapters/ClaudeAdapterV2.ts index d4715e682b2..105afef66db 100644 --- a/apps/server/src/orchestration-v2/Adapters/ClaudeAdapterV2.ts +++ b/apps/server/src/orchestration-v2/Adapters/ClaudeAdapterV2.ts @@ -1755,6 +1755,19 @@ function isClaudeActiveSteeringAbortResult(message: SDKResultMessage): boolean { return message.terminal_reason === "aborted_streaming"; } +function isClaudeProviderContinuationTurn(input: ProviderAdapterV2TurnInput): boolean { + return input.message.createdBy === "agent" && input.message.creationSource === "provider"; +} + +function isClaudeTaskNotificationOriginResult(message: SDKMessage): message is SDKResultMessage & { + readonly origin: Extract< + NonNullable, + { readonly kind: "task-notification" } + >; +} { + return message.type === "result" && message.origin?.kind === "task-notification"; +} + function providerFailureFromResult( message: SDKResultMessage, ): OrchestrationV2ProviderFailure | null { @@ -3256,6 +3269,18 @@ export function makeClaudeAdapterV2( return; } + // Task-notification-origin results can interleave during a normal + // user turn (for example a stale background stop after interrupt + // recovery). They must not finalize that turn or supply fallback + // assistant text. Provider continuation turns still consume them + // when draining buffered wake messages. + if ( + isClaudeTaskNotificationOriginResult(message) && + !isClaudeProviderContinuationTurn(context.input) + ) { + return; + } + // An is_error result's text is the error message; it belongs on the // terminal-failure item, not on a synthetic assistant message. const resultText = @@ -3539,9 +3564,7 @@ export function makeClaudeAdapterV2( // produced instead of prompting it again: drain the buffered wake // messages into this turn and let any still-streaming messages // follow live. The continuation prompt text never reaches the CLI. - const isContinuationTurn = - turnInput.message.createdBy === "agent" && - turnInput.message.creationSource === "provider"; + const isContinuationTurn = isClaudeProviderContinuationTurn(turnInput); const userMessage = isContinuationTurn ? null : yield* makeClaudeUserMessageWithAttachments({ diff --git a/apps/server/src/orchestration-v2/Adapters/CodexAdapterV2.test.ts b/apps/server/src/orchestration-v2/Adapters/CodexAdapterV2.test.ts index 6ff56268de0..e1d33093251 100644 --- a/apps/server/src/orchestration-v2/Adapters/CodexAdapterV2.test.ts +++ b/apps/server/src/orchestration-v2/Adapters/CodexAdapterV2.test.ts @@ -675,8 +675,10 @@ function makeCodexTestTurnInput(input: { function makeCodexReplayTurn(input: { readonly id: string; - readonly status: "inProgress" | "completed"; + readonly status: "inProgress" | "completed" | "interrupted" | "failed"; }): Record { + const terminal = + input.status === "completed" || input.status === "interrupted" || input.status === "failed"; return { id: input.id, items: [], @@ -684,7 +686,7 @@ function makeCodexReplayTurn(input: { status: input.status, error: null, startedAt: 1782622440, - completedAt: input.status === "completed" ? 1782622450 : null, + completedAt: terminal ? 1782622450 : null, durationMs: null, }; } @@ -1168,6 +1170,339 @@ describe("CodexAdapterV2 post-settle continuation", () => { ), ); + const INTERRUPT_SCENARIO = "codex-interrupt-mid-command"; + const INTERRUPT_NATIVE_THREAD = "native-codex-interrupt-thread"; + const INTERRUPT_NATIVE_TURN = "native-codex-interrupt-turn"; + const INTERRUPT_COMMAND_ITEM = "exec-codex-interrupt-command"; + const INTERRUPT_COMMAND_ITEM_TWO = "exec-codex-interrupt-command-two"; + const INTERRUPT_COMMAND = "bash -c 'sleep 30; echo SHOULD_NOT_FINISH_CMD_INTERRUPT_FIXTURE'"; + const INTERRUPT_COMMAND_TWO = "bash -c 'sleep 20; echo SECOND_COMMAND'"; + const INTERRUPT_PROMPT = "Run a long foreground command and wait until interrupted."; + + const interruptCommandItem = (status: "inProgress" | "completed"): Record => ({ + type: "commandExecution", + id: INTERRUPT_COMMAND_ITEM, + command: INTERRUPT_COMMAND, + cwd: "/workspace", + processId: "57680", + source: "unifiedExecStartup", + status, + commandActions: [{ type: "unknown", command: INTERRUPT_COMMAND }], + aggregatedOutput: status === "completed" ? "SHOULD_NOT_FINISH_CMD_INTERRUPT_FIXTURE\n" : null, + exitCode: status === "completed" ? 0 : null, + durationMs: status === "completed" ? 30_000 : null, + }); + + const interruptMidCommandTranscript = makeCodexReplayTranscript({ + scenario: INTERRUPT_SCENARIO, + entries: [ + ...codexReplayPreamble({ + nativeThreadId: INTERRUPT_NATIVE_THREAD, + nativeTurnId: INTERRUPT_NATIVE_TURN, + prompt: INTERRUPT_PROMPT, + }), + { + type: "emit_inbound", + label: "item/started/command", + frame: { + method: "item/started", + params: { + item: interruptCommandItem("inProgress"), + threadId: INTERRUPT_NATIVE_THREAD, + turnId: INTERRUPT_NATIVE_TURN, + startedAtMs: 1782622440500, + }, + }, + }, + { + type: "emit_inbound", + label: "item/started/command-two", + frame: { + method: "item/started", + params: { + item: { + ...interruptCommandItem("inProgress"), + id: INTERRUPT_COMMAND_ITEM_TWO, + command: INTERRUPT_COMMAND_TWO, + commandActions: [{ type: "unknown", command: INTERRUPT_COMMAND_TWO }], + }, + threadId: INTERRUPT_NATIVE_THREAD, + turnId: INTERRUPT_NATIVE_TURN, + startedAtMs: 1782622440600, + }, + }, + }, + { + type: "expect_outbound", + label: "turn/interrupt", + frame: { + id: 4, + method: "turn/interrupt", + params: { + threadId: INTERRUPT_NATIVE_THREAD, + turnId: INTERRUPT_NATIVE_TURN, + }, + }, + }, + { + type: "emit_inbound", + label: "turn/interrupt", + frame: { id: 4, result: {} }, + }, + { + type: "emit_inbound", + label: "turn/completed", + frame: { + method: "turn/completed", + params: { + threadId: INTERRUPT_NATIVE_THREAD, + turn: makeCodexReplayTurn({ + id: INTERRUPT_NATIVE_TURN, + status: "interrupted", + }), + }, + }, + }, + { + type: "emit_inbound", + label: "item/completed/command-late", + afterMs: 30_000, + frame: { + method: "item/completed", + params: { + item: interruptCommandItem("completed"), + threadId: INTERRUPT_NATIVE_THREAD, + turnId: INTERRUPT_NATIVE_TURN, + completedAtMs: 1782622465500, + }, + }, + }, + ], + }); + + it.effect( + "terminalizes running command items before turn.terminal on interrupt and ignores late completion", + () => + Effect.scoped( + Effect.gen(function* () { + const harness = yield* makeCodexReplayHarness(interruptMidCommandTranscript); + const now = yield* DateTime.now; + + yield* harness.runtime.startTurn( + makeCodexTestTurnInput({ + threadId: harness.threadId, + providerThread: harness.providerThread, + now, + attemptId: RunAttemptId.make("attempt-codex-interrupt-mid-command"), + text: INTERRUPT_PROMPT, + }), + ); + + yield* awaitUntil( + () => + harness.events.some( + (event) => + event.type === "turn_item.updated" && + event.turnItem.type === "command_execution" && + event.turnItem.status === "running", + ), + "running command item", + ); + + const providerTurnId = harness.events.find( + (event): event is Extract => + event.type === "provider_turn.updated", + )?.providerTurn.id; + assert.isDefined(providerTurnId); + + yield* harness.runtime.interruptTurn({ + providerThread: harness.providerThread, + providerTurnId, + }); + + yield* awaitUntil(() => harness.terminalEvents().length === 1, "interrupted terminal"); + assert.equal(harness.terminalEvents()[0]?.status, "interrupted"); + + const terminalIndex = harness.events.findIndex((event) => event.type === "turn.terminal"); + assert.isAtLeast(terminalIndex, 0); + + let lastCommandBeforeTerminal: + | Extract + | undefined; + for (let index = 0; index < terminalIndex; index++) { + const event = harness.events[index]; + if ( + event?.type === "turn_item.updated" && + event.turnItem.type === "command_execution" + ) { + lastCommandBeforeTerminal = event; + } + } + assert.isDefined(lastCommandBeforeTerminal); + assert.equal(lastCommandBeforeTerminal.turnItem.status, "interrupted"); + assert.isNotNull(lastCommandBeforeTerminal.turnItem.completedAt); + + const interruptedCommandsBeforeTerminal = harness.events + .slice(0, terminalIndex) + .flatMap((event) => + event.type === "turn_item.updated" && + event.turnItem.type === "command_execution" && + event.turnItem.status === "interrupted" + ? [event.turnItem.input] + : [], + ) + .sort(); + assert.deepEqual( + interruptedCommandsBeforeTerminal, + [INTERRUPT_COMMAND, INTERRUPT_COMMAND_TWO].sort(), + ); + + const interruptedCommandIndex = harness.events.findIndex( + (event, index) => + index < terminalIndex && + event.type === "turn_item.updated" && + event.turnItem.type === "command_execution" && + event.turnItem.status === "interrupted", + ); + assert.isAbove( + terminalIndex, + interruptedCommandIndex, + "command terminalization must precede turn.terminal", + ); + + assert.isFalse(yield* harness.hasPendingBackgroundWork); + assert.lengthOf(harness.continuationRequests, 0); + + // Late provider item/completed after interrupt must not revive the card + // or request a background-command wake continuation. + yield* TestClock.adjust("30 seconds"); + for (let attempt = 0; attempt < 100; attempt++) { + yield* Effect.yieldNow; + } + assert.lengthOf(harness.continuationRequests, 0); + assert.isFalse(yield* harness.hasPendingBackgroundWork); + assert.lengthOf(harness.terminalEvents(), 1); + + const postTerminalCommandUpdates = harness.events.filter( + (event, index) => + index > terminalIndex && + event.type === "turn_item.updated" && + event.turnItem.type === "command_execution", + ); + assert.lengthOf( + postTerminalCommandUpdates, + 0, + "late item/completed after interrupt must not project", + ); + + const commandUpdates = harness.events.filter( + (event): event is Extract => + event.type === "turn_item.updated" && event.turnItem.type === "command_execution", + ); + assert.isAtLeast(commandUpdates.length, 2, "start + interrupt terminalization"); + assert.equal(commandUpdates[commandUpdates.length - 1]?.turnItem.status, "interrupted"); + }).pipe(Effect.provide(Layer.merge(idAllocatorLayer, NodeServices.layer))), + ), + ); + + const FAILED_SCENARIO = "codex-failed-mid-command"; + const FAILED_NATIVE_THREAD = "native-codex-failed-thread"; + const FAILED_NATIVE_TURN = "native-codex-failed-turn"; + const FAILED_COMMAND_ITEM = "exec-codex-failed-command"; + const FAILED_COMMAND = "sleep 30"; + const FAILED_PROMPT = "Run a command that will be abandoned when the turn fails."; + + const failedMidCommandTranscript = makeCodexReplayTranscript({ + scenario: FAILED_SCENARIO, + entries: [ + ...codexReplayPreamble({ + nativeThreadId: FAILED_NATIVE_THREAD, + nativeTurnId: FAILED_NATIVE_TURN, + prompt: FAILED_PROMPT, + }), + { + type: "emit_inbound", + label: "item/started/command", + frame: { + method: "item/started", + params: { + item: { + type: "commandExecution", + id: FAILED_COMMAND_ITEM, + command: FAILED_COMMAND, + cwd: "/workspace", + processId: "99", + source: "unifiedExecStartup", + status: "inProgress", + commandActions: [{ type: "unknown", command: FAILED_COMMAND }], + aggregatedOutput: null, + exitCode: null, + durationMs: null, + }, + threadId: FAILED_NATIVE_THREAD, + turnId: FAILED_NATIVE_TURN, + startedAtMs: 1782622440500, + }, + }, + }, + { + type: "emit_inbound", + label: "turn/completed", + frame: { + method: "turn/completed", + params: { + threadId: FAILED_NATIVE_THREAD, + turn: { + ...makeCodexReplayTurn({ + id: FAILED_NATIVE_TURN, + status: "failed", + }), + error: { message: "provider failed mid-command" }, + }, + }, + }, + }, + ], + }); + + it.effect("terminalizes running command items before turn.terminal on failed turns", () => + Effect.scoped( + Effect.gen(function* () { + const harness = yield* makeCodexReplayHarness(failedMidCommandTranscript); + const now = yield* DateTime.now; + + yield* harness.runtime.startTurn( + makeCodexTestTurnInput({ + threadId: harness.threadId, + providerThread: harness.providerThread, + now, + attemptId: RunAttemptId.make("attempt-codex-failed-mid-command"), + text: FAILED_PROMPT, + }), + ); + yield* awaitUntil(() => harness.terminalEvents().length === 1, "failed terminal"); + assert.equal(harness.terminalEvents()[0]?.status, "failed"); + + const terminalIndex = harness.events.findIndex((event) => event.type === "turn.terminal"); + const failedCommandIndex = harness.events.findIndex( + (event, index) => + index < terminalIndex && + event.type === "turn_item.updated" && + event.turnItem.type === "command_execution" && + event.turnItem.status === "failed", + ); + assert.isAtLeast(failedCommandIndex, 0); + assert.isAbove( + terminalIndex, + failedCommandIndex, + "failed-turn command terminalization must precede turn.terminal", + ); + assert.isFalse(yield* harness.hasPendingBackgroundWork); + assert.lengthOf(harness.continuationRequests, 0); + }).pipe(Effect.provide(Layer.merge(idAllocatorLayer, NodeServices.layer))), + ), + ); + const RESUME_SCENARIO = "codex-resume-subagent"; const RESUME_NATIVE_THREAD = "native-codex-resume-thread"; const RESUME_NATIVE_TURN = "native-codex-resume-root-turn"; diff --git a/apps/server/src/orchestration-v2/Adapters/CodexAdapterV2.ts b/apps/server/src/orchestration-v2/Adapters/CodexAdapterV2.ts index 9adffad1e54..fcc79936fe2 100644 --- a/apps/server/src/orchestration-v2/Adapters/CodexAdapterV2.ts +++ b/apps/server/src/orchestration-v2/Adapters/CodexAdapterV2.ts @@ -759,6 +759,13 @@ interface ActiveCodexTurnContext { readonly startedAt: DateTime.Utc; } +/** Snapshot of a still-running commandExecution item for interrupt/fail terminalization. */ +interface TrackedRunningCommandItem { + readonly id: string; + readonly command: string; + readonly aggregatedOutput?: string; +} + interface CodexSubagentThreadContext { readonly parentContext: ActiveCodexTurnContext; readonly providerThread: OrchestrationV2ProviderThread; @@ -1349,7 +1356,9 @@ export function makeCodexAdapterV2(adapterOptions: CodexAdapterV2Options): Provi * keep projecting instead of being dropped. */ const settledTurns = yield* Ref.make(new Map()); - const runningCommandItemsByTurn = yield* Ref.make(new Map>()); + const runningCommandItemsByTurn = yield* Ref.make( + new Map>(), + ); const offeredContinuationItemsByTurn = yield* Ref.make(new Map>()); const emitProviderEvent = (event: ProviderAdapterV2Event) => @@ -1450,11 +1459,11 @@ export function makeCodexAdapterV2(adapterOptions: CodexAdapterV2Options): Provi return context === undefined ? undefined : ({ context, settled: false } as const); }); - const trackRunningCommandItem = (nativeTurnId: string, nativeItemId: string) => + const trackRunningCommandItem = (nativeTurnId: string, item: TrackedRunningCommandItem) => Ref.update(runningCommandItemsByTurn, (current) => { const updated = new Map(current); - const items = new Set(updated.get(nativeTurnId) ?? []); - items.add(nativeItemId); + const items = new Map(updated.get(nativeTurnId) ?? []); + items.set(item.id, item); updated.set(nativeTurnId, items); return updated; }); @@ -1466,7 +1475,7 @@ export function makeCodexAdapterV2(adapterOptions: CodexAdapterV2Options): Provi if (items === undefined || !items.has(nativeItemId)) { return [items === undefined || items.size === 0, current] as const; } - const remaining = new Set(items); + const remaining = new Map(items); remaining.delete(nativeItemId); const updated = new Map(current); if (remaining.size === 0) { @@ -1477,6 +1486,84 @@ export function makeCodexAdapterV2(adapterOptions: CodexAdapterV2Options): Provi return [remaining.size === 0, updated] as const; }); + /** + * When a turn is interrupted or failed, Codex often leaves commandExecution + * items mid-flight (no item/completed). Emit terminal turn items before + * turn.terminal so the projection never keeps a forever-running command card. + * Does not retain settled context: late completions must not wake the run. + */ + const terminalizeRunningCommandItems = ( + context: ActiveCodexTurnContext, + nativeTurnId: string, + status: "interrupted" | "failed", + completedAt: DateTime.Utc, + ) => + Effect.gen(function* () { + const items = (yield* Ref.get(runningCommandItemsByTurn)).get(nativeTurnId); + if (items === undefined || items.size === 0) { + return; + } + for (const tracked of items.values()) { + const nodeId = idAllocator.derive.nodeFromProviderItem({ + driver: CODEX_PROVIDER, + nativeItemId: tracked.id, + }); + const turnItemId = idAllocator.derive.turnItemFromProviderItem({ + driver: CODEX_PROVIDER, + nativeItemId: tracked.id, + }); + const ordinal = yield* resolveItemOrdinal(context, tracked.id); + const node: OrchestrationV2ExecutionNode = { + id: nodeId, + threadId: context.projectionThreadId, + runId: context.projectionRunId, + parentNodeId: context.itemParentNodeId, + rootNodeId: context.rootNodeId, + kind: "tool_call", + status, + countsForRun: false, + providerThreadId: context.providerThread.id, + providerTurnId: context.providerTurnId, + nativeItemRef: codexNativeItemRef(tracked.id), + runtimeRequestId: null, + checkpointScopeId: null, + startedAt: context.startedAt, + completedAt, + }; + const turnItem: OrchestrationV2TurnItem = { + id: turnItemId, + threadId: context.projectionThreadId, + runId: context.projectionRunId, + nodeId, + providerThreadId: context.providerThread.id, + providerTurnId: context.providerTurnId, + nativeItemRef: codexNativeItemRef(tracked.id), + parentItemId: null, + ordinal, + status, + title: null, + startedAt: context.startedAt, + completedAt, + updatedAt: completedAt, + type: "command_execution", + input: tracked.command, + ...(tracked.aggregatedOutput === undefined + ? {} + : { output: tracked.aggregatedOutput }), + }; + yield* emitProviderEvent({ + type: "node.updated", + driver: CODEX_PROVIDER, + node, + }); + yield* emitProviderEvent({ + type: "turn_item.updated", + driver: CODEX_PROVIDER, + turnItem, + }); + } + }); + const resolveItemOrdinal = (context: ActiveCodexTurnContext, nativeItemId: string) => Effect.gen(function* () { const existing = (yield* Ref.get(itemOrdinals)).get(nativeItemId); @@ -2926,7 +3013,14 @@ export function makeCodexAdapterV2(adapterOptions: CodexAdapterV2Options): Provi if (payload.item.type === "commandExecution") { if (!codexItemStatus(payload.item.status).completed) { - yield* trackRunningCommandItem(payload.turnId, payload.item.id); + yield* trackRunningCommandItem(payload.turnId, { + id: payload.item.id, + command: payload.item.command, + ...(payload.item.aggregatedOutput === null || + payload.item.aggregatedOutput === undefined + ? {} + : { aggregatedOutput: payload.item.aggregatedOutput }), + }); } const artifacts = yield* buildCommandExecutionArtifacts(context, payload.item); yield* emitProviderEvent({ @@ -3625,6 +3719,12 @@ export function makeCodexAdapterV2(adapterOptions: CodexAdapterV2Options): Provi completedAt, }); } + // Interrupted/failed turns leave in-flight commands without item/completed. + // Terminalize them before turn.terminal so projections never keep a + // forever-running command card. Completed turns retain background tracking. + if (status === "interrupted" || status === "failed") { + yield* terminalizeRunningCommandItems(context, payload.turn.id, status, completedAt); + } if (context.subagent === null) { const terminalStatus = providerTurnStatusToTerminal(status); yield* emitProviderEvent( diff --git a/apps/server/src/orchestration-v2/ProviderTurnStartService.ts b/apps/server/src/orchestration-v2/ProviderTurnStartService.ts index 693a00b57dd..3b6f74bf5f7 100644 --- a/apps/server/src/orchestration-v2/ProviderTurnStartService.ts +++ b/apps/server/src/orchestration-v2/ProviderTurnStartService.ts @@ -22,7 +22,7 @@ import { import { IdAllocatorV2 } from "./IdAllocator.ts"; import { ProjectionStoreV2 } from "./ProjectionStore.ts"; import { ProviderSessionManagerV2 } from "./ProviderSessionManager.ts"; -import { RunExecutionServiceV2 } from "./RunExecutionService.ts"; +import { canRouteRelatedSubagent, RunExecutionServiceV2 } from "./RunExecutionService.ts"; import { RuntimePolicyV2 } from "./RuntimePolicy.ts"; export class ProviderTurnStartError extends Schema.TaggedErrorClass()( @@ -403,6 +403,9 @@ export const layer: Layer.Layer< if (!runningWrite.committed) { return; } + const routableSubagents = projection.subagents.filter((subagent) => + canRouteRelatedSubagent(subagent.status), + ); yield* runExecution.startRootRun({ commandId: CommandId.make(`command:effect:provider-turn.start:${run.id}`), appThread: projection.thread, @@ -414,10 +417,10 @@ export const layer: Layer.Layer< providerThread: runningProviderThread, attempt: runningAttempt, attemptId: attempt.id, - relatedThreadIds: projection.subagents.flatMap((subagent) => + relatedThreadIds: routableSubagents.flatMap((subagent) => subagent.childThreadId === null ? [] : [subagent.childThreadId], ), - relatedProviderThreadIds: projection.subagents.flatMap((subagent) => + relatedProviderThreadIds: routableSubagents.flatMap((subagent) => subagent.providerThreadId === null ? [] : [subagent.providerThreadId], ), providerTurnOrdinal: diff --git a/apps/server/src/orchestration-v2/RunExecutionService.test.ts b/apps/server/src/orchestration-v2/RunExecutionService.test.ts index 2822f60851b..b290e9d96d7 100644 --- a/apps/server/src/orchestration-v2/RunExecutionService.test.ts +++ b/apps/server/src/orchestration-v2/RunExecutionService.test.ts @@ -2,14 +2,18 @@ import { assert, it } from "@effect/vitest"; import { CheckpointScopeId, CommandId, + EventId, MessageId, NodeId, type OrchestrationV2AppThread, type OrchestrationV2CheckpointScope, + type OrchestrationV2DomainEvent, type OrchestrationV2ExecutionNode, type OrchestrationV2ProviderThread, type OrchestrationV2Run, type OrchestrationV2RunAttempt, + type OrchestrationV2Subagent, + type OrchestrationV2TurnItem, ProviderDriverKind, ProviderInstanceId, ProviderSessionId, @@ -35,6 +39,8 @@ import { IdAllocatorV2, layer as idAllocatorLayer } from "./IdAllocator.ts"; import type { ProviderAdapterV2Event, ProviderAdapterV2SessionRuntime } from "./ProviderAdapter.ts"; import { ProviderEventIngestorV2 } from "./ProviderEventIngestor.ts"; import { + canRouteRelatedSubagent, + cascadeTerminalizeRunOwnedSubagents, finalProviderThreadStatus, layer as runExecutionServiceLayer, makeProviderEventRoutingState, @@ -181,6 +187,52 @@ it("does not route a superseded attempt through a reused provider thread", () => assert.isFalse(routeProviderEvent(oldTurnEvent, newAttempt, newState)[0]); }); +it("does not carry interrupted child ownership into later attempts", () => { + assert.isFalse(canRouteRelatedSubagent("interrupted")); + assert.isFalse(canRouteRelatedSubagent("failed")); + assert.isFalse(canRouteRelatedSubagent("cancelled")); + assert.isTrue(canRouteRelatedSubagent("completed")); + assert.isTrue(canRouteRelatedSubagent("running")); + + const threadId = ThreadId.make("thread:related-child:next-attempt"); + const childThreadId = ThreadId.make("thread:related-child:interrupted"); + const identity: ProviderEventRouteIdentity = { + threadId, + runId: RunId.make("run:related-child:next-attempt"), + attemptId: RunAttemptId.make("attempt:related-child:next-attempt"), + providerThreadId: ProviderThreadId.make("provider-thread:related-child:next-attempt"), + }; + const state = makeProviderEventRoutingState({ + identity, + providerTurnId: null, + relatedThreadIds: canRouteRelatedSubagent("interrupted") ? [childThreadId] : [], + }); + const childNodeId = NodeId.make("node:related-child:interrupted"); + const lateChildNode = { + type: "node.updated", + driver, + node: { + id: childNodeId, + threadId: childThreadId, + runId: null, + parentNodeId: null, + rootNodeId: childNodeId, + kind: "root_turn", + status: "completed", + countsForRun: false, + providerThreadId: null, + providerTurnId: null, + nativeItemRef: null, + runtimeRequestId: null, + checkpointScopeId: null, + startedAt: null, + completedAt: null, + }, + } satisfies ProviderAdapterV2Event; + + assert.isFalse(routeProviderEvent(lateChildNode, identity, state)[0]); +}); + it.effect("rechecks run ownership immediately before calling the provider", () => Effect.gen(function* () { const runExecution = yield* RunExecutionServiceV2; @@ -542,11 +594,730 @@ it.effect("does not pin ingestion on background items when the root turn is inte }), ); -it.effect("omits run_interrupt_result when a superseding attempt already owns the run", () => +it.effect( + "cascade-terminalizes run-owned subagent rows on interrupt before root finalization", + () => + Effect.gen(function* () { + const ids = backgroundScenarioIds("subagent-interrupt-cascade"); + const childThreadId = ids.childThreadId; + const unrelatedChildThreadId = ThreadId.make( + "thread:subagent-interrupt-cascade:unrelated-child", + ); + const providerInstanceId = ProviderInstanceId.make("codex"); + const written = yield* Ref.make>([]); + const ingested = yield* Ref.make>([]); + const ingestionDone = yield* Deferred.make(); + const testLayer = runExecutionServiceLayer.pipe( + Layer.provide( + Layer.mergeAll( + Layer.mock(CheckpointServiceV2)({ captureBaseline: () => Effect.void }), + Layer.mock(EventSinkV2)({ + write: () => Effect.succeed([]), + writeWithEffects: (input) => + Effect.gen(function* () { + yield* Ref.update(written, (current) => [...current, ...input.events]); + return []; + }), + writeIfRunCurrent: () => Effect.succeed({ committed: true, storedEvents: [] }), + }), + idAllocatorLayer, + Layer.mock(ProviderEventIngestorV2)({ + ingestNormalized: (input) => + Ref.update(ingested, (current) => [...current, input.event]).pipe(Effect.as([])), + }), + ServerSettingsService.layerTest(), + ), + ), + ); + + const runningSubagent = makeRunOwnedSubagentFixture({ + ids, + providerInstanceId, + childThreadId, + driver, + status: "running", + }); + const runningTurnItem = makeRunOwnedSubagentTurnItemFixture({ + ids, + providerInstanceId, + childThreadId, + driver, + status: "running", + }); + const runningNode = makeRunOwnedSubagentNodeFixture({ + ids, + status: "running", + }); + const runningChildNode = makeRunOwnedSubagentChildNodeFixture({ + ids, + status: "running", + }); + const unrelatedChildNode = { + ...runningChildNode, + id: NodeId.make("node:subagent-interrupt-cascade:unrelated-child"), + threadId: unrelatedChildThreadId, + runId: ids.runId, + rootNodeId: NodeId.make("node:subagent-interrupt-cascade:unrelated-child"), + }; + + yield* Effect.gen(function* () { + const runExecution = yield* RunExecutionServiceV2; + yield* runExecution.startRootRun({ + commandId: CommandId.make("command:subagent-interrupt-cascade"), + appThread: { id: ids.threadId } as OrchestrationV2AppThread, + providerSessionId: ProviderSessionId.make("session:subagent-interrupt-cascade"), + session: { + events: Stream.empty, + subscribeEvents: Effect.succeed({ + events: Stream.fromIterable([ + childThreadCreatedEvent(ids), + { + type: "subagent.updated", + driver, + subagent: runningSubagent, + }, + { + type: "node.updated", + driver, + node: runningNode, + }, + { + type: "node.updated", + driver, + node: runningChildNode, + }, + { + type: "node.updated", + driver, + node: unrelatedChildNode, + }, + { + type: "turn_item.updated", + driver, + turnItem: runningTurnItem, + }, + rootTerminalEvent(ids, "interrupted"), + // Late provider completion after interrupt must not be ingested. + { + type: "subagent.updated", + driver, + subagent: { + ...runningSubagent, + status: "completed" as const, + result: "should-not-apply", + completedAt: runningSubagent.updatedAt, + }, + }, + { + type: "turn_item.updated", + driver, + turnItem: { + ...runningTurnItem, + status: "completed" as const, + result: "should-not-apply", + completedAt: runningTurnItem.updatedAt, + updatedAt: runningTurnItem.updatedAt, + }, + }, + { + type: "node.updated", + driver, + node: { + ...runningChildNode, + status: "completed" as const, + completedAt: runningChildNode.startedAt, + }, + }, + ] satisfies ReadonlyArray), + close: Deferred.succeed(ingestionDone, undefined), + }), + startTurn: () => Effect.void, + } as unknown as ProviderAdapterV2SessionRuntime, + run: { + id: ids.runId, + threadId: ids.threadId, + ordinal: 1, + providerInstanceId, + } as OrchestrationV2Run, + rootNode: { id: ids.rootNodeId } as OrchestrationV2ExecutionNode, + checkpointScope: { + id: CheckpointScopeId.make("checkpoint-scope:subagent-interrupt-cascade"), + } as OrchestrationV2CheckpointScope, + providerThread: { + id: ids.providerThreadId, + driver, + } as OrchestrationV2ProviderThread, + attempt: { + id: ids.attemptId, + providerTurnId: ids.rootProviderTurnId, + } as OrchestrationV2RunAttempt, + attemptId: ids.attemptId, + relatedThreadIds: [unrelatedChildThreadId], + providerTurnOrdinal: 1, + message: { + messageId: MessageId.make("message:subagent-interrupt-cascade:user"), + text: "Spawn a subagent then stop.", + attachments: [], + createdBy: "user", + creationSource: "web", + }, + modelSelection: { instanceId: providerInstanceId, model: "gpt-5.4" }, + runtimePolicy: { + runtimeMode: "full-access", + interactionMode: "default", + cwd: process.cwd(), + approvalPolicy: "never", + sandboxPolicy: { + type: "readOnly", + access: { type: "fullAccess" }, + networkAccess: false, + }, + }, + }); + }).pipe(Effect.provide(testLayer)); + + const closed = yield* Deferred.await(ingestionDone).pipe(Effect.timeoutOption("2 seconds")); + assert.isTrue(Option.isSome(closed), "event ingestion fiber did not finish"); + + const events = yield* Ref.get(written); + const subagentEvents = events.flatMap((event) => + event.type === "subagent.updated" ? [event] : [], + ); + const turnItemEvents = events.flatMap((event) => { + if (event.type !== "turn-item.updated" || event.payload.type !== "subagent") { + return []; + } + return [ + { + ...event, + payload: event.payload, + }, + ]; + }); + const nodeEvents = events.flatMap((event) => + event.type === "node.updated" && + event.payload.status === "interrupted" && + (event.payload.id === runningNode.id || event.payload.id === runningChildNode.id) + ? [event] + : [], + ); + const runUpdatedIndex = events.findIndex((event) => event.type === "run.updated"); + assert.isAtLeast(runUpdatedIndex, 0, "root run.updated must be written"); + + assert.lengthOf(subagentEvents, 1); + const terminalSubagent = subagentEvents[0]; + assert.isDefined(terminalSubagent); + assert.equal(terminalSubagent.payload.status, "interrupted"); + assert.equal(terminalSubagent.payload.childThreadId, childThreadId); + assert.equal(terminalSubagent.payload.result, null); + assert.isNotNull(terminalSubagent.payload.completedAt); + + assert.lengthOf(turnItemEvents, 1); + const terminalTurnItem = turnItemEvents[0]; + assert.isDefined(terminalTurnItem); + assert.equal(terminalTurnItem.payload.status, "interrupted"); + assert.equal(terminalTurnItem.payload.childThreadId, childThreadId); + assert.equal(terminalTurnItem.payload.result, null); + + assert.lengthOf(nodeEvents, 2); + assert.isTrue( + nodeEvents.some( + (event) => event.payload.id === runningNode.id && event.payload.kind === "subagent", + ), + ); + assert.isFalse( + events.some( + (event) => + event.type === "node.updated" && + event.payload.id === unrelatedChildNode.id && + event.payload.status === "interrupted", + ), + "owned child threads without a live run-owned subagent link must not cascade", + ); + assert.isTrue( + nodeEvents.some( + (event) => + event.payload.id === runningChildNode.id && + event.runId === (runningChildNode.runId ?? ids.runId) && + event.payload.threadId === childThreadId && + event.payload.kind === "root_turn", + ), + ); + + const cascadeIndexes = events.flatMap((event, index) => { + if (event.type === "subagent.updated" && event.payload.status === "interrupted") { + return [index]; + } + if ( + event.type === "turn-item.updated" && + event.payload.type === "subagent" && + event.payload.status === "interrupted" + ) { + return [index]; + } + if (event.type === "node.updated" && event.payload.status === "interrupted") { + return event.payload.id === runningNode.id || event.payload.id === runningChildNode.id + ? [index] + : []; + } + return []; + }); + assert.isTrue( + cascadeIndexes.every((index) => index < runUpdatedIndex), + "subagent cascade must precede run.updated", + ); + assert.isFalse( + events.some((event) => { + if (event.type === "subagent.updated") { + return event.payload.status === "completed"; + } + if (event.type === "turn-item.updated" && event.payload.type === "subagent") { + return event.payload.status === "completed"; + } + return false; + }), + "late provider completion must not reopen cascaded subagent rows", + ); + assert.isFalse( + (yield* Ref.get(ingested)).some( + (event) => + (event.type === "subagent.updated" && event.subagent.status === "completed") || + (event.type === "turn_item.updated" && event.turnItem.status === "completed") || + (event.type === "node.updated" && event.node.status === "completed"), + ), + "late provider completion must not be ingested after interrupt", + ); + }), +); + +it.effect( + "cascades linked child-thread nodes after run-owned subagent and turn-item terminalize", + () => + Effect.gen(function* () { + const ids = backgroundScenarioIds("subagent-link-survives-terminal"); + const childThreadId = ids.childThreadId; + const unrelatedChildThreadId = ThreadId.make( + "thread:subagent-link-survives-terminal:unrelated-child", + ); + const providerInstanceId = ProviderInstanceId.make("codex"); + const written = yield* Ref.make>([]); + const ingestionDone = yield* Deferred.make(); + const testLayer = runExecutionServiceLayer.pipe( + Layer.provide( + Layer.mergeAll( + Layer.mock(CheckpointServiceV2)({ captureBaseline: () => Effect.void }), + Layer.mock(EventSinkV2)({ + write: () => Effect.succeed([]), + writeWithEffects: (input) => + Effect.gen(function* () { + yield* Ref.update(written, (current) => [...current, ...input.events]); + return []; + }), + writeIfRunCurrent: () => Effect.succeed({ committed: true, storedEvents: [] }), + }), + idAllocatorLayer, + Layer.mock(ProviderEventIngestorV2)({ + ingestNormalized: () => Effect.succeed([]), + }), + ServerSettingsService.layerTest(), + ), + ), + ); + + const runningSubagent = makeRunOwnedSubagentFixture({ + ids, + providerInstanceId, + childThreadId, + driver, + status: "running", + }); + const runningTurnItem = makeRunOwnedSubagentTurnItemFixture({ + ids, + providerInstanceId, + childThreadId, + driver, + status: "running", + }); + const runningNode = makeRunOwnedSubagentNodeFixture({ + ids, + status: "running", + }); + const runningChildNode = makeRunOwnedSubagentChildNodeFixture({ + ids, + status: "running", + }); + const unrelatedChildNode = { + ...runningChildNode, + id: NodeId.make("node:subagent-link-survives-terminal:unrelated-child"), + threadId: unrelatedChildThreadId, + runId: ids.runId, + rootNodeId: NodeId.make("node:subagent-link-survives-terminal:unrelated-child"), + }; + const completedAt = runningSubagent.updatedAt; + + yield* Effect.gen(function* () { + const runExecution = yield* RunExecutionServiceV2; + yield* runExecution.startRootRun({ + commandId: CommandId.make("command:subagent-link-survives-terminal"), + appThread: { id: ids.threadId } as OrchestrationV2AppThread, + providerSessionId: ProviderSessionId.make("session:subagent-link-survives-terminal"), + session: { + events: Stream.empty, + subscribeEvents: Effect.succeed({ + events: Stream.fromIterable([ + childThreadCreatedEvent(ids), + { + type: "subagent.updated", + driver, + subagent: runningSubagent, + }, + { + type: "node.updated", + driver, + node: runningNode, + }, + { + type: "node.updated", + driver, + node: runningChildNode, + }, + { + type: "node.updated", + driver, + node: unrelatedChildNode, + }, + { + type: "turn_item.updated", + driver, + turnItem: runningTurnItem, + }, + // Subagent + turn-item settle before root interrupt; linkage + // must still prove the open child-thread node is cascadeable. + { + type: "subagent.updated", + driver, + subagent: { + ...runningSubagent, + status: "completed" as const, + result: "subagent finished first", + completedAt, + }, + }, + { + type: "turn_item.updated", + driver, + turnItem: { + ...runningTurnItem, + status: "completed" as const, + result: "turn item finished first", + completedAt, + updatedAt: completedAt, + }, + }, + { + type: "node.updated", + driver, + node: { + ...runningNode, + status: "completed" as const, + completedAt, + }, + }, + rootTerminalEvent(ids, "interrupted"), + ] satisfies ReadonlyArray), + close: Deferred.succeed(ingestionDone, undefined), + }), + startTurn: () => Effect.void, + } as unknown as ProviderAdapterV2SessionRuntime, + run: { + id: ids.runId, + threadId: ids.threadId, + ordinal: 1, + providerInstanceId, + } as OrchestrationV2Run, + rootNode: { id: ids.rootNodeId } as OrchestrationV2ExecutionNode, + checkpointScope: { + id: CheckpointScopeId.make("checkpoint-scope:subagent-link-survives-terminal"), + } as OrchestrationV2CheckpointScope, + providerThread: { + id: ids.providerThreadId, + driver, + } as OrchestrationV2ProviderThread, + attempt: { + id: ids.attemptId, + providerTurnId: ids.rootProviderTurnId, + } as OrchestrationV2RunAttempt, + attemptId: ids.attemptId, + relatedThreadIds: [unrelatedChildThreadId], + providerTurnOrdinal: 1, + message: { + messageId: MessageId.make("message:subagent-link-survives-terminal:user"), + text: "Subagent settles before root interrupt.", + attachments: [], + createdBy: "user", + creationSource: "web", + }, + modelSelection: { instanceId: providerInstanceId, model: "gpt-5.4" }, + runtimePolicy: { + runtimeMode: "full-access", + interactionMode: "default", + cwd: process.cwd(), + approvalPolicy: "never", + sandboxPolicy: { + type: "readOnly", + access: { type: "fullAccess" }, + networkAccess: false, + }, + }, + }); + }).pipe(Effect.provide(testLayer)); + + const closed = yield* Deferred.await(ingestionDone).pipe(Effect.timeoutOption("2 seconds")); + assert.isTrue(Option.isSome(closed), "event ingestion fiber did not finish"); + + const events = yield* Ref.get(written); + const runUpdatedIndex = events.findIndex((event) => event.type === "run.updated"); + assert.isAtLeast(runUpdatedIndex, 0, "root run.updated must be written"); + + const cascadedChildNodeEvents = events.flatMap((event, index) => + event.type === "node.updated" && + event.payload.id === runningChildNode.id && + event.payload.status === "interrupted" + ? [{ event, index }] + : [], + ); + assert.lengthOf( + cascadedChildNodeEvents, + 1, + "open linked child-thread node must cascade after subagent/turn-item terminalize", + ); + const cascadedChild = cascadedChildNodeEvents[0]; + assert.isDefined(cascadedChild); + assert.isTrue( + cascadedChild.index < runUpdatedIndex, + "child-thread cascade must precede run.updated", + ); + assert.equal(cascadedChild.event.payload.threadId, childThreadId); + assert.equal(cascadedChild.event.payload.kind, "root_turn"); + assert.isFalse( + events.some( + (event) => + event.type === "node.updated" && + event.payload.id === unrelatedChildNode.id && + event.payload.status === "interrupted", + ), + "related but unlinked child threads must not cascade", + ); + assert.isFalse( + events.some( + (event) => event.type === "subagent.updated" && event.payload.status === "interrupted", + ), + "already-terminal subagent rows must not be re-cascaded", + ); + assert.isFalse( + events.some( + (event) => + event.type === "turn-item.updated" && + event.payload.type === "subagent" && + event.payload.status === "interrupted", + ), + "already-terminal subagent turn items must not be re-cascaded", + ); + }), +); + +it.effect("cascade helper is provider-neutral for Claude and Codex-shaped subagent rows", () => + Effect.gen(function* () { + const now = yield* DateTime.now; + let nextId = 0; + const allocateEventId = () => + Effect.sync(() => EventId.make(`event:cascade-helper:${nextId++}`)); + + for (const driverKind of [ + ProviderDriverKind.make("claudeAgent"), + ProviderDriverKind.make("codex"), + ] as const) { + const runId = RunId.make(`run:cascade-helper:${driverKind}`); + const threadId = ThreadId.make(`thread:cascade-helper:${driverKind}`); + const childThreadId = ThreadId.make(`thread:cascade-helper:${driverKind}:child`); + const subagentId = NodeId.make(`node:cascade-helper:${driverKind}:subagent`); + const childNodeId = NodeId.make(`node:cascade-helper:${driverKind}:child-root`); + const providerInstanceId = ProviderInstanceId.make(String(driverKind)); + const terminalStatus = driverKind === "claudeAgent" ? "failed" : "cancelled"; + const subagent: OrchestrationV2Subagent = { + id: subagentId, + threadId, + runId, + parentNodeId: NodeId.make(`node:cascade-helper:${driverKind}:root`), + origin: "provider_native", + createdBy: "agent", + driver: driverKind, + providerInstanceId, + providerThreadId: null, + childThreadId, + nativeTaskRef: { + driver: driverKind, + nativeId: `native-${driverKind}`, + strength: "strong", + }, + prompt: "hold", + title: "hold", + model: null, + status: "running", + progress: "partial progress", + result: "partial result", + startedAt: now, + completedAt: null, + updatedAt: now, + }; + const turnItem = { + id: TurnItemId.make(`turn-item:cascade-helper:${driverKind}`), + threadId, + runId, + nodeId: subagentId, + providerThreadId: null, + providerTurnId: null, + nativeItemRef: subagent.nativeTaskRef, + parentItemId: null, + ordinal: 3, + status: "running" as const, + title: "hold", + startedAt: now, + completedAt: null, + updatedAt: now, + type: "subagent" as const, + subagentId, + origin: "provider_native" as const, + driver: driverKind, + providerInstanceId, + childThreadId, + prompt: "hold", + progress: "partial progress", + result: "partial result", + } satisfies Extract; + const node: OrchestrationV2ExecutionNode = { + id: subagentId, + threadId, + runId, + parentNodeId: NodeId.make(`node:cascade-helper:${driverKind}:root`), + rootNodeId: NodeId.make(`node:cascade-helper:${driverKind}:root`), + kind: "subagent", + status: "running", + countsForRun: false, + providerThreadId: null, + providerTurnId: null, + nativeItemRef: subagent.nativeTaskRef, + runtimeRequestId: null, + checkpointScopeId: null, + startedAt: now, + completedAt: null, + }; + const openChildNode: OrchestrationV2ExecutionNode = { + id: childNodeId, + threadId: childThreadId, + runId: null, + parentNodeId: null, + rootNodeId: childNodeId, + kind: "root_turn", + status: "running", + countsForRun: false, + providerThreadId: null, + providerTurnId: null, + nativeItemRef: null, + runtimeRequestId: null, + checkpointScopeId: null, + startedAt: now, + completedAt: null, + }; + + const events = yield* cascadeTerminalizeRunOwnedSubagents({ + run: { + id: runId, + threadId, + ordinal: 1, + providerInstanceId, + } as OrchestrationV2Run, + open: { + subagents: new Map([[subagentId, subagent]]), + turnItems: new Map([[subagentId, turnItem]]), + nodes: new Map([[subagentId, node]]), + linkedChildThreadIds: new Set([childThreadId]), + }, + status: terminalStatus, + completedAt: now, + allocateEventId, + }); + + assert.equal(events.length, 3, `${driverKind}: subagent + node + turn item`); + const terminalSubagent = events.find((event) => event.type === "subagent.updated"); + assert.isDefined(terminalSubagent); + if (terminalSubagent?.type !== "subagent.updated") { + assert.fail("expected subagent.updated event"); + return; + } + assert.equal(terminalSubagent.payload.status, terminalStatus); + assert.equal(terminalSubagent.payload.childThreadId, childThreadId); + assert.equal(terminalSubagent.payload.progress, "partial progress"); + assert.equal(terminalSubagent.payload.result, "partial result"); + assert.equal(terminalSubagent.payload.driver, driverKind); + + const terminalItem = events.find( + (event) => event.type === "turn-item.updated" && event.payload.type === "subagent", + ); + assert.isDefined(terminalItem); + if (terminalItem?.type !== "turn-item.updated" || terminalItem.payload.type !== "subagent") { + assert.fail("expected subagent turn-item.updated event"); + return; + } + assert.equal(terminalItem.payload.status, terminalStatus); + assert.equal(terminalItem.payload.childThreadId, childThreadId); + assert.equal(terminalItem.payload.progress, "partial progress"); + assert.equal(terminalItem.payload.result, "partial result"); + + // Shared cascade path: after subagent/turn-item rows are gone, only the + // preserved linkage may prove an open child-thread node is cascadeable. + const afterTerminalLinkEvents = yield* cascadeTerminalizeRunOwnedSubagents({ + run: { + id: runId, + threadId, + ordinal: 1, + providerInstanceId, + } as OrchestrationV2Run, + open: { + subagents: new Map(), + turnItems: new Map(), + nodes: new Map([[childNodeId, openChildNode]]), + linkedChildThreadIds: new Set([childThreadId]), + }, + status: terminalStatus, + completedAt: now, + allocateEventId, + }); + assert.equal( + afterTerminalLinkEvents.length, + 1, + `${driverKind}: only open linked child node cascades after link rows terminalize`, + ); + const cascadedChild = afterTerminalLinkEvents[0]; + assert.isDefined(cascadedChild); + if (cascadedChild?.type !== "node.updated") { + assert.fail("expected node.updated for linked child thread"); + return; + } + assert.equal(cascadedChild.payload.id, childNodeId); + assert.equal(cascadedChild.payload.threadId, childThreadId); + assert.equal(cascadedChild.payload.status, terminalStatus); + assert.equal(cascadedChild.payload.kind, "root_turn"); + } + }), +); + +it.effect("omits interrupt results and subagent cascade for a superseded attempt", () => Effect.gen(function* () { const written = yield* captureInterruptTerminalTurnItems({ key: "steer-supersede", shouldFinalizeRun: () => Effect.succeed(false), + seedOpenSubagent: true, }); assert.deepEqual( written.map((item) => item.type), @@ -609,10 +1380,18 @@ function captureInterruptTerminalTurnItems(input: { readonly key: string; readonly shouldFinalizeRun: () => Effect.Effect; readonly hasUnpairedRunInterruptRequest?: () => Effect.Effect; + readonly seedOpenSubagent?: boolean; }) { return Effect.gen(function* () { const ids = backgroundScenarioIds(input.key); const providerInstanceId = ProviderInstanceId.make("codex"); + const runningSubagent = makeRunOwnedSubagentFixture({ + ids, + providerInstanceId, + childThreadId: ids.childThreadId, + driver, + status: "running", + }); const writtenItems = yield* Ref.make< ReadonlyArray<{ readonly type: string; readonly parentItemId: string | null }> >([]); @@ -668,7 +1447,30 @@ function captureInterruptTerminalTurnItems(input: { session: { events: Stream.empty, subscribeEvents: Effect.succeed({ - events: Stream.fromIterable([rootTerminalEvent(ids, "interrupted")]), + events: Stream.fromIterable([ + ...(input.seedOpenSubagent + ? [ + { type: "subagent.updated", driver, subagent: runningSubagent } as const, + { + type: "node.updated", + driver, + node: makeRunOwnedSubagentNodeFixture({ ids, status: "running" }), + } as const, + { + type: "turn_item.updated", + driver, + turnItem: makeRunOwnedSubagentTurnItemFixture({ + ids, + providerInstanceId, + childThreadId: ids.childThreadId, + driver, + status: "running", + }), + } as const, + ] + : []), + rootTerminalEvent(ids, "interrupted"), + ] satisfies ReadonlyArray), close: Deferred.succeed(ingestionDone, undefined), }), startTurn: () => Effect.void, @@ -831,6 +1633,131 @@ function subagentEvent( } as ProviderAdapterV2Event; } +function makeRunOwnedSubagentFixture(input: { + readonly ids: BackgroundScenarioIds; + readonly providerInstanceId: ProviderInstanceId; + readonly childThreadId: ThreadId; + readonly driver: typeof driver; + readonly status: "running" | "interrupted"; +}): OrchestrationV2Subagent { + const now = DateTime.makeUnsafe("2026-07-21T12:00:00.000Z"); + return { + id: input.ids.subagentNodeId, + threadId: input.ids.threadId, + runId: input.ids.runId, + parentNodeId: input.ids.rootNodeId, + origin: "provider_native", + createdBy: "agent", + driver: input.driver, + providerInstanceId: input.providerInstanceId, + providerThreadId: null, + childThreadId: input.childThreadId, + nativeTaskRef: { + driver: input.driver, + nativeId: `task:${input.ids.subagentNodeId}`, + strength: "strong", + }, + prompt: "hold", + title: "Live-test subagent hold", + model: null, + status: input.status, + result: null, + startedAt: now, + completedAt: null, + updatedAt: now, + }; +} + +function makeRunOwnedSubagentTurnItemFixture(input: { + readonly ids: BackgroundScenarioIds; + readonly providerInstanceId: ProviderInstanceId; + readonly childThreadId: ThreadId; + readonly driver: typeof driver; + readonly status: "running" | "interrupted"; +}): Extract { + const now = DateTime.makeUnsafe("2026-07-21T12:00:00.000Z"); + return { + id: input.ids.itemId, + threadId: input.ids.threadId, + runId: input.ids.runId, + nodeId: input.ids.subagentNodeId, + providerThreadId: input.ids.providerThreadId, + providerTurnId: input.ids.rootProviderTurnId, + nativeItemRef: { + driver: input.driver, + nativeId: `task:${input.ids.subagentNodeId}`, + strength: "strong", + }, + parentItemId: null, + ordinal: 3, + status: input.status, + title: "Live-test subagent hold", + startedAt: now, + completedAt: null, + updatedAt: now, + type: "subagent", + subagentId: input.ids.subagentNodeId, + origin: "provider_native", + driver: input.driver, + providerInstanceId: input.providerInstanceId, + childThreadId: input.childThreadId, + prompt: "hold", + result: null, + }; +} + +function makeRunOwnedSubagentNodeFixture(input: { + readonly ids: BackgroundScenarioIds; + readonly status: "running" | "interrupted"; +}): OrchestrationV2ExecutionNode { + const now = DateTime.makeUnsafe("2026-07-21T12:00:00.000Z"); + return { + id: input.ids.subagentNodeId, + threadId: input.ids.threadId, + runId: input.ids.runId, + parentNodeId: input.ids.rootNodeId, + rootNodeId: input.ids.rootNodeId, + kind: "subagent", + status: input.status, + countsForRun: false, + providerThreadId: input.ids.providerThreadId, + providerTurnId: input.ids.rootProviderTurnId, + nativeItemRef: { + driver, + nativeId: `task:${input.ids.subagentNodeId}`, + strength: "strong", + }, + runtimeRequestId: null, + checkpointScopeId: null, + startedAt: now, + completedAt: null, + }; +} + +function makeRunOwnedSubagentChildNodeFixture(input: { + readonly ids: BackgroundScenarioIds; + readonly status: "running" | "interrupted"; +}): OrchestrationV2ExecutionNode { + const now = DateTime.makeUnsafe("2026-07-21T12:00:00.000Z"); + return { + id: NodeId.make(`${input.ids.subagentNodeId}:child-root`), + threadId: input.ids.childThreadId, + runId: null, + parentNodeId: null, + rootNodeId: NodeId.make(`${input.ids.subagentNodeId}:child-root`), + kind: "root_turn", + status: input.status, + countsForRun: false, + providerThreadId: null, + providerTurnId: null, + nativeItemRef: null, + runtimeRequestId: null, + checkpointScopeId: null, + startedAt: now, + completedAt: null, + }; +} + function rootTerminalEvent( ids: BackgroundScenarioIds, status: "completed" | "interrupted", diff --git a/apps/server/src/orchestration-v2/RunExecutionService.ts b/apps/server/src/orchestration-v2/RunExecutionService.ts index c833fdf7b36..6109c56f248 100644 --- a/apps/server/src/orchestration-v2/RunExecutionService.ts +++ b/apps/server/src/orchestration-v2/RunExecutionService.ts @@ -1,9 +1,11 @@ import { CommandId, + type EventId, type ModelSelection, type NodeId, type OrchestrationV2AppThread, type OrchestrationV2CheckpointScope, + type OrchestrationV2DomainEvent, type OrchestrationV2ExecutionNode, type OrchestrationV2ProviderFailure, type OrchestrationV2ProviderThread, @@ -31,7 +33,11 @@ import * as Stream from "effect/Stream"; import { ServerSettingsService } from "../serverSettings.ts"; import { CheckpointServiceV2 } from "./CheckpointService.ts"; import { EventSinkV2 } from "./EventSink.ts"; -import { IdAllocatorV2, type IdAllocatorV2Shape } from "./IdAllocator.ts"; +import { + IdAllocatorV2, + type IdAllocatorV2AllocationError, + type IdAllocatorV2Shape, +} from "./IdAllocator.ts"; import type { ProviderAdapterV2Event, ProviderAdapterV2RuntimePolicy, @@ -94,6 +100,147 @@ function isTerminalTurnItemStatus(status: OrchestrationV2TurnItem["status"]): bo ); } +type SubagentTurnItem = Extract; + +type OpenRunOwnedSubagentProjection = { + readonly subagents: ReadonlyMap; + readonly turnItems: ReadonlyMap; + readonly nodes: ReadonlyMap; + /** Child threads once linked by a root-run subagent row; kept for cascade. */ + readonly linkedChildThreadIds: ReadonlySet; +}; + +type RunOwnedSubagentTerminalStatus = Extract< + OrchestrationV2Subagent["status"], + "interrupted" | "failed" | "cancelled" +>; + +function isOpenExecutionNodeStatus(status: OrchestrationV2ExecutionNode["status"]): boolean { + return status === "pending" || status === "running" || status === "waiting"; +} + +function isRunOwnedSubagentTerminalStatus( + status: ProviderTerminalEvent["status"], +): status is RunOwnedSubagentTerminalStatus { + return status === "interrupted" || status === "failed" || status === "cancelled"; +} + +export function canRouteRelatedSubagent(status: OrchestrationV2Subagent["status"]): boolean { + return status !== "interrupted" && status !== "failed" && status !== "cancelled"; +} + +function emptyOpenRunOwnedSubagentProjection(): OpenRunOwnedSubagentProjection { + return { + subagents: new Map(), + turnItems: new Map(), + nodes: new Map(), + linkedChildThreadIds: new Set(), + }; +} + +function withLinkedChildThreadId( + current: OpenRunOwnedSubagentProjection, + childThreadId: ThreadId | null, +): OpenRunOwnedSubagentProjection { + if (childThreadId === null || current.linkedChildThreadIds.has(childThreadId)) { + return current; + } + const linkedChildThreadIds = new Set(current.linkedChildThreadIds); + linkedChildThreadIds.add(childThreadId); + return { ...current, linkedChildThreadIds }; +} + +export function cascadeTerminalizeRunOwnedSubagents(input: { + readonly run: OrchestrationV2Run; + readonly open: OpenRunOwnedSubagentProjection; + readonly status: RunOwnedSubagentTerminalStatus; + readonly completedAt: DateTime.Utc; + readonly allocateEventId: () => Effect.Effect; +}): Effect.Effect, IdAllocatorV2AllocationError> { + return Effect.gen(function* () { + const events: Array = []; + // Prefer lifetime linkage over currently-open rows: subagent/turn-item + // snapshots may terminalize before the linked child-thread node settles. + const childThreadIds = new Set(input.open.linkedChildThreadIds); + for (const item of [...input.open.subagents.values(), ...input.open.turnItems.values()]) { + if (item.childThreadId !== null) { + childThreadIds.add(item.childThreadId); + } + } + const keys = new Set([ + ...input.open.subagents.keys(), + ...input.open.turnItems.keys(), + ...input.open.nodes.keys(), + ]); + for (const key of keys) { + const subagent = input.open.subagents.get(key); + if (subagent !== undefined && !isTerminalSubagentStatus(subagent.status)) { + events.push({ + id: yield* input.allocateEventId(), + type: "subagent.updated", + threadId: subagent.threadId, + runId: input.run.id, + nodeId: subagent.id, + driver: subagent.driver, + providerInstanceId: subagent.providerInstanceId, + occurredAt: input.completedAt, + payload: { + ...subagent, + status: input.status, + completedAt: input.completedAt, + updatedAt: input.completedAt, + }, + }); + } + const node = input.open.nodes.get(key); + if ( + node !== undefined && + ((node.threadId === input.run.threadId && node.runId === input.run.id) || + childThreadIds.has(node.threadId)) && + isOpenExecutionNodeStatus(node.status) + ) { + events.push({ + id: yield* input.allocateEventId(), + type: "node.updated", + threadId: node.threadId, + runId: node.runId ?? input.run.id, + nodeId: node.id, + providerInstanceId: input.run.providerInstanceId, + occurredAt: input.completedAt, + payload: { + ...node, + status: input.status, + completedAt: input.completedAt, + }, + }); + } + const turnItem = input.open.turnItems.get(key); + if ( + turnItem !== undefined && + turnItem.runId === input.run.id && + !isTerminalTurnItemStatus(turnItem.status) + ) { + events.push({ + id: yield* input.allocateEventId(), + type: "turn-item.updated", + threadId: turnItem.threadId, + runId: input.run.id, + ...(turnItem.nodeId === null ? {} : { nodeId: turnItem.nodeId }), + providerInstanceId: input.run.providerInstanceId, + occurredAt: input.completedAt, + payload: { + ...turnItem, + status: input.status, + completedAt: input.completedAt, + updatedAt: input.completedAt, + }, + }); + } + } + return events; + }); +} + export function finalProviderThreadStatus( disposition: ProviderTerminalEvent["threadDisposition"], ): OrchestrationV2ProviderThread["status"] { @@ -325,6 +472,7 @@ export const layer: Layer.Layer< readonly attempt: OrchestrationV2RunAttempt; readonly shouldFinalizeRun?: () => Effect.Effect; readonly hasUnpairedRunInterruptRequest?: () => Effect.Effect; + readonly openRunOwnedSubagents?: OpenRunOwnedSubagentProjection; readonly terminal: ProviderTerminalEvent; readonly failureItemPersisted: boolean; }) => @@ -372,6 +520,20 @@ export const layer: Layer.Layer< } return; } + const allocateEventId = () => idAllocator.allocate.event({ threadId: input.run.threadId }); + const open = input.openRunOwnedSubagents ?? emptyOpenRunOwnedSubagentProjection(); + const hasOpenSubagentProjection = + open.subagents.size > 0 || open.turnItems.size > 0 || open.nodes.size > 0; + const cascadedSubagentEvents = + isRunOwnedSubagentTerminalStatus(input.terminal.status) && hasOpenSubagentProjection + ? yield* cascadeTerminalizeRunOwnedSubagents({ + run: input.run, + open, + status: input.terminal.status, + completedAt, + allocateEventId, + }) + : []; const persistedStatus = input.terminal.status === "completed" ? "waiting" : input.terminal.status; const finalizedRun: OrchestrationV2Run = { @@ -390,11 +552,9 @@ export const layer: Layer.Layer< status: finalProviderThreadStatus(input.terminal.threadDisposition), updatedAt: completedAt, }; - const runEventId = yield* idAllocator.allocate.event({ threadId: input.run.threadId }); - const nodeEventId = yield* idAllocator.allocate.event({ threadId: input.run.threadId }); - const providerThreadEventId = yield* idAllocator.allocate.event({ - threadId: input.run.threadId, - }); + const runEventId = yield* allocateEventId(); + const nodeEventId = yield* allocateEventId(); + const providerThreadEventId = yield* allocateEventId(); const checkpointCaptureCommandId = CommandId.make( `command:effect:checkpoint.capture:${input.run.id}`, ); @@ -415,11 +575,14 @@ export const layer: Layer.Layer< ] : [], events: [ + // Terminalize open run-owned subagent rows before the root run + // settles so projections never keep a forever-running subagent card. + ...cascadedSubagentEvents, ...(finalizedAttempt === null ? [] : [ { - id: yield* idAllocator.allocate.event({ threadId: input.run.threadId }), + id: yield* allocateEventId(), type: "run-attempt.updated" as const, threadId: input.run.threadId, runId: input.run.id, @@ -432,7 +595,7 @@ export const layer: Layer.Layer< ...(input.terminal.status === "interrupted" ? [ { - id: yield* idAllocator.allocate.event({ threadId: input.run.threadId }), + id: yield* allocateEventId(), type: "turn-item.updated" as const, threadId: input.run.threadId, runId: input.run.id, @@ -452,7 +615,7 @@ export const layer: Layer.Layer< ...(input.terminal.status === "failed" && !input.failureItemPersisted ? [ { - id: yield* idAllocator.allocate.event({ threadId: input.run.threadId }), + id: yield* allocateEventId(), type: "turn-item.updated" as const, threadId: input.run.threadId, runId: input.run.id, @@ -588,12 +751,14 @@ export const layer: Layer.Layer< const activeBackgroundTurnItems = yield* Ref.make< ReadonlySet >(new Set()); + const openRunOwnedSubagents = yield* Ref.make(emptyOpenRunOwnedSubagentProjection()); const finalizeRootRun = (terminal: ProviderTerminalEvent) => Effect.gen(function* () { if (yield* Ref.get(rootRunFinalized)) { return; } const providerThread = yield* Ref.get(latestProviderThread); + const openSubagents = yield* Ref.get(openRunOwnedSubagents); yield* writeFinalRunEvents({ run: input.run, rootNode: input.rootNode, @@ -608,6 +773,7 @@ export const layer: Layer.Layer< : { hasUnpairedRunInterruptRequest: input.hasUnpairedRunInterruptRequest, }), + openRunOwnedSubagents: openSubagents, terminal, failureItemPersisted: terminal.status === "failed", }).pipe( @@ -615,6 +781,9 @@ export const layer: Layer.Layer< (cause) => new RunExecutionIngestError({ runId: input.run.id, cause }), ), ); + if (isRunOwnedSubagentTerminalStatus(terminal.status)) { + yield* Ref.set(openRunOwnedSubagents, emptyOpenRunOwnedSubagentProjection()); + } yield* Ref.set(rootRunFinalized, true); }); const trackChildLifecycle = (event: ProviderAdapterV2Event) => @@ -652,6 +821,41 @@ export const layer: Layer.Layer< return next; }); } + // Snapshot run-owned subagents for interrupt cascade. + // Preserve childThreadId linkage for the root-run lifetime even + // after the subagent row terminalizes, so open child-thread + // nodes can still be proven linked on a later root interrupt. + if (belongsToRootRun) { + yield* Ref.update(openRunOwnedSubagents, (current) => { + const withLink = withLinkedChildThreadId(current, event.subagent.childThreadId); + const subagents = new Map(withLink.subagents); + if (isTerminalSubagentStatus(event.subagent.status)) { + subagents.delete(event.subagent.id); + } else { + subagents.set(event.subagent.id, event.subagent); + } + return { ...withLink, subagents }; + }); + } + } + if (event.type === "node.updated") { + const belongsToRootSubagent = + event.node.kind === "subagent" && event.node.runId === input.run.id; + const belongsToOwnedChildThread = + event.node.threadId !== input.run.threadId && + routing.ownedThreadIds.has(event.node.threadId); + if (!belongsToRootSubagent && !belongsToOwnedChildThread) { + return; + } + yield* Ref.update(openRunOwnedSubagents, (current) => { + const nodes = new Map(current.nodes); + if (isOpenExecutionNodeStatus(event.node.status)) { + nodes.set(event.node.id, event.node); + } else { + nodes.delete(event.node.id); + } + return { ...current, nodes }; + }); } if ( event.type === "turn_item.updated" && @@ -672,12 +876,29 @@ export const layer: Layer.Layer< return next; }); } + if (belongsToRootRun && event.turnItem.type === "subagent") { + const subagentItem = event.turnItem; + yield* Ref.update(openRunOwnedSubagents, (current) => { + const withLink = withLinkedChildThreadId(current, subagentItem.childThreadId); + const turnItems = new Map(withLink.turnItems); + if (isTerminalTurnItemStatus(subagentItem.status)) { + turnItems.delete(subagentItem.subagentId); + } else { + turnItems.set(subagentItem.subagentId, subagentItem); + } + return { ...withLink, turnItems }; + }); + } } }); const shouldStopProviderEventIngestion = Effect.gen(function* () { if (!(yield* Ref.get(rootTerminalSeen))) { return false; } + const terminal = yield* Ref.get(terminalEvent); + if (terminal !== null && terminal.status !== "completed") { + return true; + } const childProviderTurns = yield* Ref.get(activeChildProviderTurns); if (childProviderTurns.size > 0) { return false; @@ -694,12 +915,8 @@ export const layer: Layer.Layer< // rather than pinning the stream open. Assumes adapters emit an // item's non-terminal event before the root terminal; an item // first seen after the terminal is not pinned. - const terminal = yield* Ref.get(terminalEvent); - if (terminal !== null && terminal.status === "completed") { - const backgroundItems = yield* Ref.get(activeBackgroundTurnItems); - return backgroundItems.size === 0; - } - return true; + const backgroundItems = yield* Ref.get(activeBackgroundTurnItems); + return backgroundItems.size === 0; }); const eventSubscription = input.session.subscribeEvents === undefined @@ -781,30 +998,35 @@ export const layer: Layer.Layer< Effect.flatMap((providerThread) => Ref.get(latestTurnItemOrdinal).pipe( Effect.flatMap((latestItemOrdinal) => - writeFinalRunEvents({ - run: input.run, - rootNode: input.rootNode, - checkpointScope: input.checkpointScope, - providerThread, - attempt: input.attempt, - ...(input.shouldFinalizeRun === undefined - ? {} - : { shouldFinalizeRun: input.shouldFinalizeRun }), - ...(input.hasUnpairedRunInterruptRequest === undefined - ? {} - : { - hasUnpairedRunInterruptRequest: - input.hasUnpairedRunInterruptRequest, - }), - terminal: makeFailedTerminalEvent( - makeProviderFailure({ - cause: Cause.squash(cause), - class: "unknown", + Ref.get(openRunOwnedSubagents).pipe( + Effect.flatMap((openSubagents) => + writeFinalRunEvents({ + run: input.run, + rootNode: input.rootNode, + checkpointScope: input.checkpointScope, + providerThread, + attempt: input.attempt, + ...(input.shouldFinalizeRun === undefined + ? {} + : { shouldFinalizeRun: input.shouldFinalizeRun }), + ...(input.hasUnpairedRunInterruptRequest === undefined + ? {} + : { + hasUnpairedRunInterruptRequest: + input.hasUnpairedRunInterruptRequest, + }), + openRunOwnedSubagents: openSubagents, + terminal: makeFailedTerminalEvent( + makeProviderFailure({ + cause: Cause.squash(cause), + class: "unknown", + }), + latestItemOrdinal + 1, + ), + failureItemPersisted: false, }), - latestItemOrdinal + 1, ), - failureItemPersisted: false, - }), + ), ), ), ), @@ -858,30 +1080,35 @@ export const layer: Layer.Layer< Effect.flatMap((providerThread) => Ref.get(latestTurnItemOrdinal).pipe( Effect.flatMap((latestItemOrdinal) => - writeFinalRunEvents({ - run: input.run, - rootNode: input.rootNode, - checkpointScope: input.checkpointScope, - providerThread, - attempt: input.attempt, - ...(input.shouldFinalizeRun === undefined - ? {} - : { shouldFinalizeRun: input.shouldFinalizeRun }), - ...(input.hasUnpairedRunInterruptRequest === undefined - ? {} - : { - hasUnpairedRunInterruptRequest: - input.hasUnpairedRunInterruptRequest, - }), - terminal: makeFailedTerminalEvent( - makeProviderFailure({ - cause: Cause.squash(cause), - class: "provider_error", + Ref.get(openRunOwnedSubagents).pipe( + Effect.flatMap((openSubagents) => + writeFinalRunEvents({ + run: input.run, + rootNode: input.rootNode, + checkpointScope: input.checkpointScope, + providerThread, + attempt: input.attempt, + ...(input.shouldFinalizeRun === undefined + ? {} + : { shouldFinalizeRun: input.shouldFinalizeRun }), + ...(input.hasUnpairedRunInterruptRequest === undefined + ? {} + : { + hasUnpairedRunInterruptRequest: + input.hasUnpairedRunInterruptRequest, + }), + openRunOwnedSubagents: openSubagents, + terminal: makeFailedTerminalEvent( + makeProviderFailure({ + cause: Cause.squash(cause), + class: "provider_error", + }), + latestItemOrdinal + 1, + ), + failureItemPersisted: false, }), - latestItemOrdinal + 1, ), - failureItemPersisted: false, - }), + ), ), ), ), diff --git a/apps/server/src/orchestration-v2/testkit/fixtures/turn_interrupt_mid_tool/codex_output.ts b/apps/server/src/orchestration-v2/testkit/fixtures/turn_interrupt_mid_tool/codex_output.ts index f3f98b53104..47ea4dcf426 100644 --- a/apps/server/src/orchestration-v2/testkit/fixtures/turn_interrupt_mid_tool/codex_output.ts +++ b/apps/server/src/orchestration-v2/testkit/fixtures/turn_interrupt_mid_tool/codex_output.ts @@ -80,7 +80,10 @@ export function assertTurnInterruptMidToolCodexOutput( assert.isDefined(commandItem); assert.isDefined(interruptRequest); assert.isDefined(interruptResult); - assert.include(["running", "completed", "failed"], commandItem.status); + // Interrupted Codex turns must terminalize mid-flight commandExecution items so + // the projected card is never left running forever after turn.terminal. + assert.equal(commandItem.status, "interrupted"); + assert.isNotNull(commandItem.completedAt); assert.include(commandItem.input, "node -e"); assert.equal(interruptRequest.status, "completed"); assert.equal(interruptResult.status, "interrupted"); @@ -91,4 +94,15 @@ export function assertTurnInterruptMidToolCodexOutput( ); assert.equal(projection.providerThreads[0]?.status, "idle"); assert.include(["interrupted", "cancelled"], projection.providerTurns[0]?.status); + + const runningCommands = projection.turnItems.filter( + (item) => item.type === "command_execution" && item.status === "running", + ); + assert.lengthOf(runningCommands, 0, "interrupted turn must not leave running command items"); + + const toolNodes = projection.nodes.filter((node) => node.kind === "tool_call"); + for (const node of toolNodes) { + assert.notEqual(node.status, "running", "interrupted turn tool nodes must be terminal"); + assert.isNotNull(node.completedAt); + } }