diff --git a/apps/server/src/orchestration/Layers/ProviderCommandReactor.ts b/apps/server/src/orchestration/Layers/ProviderCommandReactor.ts index 84d472089d8d..05670b0903c9 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, @@ -775,6 +776,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; @@ -826,6 +828,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 } : {}), @@ -1213,6 +1218,7 @@ const make = Effect.gen(function* () { const sendTurnRequest = yield* buildSendTurnRequestForThread({ threadId: event.payload.threadId, + turnStartMessageId: event.payload.messageId, messageText: message.text, ...(message.attachments !== undefined ? { attachments: message.attachments } : {}), ...(event.payload.modelSelection !== undefined diff --git a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts index 26332f9f8c9c..39bb2a2da93b 100644 --- a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts +++ b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts @@ -259,6 +259,9 @@ describe("ProviderRuntimeIngestion", () => { const engine = await runtime.runPromise(Effect.service(OrchestrationEngineService)); const snapshotQuery = await runtime.runPromise(Effect.service(ProjectionSnapshotQuery)); const ingestion = await runtime.runPromise(Effect.service(ProviderRuntimeIngestionService)); + const planProgress = await runtime.runPromise( + Effect.service(ThreadPlanProgress.ThreadPlanProgressService), + ); scope = await Effect.runPromise(Scope.make("sequential")); await Effect.runPromise(ingestion.start().pipe(Scope.provide(scope))); const drain = () => Effect.runPromise(ingestion.drain); @@ -324,6 +327,7 @@ describe("ProviderRuntimeIngestion", () => { emit: provider.emit, setProviderSession: provider.setSession, drain, + planProgress, }; } @@ -369,6 +373,368 @@ describe("ProviderRuntimeIngestion", () => { expect(thread.session?.lastError).toBe("turn failed"); }); + it("settles the durable session when a turn is aborted", async () => { + const harness = await createHarness(); + const now = "2026-01-01T00:00:00.000Z"; + const turnId = asTurnId("turn-aborted"); + + harness.emit({ + type: "turn.started", + eventId: asEventId("evt-turn-aborted-started"), + provider: ProviderDriverKind.make("opencode"), + threadId: asThreadId("thread-1"), + createdAt: now, + turnId, + }); + + 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?.activeTurnId).toBe(turnId); + + harness.emit({ + type: "turn.aborted", + eventId: asEventId("evt-turn-aborted"), + provider: ProviderDriverKind.make("opencode"), + threadId: asThreadId("thread-1"), + 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 === asThreadId("thread-1"), + ); + 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: "turn.started", + eventId: asEventId("evt-turn-late-after-abort-started"), + provider: ProviderDriverKind.make("opencode"), + threadId, + createdAt: now, + turnId, + }); + await harness.drain(); + + 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: "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(); + + 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("persists an OpenCode prompt failure carried by an aborted turn", async () => { + const harness = await createHarness(); + 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-aborted-prompt-failure-started"), + provider: ProviderDriverKind.make("opencode"), + threadId: asThreadId("thread-1"), + createdAt: "2026-01-01T00:00:00.000Z", + turnId, + }); + await harness.drain(); + + harness.emit({ + 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", + turnId, + payload: { reason: failureReason }, + }); + await harness.drain(); + + const thread = (await harness.readModel()).threads.find( + (entry) => entry.id === asThreadId("thread-1"), + ); + expect(thread?.session?.status).toBe("error"); + expect(thread?.session?.activeTurnId).toBeNull(); + expect(thread?.session?.lastError).toBe(failureReason); + expect(thread?.latestTurn?.state).toBe("error"); + }); + + 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."; + + 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(); + + 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"); + }); + + 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"; @@ -996,6 +1362,709 @@ describe("ProviderRuntimeIngestion", () => { 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); + }); + + 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-abort-untargeted-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-untargeted"), + provider: ProviderDriverKind.make("opencode"), + createdAt: "2026-01-01T00:00:01.000Z", + 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("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", + activeTurnId: null, + updatedAt: seededAt, + lastError: null, + }, + createdAt: seededAt, + }); + + 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, + 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(); + + let thread = (await harness.readModel()).threads.find((entry) => entry.id === threadId); + expect(thread?.session?.status).toBe("starting"); + expect(thread?.session?.activeTurnId).toBeNull(); + + // 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 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.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(); + + let 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, + }); + harness.emit({ + 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(); + + 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.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(); + + thread = (await harness.readModel()).threads.find((entry) => entry.id === threadId); + expect(thread?.session?.status).toBe("interrupted"); + expect(thread?.session?.activeTurnId).toBeNull(); + }); + + it("rejects a targeted abort that does not match the pending provider turn", async () => { + const harness = await createHarness(); + 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 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.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(); + + 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.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(); + + thread = (await harness.readModel()).threads.find((entry) => entry.id === threadId); + expect(thread?.session?.status).toBe("interrupted"); + expect(thread?.session?.activeTurnId).toBeNull(); + }); + + 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.setProviderSession({ + provider: ProviderDriverKind.make("opencode"), + status: "running", + runtimeMode: "approval-required", + threadId, + createdAt: now, + updatedAt: now, + activeTurnId: turnId, + activeTurnStartMessageId: originalMessageId, + }); + harness.emit({ + type: "turn.started", + eventId: asEventId("evt-turn-steer-abort-started"), + provider: ProviderDriverKind.make("opencode"), + threadId, + createdAt: now, + turnId, + }); + await harness.drain(); + + await harness.dispatch({ + 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: [], + }, + 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.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 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("rejects an old abort when a newer provider turn is active", async () => { + const harness = await createHarness(); + 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(); + + await harness.dispatch({ + type: "thread.turn.start", + commandId: CommandId.make("cmd-turn-start-new-bound"), + threadId, + message: { + messageId: pendingMessageId, + role: "user", + text: "start the newer turn", + attachments: [], + }, + 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.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(); + + 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("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.proposed.completed", + eventId: asEventId("evt-clear-before-emit-source-plan"), + provider: ProviderDriverKind.make("codex"), + threadId, + createdAt: now, + turnId: sourceTurnId, + payload: { planMarkdown: "# Source plan" }, + }); + const threadWithSourcePlan = await waitForThread( + harness.readModel, + (thread) => + 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."); + } + + 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, + }); + + harness.emit({ + 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(); + + const thread = (await harness.readModel()).threads.find((entry) => entry.id === threadId); + expect(thread?.session?.status).toBe("running"); + expect(thread?.session?.activeTurnId).toBe(oldTurnId); + + 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(); + + 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 () => { const harness = await createHarness(); const now = "2026-01-01T00:00:00.000Z"; diff --git a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts index a90010f0b6e2..ce78b049d388 100644 --- a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts +++ b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts @@ -271,6 +271,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" { @@ -1443,13 +1461,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", @@ -1458,8 +1485,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; } @@ -1525,17 +1552,48 @@ const make = Effect.gen(function* () { event.type === "session.exited" || event.type === "thread.started" || event.type === "turn.started" || - event.type === "turn.completed" + event.type === "turn.completed" || + event.type === "turn.aborted" ? yield* projectionTurnRepository.getPendingTurnStartByThreadId({ 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 @@ -1545,7 +1603,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; @@ -1560,22 +1619,54 @@ 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); } - // No active turn tracked: accept only completions that name their - // turn (covers a real completion whose turn.started was lost). An - // untargeted completion cannot prove it belongs to any turn this - // thread ran — the known emitter was the Claude resume handshake - // (system/init + result(num_turns: 0)), which is not a turn at - // all — and applying it here stomps the "starting" lifecycle - // state while a turn start is pending. + // No active turn tracked: accept only terminal events that name + // their turn (covers a real completion/abort whose turn.started + // was lost). An untargeted terminal event cannot prove it belongs + // to any turn this thread ran — the known emitter was the Claude + // resume handshake (system/init + result(num_turns: 0)), which is + // not a turn at all — and applying it here stomps the "starting" + // lifecycle state while a turn start is pending. return eventTurnId !== undefined; default: return true; @@ -1592,13 +1683,16 @@ const make = Effect.gen(function* () { event.type === "session.exited" || event.type === "thread.started" || event.type === "turn.started" || - event.type === "turn.completed" + event.type === "turn.completed" || + event.type === "turn.aborted" ) { const status = (() => { 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"; @@ -1608,17 +1702,25 @@ const make = Effect.gen(function* () { 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 = event.type === "turn.started" ? (eventTurnId ?? null) - : event.type === "turn.completed" || event.type === "session.exited" + : event.type === "turn.completed" || + event.type === "turn.aborted" || + event.type === "session.exited" ? null : event.type === "session.state.changed" && !sessionStatusAllowsActiveTurn( @@ -1632,9 +1734,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" - ? 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) { @@ -1859,7 +1963,10 @@ const make = Effect.gen(function* () { }); } - if (event.type === "turn.completed") { + if ( + (event.type === "turn.completed" || event.type === "turn.aborted") && + shouldApplyThreadLifecycle + ) { const detailedThread = yield* getLoadedThreadDetail(); const messages = detailedThread?.messages ?? []; const proposedPlans = detailedThread?.proposedPlans ?? []; @@ -1989,7 +2096,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 (event.type === "turn.completed" || event.type === "turn.aborted") { diff --git a/apps/server/src/provider/Layers/OpenCodeAdapter.test.ts b/apps/server/src/provider/Layers/OpenCodeAdapter.test.ts index d297360e6d34..f0af1d0627e2 100644 --- a/apps/server/src/provider/Layers/OpenCodeAdapter.test.ts +++ b/apps/server/src/provider/Layers/OpenCodeAdapter.test.ts @@ -19,6 +19,7 @@ import type { PermissionRequest, QuestionRequest } from "@opencode-ai/sdk/v2"; import { ApprovalRequestId, + MessageId, OpenCodeSettings, ProviderDriverKind, ProviderInstanceId, @@ -1269,6 +1270,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", @@ -1290,6 +1292,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"), + ); }), ); @@ -1332,6 +1339,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; @@ -1376,6 +1533,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(); @@ -1422,6 +1580,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", @@ -1444,6 +1603,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, @@ -3302,6 +3467,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", @@ -3335,6 +3501,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 d0b4f0de78ce..eb3048b0258a 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, @@ -338,6 +339,7 @@ interface OpenCodeSessionContext { readonly completedAssistantPartIds: Set; readonly turns: Array; activeTurnId: TurnId | undefined; + activeTurnStartMessageId: MessageId | undefined; activeAgent: string | undefined; activeVariant: string | undefined; cancellation: OpenCodeCancellation | undefined; @@ -650,6 +652,7 @@ function updateProviderSession( patch: Partial, options?: { readonly clearActiveTurnId?: boolean; + readonly clearActiveTurnStartMessageId?: boolean; readonly clearLastError?: boolean; }, ): Effect.Effect { @@ -664,6 +667,7 @@ function applyProviderSessionUpdate( options: | { readonly clearActiveTurnId?: boolean; + readonly clearActiveTurnStartMessageId?: boolean; readonly clearLastError?: boolean; } | undefined, @@ -678,6 +682,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; } @@ -1049,6 +1069,7 @@ export function makeOpenCodeAdapter( context.pendingIdleReconciliation = undefined; } context.activeTurnId = undefined; + context.activeTurnStartMessageId = undefined; context.activeAgent = undefined; context.activeVariant = undefined; context.interruptedTurnId = undefined; @@ -1057,7 +1078,7 @@ export function makeOpenCodeAdapter( applyProviderSessionUpdate( context, { status: "ready" }, - { clearActiveTurnId: true }, + { clearActiveTurnId: true, clearActiveTurnStartMessageId: true }, updatedAt, ); if (pendingIdleReconciliation?.fiber) { @@ -1211,6 +1232,7 @@ export function makeOpenCodeAdapter( } context.promptAdmission = undefined; context.activeTurnId = undefined; + context.activeTurnStartMessageId = undefined; context.activeAgent = undefined; context.activeVariant = undefined; context.awaitingBusyAfterInterruption = false; @@ -1218,7 +1240,7 @@ export function makeOpenCodeAdapter( yield* updateProviderSession( context, { status: "error", lastError: detail }, - { clearActiveTurnId: true }, + { clearActiveTurnId: true, clearActiveTurnStartMessageId: true }, ); yield* emit({ ...(yield* buildEventBase({ @@ -1408,13 +1430,25 @@ export function makeOpenCodeAdapter( context.cancellation = undefined; } if (context.activeTurnId === turnId) { + 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* emit({ @@ -2188,6 +2222,7 @@ export function makeOpenCodeAdapter( yield* updateProviderSession(context, { status: "running", activeTurnId: turnId, + activeTurnStartMessageId: context.activeTurnStartMessageId, }); } @@ -2260,6 +2295,7 @@ export function makeOpenCodeAdapter( terminalCancellation.acknowledged = true; } context.activeTurnId = undefined; + context.activeTurnStartMessageId = undefined; context.activeAgent = undefined; context.activeVariant = undefined; context.reconcileIdleStatus = false; @@ -2269,7 +2305,7 @@ export function makeOpenCodeAdapter( status: "error", lastError: message, }, - { clearActiveTurnId: true }, + { clearActiveTurnId: true, clearActiveTurnStartMessageId: true }, ); if (activeTurnId) { yield* emit({ @@ -2564,6 +2600,7 @@ export function makeOpenCodeAdapter( completedAssistantPartIds: new Set(), turns: [], activeTurnId: undefined, + activeTurnStartMessageId: undefined, activeAgent: undefined, activeVariant: undefined, cancellation: undefined, @@ -2701,6 +2738,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"); @@ -2735,6 +2776,7 @@ export function makeOpenCodeAdapter( context.promptAdmission = promptAdmission; context.activeTurnId = turnId; + context.activeTurnStartMessageId = turnStartMessageId; context.activeAgent = agent ?? (input.interactionMode === "plan" ? "plan" : undefined); context.activeVariant = variant; if (steeringTurnId === undefined) { @@ -2748,6 +2790,7 @@ export function makeOpenCodeAdapter( { status: "running", activeTurnId: turnId, + activeTurnStartMessageId: turnStartMessageId, model: modelSelection?.model ?? context.session.model, }, { clearLastError: true }, @@ -2821,6 +2864,7 @@ export function makeOpenCodeAdapter( } context.promptAdmission = undefined; context.activeTurnId = undefined; + context.activeTurnStartMessageId = undefined; context.activeAgent = undefined; context.activeVariant = undefined; yield* updateProviderSession( @@ -2829,8 +2873,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 })), @@ -2865,6 +2913,7 @@ export function makeOpenCodeAdapter( } context.promptAdmission = undefined; context.activeTurnId = undefined; + context.activeTurnStartMessageId = undefined; context.activeAgent = undefined; context.activeVariant = undefined; context.awaitingBusyAfterInterruption = false; @@ -2875,8 +2924,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 42a943923037..1ea3bdfe56d0 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, @@ -44,6 +45,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), @@ -75,6 +81,7 @@ export const ProviderSendTurnInput = Schema.Struct({ ), modelSelection: Schema.optional(ModelSelection), interactionMode: Schema.optional(ProviderInteractionMode), + turnStartMessageId: Schema.optional(MessageId), }); export type ProviderSendTurnInput = typeof ProviderSendTurnInput.Type;