From 7e666bf3bed4e303caef978c246e731a68c1c0de Mon Sep 17 00:00:00 2001 From: dlowzzxx Date: Sun, 13 Sep 2026 09:51:35 +0200 Subject: [PATCH] fix(orchestration): correlate steered and stale aborts (#8939) --- .../Layers/ProviderCommandReactor.ts | 6 + .../Layers/ProviderRuntimeIngestion.test.ts | 2063 +++++++++++++---- .../Layers/ProviderRuntimeIngestion.ts | 143 +- .../provider/Layers/OpenCodeAdapter.test.ts | 171 ++ .../src/provider/Layers/OpenCodeAdapter.ts | 67 +- packages/contracts/src/provider.ts | 7 + 6 files changed, 1931 insertions(+), 526 deletions(-) diff --git a/apps/server/src/orchestration/Layers/ProviderCommandReactor.ts b/apps/server/src/orchestration/Layers/ProviderCommandReactor.ts index cdaadba1a96d..f34809b623c4 100644 --- a/apps/server/src/orchestration/Layers/ProviderCommandReactor.ts +++ b/apps/server/src/orchestration/Layers/ProviderCommandReactor.ts @@ -2,6 +2,7 @@ import { type ChatAttachment, CommandId, EventId, + type MessageId, type ModelSelection, type OrchestrationEvent, ProviderDriverKind, @@ -930,6 +931,7 @@ const make = Effect.gen(function* () { const buildSendTurnRequestForThread = Effect.fnUntraced(function* (input: { readonly threadId: ThreadId; + readonly turnStartMessageId?: MessageId; readonly messageText: string; readonly attachments?: ReadonlyArray; readonly modelSelection?: ModelSelection; @@ -981,6 +983,9 @@ const make = Effect.gen(function* () { return { threadId: input.threadId, + ...(input.turnStartMessageId !== undefined + ? { turnStartMessageId: input.turnStartMessageId } + : {}), ...(normalizedInput ? { input: normalizedInput } : {}), ...(normalizedAttachments.length > 0 ? { attachments: normalizedAttachments } : {}), ...(modelForTurn !== undefined ? { modelSelection: modelForTurn } : {}), @@ -1546,6 +1551,7 @@ const make = Effect.gen(function* () { } const sendTurnRequest = yield* buildSendTurnRequestForThread({ threadId: event.payload.threadId, + turnStartMessageId: event.payload.messageId, messageText: projectComposerContextForProvider({ text: message.text, records: message.context?.records ?? [], diff --git a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts index 1094ab48b7ac..dd004dbc30c3 100644 --- a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts +++ b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts @@ -318,6 +318,9 @@ describe("ProviderRuntimeIngestion", () => { const engine = await testRuntime.runPromise(Effect.service(OrchestrationEngineService)); const snapshotQuery = await testRuntime.runPromise(Effect.service(ProjectionSnapshotQuery)); const ingestion = await testRuntime.runPromise(Effect.service(ProviderRuntimeIngestionService)); + const planProgress = await testRuntime.runPromise( + Effect.service(ThreadPlanProgress.ThreadPlanProgressService), + ); scope = await Effect.runPromise(Scope.make("sequential")); await testRuntime.runPromise(ingestion.start().pipe(Scope.provide(scope))); const drain = () => testRuntime.runPromise(ingestion.drain); @@ -396,6 +399,7 @@ describe("ProviderRuntimeIngestion", () => { sqlCount: sqlCounter.count, setProviderSession: provider.setSession, drain, + planProgress, }; } @@ -714,631 +718,1696 @@ describe("ProviderRuntimeIngestion", () => { }); }); - it("applies provider session.state.changed transitions directly", async () => { + it("settles the durable session when a turn is aborted", async () => { const harness = await createHarness(); - const waitingAt = "2026-01-01T00:00:00.000Z"; + const now = "2026-01-01T00:00:00.000Z"; + const turnId = asTurnId("turn-aborted"); harness.emit({ - type: "session.state.changed", - eventId: asEventId("evt-session-state-waiting"), - provider: ProviderDriverKind.make("codex"), + type: "turn.started", + eventId: asEventId("evt-turn-aborted-started"), + provider: ProviderDriverKind.make("opencode"), threadId: asThreadId("thread-1"), - createdAt: waitingAt, - payload: { - state: "waiting", - reason: "awaiting approval", - }, + createdAt: now, + turnId, }); - let thread = await waitForThread( - harness.readModel, - (entry) => entry.session?.status === "running" && entry.session?.activeTurnId === null, + await harness.drain(); + let thread = (await harness.readModel()).threads.find( + (entry) => entry.id === asThreadId("thread-1"), ); - expect(thread.session?.status).toBe("running"); - expect(thread.session?.lastError).toBeNull(); + expect(thread?.session?.status).toBe("running"); + expect(thread?.session?.activeTurnId).toBe(turnId); harness.emit({ - type: "session.state.changed", - eventId: asEventId("evt-session-state-error"), - provider: ProviderDriverKind.make("codex"), + type: "turn.aborted", + eventId: asEventId("evt-turn-aborted"), + provider: ProviderDriverKind.make("opencode"), threadId: asThreadId("thread-1"), - createdAt: "2026-01-01T00:00:00.000Z", + createdAt: "2026-01-01T00:00:01.000Z", + turnId, payload: { - state: "error", - reason: "provider crashed", + reason: "Interrupted by user.", }, }); - thread = await waitForThread( - harness.readModel, - (entry) => - entry.session?.status === "error" && - entry.session?.activeTurnId === null && - entry.session?.lastError === "provider crashed", + await harness.drain(); + thread = (await harness.readModel()).threads.find( + (entry) => entry.id === asThreadId("thread-1"), ); - expect(thread.session?.status).toBe("error"); - expect(thread.session?.lastError).toBe("provider crashed"); + expect(thread?.session?.status).toBe("interrupted"); + expect(thread?.session?.activeTurnId).toBeNull(); + expect(thread?.session?.lastError).toBeNull(); + expect(thread?.session?.providerName).toBe("opencode"); + expect(thread?.session?.runtimeMode).toBe("approval-required"); + expect(thread?.session?.updatedAt).toBe("2026-01-01T00:00:01.000Z"); + expect(thread?.latestTurn?.state).toBe("interrupted"); + expect(thread?.latestTurn?.completedAt).toBe("2026-01-01T00:00:01.000Z"); + }); + + it("ignores a late turn start for a retained aborted turn", async () => { + const harness = await createHarness(); + const threadId = asThreadId("thread-1"); + const turnId = asTurnId("turn-late-after-abort"); + const turnStartMessageId = asMessageId("message-late-after-abort"); + const now = "2026-01-01T00:00:00.000Z"; + harness.setProviderSession({ + provider: ProviderDriverKind.make("opencode"), + status: "running", + runtimeMode: "approval-required", + threadId, + createdAt: now, + updatedAt: now, + activeTurnId: turnId, + activeTurnStartMessageId: turnStartMessageId, + }); harness.emit({ - type: "session.state.changed", - eventId: asEventId("evt-session-state-stopped"), - provider: ProviderDriverKind.make("codex"), - threadId: asThreadId("thread-1"), - createdAt: "2026-01-01T00:00:00.000Z", - payload: { - state: "stopped", - }, + type: "turn.started", + eventId: asEventId("evt-turn-late-after-abort-started"), + provider: ProviderDriverKind.make("opencode"), + threadId, + createdAt: now, + turnId, }); + await harness.drain(); - thread = await waitForThread( - harness.readModel, - (entry) => - entry.session?.status === "stopped" && - entry.session?.activeTurnId === null && - entry.session?.lastError === "provider crashed", - ); - expect(thread.session?.status).toBe("stopped"); - expect(thread.session?.lastError).toBe("provider crashed"); + harness.emit({ + type: "turn.aborted", + eventId: asEventId("evt-turn-late-after-abort-aborted"), + provider: ProviderDriverKind.make("opencode"), + threadId, + createdAt: "2026-01-01T00:00:01.000Z", + turnId, + payload: { reason: "Interrupted by user." }, + }); + await harness.drain(); + harness.setProviderSession({ + provider: ProviderDriverKind.make("opencode"), + status: "ready", + runtimeMode: "approval-required", + threadId, + createdAt: now, + updatedAt: "2026-01-01T00:00:01.000Z", + lastAbortedTurnId: turnId, + lastAbortedMessageId: turnStartMessageId, + }); harness.emit({ - type: "session.state.changed", - eventId: asEventId("evt-session-state-ready"), - provider: ProviderDriverKind.make("codex"), - threadId: asThreadId("thread-1"), - createdAt: "2026-01-01T00:00:00.000Z", - payload: { - state: "ready", - }, + type: "turn.started", + eventId: asEventId("evt-turn-late-after-abort-late-start"), + provider: ProviderDriverKind.make("opencode"), + threadId, + createdAt: "2026-01-01T00:00:02.000Z", + turnId, }); + await harness.drain(); - thread = await waitForThread( - harness.readModel, - (entry) => - entry.session?.status === "ready" && - entry.session?.activeTurnId === null && - entry.session?.lastError === null, - ); - expect(thread.session?.status).toBe("ready"); - expect(thread.session?.lastError).toBeNull(); + const thread = (await harness.readModel()).threads.find((entry) => entry.id === threadId); + expect(thread?.session?.status).toBe("interrupted"); + expect(thread?.session?.activeTurnId).toBeNull(); + expect(thread?.latestTurn?.state).toBe("interrupted"); }); - it("clears active turn when provider session becomes ready", async () => { + it("persists an OpenCode prompt failure carried by an aborted turn", async () => { const harness = await createHarness(); - const now = "2026-01-01T00:00:00.000Z"; + const turnId = asTurnId("turn-aborted-prompt-failure"); + const failureReason = "OpenCode prompt submission did not complete within 10 seconds."; harness.emit({ type: "turn.started", - eventId: asEventId("evt-turn-started-session-ready"), - provider: ProviderDriverKind.make("codex"), - createdAt: now, + eventId: asEventId("evt-turn-aborted-prompt-failure-started"), + provider: ProviderDriverKind.make("opencode"), threadId: asThreadId("thread-1"), - turnId: asTurnId("turn-session-ready"), + createdAt: "2026-01-01T00:00:00.000Z", + turnId, }); - - await waitForThread( - harness.readModel, - (thread) => - thread.session?.status === "running" && - thread.session?.activeTurnId === "turn-session-ready", - 10_000, - ); + await harness.drain(); harness.emit({ - type: "session.state.changed", - eventId: asEventId("evt-session-state-ready-with-active-turn"), - provider: ProviderDriverKind.make("codex"), + type: "turn.aborted", + eventId: asEventId("evt-turn-aborted-prompt-failure"), + provider: ProviderDriverKind.make("opencode"), threadId: asThreadId("thread-1"), createdAt: "2026-01-01T00:00:01.000Z", - payload: { - state: "ready", - }, + turnId, + payload: { reason: failureReason }, }); + await harness.drain(); - const thread = await waitForThread( - harness.readModel, - (entry) => - entry.session?.status === "ready" && - entry.session?.activeTurnId === null && - entry.session?.lastError === null, - 10_000, + const thread = (await harness.readModel()).threads.find( + (entry) => entry.id === asThreadId("thread-1"), ); - expect(thread.session?.status).toBe("ready"); - expect(thread.session?.activeTurnId).toBeNull(); - expect(thread.session?.lastError).toBeNull(); + expect(thread?.session?.status).toBe("error"); + expect(thread?.session?.activeTurnId).toBeNull(); + expect(thread?.session?.lastError).toBe(failureReason); + expect(thread?.latestTurn?.state).toBe("error"); }); - effectIt.effect( - "keeps a reconnecting pending turn starting while ready clears stale active state", - () => - Effect.gen(function* () { - const harness = yield* Effect.promise(() => createHarness()); - const threadId = asThreadId("thread-1"); - const staleTurnId = asTurnId("turn-stale-before-reconnect"); + it("accepts a named Codex abort during a pending turn start without a message token", async () => { + const harness = await createHarness(); + const threadId = asThreadId("thread-1"); + const turnId = asTurnId("turn-codex-aborted-pending"); + const messageId = asMessageId("message-codex-aborted-pending"); + const createdAt = "2026-01-01T00:00:00.000Z"; + const abortReason = "Codex cancelled the turn after the user stopped it."; - yield* harness.engine.dispatch({ - type: "thread.turn.start", - commandId: CommandId.make("cmd-turn-start-pending-reconnect"), - threadId, - message: { - messageId: MessageId.make("message-pending-reconnect"), - role: "user", - text: "resume after reconnect", - attachments: [], - }, - interactionMode: DEFAULT_PROVIDER_INTERACTION_MODE, - runtimeMode: "approval-required", - createdAt: "2026-01-01T00:00:01.000Z", - }); - yield* harness.engine.dispatch({ - type: "thread.session.set", - commandId: CommandId.make("cmd-session-starting-pending-reconnect"), - threadId, - session: { - threadId, - status: "starting", - providerName: "codex", - runtimeMode: "approval-required", - activeTurnId: staleTurnId, - lastError: null, - updatedAt: "2026-01-01T00:00:01.000Z", - }, - createdAt: "2026-01-01T00:00:01.000Z", - }); + await harness.dispatch({ + type: "thread.turn.start", + commandId: CommandId.make("cmd-turn-start-codex-abort"), + threadId, + message: { + messageId, + role: "user", + text: "start the Codex turn", + attachments: [], + }, + interactionMode: DEFAULT_PROVIDER_INTERACTION_MODE, + runtimeMode: "approval-required", + createdAt, + }); + await harness.dispatch({ + type: "thread.session.set", + commandId: CommandId.make("cmd-session-starting-codex-abort"), + threadId, + session: { + threadId, + status: "starting", + providerName: "codex", + runtimeMode: "approval-required", + activeTurnId: null, + lastError: "previous provider error", + updatedAt: createdAt, + }, + createdAt, + }); + harness.setProviderSession({ + provider: ProviderDriverKind.make("codex"), + status: "running", + runtimeMode: "approval-required", + threadId, + createdAt, + updatedAt: createdAt, + activeTurnId: turnId, + }); + harness.emit({ + type: "turn.aborted", + eventId: asEventId("evt-turn-abort-codex-pending"), + provider: ProviderDriverKind.make("codex"), + threadId, + createdAt: "2026-01-01T00:00:01.000Z", + turnId, + payload: { reason: abortReason }, + }); + await harness.drain(); - harness.emit({ - type: "session.state.changed", - eventId: asEventId("evt-session-ready-pending-reconnect"), - provider: ProviderDriverKind.make("codex"), - threadId, - createdAt: "2026-01-01T00:00:02.000Z", - payload: { state: "ready" }, - }); + const thread = (await harness.readModel()).threads.find((entry) => entry.id === threadId); + expect(thread?.session?.status).toBe("interrupted"); + expect(thread?.session?.activeTurnId).toBeNull(); + expect(thread?.session?.lastError).toBeNull(); + expect(thread?.session?.updatedAt).toBe("2026-01-01T00:00:01.000Z"); + }); - let thread = yield* Effect.promise(() => - waitForThread( - harness.readModel, - (entry) => entry.session?.status === "starting" && entry.session.activeTurnId === null, - ), - ); - expect(thread.session?.status).toBe("starting"); - expect(thread.session?.activeTurnId).toBeNull(); + it("finalizes buffered assistant text and proposed plans when a turn is aborted", async () => { + const harness = await createHarness(); + const threadId = asThreadId("thread-1"); + const turnId = asTurnId("turn-abort-buffered-cleanup"); + const now = "2026-01-01T00:00:00.000Z"; + + harness.emit({ + type: "turn.started", + eventId: asEventId("evt-turn-abort-buffered-cleanup-started"), + provider: ProviderDriverKind.make("opencode"), + threadId, + createdAt: now, + turnId, + }); + await harness.drain(); + + harness.emit({ + type: "turn.plan.updated", + eventId: asEventId("evt-turn-abort-buffered-cleanup-plan"), + provider: ProviderDriverKind.make("opencode"), + threadId, + createdAt: now, + turnId, + payload: { + explanation: "Working", + plan: [ + { step: "Inspect", status: "completed" }, + { step: "Apply", status: "in_progress" }, + ], + }, + }); + await harness.drain(); + expect(harness.planProgress.getThreadPlanProgress(threadId)).toMatchObject({ + step: "Apply", + completedSteps: 1, + totalSteps: 2, + }); + + harness.emit({ + type: "turn.aborted", + eventId: asEventId("evt-turn-abort-buffered-cleanup-untargeted"), + provider: ProviderDriverKind.make("opencode"), + threadId, + createdAt: now, + payload: { reason: "Interrupted by user." }, + }); + await harness.drain(); + expect(harness.planProgress.getThreadPlanProgress(threadId)).not.toBeNull(); + + harness.emit({ + type: "content.delta", + eventId: asEventId("evt-turn-abort-buffered-cleanup-message"), + provider: ProviderDriverKind.make("opencode"), + threadId, + createdAt: now, + turnId, + itemId: asItemId("item-abort-buffered-cleanup"), + payload: { streamKind: "assistant_text", delta: "partial answer" }, + }); + harness.emit({ + type: "turn.proposed.delta", + eventId: asEventId("evt-turn-abort-buffered-cleanup-proposed-plan"), + provider: ProviderDriverKind.make("opencode"), + threadId, + createdAt: now, + turnId, + payload: { delta: "# Proposed plan\n\n- Keep the useful work" }, + }); + await harness.drain(); + + harness.emit({ + type: "turn.aborted", + eventId: asEventId("evt-turn-abort-buffered-cleanup"), + provider: ProviderDriverKind.make("opencode"), + threadId, + createdAt: "2026-01-01T00:00:01.000Z", + turnId, + payload: { reason: "Interrupted by user." }, + }); + await harness.drain(); + + const thread = (await harness.readModel()).threads.find((entry) => entry.id === threadId); + const message = thread?.messages.find( + (entry: ProviderRuntimeTestMessage) => entry.id === "assistant:item-abort-buffered-cleanup", + ); + expect(message?.text).toBe("partial answer"); + expect(message?.streaming).toBe(false); + expect(thread?.proposedPlans).toEqual( + expect.arrayContaining([ + expect.objectContaining({ + id: "plan:thread-1:turn:turn-abort-buffered-cleanup", + planMarkdown: "# Proposed plan\n\n- Keep the useful work", + }), + ]), + ); + expect(harness.planProgress.getThreadPlanProgress(threadId)).toBeNull(); + }); + + it("completes streamed assistant messages when a turn is aborted", async () => { + const harness = await createHarness({ serverSettings: { enableLegacyTokenStreaming: true } }); + const threadId = asThreadId("thread-1"); + const turnId = asTurnId("turn-abort-streaming-cleanup"); + const now = "2026-01-01T00:00:00.000Z"; + + harness.emit({ + type: "turn.started", + eventId: asEventId("evt-turn-abort-streaming-cleanup-started"), + provider: ProviderDriverKind.make("opencode"), + threadId, + createdAt: now, + turnId, + }); + await harness.drain(); + harness.emit({ + type: "content.delta", + eventId: asEventId("evt-turn-abort-streaming-cleanup-message"), + provider: ProviderDriverKind.make("opencode"), + threadId, + createdAt: now, + turnId, + itemId: asItemId("item-abort-streaming-cleanup"), + payload: { streamKind: "assistant_text", delta: "streamed answer" }, + }); + await harness.drain(); + + let thread = (await harness.readModel()).threads.find((entry) => entry.id === threadId); + expect( + thread?.messages.find( + (entry: ProviderRuntimeTestMessage) => + entry.id === "assistant:item-abort-streaming-cleanup", + )?.streaming, + ).toBe(true); + + harness.emit({ + type: "turn.aborted", + eventId: asEventId("evt-turn-abort-streaming-cleanup"), + provider: ProviderDriverKind.make("opencode"), + threadId, + createdAt: "2026-01-01T00:00:01.000Z", + turnId, + payload: { reason: "Interrupted by user." }, + }); + await harness.drain(); + + thread = (await harness.readModel()).threads.find((entry) => entry.id === threadId); + const message = thread?.messages.find( + (entry: ProviderRuntimeTestMessage) => entry.id === "assistant:item-abort-streaming-cleanup", + ); + expect(message?.text).toBe("streamed answer"); + expect(message?.streaming).toBe(false); + }); + + it("applies provider session.state.changed transitions directly", async () => { + const harness = await createHarness(); + const waitingAt = "2026-01-01T00:00:00.000Z"; + + harness.emit({ + type: "session.state.changed", + eventId: asEventId("evt-session-state-waiting"), + provider: ProviderDriverKind.make("codex"), + threadId: asThreadId("thread-1"), + createdAt: waitingAt, + payload: { + state: "waiting", + reason: "awaiting approval", + }, + }); + + let thread = await waitForThread( + harness.readModel, + (entry) => entry.session?.status === "running" && entry.session?.activeTurnId === null, + ); + expect(thread.session?.status).toBe("running"); + expect(thread.session?.lastError).toBeNull(); + + harness.emit({ + type: "session.state.changed", + eventId: asEventId("evt-session-state-error"), + provider: ProviderDriverKind.make("codex"), + threadId: asThreadId("thread-1"), + createdAt: "2026-01-01T00:00:00.000Z", + payload: { + state: "error", + reason: "provider crashed", + }, + }); + + thread = await waitForThread( + harness.readModel, + (entry) => + entry.session?.status === "error" && + entry.session?.activeTurnId === null && + entry.session?.lastError === "provider crashed", + ); + expect(thread.session?.status).toBe("error"); + expect(thread.session?.lastError).toBe("provider crashed"); + + harness.emit({ + type: "session.state.changed", + eventId: asEventId("evt-session-state-stopped"), + provider: ProviderDriverKind.make("codex"), + threadId: asThreadId("thread-1"), + createdAt: "2026-01-01T00:00:00.000Z", + payload: { + state: "stopped", + }, + }); + + thread = await waitForThread( + harness.readModel, + (entry) => + entry.session?.status === "stopped" && + entry.session?.activeTurnId === null && + entry.session?.lastError === "provider crashed", + ); + expect(thread.session?.status).toBe("stopped"); + expect(thread.session?.lastError).toBe("provider crashed"); + + harness.emit({ + type: "session.state.changed", + eventId: asEventId("evt-session-state-ready"), + provider: ProviderDriverKind.make("codex"), + threadId: asThreadId("thread-1"), + createdAt: "2026-01-01T00:00:00.000Z", + payload: { + state: "ready", + }, + }); + + thread = await waitForThread( + harness.readModel, + (entry) => + entry.session?.status === "ready" && + entry.session?.activeTurnId === null && + entry.session?.lastError === null, + ); + expect(thread.session?.status).toBe("ready"); + expect(thread.session?.lastError).toBeNull(); + }); + + it("clears active turn when provider session becomes ready", async () => { + const harness = await createHarness(); + const now = "2026-01-01T00:00:00.000Z"; + + harness.emit({ + type: "turn.started", + eventId: asEventId("evt-turn-started-session-ready"), + provider: ProviderDriverKind.make("codex"), + createdAt: now, + threadId: asThreadId("thread-1"), + turnId: asTurnId("turn-session-ready"), + }); + + await waitForThread( + harness.readModel, + (thread) => + thread.session?.status === "running" && + thread.session?.activeTurnId === "turn-session-ready", + 10_000, + ); + + harness.emit({ + type: "session.state.changed", + eventId: asEventId("evt-session-state-ready-with-active-turn"), + provider: ProviderDriverKind.make("codex"), + threadId: asThreadId("thread-1"), + createdAt: "2026-01-01T00:00:01.000Z", + payload: { + state: "ready", + }, + }); + + const thread = await waitForThread( + harness.readModel, + (entry) => + entry.session?.status === "ready" && + entry.session?.activeTurnId === null && + entry.session?.lastError === null, + 10_000, + ); + expect(thread.session?.status).toBe("ready"); + expect(thread.session?.activeTurnId).toBeNull(); + expect(thread.session?.lastError).toBeNull(); + }); + + effectIt.effect( + "keeps a reconnecting pending turn starting while ready clears stale active state", + () => + Effect.gen(function* () { + const harness = yield* Effect.promise(() => createHarness()); + const threadId = asThreadId("thread-1"); + const staleTurnId = asTurnId("turn-stale-before-reconnect"); + + yield* harness.engine.dispatch({ + type: "thread.turn.start", + commandId: CommandId.make("cmd-turn-start-pending-reconnect"), + threadId, + message: { + messageId: MessageId.make("message-pending-reconnect"), + role: "user", + text: "resume after reconnect", + attachments: [], + }, + interactionMode: DEFAULT_PROVIDER_INTERACTION_MODE, + runtimeMode: "approval-required", + createdAt: "2026-01-01T00:00:01.000Z", + }); + yield* harness.engine.dispatch({ + type: "thread.session.set", + commandId: CommandId.make("cmd-session-starting-pending-reconnect"), + threadId, + session: { + threadId, + status: "starting", + providerName: "codex", + runtimeMode: "approval-required", + activeTurnId: staleTurnId, + lastError: null, + updatedAt: "2026-01-01T00:00:01.000Z", + }, + createdAt: "2026-01-01T00:00:01.000Z", + }); + + harness.emit({ + type: "session.state.changed", + eventId: asEventId("evt-session-ready-pending-reconnect"), + provider: ProviderDriverKind.make("codex"), + threadId, + createdAt: "2026-01-01T00:00:02.000Z", + payload: { state: "ready" }, + }); + + let thread = yield* Effect.promise(() => + waitForThread( + harness.readModel, + (entry) => entry.session?.status === "starting" && entry.session.activeTurnId === null, + ), + ); + expect(thread.session?.status).toBe("starting"); + expect(thread.session?.activeTurnId).toBeNull(); + + harness.emit({ + type: "session.started", + eventId: asEventId("evt-session-started-pending-reconnect"), + provider: ProviderDriverKind.make("codex"), + threadId, + createdAt: "2026-01-01T00:00:03.000Z", + }); + yield* Effect.promise(() => harness.drain()); + thread = (yield* Effect.promise(() => harness.readModel())).threads.find( + (entry) => entry.id === threadId, + )!; + expect(thread.session?.status).toBe("starting"); + expect(thread.session?.activeTurnId).toBeNull(); + + harness.emit({ + type: "turn.started", + eventId: asEventId("evt-turn-started-pending-reconnect"), + provider: ProviderDriverKind.make("codex"), + threadId, + turnId: asTurnId("turn-after-reconnect"), + createdAt: "2026-01-01T00:00:04.000Z", + }); + thread = yield* Effect.promise(() => + waitForThread( + harness.readModel, + (entry) => + entry.session?.status === "running" && + entry.session.activeTurnId === asTurnId("turn-after-reconnect"), + ), + ); + expect(thread.session?.status).toBe("running"); + + harness.emit({ + type: "session.started", + eventId: asEventId("evt-session-started-duplicate-midturn"), + provider: ProviderDriverKind.make("codex"), + threadId, + createdAt: "2026-01-01T00:00:05.000Z", + }); + yield* Effect.promise(() => harness.drain()); + thread = (yield* Effect.promise(() => harness.readModel())).threads.find( + (entry) => entry.id === threadId, + )!; + expect(thread.session?.status).toBe("running"); + expect(thread.session?.activeTurnId).toBe(asTurnId("turn-after-reconnect")); + }), + ); + + effectIt.effect("keeps an aborted pending start stopped across duplicate exit events", () => + Effect.gen(function* () { + const harness = yield* Effect.promise(() => createHarness()); + const threadId = asThreadId("thread-1"); + const stoppedAt = "2026-01-01T00:00:02.000Z"; + + yield* harness.engine.dispatch({ + type: "thread.turn.start", + commandId: CommandId.make("cmd-turn-start-before-stop"), + threadId, + message: { + messageId: MessageId.make("message-before-stop"), + role: "user", + text: "stop this startup", + attachments: [], + }, + interactionMode: DEFAULT_PROVIDER_INTERACTION_MODE, + runtimeMode: "approval-required", + createdAt: "2026-01-01T00:00:01.000Z", + }); + yield* harness.engine.dispatch({ + type: "thread.session.set", + commandId: CommandId.make("cmd-session-starting-before-stop"), + threadId, + session: { + threadId, + status: "starting", + providerName: "codex", + runtimeMode: "approval-required", + activeTurnId: null, + lastError: null, + updatedAt: "2026-01-01T00:00:01.000Z", + }, + createdAt: "2026-01-01T00:00:01.000Z", + }); + yield* harness.engine.dispatch({ + type: "thread.session.set", + commandId: CommandId.make("cmd-session-stop-pending-start"), + threadId, + session: { + threadId, + status: "stopped", + providerName: "codex", + runtimeMode: "approval-required", + activeTurnId: null, + lastError: null, + updatedAt: stoppedAt, + }, + createdAt: stoppedAt, + }); + + harness.emit({ + type: "session.exited", + eventId: asEventId("evt-session-exited-after-stop"), + provider: ProviderDriverKind.make("codex"), + threadId, + createdAt: "2026-01-01T00:00:03.000Z", + }); + harness.emit({ + type: "session.exited", + eventId: asEventId("evt-duplicate-session-exited-after-stop"), + provider: ProviderDriverKind.make("codex"), + threadId, + createdAt: "2026-01-01T00:00:04.000Z", + }); + + yield* Effect.promise(() => harness.drain()); + const thread = (yield* Effect.promise(() => harness.readModel())).threads.find( + (entry) => entry.id === threadId, + ); + expect(thread?.session?.status).toBe("stopped"); + expect(thread?.session?.activeTurnId).toBeNull(); + }), + ); + + it("does not clear active turn when session/thread started arrives mid-turn", async () => { + const harness = await createHarness(); + const now = "2026-01-01T00:00:00.000Z"; + + harness.emit({ + type: "turn.started", + eventId: asEventId("evt-turn-started-midturn-lifecycle"), + provider: ProviderDriverKind.make("codex"), + createdAt: now, + threadId: asThreadId("thread-1"), + turnId: asTurnId("turn-midturn-lifecycle"), + }); + + await waitForThread( + harness.readModel, + (thread) => + thread.session?.status === "running" && + thread.session?.activeTurnId === "turn-midturn-lifecycle", + 10_000, + ); + + harness.emit({ + type: "thread.started", + eventId: asEventId("evt-thread-started-midturn-lifecycle"), + provider: ProviderDriverKind.make("codex"), + createdAt: "2026-01-01T00:00:00.000Z", + threadId: asThreadId("thread-1"), + }); + harness.emit({ + type: "session.started", + eventId: asEventId("evt-session-started-midturn-lifecycle"), + provider: ProviderDriverKind.make("codex"), + createdAt: "2026-01-01T00:00:00.000Z", + threadId: asThreadId("thread-1"), + }); + + await harness.drain(); + const midReadModel = await harness.readModel(); + const midThread = midReadModel.threads.find((entry) => entry.id === ThreadId.make("thread-1")); + expect(midThread?.session?.status).toBe("running"); + expect(midThread?.session?.activeTurnId).toBe("turn-midturn-lifecycle"); + + harness.emit({ + type: "turn.completed", + eventId: asEventId("evt-turn-completed-midturn-lifecycle"), + provider: ProviderDriverKind.make("codex"), + createdAt: "2026-01-01T00:00:00.000Z", + threadId: asThreadId("thread-1"), + turnId: asTurnId("turn-midturn-lifecycle"), + status: "completed", + }); + + await waitForThread( + harness.readModel, + (thread) => thread.session?.status === "ready" && thread.session?.activeTurnId === null, + 10_000, + ); + }); + + it("accepts claude turn lifecycle when seeded thread id is a synthetic placeholder", async () => { + const harness = await createHarness(); + const seededAt = "2026-01-01T00:00:00.000Z"; + + await Effect.runPromise( + harness.engine.dispatch({ + type: "thread.session.set", + commandId: CommandId.make("cmd-session-seed-claude-placeholder"), + threadId: ThreadId.make("thread-1"), + session: { + threadId: ThreadId.make("thread-1"), + status: "ready", + providerName: "claudeAgent", + runtimeMode: "approval-required", + activeTurnId: null, + updatedAt: seededAt, + lastError: null, + }, + createdAt: seededAt, + }), + ); + + harness.emit({ + type: "turn.started", + eventId: asEventId("evt-turn-started-claude-placeholder"), + provider: ProviderDriverKind.make("claudeAgent"), + createdAt: "2026-01-01T00:00:00.000Z", + threadId: asThreadId("thread-1"), + turnId: asTurnId("turn-claude-placeholder"), + }); + + await waitForThread( + harness.readModel, + (thread) => + thread.session?.status === "running" && + thread.session?.activeTurnId === "turn-claude-placeholder", + ); + + harness.emit({ + type: "turn.completed", + eventId: asEventId("evt-turn-completed-claude-placeholder"), + provider: ProviderDriverKind.make("claudeAgent"), + createdAt: "2026-01-01T00:00:00.000Z", + threadId: asThreadId("thread-1"), + turnId: asTurnId("turn-claude-placeholder"), + status: "completed", + }); + + await waitForThread( + harness.readModel, + (thread) => thread.session?.status === "ready" && thread.session?.activeTurnId === null, + ); + }); + + it("ignores auxiliary turn completions from a different provider thread", async () => { + const harness = await createHarness(); + const now = "2026-01-01T00:00:00.000Z"; + + harness.emit({ + type: "turn.started", + eventId: asEventId("evt-turn-started-primary"), + provider: ProviderDriverKind.make("codex"), + createdAt: now, + threadId: asThreadId("thread-1"), + turnId: asTurnId("turn-primary"), + }); + + await waitForThread( + harness.readModel, + (thread) => + thread.session?.status === "running" && thread.session?.activeTurnId === "turn-primary", + ); + + harness.emit({ + type: "turn.completed", + eventId: asEventId("evt-turn-completed-aux"), + provider: ProviderDriverKind.make("codex"), + createdAt: "2026-01-01T00:00:00.000Z", + threadId: asThreadId("thread-1"), + turnId: asTurnId("turn-aux"), + status: "completed", + }); + + await harness.drain(); + const midReadModel = await harness.readModel(); + const midThread = midReadModel.threads.find((entry) => entry.id === ThreadId.make("thread-1")); + expect(midThread?.session?.status).toBe("running"); + expect(midThread?.session?.activeTurnId).toBe("turn-primary"); + + harness.emit({ + type: "turn.completed", + eventId: asEventId("evt-turn-completed-primary"), + provider: ProviderDriverKind.make("codex"), + createdAt: "2026-01-01T00:00:00.000Z", + threadId: asThreadId("thread-1"), + turnId: asTurnId("turn-primary"), + status: "completed", + }); + + await waitForThread( + harness.readModel, + (thread) => thread.session?.status === "ready" && thread.session?.activeTurnId === null, + ); + }); + + it("rejects an untargeted turn.completed when no turn is active", async () => { + const harness = await createHarness(); + const seededAt = "2026-01-01T00:00:00.000Z"; + + // A turn start is pending: the session reads "starting" with no active + // turn tracked yet. This is the window the Claude resume handshake's + // phantom (turn.completed with no turnId) used to slip through, stomping + // "starting" back to "ready" for a turn that never existed. + await harness.dispatch({ + type: "thread.session.set", + commandId: CommandId.make("cmd-session-seed-untargeted-completion"), + threadId: ThreadId.make("thread-1"), + session: { + threadId: ThreadId.make("thread-1"), + status: "starting", + providerName: "claudeAgent", + runtimeMode: "approval-required", + activeTurnId: null, + updatedAt: seededAt, + lastError: null, + }, + createdAt: seededAt, + }); + + harness.emit({ + type: "turn.completed", + eventId: asEventId("evt-turn-completed-untargeted"), + provider: ProviderDriverKind.make("claudeAgent"), + createdAt: seededAt, + threadId: asThreadId("thread-1"), + status: "completed", + }); + + await harness.drain(); + const readModel = await harness.readModel(); + const thread = readModel.threads.find((entry) => entry.id === ThreadId.make("thread-1")); + expect(thread?.session?.status).toBe("starting"); + expect(thread?.session?.activeTurnId).toBeNull(); + }); + + it("accepts a targeted turn.completed when no turn is active", async () => { + const harness = await createHarness(); + const seededAt = "2026-01-01T00:00:00.000Z"; + + // A completion that names its turn still lands even when no active turn + // is tracked (e.g. its turn.started was lost). Only untargeted + // completions are rejected. + await harness.dispatch({ + type: "thread.session.set", + commandId: CommandId.make("cmd-session-seed-targeted-completion"), + threadId: ThreadId.make("thread-1"), + session: { + threadId: ThreadId.make("thread-1"), + status: "starting", + providerName: "claudeAgent", + runtimeMode: "approval-required", + activeTurnId: null, + updatedAt: seededAt, + lastError: null, + }, + createdAt: seededAt, + }); + + harness.emit({ + type: "turn.completed", + eventId: asEventId("evt-turn-completed-targeted-late"), + provider: ProviderDriverKind.make("claudeAgent"), + createdAt: seededAt, + threadId: asThreadId("thread-1"), + turnId: asTurnId("turn-late"), + status: "completed", + }); + + await waitForThread(harness.readModel, (thread) => thread.session?.status === "ready"); + }); + + it("ignores non-active turn completion when runtime omits thread id", async () => { + const harness = await createHarness(); + const now = "2026-01-01T00:00:00.000Z"; + + harness.emit({ + type: "turn.started", + eventId: asEventId("evt-turn-started-guarded"), + provider: ProviderDriverKind.make("codex"), + createdAt: now, + threadId: asThreadId("thread-1"), + turnId: asTurnId("turn-guarded-main"), + }); + + await waitForThread( + harness.readModel, + (thread) => + thread.session?.status === "running" && + thread.session?.activeTurnId === "turn-guarded-main", + ); + + harness.emit({ + type: "turn.completed", + eventId: asEventId("evt-turn-completed-guarded-other"), + provider: ProviderDriverKind.make("codex"), + createdAt: "2026-01-01T00:00:00.000Z", + threadId: asThreadId("thread-1"), + turnId: asTurnId("turn-guarded-other"), + status: "completed", + }); + + await harness.drain(); + const midReadModel = await harness.readModel(); + const midThread = midReadModel.threads.find((entry) => entry.id === ThreadId.make("thread-1")); + expect(midThread?.session?.status).toBe("running"); + expect(midThread?.session?.activeTurnId).toBe("turn-guarded-main"); + + harness.emit({ + type: "turn.completed", + eventId: asEventId("evt-turn-completed-guarded-main"), + provider: ProviderDriverKind.make("codex"), + createdAt: "2026-01-01T00:00:00.000Z", + threadId: asThreadId("thread-1"), + turnId: asTurnId("turn-guarded-main"), + status: "completed", + }); + + await waitForThread( + harness.readModel, + (thread) => thread.session?.status === "ready" && thread.session?.activeTurnId === null, + ); + }); + + it("ignores provider content deltas that cannot change thread state", async () => { + const harness = await createHarness(); + const initial = await harness.readModel(); + + for (const streamKind of ["reasoning_text", "command_output", "file_change_output"] as const) { + harness.emit({ + type: "content.delta", + eventId: asEventId(`evt-ignored-${streamKind}`), + provider: ProviderDriverKind.make("codex"), + createdAt: "2026-01-01T00:00:00.000Z", + threadId: asThreadId("thread-1"), + turnId: asTurnId("turn-ignored"), + payload: { + streamKind, + delta: "ignored output", + }, + }); + } + + await harness.drain(); + expect(await harness.readModel()).toEqual(initial); + }); + + it("ignores an aborted event for a different active turn", async () => { + const harness = await createHarness(); + const now = "2026-01-01T00:00:00.000Z"; + const activeTurnId = asTurnId("turn-abort-guarded-main"); + + harness.emit({ + type: "turn.started", + eventId: asEventId("evt-turn-abort-guarded-started"), + provider: ProviderDriverKind.make("opencode"), + createdAt: now, + threadId: asThreadId("thread-1"), + turnId: activeTurnId, + }); + + await waitForThread( + harness.readModel, + (thread) => + thread.session?.status === "running" && thread.session?.activeTurnId === activeTurnId, + ); + + harness.emit({ + type: "turn.aborted", + eventId: asEventId("evt-turn-abort-guarded-stale"), + provider: ProviderDriverKind.make("opencode"), + createdAt: "2026-01-01T00:00:01.000Z", + threadId: asThreadId("thread-1"), + turnId: asTurnId("turn-abort-guarded-stale"), + payload: { + reason: "Interrupted by user.", + }, + }); + + await harness.drain(); + const readModel = await harness.readModel(); + const thread = readModel.threads.find((entry) => entry.id === asThreadId("thread-1")); + expect(thread?.session?.status).toBe("running"); + expect(thread?.session?.activeTurnId).toBe(activeTurnId); + }); - harness.emit({ - type: "session.started", - eventId: asEventId("evt-session-started-pending-reconnect"), - provider: ProviderDriverKind.make("codex"), - threadId, - createdAt: "2026-01-01T00:00:03.000Z", - }); - yield* Effect.promise(() => harness.drain()); - thread = (yield* Effect.promise(() => harness.readModel())).threads.find( - (entry) => entry.id === threadId, - )!; - expect(thread.session?.status).toBe("starting"); - expect(thread.session?.activeTurnId).toBeNull(); + it("ignores an untargeted aborted event while a turn is active", async () => { + const harness = await createHarness(); + const now = "2026-01-01T00:00:00.000Z"; + const activeTurnId = asTurnId("turn-abort-untargeted-active"); - harness.emit({ - type: "turn.started", - eventId: asEventId("evt-turn-started-pending-reconnect"), - provider: ProviderDriverKind.make("codex"), - threadId, - turnId: asTurnId("turn-after-reconnect"), - createdAt: "2026-01-01T00:00:04.000Z", - }); - thread = yield* Effect.promise(() => - waitForThread( - harness.readModel, - (entry) => - entry.session?.status === "running" && - entry.session.activeTurnId === asTurnId("turn-after-reconnect"), - ), - ); - expect(thread.session?.status).toBe("running"); + harness.emit({ + type: "turn.started", + eventId: asEventId("evt-turn-abort-untargeted-started"), + provider: ProviderDriverKind.make("opencode"), + createdAt: now, + threadId: asThreadId("thread-1"), + turnId: activeTurnId, + }); - harness.emit({ - type: "session.started", - eventId: asEventId("evt-session-started-duplicate-midturn"), - provider: ProviderDriverKind.make("codex"), - threadId, - createdAt: "2026-01-01T00:00:05.000Z", - }); - yield* Effect.promise(() => harness.drain()); - thread = (yield* Effect.promise(() => harness.readModel())).threads.find( - (entry) => entry.id === threadId, - )!; - expect(thread.session?.status).toBe("running"); - expect(thread.session?.activeTurnId).toBe(asTurnId("turn-after-reconnect")); - }), - ); + await waitForThread( + harness.readModel, + (thread) => + thread.session?.status === "running" && thread.session?.activeTurnId === activeTurnId, + ); - effectIt.effect("keeps an aborted pending start stopped across duplicate exit events", () => - Effect.gen(function* () { - const harness = yield* Effect.promise(() => createHarness()); - const threadId = asThreadId("thread-1"); - const stoppedAt = "2026-01-01T00:00:02.000Z"; + harness.emit({ + type: "turn.aborted", + eventId: asEventId("evt-turn-abort-untargeted"), + provider: ProviderDriverKind.make("opencode"), + createdAt: "2026-01-01T00:00:01.000Z", + threadId: asThreadId("thread-1"), + payload: { + reason: "Interrupted by user.", + }, + }); - yield* harness.engine.dispatch({ - type: "thread.turn.start", - commandId: CommandId.make("cmd-turn-start-before-stop"), - threadId, - message: { - messageId: MessageId.make("message-before-stop"), - role: "user", - text: "stop this startup", - attachments: [], - }, - interactionMode: DEFAULT_PROVIDER_INTERACTION_MODE, + await harness.drain(); + const readModel = await harness.readModel(); + const thread = readModel.threads.find((entry) => entry.id === asThreadId("thread-1")); + expect(thread?.session?.status).toBe("running"); + expect(thread?.session?.activeTurnId).toBe(activeTurnId); + }); + + it("rejects an untargeted aborted event while a turn start is pending", async () => { + const harness = await createHarness(); + const seededAt = "2026-01-01T00:00:00.000Z"; + + await harness.dispatch({ + type: "thread.session.set", + commandId: CommandId.make("cmd-session-seed-untargeted-abort"), + threadId: ThreadId.make("thread-1"), + session: { + threadId: ThreadId.make("thread-1"), + status: "starting", + providerName: "opencode", runtimeMode: "approval-required", - createdAt: "2026-01-01T00:00:01.000Z", - }); - yield* harness.engine.dispatch({ - type: "thread.session.set", - commandId: CommandId.make("cmd-session-starting-before-stop"), - threadId, - session: { - threadId, - status: "starting", - providerName: "codex", - runtimeMode: "approval-required", - activeTurnId: null, - lastError: null, - updatedAt: "2026-01-01T00:00:01.000Z", - }, - createdAt: "2026-01-01T00:00:01.000Z", - }); - yield* harness.engine.dispatch({ - type: "thread.session.set", - commandId: CommandId.make("cmd-session-stop-pending-start"), - threadId, - session: { - threadId, - status: "stopped", - providerName: "codex", - runtimeMode: "approval-required", - activeTurnId: null, - lastError: null, - updatedAt: stoppedAt, - }, - createdAt: stoppedAt, - }); + activeTurnId: null, + updatedAt: seededAt, + lastError: null, + }, + createdAt: seededAt, + }); - harness.emit({ - type: "session.exited", - eventId: asEventId("evt-session-exited-after-stop"), - provider: ProviderDriverKind.make("codex"), - threadId, - createdAt: "2026-01-01T00:00:03.000Z", - }); - harness.emit({ - type: "session.exited", - eventId: asEventId("evt-duplicate-session-exited-after-stop"), - provider: ProviderDriverKind.make("codex"), + harness.emit({ + type: "turn.aborted", + eventId: asEventId("evt-turn-abort-untargeted-pending"), + provider: ProviderDriverKind.make("opencode"), + createdAt: seededAt, + threadId: asThreadId("thread-1"), + payload: { + reason: "Interrupted by user.", + }, + }); + + await harness.drain(); + const readModel = await harness.readModel(); + const thread = readModel.threads.find((entry) => entry.id === asThreadId("thread-1")); + expect(thread?.session?.status).toBe("starting"); + expect(thread?.session?.activeTurnId).toBeNull(); + }); + + it("rejects a named stale abort while a newer turn start is pending", async () => { + const harness = await createHarness(); + const threadId = asThreadId("thread-1"); + const pendingTurnId = asTurnId("turn-pending-abort"); + const staleTurnId = asTurnId("turn-stale-abort"); + const createdAt = "2026-01-01T00:00:00.000Z"; + const staleAbortAt = "2025-12-31T23:59:59.000Z"; + + harness.setProviderSession({ + provider: ProviderDriverKind.make("opencode"), + status: "ready", + runtimeMode: "approval-required", + threadId, + createdAt: staleAbortAt, + updatedAt: staleAbortAt, + lastAbortedTurnId: staleTurnId, + lastAbortedMessageId: asMessageId("message-old-abort"), + }); + + await harness.dispatch({ + type: "thread.turn.start", + commandId: CommandId.make("cmd-turn-start-pending-abort-guard"), + threadId, + message: { + messageId: asMessageId("message-pending-abort-guard"), + role: "user", + text: "start the pending turn", + attachments: [], + }, + interactionMode: DEFAULT_PROVIDER_INTERACTION_MODE, + runtimeMode: "approval-required", + createdAt, + }); + await harness.dispatch({ + type: "thread.session.set", + commandId: CommandId.make("cmd-session-starting-abort-guard"), + threadId, + session: { threadId, - createdAt: "2026-01-01T00:00:04.000Z", - }); + status: "starting", + providerName: "opencode", + runtimeMode: "approval-required", + activeTurnId: null, + lastError: null, + updatedAt: createdAt, + }, + createdAt, + }); + harness.emit({ + type: "turn.aborted", + eventId: asEventId("evt-turn-abort-stale-pending"), + provider: ProviderDriverKind.make("opencode"), + threadId, + createdAt: "2026-01-01T00:00:01.000Z", + turnId: staleTurnId, + payload: { reason: "Interrupted by user." }, + }); + await harness.drain(); - yield* Effect.promise(() => harness.drain()); - const thread = (yield* Effect.promise(() => harness.readModel())).threads.find( - (entry) => entry.id === threadId, - ); - expect(thread?.session?.status).toBe("stopped"); - expect(thread?.session?.activeTurnId).toBeNull(); - }), - ); + let thread = (await harness.readModel()).threads.find((entry) => entry.id === threadId); + expect(thread?.session?.status).toBe("starting"); + expect(thread?.session?.activeTurnId).toBeNull(); - it("does not clear active turn when session/thread started arrives mid-turn", async () => { + // OpenCode clears activeTurnId before its queued turn.aborted event is + // ingested. The terminal ID keeps a legitimate abort correlated to the + // pending durable turn without reopening the stale-abort race. + harness.setProviderSession({ + provider: ProviderDriverKind.make("opencode"), + status: "ready", + runtimeMode: "approval-required", + threadId, + createdAt, + updatedAt: "2026-01-01T00:00:02.000Z", + lastAbortedTurnId: pendingTurnId, + lastAbortedMessageId: asMessageId("message-pending-abort-guard"), + }); + harness.emit({ + type: "turn.aborted", + eventId: asEventId("evt-turn-abort-pending"), + provider: ProviderDriverKind.make("opencode"), + threadId, + createdAt: "2026-01-01T00:00:01.000Z", + turnId: pendingTurnId, + payload: { reason: "Interrupted by user." }, + }); + await harness.drain(); + + thread = (await harness.readModel()).threads.find((entry) => entry.id === threadId); + expect(thread?.session?.status).toBe("interrupted"); + expect(thread?.session?.activeTurnId).toBeNull(); + }); + + it("requires the pending start token for active provider aborts", async () => { const harness = await createHarness(); - const now = "2026-01-01T00:00:00.000Z"; + const threadId = asThreadId("thread-1"); + const pendingTurnId = asTurnId("turn-pending-active-token"); + const oldTurnId = asTurnId("turn-old-active-token"); + const pendingMessageId = asMessageId("message-pending-active-token"); + const oldMessageId = asMessageId("message-old-active-token"); + const createdAt = "2026-01-01T00:00:00.000Z"; + + await harness.dispatch({ + type: "thread.turn.start", + commandId: CommandId.make("cmd-turn-start-active-token-guard"), + threadId, + message: { + messageId: pendingMessageId, + role: "user", + text: "start the pending turn", + attachments: [], + }, + interactionMode: DEFAULT_PROVIDER_INTERACTION_MODE, + runtimeMode: "approval-required", + createdAt, + }); + await harness.dispatch({ + type: "thread.session.set", + commandId: CommandId.make("cmd-session-starting-active-token-guard"), + threadId, + session: { + threadId, + status: "starting", + providerName: "opencode", + runtimeMode: "approval-required", + activeTurnId: null, + lastError: null, + updatedAt: createdAt, + }, + createdAt, + }); + harness.setProviderSession({ + provider: ProviderDriverKind.make("opencode"), + status: "running", + runtimeMode: "approval-required", + threadId, + createdAt, + updatedAt: createdAt, + activeTurnId: oldTurnId, + activeTurnStartMessageId: oldMessageId, + }); harness.emit({ - type: "turn.started", - eventId: asEventId("evt-turn-started-midturn-lifecycle"), - provider: ProviderDriverKind.make("codex"), - createdAt: now, - threadId: asThreadId("thread-1"), - turnId: asTurnId("turn-midturn-lifecycle"), + type: "turn.aborted", + eventId: asEventId("evt-turn-abort-active-old-token"), + provider: ProviderDriverKind.make("opencode"), + threadId, + createdAt: "2026-01-01T00:00:01.000Z", + turnId: oldTurnId, + payload: { reason: "Interrupted by user." }, }); + await harness.drain(); - await waitForThread( - harness.readModel, - (thread) => - thread.session?.status === "running" && - thread.session?.activeTurnId === "turn-midturn-lifecycle", - 10_000, - ); + let thread = (await harness.readModel()).threads.find((entry) => entry.id === threadId); + expect(thread?.session?.status).toBe("starting"); + expect(thread?.session?.activeTurnId).toBeNull(); - harness.emit({ - type: "thread.started", - eventId: asEventId("evt-thread-started-midturn-lifecycle"), - provider: ProviderDriverKind.make("codex"), - createdAt: "2026-01-01T00:00:00.000Z", - threadId: asThreadId("thread-1"), + harness.setProviderSession({ + provider: ProviderDriverKind.make("opencode"), + status: "running", + runtimeMode: "approval-required", + threadId, + createdAt, + updatedAt: createdAt, + activeTurnId: pendingTurnId, }); harness.emit({ - type: "session.started", - eventId: asEventId("evt-session-started-midturn-lifecycle"), - provider: ProviderDriverKind.make("codex"), - createdAt: "2026-01-01T00:00:00.000Z", - threadId: asThreadId("thread-1"), + type: "turn.aborted", + eventId: asEventId("evt-turn-abort-active-missing-token"), + provider: ProviderDriverKind.make("opencode"), + threadId, + createdAt: "2026-01-01T00:00:02.000Z", + turnId: pendingTurnId, + payload: { reason: "Interrupted by user." }, }); - await harness.drain(); - const midReadModel = await harness.readModel(); - const midThread = midReadModel.threads.find((entry) => entry.id === ThreadId.make("thread-1")); - expect(midThread?.session?.status).toBe("running"); - expect(midThread?.session?.activeTurnId).toBe("turn-midturn-lifecycle"); + thread = (await harness.readModel()).threads.find((entry) => entry.id === threadId); + expect(thread?.session?.status).toBe("starting"); + expect(thread?.session?.activeTurnId).toBeNull(); + + harness.setProviderSession({ + provider: ProviderDriverKind.make("opencode"), + status: "running", + runtimeMode: "approval-required", + threadId, + createdAt, + updatedAt: createdAt, + activeTurnId: pendingTurnId, + activeTurnStartMessageId: pendingMessageId, + }); harness.emit({ - type: "turn.completed", - eventId: asEventId("evt-turn-completed-midturn-lifecycle"), - provider: ProviderDriverKind.make("codex"), - createdAt: "2026-01-01T00:00:00.000Z", - threadId: asThreadId("thread-1"), - turnId: asTurnId("turn-midturn-lifecycle"), - status: "completed", + type: "turn.aborted", + eventId: asEventId("evt-turn-abort-active-current-token"), + provider: ProviderDriverKind.make("opencode"), + threadId, + createdAt: "2026-01-01T00:00:03.000Z", + turnId: pendingTurnId, + payload: { reason: "Interrupted by user." }, }); + await harness.drain(); - await waitForThread( - harness.readModel, - (thread) => thread.session?.status === "ready" && thread.session?.activeTurnId === null, - 10_000, - ); + thread = (await harness.readModel()).threads.find((entry) => entry.id === threadId); + expect(thread?.session?.status).toBe("interrupted"); + expect(thread?.session?.activeTurnId).toBeNull(); }); - it("accepts claude turn lifecycle when seeded thread id is a synthetic placeholder", async () => { + it("rejects a targeted abort that does not match the pending provider turn", async () => { const harness = await createHarness(); - const seededAt = "2026-01-01T00:00:00.000Z"; + const threadId = asThreadId("thread-1"); + const liveTurnId = asTurnId("turn-pending-provider-live"); + const staleTurnId = asTurnId("turn-pending-provider-stale"); + const pendingMessageId = asMessageId("message-pending-provider-match"); + const liveMessageId = asMessageId("message-live-provider-turn"); + const createdAt = "2026-01-01T00:00:00.000Z"; - await Effect.runPromise( - harness.engine.dispatch({ - type: "thread.session.set", - commandId: CommandId.make("cmd-session-seed-claude-placeholder"), - threadId: ThreadId.make("thread-1"), - session: { - threadId: ThreadId.make("thread-1"), - status: "ready", - providerName: "claudeAgent", - runtimeMode: "approval-required", - activeTurnId: null, - updatedAt: seededAt, - lastError: null, - }, - createdAt: seededAt, - }), - ); + await harness.dispatch({ + type: "thread.turn.start", + commandId: CommandId.make("cmd-turn-start-pending-provider-match"), + threadId, + message: { + messageId: pendingMessageId, + role: "user", + text: "start the pending turn", + attachments: [], + }, + interactionMode: DEFAULT_PROVIDER_INTERACTION_MODE, + runtimeMode: "approval-required", + createdAt, + }); + await harness.dispatch({ + type: "thread.session.set", + commandId: CommandId.make("cmd-session-starting-pending-provider-match"), + threadId, + session: { + threadId, + status: "starting", + providerName: "opencode", + runtimeMode: "approval-required", + activeTurnId: null, + lastError: null, + updatedAt: createdAt, + }, + createdAt, + }); + // The provider is already running a different turn than the stale abort + // names, so the abort must not settle the pending start. + harness.setProviderSession({ + provider: ProviderDriverKind.make("opencode"), + status: "running", + runtimeMode: "approval-required", + threadId, + createdAt, + updatedAt: createdAt, + activeTurnId: liveTurnId, + activeTurnStartMessageId: liveMessageId, + }); harness.emit({ - type: "turn.started", - eventId: asEventId("evt-turn-started-claude-placeholder"), - provider: ProviderDriverKind.make("claudeAgent"), - createdAt: "2026-01-01T00:00:00.000Z", - threadId: asThreadId("thread-1"), - turnId: asTurnId("turn-claude-placeholder"), + type: "turn.aborted", + eventId: asEventId("evt-turn-abort-pending-provider-mismatch"), + provider: ProviderDriverKind.make("opencode"), + threadId, + createdAt: "2026-01-01T00:00:01.000Z", + turnId: staleTurnId, + payload: { reason: "Interrupted by user." }, }); + await harness.drain(); - await waitForThread( - harness.readModel, - (thread) => - thread.session?.status === "running" && - thread.session?.activeTurnId === "turn-claude-placeholder", - ); + let thread = (await harness.readModel()).threads.find((entry) => entry.id === threadId); + expect(thread?.session?.status).toBe("starting"); + expect(thread?.session?.activeTurnId).toBeNull(); + // The same pending start is settled once the abort names the provider's + // pending turn with the pending message token. + harness.setProviderSession({ + provider: ProviderDriverKind.make("opencode"), + status: "running", + runtimeMode: "approval-required", + threadId, + createdAt, + updatedAt: "2026-01-01T00:00:02.000Z", + activeTurnId: liveTurnId, + activeTurnStartMessageId: pendingMessageId, + }); harness.emit({ - type: "turn.completed", - eventId: asEventId("evt-turn-completed-claude-placeholder"), - provider: ProviderDriverKind.make("claudeAgent"), - createdAt: "2026-01-01T00:00:00.000Z", - threadId: asThreadId("thread-1"), - turnId: asTurnId("turn-claude-placeholder"), - status: "completed", + type: "turn.aborted", + eventId: asEventId("evt-turn-abort-pending-provider-match"), + provider: ProviderDriverKind.make("opencode"), + threadId, + createdAt: "2026-01-01T00:00:03.000Z", + turnId: liveTurnId, + payload: { reason: "Interrupted by user." }, }); + await harness.drain(); - await waitForThread( - harness.readModel, - (thread) => thread.session?.status === "ready" && thread.session?.activeTurnId === null, - ); + thread = (await harness.readModel()).threads.find((entry) => entry.id === threadId); + expect(thread?.session?.status).toBe("interrupted"); + expect(thread?.session?.activeTurnId).toBeNull(); }); - it("ignores auxiliary turn completions from a different provider thread", async () => { + it("accepts an OpenCode abort after a steer with the pending message token", async () => { const harness = await createHarness(); + const threadId = asThreadId("thread-1"); + const turnId = asTurnId("turn-steer-abort"); + const originalMessageId = asMessageId("message-original-steer-abort"); + const steeringMessageId = asMessageId("message-steering-abort"); const now = "2026-01-01T00:00:00.000Z"; - harness.emit({ - type: "turn.started", - eventId: asEventId("evt-turn-started-primary"), - provider: ProviderDriverKind.make("codex"), + harness.setProviderSession({ + provider: ProviderDriverKind.make("opencode"), + status: "running", + runtimeMode: "approval-required", + threadId, createdAt: now, - threadId: asThreadId("thread-1"), - turnId: asTurnId("turn-primary"), - }); - - await waitForThread( - harness.readModel, - (thread) => - thread.session?.status === "running" && thread.session?.activeTurnId === "turn-primary", - ); - - harness.emit({ - type: "turn.completed", - eventId: asEventId("evt-turn-completed-aux"), - provider: ProviderDriverKind.make("codex"), - createdAt: "2026-01-01T00:00:00.000Z", - threadId: asThreadId("thread-1"), - turnId: asTurnId("turn-aux"), - status: "completed", + updatedAt: now, + activeTurnId: turnId, + activeTurnStartMessageId: originalMessageId, }); - - await harness.drain(); - const midReadModel = await harness.readModel(); - const midThread = midReadModel.threads.find((entry) => entry.id === ThreadId.make("thread-1")); - expect(midThread?.session?.status).toBe("running"); - expect(midThread?.session?.activeTurnId).toBe("turn-primary"); - harness.emit({ - type: "turn.completed", - eventId: asEventId("evt-turn-completed-primary"), - provider: ProviderDriverKind.make("codex"), - createdAt: "2026-01-01T00:00:00.000Z", - threadId: asThreadId("thread-1"), - turnId: asTurnId("turn-primary"), - status: "completed", - }); - - await waitForThread( - harness.readModel, - (thread) => thread.session?.status === "ready" && thread.session?.activeTurnId === null, - ); - }); - - it("rejects an untargeted turn.completed when no turn is active", async () => { - const harness = await createHarness(); - const seededAt = "2026-01-01T00:00:00.000Z"; + type: "turn.started", + eventId: asEventId("evt-turn-steer-abort-started"), + provider: ProviderDriverKind.make("opencode"), + threadId, + createdAt: now, + turnId, + }); + await harness.drain(); - // A turn start is pending: the session reads "starting" with no active - // turn tracked yet. This is the window the Claude resume handshake's - // phantom (turn.completed with no turnId) used to slip through, stomping - // "starting" back to "ready" for a turn that never existed. await harness.dispatch({ - type: "thread.session.set", - commandId: CommandId.make("cmd-session-seed-untargeted-completion"), - threadId: ThreadId.make("thread-1"), - session: { - threadId: ThreadId.make("thread-1"), - status: "starting", - providerName: "claudeAgent", - runtimeMode: "approval-required", - activeTurnId: null, - updatedAt: seededAt, - lastError: null, + type: "thread.turn.start", + commandId: CommandId.make("cmd-turn-start-steer-abort"), + threadId, + message: { + messageId: steeringMessageId, + role: "user", + text: "actually do the other thing", + attachments: [], }, - createdAt: seededAt, + interactionMode: DEFAULT_PROVIDER_INTERACTION_MODE, + runtimeMode: "approval-required", + createdAt: "2026-01-01T00:00:01.000Z", + }); + harness.setProviderSession({ + provider: ProviderDriverKind.make("opencode"), + status: "ready", + runtimeMode: "approval-required", + threadId, + createdAt: now, + updatedAt: "2026-01-01T00:00:01.000Z", + lastAbortedTurnId: turnId, + lastAbortedMessageId: steeringMessageId, }); harness.emit({ - type: "turn.completed", - eventId: asEventId("evt-turn-completed-untargeted"), - provider: ProviderDriverKind.make("claudeAgent"), - createdAt: seededAt, - threadId: asThreadId("thread-1"), - status: "completed", + type: "turn.aborted", + eventId: asEventId("evt-turn-steer-abort-aborted"), + provider: ProviderDriverKind.make("opencode"), + threadId, + createdAt: "2026-01-01T00:00:02.000Z", + turnId, + payload: { reason: "Interrupted by user." }, }); - await harness.drain(); - const readModel = await harness.readModel(); - const thread = readModel.threads.find((entry) => entry.id === ThreadId.make("thread-1")); - expect(thread?.session?.status).toBe("starting"); + + const thread = (await harness.readModel()).threads.find((entry) => entry.id === threadId); + expect(thread?.session?.status).toBe("interrupted"); expect(thread?.session?.activeTurnId).toBeNull(); + expect(thread?.latestTurn?.state).toBe("interrupted"); }); - it("accepts a targeted turn.completed when no turn is active", async () => { + it("rejects an old abort when a newer provider turn is active", async () => { const harness = await createHarness(); - const seededAt = "2026-01-01T00:00:00.000Z"; + const threadId = asThreadId("thread-1"); + const oldTurnId = asTurnId("turn-old-before-new-bind"); + const newTurnId = asTurnId("turn-new-bound"); + const pendingMessageId = asMessageId("message-new-bound"); + const newMessageId = asMessageId("message-new-bound-provider"); + const now = "2026-01-01T00:00:00.000Z"; + + harness.setProviderSession({ + provider: ProviderDriverKind.make("opencode"), + status: "running", + runtimeMode: "approval-required", + threadId, + createdAt: now, + updatedAt: now, + activeTurnId: oldTurnId, + }); + harness.emit({ + type: "turn.started", + eventId: asEventId("evt-turn-old-before-new-bind-started"), + provider: ProviderDriverKind.make("opencode"), + threadId, + createdAt: now, + turnId: oldTurnId, + }); + await harness.drain(); - // A completion that names its turn still lands even when no active turn - // is tracked (e.g. its turn.started was lost). Only untargeted - // completions are rejected. await harness.dispatch({ - type: "thread.session.set", - commandId: CommandId.make("cmd-session-seed-targeted-completion"), - threadId: ThreadId.make("thread-1"), - session: { - threadId: ThreadId.make("thread-1"), - status: "starting", - providerName: "claudeAgent", - runtimeMode: "approval-required", - activeTurnId: null, - updatedAt: seededAt, - lastError: null, + type: "thread.turn.start", + commandId: CommandId.make("cmd-turn-start-new-bound"), + threadId, + message: { + messageId: pendingMessageId, + role: "user", + text: "start the newer turn", + attachments: [], }, - createdAt: seededAt, + interactionMode: DEFAULT_PROVIDER_INTERACTION_MODE, + runtimeMode: "approval-required", + createdAt: "2026-01-01T00:00:01.000Z", + }); + harness.setProviderSession({ + provider: ProviderDriverKind.make("opencode"), + status: "running", + runtimeMode: "approval-required", + threadId, + createdAt: now, + updatedAt: "2026-01-01T00:00:01.000Z", + activeTurnId: newTurnId, + activeTurnStartMessageId: newMessageId, }); harness.emit({ - type: "turn.completed", - eventId: asEventId("evt-turn-completed-targeted-late"), - provider: ProviderDriverKind.make("claudeAgent"), - createdAt: seededAt, - threadId: asThreadId("thread-1"), - turnId: asTurnId("turn-late"), - status: "completed", + type: "turn.aborted", + eventId: asEventId("evt-turn-old-before-new-bind-aborted"), + provider: ProviderDriverKind.make("opencode"), + threadId, + createdAt: "2026-01-01T00:00:02.000Z", + turnId: oldTurnId, + payload: { reason: "Interrupted by user." }, }); + await harness.drain(); - await waitForThread(harness.readModel, (thread) => thread.session?.status === "ready"); + const thread = (await harness.readModel()).threads.find((entry) => entry.id === threadId); + expect(thread?.session?.status).toBe("running"); + expect(thread?.session?.activeTurnId).toBe(oldTurnId); + expect( + thread?.messages.find( + (message: ProviderRuntimeTestMessage) => message.id === pendingMessageId, + )?.text, + ).toBe("start the newer turn"); }); - it("ignores non-active turn completion when runtime omits thread id", async () => { + it("preserves a newer pending start when a clear-before-emit abort is stale", async () => { const harness = await createHarness(); + const threadId = asThreadId("thread-1"); + const oldTurnId = asTurnId("turn-clear-before-emit-old"); + const oldMessageId = asMessageId("message-clear-before-emit-old"); + const sourceTurnId = asTurnId("turn-clear-before-emit-source-plan"); + const pendingTurnId = asTurnId("turn-clear-before-emit-new"); + const pendingMessageId = asMessageId("message-clear-before-emit-new"); const now = "2026-01-01T00:00:00.000Z"; harness.emit({ - type: "turn.started", - eventId: asEventId("evt-turn-started-guarded"), + type: "turn.proposed.completed", + eventId: asEventId("evt-clear-before-emit-source-plan"), provider: ProviderDriverKind.make("codex"), + threadId, createdAt: now, - threadId: asThreadId("thread-1"), - turnId: asTurnId("turn-guarded-main"), + turnId: sourceTurnId, + payload: { planMarkdown: "# Source plan" }, }); - - await waitForThread( + const threadWithSourcePlan = await waitForThread( harness.readModel, (thread) => - thread.session?.status === "running" && - thread.session?.activeTurnId === "turn-guarded-main", + thread.proposedPlans.some( + (proposedPlan: ProviderRuntimeTestProposedPlan) => + proposedPlan.id === `plan:${threadId}:turn:${sourceTurnId}` && + proposedPlan.implementedAt === null, + ), + 2_000, + threadId, + ); + const sourcePlan = threadWithSourcePlan.proposedPlans.find( + (entry: ProviderRuntimeTestProposedPlan) => + entry.id === `plan:${threadId}:turn:${sourceTurnId}`, ); + expect(sourcePlan).toBeDefined(); + if (!sourcePlan) { + throw new Error("Expected source plan to exist."); + } - harness.emit({ - type: "turn.completed", - eventId: asEventId("evt-turn-completed-guarded-other"), - provider: ProviderDriverKind.make("codex"), - createdAt: "2026-01-01T00:00:00.000Z", - threadId: asThreadId("thread-1"), - turnId: asTurnId("turn-guarded-other"), - status: "completed", + await harness.dispatch({ + type: "thread.session.set", + commandId: CommandId.make("cmd-session-running-clear-before-emit"), + threadId, + session: { + threadId, + status: "running", + providerName: "opencode", + runtimeMode: "approval-required", + activeTurnId: oldTurnId, + lastError: null, + updatedAt: now, + }, + createdAt: now, + }); + await harness.dispatch({ + type: "thread.turn.start", + commandId: CommandId.make("cmd-turn-start-clear-before-emit"), + threadId, + message: { + messageId: pendingMessageId, + role: "user", + text: "start the newer turn", + attachments: [], + }, + sourceProposedPlan: { + threadId, + planId: sourcePlan.id, + }, + interactionMode: DEFAULT_PROVIDER_INTERACTION_MODE, + runtimeMode: "approval-required", + createdAt: "2026-01-01T00:00:01.000Z", + }); + harness.setProviderSession({ + provider: ProviderDriverKind.make("opencode"), + status: "ready", + runtimeMode: "approval-required", + threadId, + createdAt: now, + updatedAt: "2026-01-01T00:00:01.000Z", + lastAbortedTurnId: oldTurnId, + lastAbortedMessageId: oldMessageId, }); - - await harness.drain(); - const midReadModel = await harness.readModel(); - const midThread = midReadModel.threads.find((entry) => entry.id === ThreadId.make("thread-1")); - expect(midThread?.session?.status).toBe("running"); - expect(midThread?.session?.activeTurnId).toBe("turn-guarded-main"); harness.emit({ - type: "turn.completed", - eventId: asEventId("evt-turn-completed-guarded-main"), - provider: ProviderDriverKind.make("codex"), - createdAt: "2026-01-01T00:00:00.000Z", - threadId: asThreadId("thread-1"), - turnId: asTurnId("turn-guarded-main"), - status: "completed", + type: "turn.aborted", + eventId: asEventId("evt-clear-before-emit-stale-abort"), + provider: ProviderDriverKind.make("opencode"), + threadId, + createdAt: "2026-01-01T00:00:02.000Z", + turnId: oldTurnId, + payload: { reason: "Interrupted by user." }, }); + await harness.drain(); - await waitForThread( - harness.readModel, - (thread) => thread.session?.status === "ready" && thread.session?.activeTurnId === null, - ); - }); - - it("ignores provider content deltas that cannot change thread state", async () => { - const harness = await createHarness(); - const initial = await harness.readModel(); + const thread = (await harness.readModel()).threads.find((entry) => entry.id === threadId); + expect(thread?.session?.status).toBe("running"); + expect(thread?.session?.activeTurnId).toBe(oldTurnId); - for (const streamKind of ["reasoning_text", "command_output", "file_change_output"] as const) { - harness.emit({ - type: "content.delta", - eventId: asEventId(`evt-ignored-${streamKind}`), - provider: ProviderDriverKind.make("codex"), - createdAt: "2026-01-01T00:00:00.000Z", - threadId: asThreadId("thread-1"), - turnId: asTurnId("turn-ignored"), - payload: { - streamKind, - delta: "ignored output", - }, - }); - } + const pendingMessage = thread?.messages.find( + (message: ProviderRuntimeTestMessage) => message.id === pendingMessageId, + ); + expect(pendingMessage?.text).toBe("start the newer turn"); + expect(thread?.proposedPlans.find((entry) => entry.id === sourcePlan.id)).toMatchObject({ + implementedAt: null, + implementationThreadId: null, + }); + harness.setProviderSession({ + provider: ProviderDriverKind.make("opencode"), + status: "running", + runtimeMode: "approval-required", + threadId, + createdAt: now, + updatedAt: "2026-01-01T00:00:02.000Z", + activeTurnId: pendingTurnId, + activeTurnStartMessageId: pendingMessageId, + }); + harness.emit({ + type: "turn.started", + eventId: asEventId("evt-clear-before-emit-new-started"), + provider: ProviderDriverKind.make("opencode"), + threadId, + createdAt: "2026-01-01T00:00:02.000Z", + turnId: pendingTurnId, + }); await harness.drain(); - expect(await harness.readModel()).toEqual(initial); + + const threadAfterPendingStart = (await harness.readModel()).threads.find( + (entry) => entry.id === threadId, + ); + expect(threadAfterPendingStart?.session?.status).toBe("running"); + expect(threadAfterPendingStart?.session?.activeTurnId).toBe(pendingTurnId); + expect(threadAfterPendingStart?.latestTurn?.sourceProposedPlan).toEqual({ + threadId, + planId: sourcePlan.id, + }); + expect( + threadAfterPendingStart?.proposedPlans.find((entry) => entry.id === sourcePlan.id), + ).toMatchObject({ + implementedAt: "2026-01-01T00:00:02.000Z", + implementationThreadId: threadId, + }); }); it("maps canonical content delta/item completed into finalized assistant messages", async () => { diff --git a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts index 964f60d3a306..0efafcf9a98d 100644 --- a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts +++ b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts @@ -259,6 +259,24 @@ function normalizeRuntimeTurnState( } } +function isOpenCodeUserAbort( + event: ProviderRuntimeEvent, +): event is Extract { + return ( + event.type === "turn.aborted" && + event.provider === "opencode" && + event.payload.reason === "Interrupted by user." + ); +} + +function isOpenCodePromptFailureAbort( + event: ProviderRuntimeEvent, +): event is Extract { + return ( + event.type === "turn.aborted" && event.provider === "opencode" && !isOpenCodeUserAbort(event) + ); +} + function orchestrationSessionStatusFromRuntimeState( state: "starting" | "running" | "waiting" | "ready" | "interrupted" | "stopped" | "error", ): "starting" | "running" | "ready" | "interrupted" | "stopped" | "error" { @@ -1413,13 +1431,22 @@ const make = Effect.gen(function* () { } as const; }); - const getExpectedProviderTurnIdForThread = Effect.fn("getExpectedProviderTurnIdForThread")( - function* (threadId: ThreadId) { - const sessions = yield* providerService.listSessions(); - const session = sessions.find((entry) => entry.threadId === threadId); - return session?.activeTurnId; - }, - ); + const getExpectedProviderTurnForThread = Effect.fn("getExpectedProviderTurnForThread")(function* ( + threadId: ThreadId, + ) { + const sessions = yield* providerService.listSessions(); + const session = sessions.find((entry) => entry.threadId === threadId); + return { + provider: session?.provider, + activeTurnId: session?.activeTurnId, + activeTurnStartMessageId: + session?.activeTurnId === undefined ? undefined : session?.activeTurnStartMessageId, + lastAbortedTurnId: + session?.activeTurnId === undefined ? session?.lastAbortedTurnId : undefined, + lastAbortedMessageId: + session?.activeTurnId === undefined ? session?.lastAbortedMessageId : undefined, + }; + }); const getSourceProposedPlanReferenceForAcceptedTurnStart = Effect.fn( "getSourceProposedPlanReferenceForAcceptedTurnStart", @@ -1428,8 +1455,8 @@ const make = Effect.gen(function* () { return null; } - const expectedTurnId = yield* getExpectedProviderTurnIdForThread(threadId); - if (!sameId(expectedTurnId, eventTurnId)) { + const expectedTurn = yield* getExpectedProviderTurnForThread(threadId); + if (!sameId(expectedTurn.activeTurnId, eventTurnId)) { return null; } @@ -1502,12 +1529,42 @@ const make = Effect.gen(function* () { threadId: thread.id, }) : Option.none(); - const hasPendingTurnStart = - Option.isSome(pendingTurnStart) && thread.session?.status === "starting"; + const hasPendingTurnStart = Option.isSome(pendingTurnStart); + const hasPendingTurnStartWhileStarting = + hasPendingTurnStart && thread.session?.status === "starting"; + const expectedProviderTurn = + (event.type === "turn.aborted" && hasPendingTurnStart) || + (event.type === "turn.started" && activeTurnId === null) + ? yield* getExpectedProviderTurnForThread(thread.id) + : undefined; + + const turnStartedMatchesRetainedAbort = + event.type === "turn.started" && + activeTurnId === null && + eventTurnId !== undefined && + expectedProviderTurn !== undefined && + sameId(expectedProviderTurn.lastAbortedTurnId, eventTurnId); + const pendingTurnStartMatchesRetainedAbort = + event.type === "turn.aborted" && + Option.isSome(pendingTurnStart) && + expectedProviderTurn !== undefined && + expectedProviderTurn.activeTurnId === undefined && + eventTurnId !== undefined && + sameId(expectedProviderTurn.lastAbortedTurnId, eventTurnId) && + !sameId(expectedProviderTurn.lastAbortedMessageId, pendingTurnStart.value.messageId); + const pendingTurnStartConflictsWithActiveProviderTurn = + event.type === "turn.aborted" && + Option.isSome(pendingTurnStart) && + expectedProviderTurn?.activeTurnId !== undefined && + !sameId(expectedProviderTurn.activeTurnId, eventTurnId); const conflictsWithActiveTurn = activeTurnId !== null && eventTurnId !== undefined && !sameId(activeTurnId, eventTurnId); const missingTurnForActiveTurn = activeTurnId !== null && eventTurnId === undefined; + const expectedProviderTurnForStart = + event.type === "turn.started" && conflictsWithActiveTurn + ? yield* getExpectedProviderTurnForThread(thread.id) + : undefined; // A turn.started that conflicts with the active turn is legitimate when // the server itself has a turn start pending for this thread AND the @@ -1517,7 +1574,8 @@ const make = Effect.gen(function* () { // turn.started for some other turn id still gets rejected. const conflictingTurnStartIsPendingTurnStart = event.type === "turn.started" && conflictsWithActiveTurn - ? sameId(yield* getExpectedProviderTurnIdForThread(thread.id), eventTurnId) && + ? expectedProviderTurnForStart !== undefined && + sameId(expectedProviderTurnForStart.activeTurnId, eventTurnId) && Option.isSome(pendingTurnStart) : false; @@ -1532,12 +1590,43 @@ const make = Effect.gen(function* () { case "thread.started": return true; case "turn.started": - return !conflictsWithActiveTurn || conflictingTurnStartIsPendingTurnStart; + return ( + !turnStartedMatchesRetainedAbort && + (!conflictsWithActiveTurn || conflictingTurnStartIsPendingTurnStart) + ); case "turn.completed": case "turn.aborted": if (conflictsWithActiveTurn || missingTurnForActiveTurn) { return false; } + if (pendingTurnStartConflictsWithActiveProviderTurn) { + return false; + } + if (pendingTurnStartMatchesRetainedAbort) { + return false; + } + // A named abort may arrive after the server has requested a new + // turn but before its turn.started event. In that window there + // is no projected active turn to compare against, so only the + // provider's currently active turn can settle the pending start. + if (event.type === "turn.aborted" && hasPendingTurnStart && activeTurnId === null) { + return ( + eventTurnId !== undefined && + expectedProviderTurn !== undefined && + sameId( + expectedProviderTurn.activeTurnId ?? expectedProviderTurn.lastAbortedTurnId, + eventTurnId, + ) && + Option.isSome(pendingTurnStart) && + (expectedProviderTurn.provider !== "opencode" || + sameId( + expectedProviderTurn.activeTurnId !== undefined + ? expectedProviderTurn.activeTurnStartMessageId + : expectedProviderTurn.lastAbortedMessageId, + pendingTurnStart.value.messageId, + )) + ); + } // Only the active turn may close the lifecycle state. if (activeTurnId !== null && eventTurnId !== undefined) { return sameId(activeTurnId, eventTurnId); @@ -1567,23 +1656,29 @@ const make = Effect.gen(function* () { switch (event.type) { case "session.state.changed": { const runtimeStatus = orchestrationSessionStatusFromRuntimeState(event.payload.state); - return hasPendingTurnStart && runtimeStatus === "ready" ? "starting" : runtimeStatus; + return hasPendingTurnStartWhileStarting && runtimeStatus === "ready" + ? "starting" + : runtimeStatus; } case "turn.started": return "running"; case "session.exited": return "stopped"; - case "turn.aborted": - return "interrupted"; case "turn.completed": return normalizeRuntimeTurnState(event.payload.state) === "failed" ? "error" : "ready"; + case "turn.aborted": + return isOpenCodePromptFailureAbort(event) ? "error" : "interrupted"; case "session.started": case "thread.started": // Provider thread/session start notifications can arrive during an // active or pending turn; preserve that lifecycle state. - return activeTurnId !== null ? "running" : hasPendingTurnStart ? "starting" : "ready"; + return activeTurnId !== null + ? "running" + : hasPendingTurnStartWhileStarting + ? "starting" + : "ready"; } })(); const nextActiveTurnId = @@ -1603,9 +1698,11 @@ const make = Effect.gen(function* () { : event.type === "turn.completed" && normalizeRuntimeTurnState(event.payload.state) === "failed" ? (event.payload.errorMessage ?? thread.session?.lastError ?? "Turn failed") - : status === "ready" || status === "interrupted" - ? null - : (thread.session?.lastError ?? null); + : isOpenCodePromptFailureAbort(event) + ? event.payload.reason + : status === "ready" || status === "interrupted" + ? null + : (thread.session?.lastError ?? null); if (shouldApplyThreadLifecycle) { if (event.type === "turn.started" && acceptedTurnStartedSourcePlan !== null) { @@ -1840,7 +1937,7 @@ const make = Effect.gen(function* () { }); } - if (isTerminalTurn) { + if (isTerminalTurn && shouldApplyThreadLifecycle) { const turnId = toTurnId(event.turnId); if (turnId) { const userInputActivities = @@ -2016,7 +2113,9 @@ const make = Effect.gen(function* () { // active turn's progress; session.exited always clears. if (event.type === "session.exited") { threadPlanProgress.clearThreadPlanProgress(thread.id); - } else if (!conflictsWithActiveTurn) { + } else if ( + event.type === "turn.aborted" ? shouldApplyThreadLifecycle : !conflictsWithActiveTurn + ) { if (event.type === "turn.plan.updated") { threadPlanProgress.recordPlanProgress(thread.id, event.payload.plan); } else if (isTerminalTurn && shouldApplyThreadLifecycle) { diff --git a/apps/server/src/provider/Layers/OpenCodeAdapter.test.ts b/apps/server/src/provider/Layers/OpenCodeAdapter.test.ts index a6636fcd69d0..c6a987f96588 100644 --- a/apps/server/src/provider/Layers/OpenCodeAdapter.test.ts +++ b/apps/server/src/provider/Layers/OpenCodeAdapter.test.ts @@ -26,6 +26,7 @@ import type { import { ApprovalRequestId, + MessageId, OpenCodeSettings, ProviderDriverKind, ProviderInstanceId, @@ -1457,6 +1458,7 @@ it.layer(OpenCodeAdapterTestLayer)("OpenCodeAdapterLive", (it) => { .sendTurn({ threadId: asThreadId("thread-send-turn-failure"), input: "Fix it", + turnStartMessageId: MessageId.make("message-send-turn-failure"), modelSelection: { instanceId: ProviderInstanceId.make("opencode"), model: "openai/gpt-5", @@ -1478,6 +1480,11 @@ it.layer(OpenCodeAdapterTestLayer)("OpenCodeAdapterLive", (it) => { NodeAssert.equal(sessions[0]?.status, "ready"); NodeAssert.equal(sessions[0]?.activeTurnId, undefined); NodeAssert.equal(sessions[0]?.lastError, "prompt failed"); + NodeAssert.equal(typeof sessions[0]?.lastAbortedTurnId, "string"); + NodeAssert.equal( + sessions[0]?.lastAbortedMessageId, + MessageId.make("message-send-turn-failure"), + ); }), ); @@ -1520,6 +1527,156 @@ it.layer(OpenCodeAdapterTestLayer)("OpenCodeAdapterLive", (it) => { }), ); + it.effect("uses the steering turn token through steering and abort", () => + Effect.gen(function* () { + const adapter = yield* OpenCodeAdapter; + const threadId = asThreadId("thread-steer-abort-token"); + const originalTurnStartMessageId = MessageId.make("message-original-turn"); + const steeringTurnStartMessageId = MessageId.make("message-steering-turn"); + yield* adapter.startSession({ + provider: ProviderDriverKind.make("opencode"), + threadId, + runtimeMode: "full-access", + }); + + const turn = yield* adapter.sendTurn({ + threadId, + input: "run 5 commands", + turnStartMessageId: originalTurnStartMessageId, + modelSelection: { + instanceId: ProviderInstanceId.make("opencode"), + model: "openai/gpt-5", + }, + }); + const sessionAfterTurn = (yield* adapter.listSessions()).find( + (entry) => entry.threadId === threadId, + ); + NodeAssert.equal(sessionAfterTurn?.activeTurnStartMessageId, originalTurnStartMessageId); + + const steeredTurn = yield* adapter.sendTurn({ + threadId, + input: "actually run 15", + turnStartMessageId: steeringTurnStartMessageId, + modelSelection: { + instanceId: ProviderInstanceId.make("opencode"), + model: "openai/gpt-5", + }, + }); + NodeAssert.equal(String(steeredTurn.turnId), String(turn.turnId)); + + const sessionAfterSteer = (yield* adapter.listSessions()).find( + (entry) => entry.threadId === threadId, + ); + NodeAssert.equal(sessionAfterSteer?.activeTurnStartMessageId, steeringTurnStartMessageId); + + yield* adapter.interruptTurn(threadId, turn.turnId); + const sessionAfterAbort = (yield* adapter.listSessions()).find( + (entry) => entry.threadId === threadId, + ); + NodeAssert.equal(sessionAfterAbort?.lastAbortedTurnId, turn.turnId); + NodeAssert.equal(sessionAfterAbort?.lastAbortedMessageId, steeringTurnStartMessageId); + }), + ); + + it.effect("uses the steering turn token through a steering timeout", () => + Effect.gen(function* () { + const adapter = yield* OpenCodeAdapter; + const threadId = asThreadId("thread-steer-timeout-token"); + const originalTurnStartMessageId = MessageId.make("message-original-timeout-turn"); + const steeringTurnStartMessageId = MessageId.make("message-steering-timeout-turn"); + const promptStarted = promiseWithResolvers(); + yield* adapter.startSession({ + provider: ProviderDriverKind.make("opencode"), + threadId, + runtimeMode: "full-access", + }); + + const turn = yield* adapter.sendTurn({ + threadId, + input: "run 5 commands", + turnStartMessageId: originalTurnStartMessageId, + modelSelection: { + instanceId: ProviderInstanceId.make("opencode"), + model: "openai/gpt-5", + }, + }); + runtimeMock.state.promptAsyncImplementation = async () => { + promptStarted.resolve(undefined); + await new Promise(() => {}); + }; + + const steeringFiber = yield* adapter + .sendTurn({ + threadId, + input: "actually run 15", + turnStartMessageId: steeringTurnStartMessageId, + modelSelection: { + instanceId: ProviderInstanceId.make("opencode"), + model: "openai/gpt-5", + }, + }) + .pipe(Effect.exit, Effect.forkChild); + yield* Effect.promise(() => promptStarted.promise); + yield* advanceTestClock(10_000); + + const steeringResult = yield* Fiber.join(steeringFiber); + NodeAssert.equal(steeringResult._tag, "Failure"); + const session = (yield* adapter.listSessions()).find((entry) => entry.threadId === threadId); + NodeAssert.equal(session?.lastAbortedTurnId, turn.turnId); + NodeAssert.equal(session?.lastAbortedMessageId, steeringTurnStartMessageId); + }), + ); + + it.effect("preserves the original turn token through a tokenless steering timeout", () => + Effect.gen(function* () { + const adapter = yield* OpenCodeAdapter; + const threadId = asThreadId("thread-steer-timeout-no-token"); + const originalTurnStartMessageId = MessageId.make("message-original-timeout-no-token"); + const promptStarted = promiseWithResolvers(); + yield* adapter.startSession({ + provider: ProviderDriverKind.make("opencode"), + threadId, + runtimeMode: "full-access", + }); + + const turn = yield* adapter.sendTurn({ + threadId, + input: "run 5 commands", + turnStartMessageId: originalTurnStartMessageId, + modelSelection: { + instanceId: ProviderInstanceId.make("opencode"), + model: "openai/gpt-5", + }, + }); + runtimeMock.state.promptAsyncImplementation = async () => { + promptStarted.resolve(undefined); + await new Promise(() => {}); + }; + + // A steer without its own turn token must keep the active turn's + // original token, so a later prompt timeout still correlates the abort + // to the running turn instead of persisting a missing token. + const steeringFiber = yield* adapter + .sendTurn({ + threadId, + input: "actually run 15", + modelSelection: { + instanceId: ProviderInstanceId.make("opencode"), + model: "openai/gpt-5", + }, + }) + .pipe(Effect.exit, Effect.forkChild); + yield* Effect.promise(() => promptStarted.promise); + yield* advanceTestClock(10_000); + + const steeringResult = yield* Fiber.join(steeringFiber); + NodeAssert.equal(steeringResult._tag, "Failure"); + const session = (yield* adapter.listSessions()).find((entry) => entry.threadId === threadId); + NodeAssert.equal(session?.lastAbortedTurnId, turn.turnId); + NodeAssert.equal(session?.lastAbortedMessageId, originalTurnStartMessageId); + }), + ); + it.effect("keeps the running turn when a steer prompt fails", () => Effect.gen(function* () { const adapter = yield* OpenCodeAdapter; @@ -1564,6 +1721,7 @@ it.layer(OpenCodeAdapterTestLayer)("OpenCodeAdapterLive", (it) => { Effect.gen(function* () { const adapter = yield* OpenCodeAdapter; const threadId = asThreadId("thread-steer-idle-admission"); + const activeTurnStartMessageId = MessageId.make("message-steer-idle-admission"); const busyBeforeSteer = promiseWithResolvers(); const idleBeforeSteer = promiseWithResolvers(); const idleAfterSteer = promiseWithResolvers(); @@ -1610,6 +1768,7 @@ it.layer(OpenCodeAdapterTestLayer)("OpenCodeAdapterLive", (it) => { const activeTurn = yield* adapter.sendTurn({ threadId, input: "Start the next turn", + turnStartMessageId: activeTurnStartMessageId, modelSelection: createModelSelection( ProviderInstanceId.make("opencode"), "opencode/kimi-k3", @@ -1632,6 +1791,12 @@ it.layer(OpenCodeAdapterTestLayer)("OpenCodeAdapterLive", (it) => { }, }); yield* Effect.promise(() => statusStarted.promise); + yield* Effect.yieldNow; + const sessionAfterBusy = (yield* adapter.listSessions()).find( + (candidate) => candidate.threadId === threadId, + ); + NodeAssert.equal(sessionAfterBusy?.activeTurnId, activeTurn.turnId); + NodeAssert.equal(sessionAfterBusy?.activeTurnStartMessageId, activeTurnStartMessageId); const steerFiber = yield* adapter .sendTurn({ threadId, @@ -4850,6 +5015,7 @@ it.layer(OpenCodeAdapterTestLayer)("OpenCodeAdapterLive", (it) => { const turn = yield* adapter.sendTurn({ threadId, input: "Keep working", + turnStartMessageId: MessageId.make("message-interrupt-idle-race"), modelSelection: createModelSelection( ProviderInstanceId.make("opencode"), "opencode/kimi-k3", @@ -4883,6 +5049,11 @@ it.layer(OpenCodeAdapterTestLayer)("OpenCodeAdapterLive", (it) => { const session = sessions.find((candidate) => candidate.threadId === threadId); NodeAssert.equal(session?.status, "ready"); NodeAssert.equal(session?.activeTurnId, undefined); + NodeAssert.equal(session?.lastAbortedTurnId, turn.turnId); + NodeAssert.equal( + session?.lastAbortedMessageId, + MessageId.make("message-interrupt-idle-race"), + ); yield* adapter.stopSession(threadId); }), diff --git a/apps/server/src/provider/Layers/OpenCodeAdapter.ts b/apps/server/src/provider/Layers/OpenCodeAdapter.ts index 3d216bb1167b..1124705584f8 100644 --- a/apps/server/src/provider/Layers/OpenCodeAdapter.ts +++ b/apps/server/src/provider/Layers/OpenCodeAdapter.ts @@ -1,5 +1,6 @@ import { EventId, + type MessageId, type OpenCodeSettings, ProviderDriverKind, ProviderInstanceId, @@ -352,6 +353,7 @@ interface OpenCodeSessionContext { readonly textPartsByMessageId: Map>; turnTokenUsage: OpenCodeTurnTokenUsageAccumulator | undefined; activeTurnId: TurnId | undefined; + activeTurnStartMessageId: MessageId | undefined; activeAgent: string | undefined; activeVariant: string | undefined; cancellation: OpenCodeCancellation | undefined; @@ -724,6 +726,7 @@ function updateProviderSession( patch: Partial, options?: { readonly clearActiveTurnId?: boolean; + readonly clearActiveTurnStartMessageId?: boolean; readonly clearLastError?: boolean; }, ): Effect.Effect { @@ -738,6 +741,7 @@ function applyProviderSessionUpdate( options: | { readonly clearActiveTurnId?: boolean; + readonly clearActiveTurnStartMessageId?: boolean; readonly clearLastError?: boolean; } | undefined, @@ -752,6 +756,22 @@ function applyProviderSessionUpdate( if (options?.clearActiveTurnId) { delete mutableSession.activeTurnId; } + if (options?.clearActiveTurnStartMessageId) { + delete mutableSession.activeTurnStartMessageId; + } + if (patch.activeTurnId !== undefined) { + delete mutableSession.lastAbortedTurnId; + delete mutableSession.lastAbortedMessageId; + delete mutableSession.activeTurnStartMessageId; + if (patch.activeTurnStartMessageId !== undefined) { + mutableSession.activeTurnStartMessageId = patch.activeTurnStartMessageId; + } + } + if (patch.lastAbortedTurnId !== undefined) { + if (patch.lastAbortedMessageId === undefined) { + delete mutableSession.lastAbortedMessageId; + } + } if (options?.clearLastError) { delete mutableSession.lastError; } @@ -1130,6 +1150,7 @@ export function makeOpenCodeAdapter( } const tokenUsage = takeOpenCodeTurnTokenUsage(context, true); context.activeTurnId = undefined; + context.activeTurnStartMessageId = undefined; context.activeAgent = undefined; context.activeVariant = undefined; context.interruptedTurnId = undefined; @@ -1142,7 +1163,7 @@ export function makeOpenCodeAdapter( applyProviderSessionUpdate( context, { status: "ready" }, - { clearActiveTurnId: true }, + { clearActiveTurnId: true, clearActiveTurnStartMessageId: true }, updatedAt, ); if (pendingIdleReconciliation?.fiber) { @@ -1299,6 +1320,7 @@ export function makeOpenCodeAdapter( const tokenUsage = takeOpenCodeTurnTokenUsage(context, false); context.promptAdmission = undefined; context.activeTurnId = undefined; + context.activeTurnStartMessageId = undefined; context.activeAgent = undefined; context.activeVariant = undefined; context.awaitingBusyAfterInterruption = false; @@ -1306,7 +1328,7 @@ export function makeOpenCodeAdapter( yield* updateProviderSession( context, { status: "error", lastError: detail }, - { clearActiveTurnId: true }, + { clearActiveTurnId: true, clearActiveTurnStartMessageId: true }, ); yield* emit({ ...(yield* buildEventBase({ @@ -1504,13 +1526,25 @@ export function makeOpenCodeAdapter( }; if (context.activeTurnId === turnId) { tokenUsage = takeOpenCodeTurnTokenUsage(context, false); + const turnStartMessageId = context.activeTurnStartMessageId; context.activeTurnId = undefined; + context.activeTurnStartMessageId = undefined; context.activeAgent = undefined; context.activeVariant = undefined; yield* updateProviderSession( context, - { status: "ready" }, - { clearActiveTurnId: true, clearLastError: true }, + { + status: "ready", + lastAbortedTurnId: turnId, + ...(turnStartMessageId !== undefined + ? { lastAbortedMessageId: turnStartMessageId } + : {}), + }, + { + clearActiveTurnId: true, + clearActiveTurnStartMessageId: true, + clearLastError: true, + }, ); } yield* clearPendingOpenCodeRequests(context, { type: "session.abort" }); @@ -2577,6 +2611,7 @@ export function makeOpenCodeAdapter( yield* updateProviderSession(context, { status: "running", activeTurnId: turnId, + activeTurnStartMessageId: context.activeTurnStartMessageId, }); } @@ -2650,6 +2685,7 @@ export function makeOpenCodeAdapter( } const tokenUsage = activeTurnId ? takeOpenCodeTurnTokenUsage(context, false) : undefined; context.activeTurnId = undefined; + context.activeTurnStartMessageId = undefined; context.activeAgent = undefined; context.activeVariant = undefined; context.reconcileIdleStatus = false; @@ -2660,7 +2696,7 @@ export function makeOpenCodeAdapter( status: "error", lastError: message, }, - { clearActiveTurnId: true }, + { clearActiveTurnId: true, clearActiveTurnStartMessageId: true }, ); if (activeTurnId) { yield* emit({ @@ -2996,6 +3032,7 @@ export function makeOpenCodeAdapter( messageRoleById: new Map(), turnTokenUsage: undefined, activeTurnId: undefined, + activeTurnStartMessageId: undefined, activeAgent: undefined, activeVariant: undefined, cancellation: undefined, @@ -3133,6 +3170,10 @@ export function makeOpenCodeAdapter( // A sendTurn while a turn is active is a steer. OpenCode queues the // prompt into the running session, so the active turn id is reused. const steeringTurnId = context.activeTurnId; + const turnStartMessageId = + steeringTurnId === undefined + ? input.turnStartMessageId + : (input.turnStartMessageId ?? context.activeTurnStartMessageId); const turnId = steeringTurnId ?? freshTurnId; const agent = getModelSelectionStringOptionValue(modelSelection, "agent"); const variant = getModelSelectionStringOptionValue(modelSelection, "variant"); @@ -3171,6 +3212,7 @@ export function makeOpenCodeAdapter( context.turnTokenUsage = makeOpenCodeTurnTokenUsageAccumulator(); } context.turnTokenUsage?.promptMessageIds.add(messageId); + context.activeTurnStartMessageId = turnStartMessageId; context.activeAgent = agent ?? (input.interactionMode === "plan" ? "plan" : undefined); context.activeVariant = variant; if (steeringTurnId === undefined) { @@ -3184,6 +3226,7 @@ export function makeOpenCodeAdapter( { status: "running", activeTurnId: turnId, + activeTurnStartMessageId: turnStartMessageId, model: modelSelection?.model ?? context.session.model, }, { clearLastError: true }, @@ -3262,6 +3305,7 @@ export function makeOpenCodeAdapter( const tokenUsage = takeOpenCodeTurnTokenUsage(context, false); context.promptAdmission = undefined; context.activeTurnId = undefined; + context.activeTurnStartMessageId = undefined; context.activeAgent = undefined; context.activeVariant = undefined; yield* updateProviderSession( @@ -3270,8 +3314,12 @@ export function makeOpenCodeAdapter( status: "ready", model: modelSelection?.model ?? context.session.model, lastError: requestError.detail, + lastAbortedTurnId: turnId, + ...(turnStartMessageId !== undefined + ? { lastAbortedMessageId: turnStartMessageId } + : {}), }, - { clearActiveTurnId: true }, + { clearActiveTurnId: true, clearActiveTurnStartMessageId: true }, ); yield* emit({ ...(yield* buildEventBase({ threadId: input.threadId, turnId })), @@ -3310,6 +3358,7 @@ export function makeOpenCodeAdapter( const tokenUsage = takeOpenCodeTurnTokenUsage(context, false); context.promptAdmission = undefined; context.activeTurnId = undefined; + context.activeTurnStartMessageId = undefined; context.activeAgent = undefined; context.activeVariant = undefined; context.awaitingBusyAfterInterruption = false; @@ -3320,8 +3369,12 @@ export function makeOpenCodeAdapter( status: "ready", model: modelSelection?.model ?? context.session.model, lastError: requestError.detail, + lastAbortedTurnId: turnId, + ...(turnStartMessageId !== undefined + ? { lastAbortedMessageId: turnStartMessageId } + : {}), }, - { clearActiveTurnId: true }, + { clearActiveTurnId: true, clearActiveTurnStartMessageId: true }, ); yield* emit({ ...(yield* buildEventBase({ diff --git a/packages/contracts/src/provider.ts b/packages/contracts/src/provider.ts index df91839e8520..33fc6dc86f38 100644 --- a/packages/contracts/src/provider.ts +++ b/packages/contracts/src/provider.ts @@ -4,6 +4,7 @@ import { ApprovalRequestId, EventId, IsoDateTime, + MessageId, ProviderItemId, ThreadId, TurnId, @@ -45,6 +46,11 @@ export const ProviderSession = Schema.Struct({ threadId: ThreadId, resumeCursor: Schema.optional(Schema.Unknown), activeTurnId: Schema.optional(TurnId), + activeTurnStartMessageId: Schema.optional(MessageId), + // Retained after a provider abort clears its active turn so consumers can + // correlate the terminal event with a pending durable turn start. + lastAbortedTurnId: Schema.optional(TurnId), + lastAbortedMessageId: Schema.optional(MessageId), createdAt: IsoDateTime, updatedAt: IsoDateTime, lastError: Schema.optional(TrimmedNonEmptyString), @@ -79,6 +85,7 @@ export const ProviderSendTurnInput = Schema.Struct({ ), modelSelection: Schema.optional(ModelSelection), interactionMode: Schema.optional(ProviderInteractionMode), + turnStartMessageId: Schema.optional(MessageId), }); export type ProviderSendTurnInput = typeof ProviderSendTurnInput.Type;