diff --git a/apps/server/src/orchestration/Layers/ProjectionPipeline.ts b/apps/server/src/orchestration/Layers/ProjectionPipeline.ts index f12df850941..20876c99630 100644 --- a/apps/server/src/orchestration/Layers/ProjectionPipeline.ts +++ b/apps/server/src/orchestration/Layers/ProjectionPipeline.ts @@ -819,6 +819,9 @@ const makeOrchestrationProjectionPipeline = Effect.fn("makeOrchestrationProjecti const nextText = Option.match(existingMessage, { onNone: () => event.payload.text, onSome: (message) => { + if (event.payload.reset === true) { + return event.payload.text; + } if (event.payload.streaming) { return `${message.text}${event.payload.text}`; } diff --git a/apps/server/src/orchestration/Layers/ProviderCommandReactor.test.ts b/apps/server/src/orchestration/Layers/ProviderCommandReactor.test.ts index e9383ce204a..3427f3afbf6 100644 --- a/apps/server/src/orchestration/Layers/ProviderCommandReactor.test.ts +++ b/apps/server/src/orchestration/Layers/ProviderCommandReactor.test.ts @@ -1703,7 +1703,93 @@ describe("ProviderCommandReactor", () => { await waitFor(() => harness.interruptTurn.mock.calls.length === 1); expect(harness.interruptTurn.mock.calls[0]?.[0]).toEqual({ threadId: "thread-1", + turnId: "turn-1", }); + + const readModel = await harness.readModel(); + const thread = readModel.threads.find((entry) => entry.id === ThreadId.make("thread-1")); + expect(thread?.session?.status).toBe("ready"); + expect(thread?.session?.activeTurnId).toBeNull(); + }); + + it("settles a stuck running session when stop is pressed after the turn is already interrupted", async () => { + const harness = await createHarness(); + const firstInterruptAt = "2026-01-01T00:00:00.000Z"; + const repeatInterruptAt = "2026-01-01T00:00:01.000Z"; + const threadId = ThreadId.make("thread-1"); + const turnId = asTurnId("turn-1"); + + await Effect.runPromise( + harness.engine.dispatch({ + type: "thread.session.set", + commandId: CommandId.make("cmd-session-set-stuck"), + threadId, + session: { + threadId, + status: "running", + providerName: "codex", + runtimeMode: "approval-required", + activeTurnId: turnId, + lastError: null, + updatedAt: firstInterruptAt, + }, + createdAt: firstInterruptAt, + }), + ); + + await Effect.runPromise( + harness.engine.dispatch({ + type: "thread.turn.interrupt", + commandId: CommandId.make("cmd-turn-interrupt-stuck"), + threadId, + turnId, + createdAt: firstInterruptAt, + }), + ); + + await waitFor(() => harness.interruptTurn.mock.calls.length === 1); + + harness.interruptTurn.mockClear(); + + await Effect.runPromise( + harness.engine.dispatch({ + type: "thread.session.set", + commandId: CommandId.make("cmd-session-set-stuck-again"), + threadId, + session: { + threadId, + status: "running", + providerName: "codex", + runtimeMode: "approval-required", + activeTurnId: turnId, + lastError: null, + updatedAt: repeatInterruptAt, + }, + createdAt: repeatInterruptAt, + }), + ); + + await Effect.runPromise( + harness.engine.dispatch({ + type: "thread.turn.interrupt", + commandId: CommandId.make("cmd-turn-interrupt-stuck-repeat"), + threadId, + turnId, + createdAt: repeatInterruptAt, + }), + ); + + await waitFor(async () => { + const readModel = await harness.readModel(); + const thread = readModel.threads.find((entry) => entry.id === threadId); + return thread?.session?.status === "ready" && thread?.session?.activeTurnId === null; + }); + + expect(harness.interruptTurn.mock.calls.length).toBe(0); + const readModel = await harness.readModel(); + const thread = readModel.threads.find((entry) => entry.id === threadId); + expect(thread?.session?.status).toBe("ready"); + expect(thread?.session?.activeTurnId).toBeNull(); }); it("starts a fresh session when only projected session state exists", async () => { diff --git a/apps/server/src/orchestration/Layers/ProviderCommandReactor.ts b/apps/server/src/orchestration/Layers/ProviderCommandReactor.ts index 6a71a4735d8..3c6ca52af7b 100644 --- a/apps/server/src/orchestration/Layers/ProviderCommandReactor.ts +++ b/apps/server/src/orchestration/Layers/ProviderCommandReactor.ts @@ -870,8 +870,8 @@ const make = Effect.gen(function* () { if (!thread) { return; } - const hasSession = thread.session && thread.session.status !== "stopped"; - if (!hasSession) { + const session = thread.session; + if (!session || session.status === "stopped") { return yield* appendProviderFailureActivity({ threadId: event.payload.threadId, kind: "provider.turn.interrupt.failed", @@ -882,8 +882,63 @@ const make = Effect.gen(function* () { }); } - // Orchestration turn ids are not provider turn ids, so interrupt by session. - yield* providerService.interruptTurn({ threadId: event.payload.threadId }); + const now = event.payload.createdAt; + const turnId = event.payload.turnId; + const latestTurn = thread.latestTurn; + const sessionStillShowsActiveTurn = + session.status === "running" || + session.status === "starting" || + session.activeTurnId !== null; + const turnWasSettledBeforeThisInterrupt = + turnId !== undefined && + latestTurn?.turnId === turnId && + latestTurn.state !== "running" && + latestTurn.completedAt !== null && + latestTurn.completedAt < now; + const shouldSignalProvider = + !turnWasSettledBeforeThisInterrupt && + (session.status === "running" || + session.status === "starting" || + session.activeTurnId !== null); + + if (shouldSignalProvider) { + // Orchestration turn ids are not provider turn ids, so interrupt by session. + yield* providerService + .interruptTurn({ + threadId: event.payload.threadId, + ...(turnId !== undefined ? { turnId } : {}), + }) + .pipe( + Effect.catchCause((cause) => + appendProviderFailureActivity({ + threadId: event.payload.threadId, + kind: "provider.turn.interrupt.failed", + summary: "Provider turn interrupt failed", + detail: Cause.pretty(cause), + turnId: turnId ?? null, + createdAt: now, + }), + ), + ); + } + + // Projection marks the turn interrupted, but only provider runtime events + // (or an explicit session set) clear session.activeTurnId. Without this, + // stop leaves the thread "working" forever — especially when the provider + // process is already gone and interrupt could not reach it. + if (sessionStillShowsActiveTurn) { + yield* setThreadSession({ + threadId: thread.id, + session: { + ...session, + status: "ready", + activeTurnId: null, + lastError: null, + updatedAt: now, + }, + createdAt: now, + }); + } }); const processApprovalResponseRequested = Effect.fn("processApprovalResponseRequested")(function* ( diff --git a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts index 001ba388949..88ee2c5d243 100644 --- a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts +++ b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts @@ -2072,6 +2072,119 @@ describe("ProviderRuntimeIngestion", () => { ).toBe(" after approval"); }); + it("resets a completed assistant segment when item.started arrives again after resume", async () => { + const harness = await createHarness({ serverSettings: { enableAssistantStreaming: true } }); + const startedAt = "2026-03-28T08:00:00.000Z"; + const resumedAt = "2026-03-28T08:00:10.000Z"; + + harness.emit({ + type: "turn.started", + eventId: asEventId("evt-turn-started-resume-assistant-segment"), + provider: ProviderDriverKind.make("cursor"), + createdAt: startedAt, + threadId: asThreadId("thread-1"), + turnId: asTurnId("turn-resume-assistant-segment"), + }); + await waitForThread( + harness.readModel, + (thread) => + thread.session?.status === "running" && + thread.session?.activeTurnId === "turn-resume-assistant-segment", + ); + + const resumedItemId = asItemId("assistant:session-1:segment:3"); + const resumedMessageId = "assistant:assistant:session-1:segment:3"; + + harness.emit({ + type: "content.delta", + eventId: asEventId("evt-message-delta-resume-assistant-segment-initial"), + provider: ProviderDriverKind.make("cursor"), + createdAt: startedAt, + threadId: asThreadId("thread-1"), + turnId: asTurnId("turn-resume-assistant-segment"), + itemId: resumedItemId, + payload: { + streamKind: "assistant_text", + delta: "stale summary before restart", + }, + }); + harness.emit({ + type: "item.completed", + eventId: asEventId("evt-message-completed-resume-assistant-segment-initial"), + provider: ProviderDriverKind.make("cursor"), + createdAt: startedAt, + threadId: asThreadId("thread-1"), + turnId: asTurnId("turn-resume-assistant-segment"), + itemId: resumedItemId, + payload: { + itemType: "assistant_message", + status: "completed", + }, + }); + + await waitForThread(harness.readModel, (entry) => + entry.messages.some( + (message: ProviderRuntimeTestMessage) => + message.id === resumedMessageId && + !message.streaming && + message.text === "stale summary before restart", + ), + ); + + harness.emit({ + type: "item.started", + eventId: asEventId("evt-message-started-resume-assistant-segment"), + provider: ProviderDriverKind.make("cursor"), + createdAt: resumedAt, + threadId: asThreadId("thread-1"), + turnId: asTurnId("turn-resume-assistant-segment"), + itemId: resumedItemId, + payload: { + itemType: "assistant_message", + status: "inProgress", + }, + }); + harness.emit({ + type: "content.delta", + eventId: asEventId("evt-message-delta-resume-assistant-segment-after-resume"), + provider: ProviderDriverKind.make("cursor"), + createdAt: resumedAt, + threadId: asThreadId("thread-1"), + turnId: asTurnId("turn-resume-assistant-segment"), + itemId: resumedItemId, + payload: { + streamKind: "assistant_text", + delta: "fresh answer after resume", + }, + }); + harness.emit({ + type: "item.completed", + eventId: asEventId("evt-message-completed-resume-assistant-segment-after-resume"), + provider: ProviderDriverKind.make("cursor"), + createdAt: resumedAt, + threadId: asThreadId("thread-1"), + turnId: asTurnId("turn-resume-assistant-segment"), + itemId: resumedItemId, + payload: { + itemType: "assistant_message", + status: "completed", + }, + }); + + const thread = await waitForThread(harness.readModel, (entry) => + entry.messages.some( + (message: ProviderRuntimeTestMessage) => + message.id === resumedMessageId && + !message.streaming && + message.text === "fresh answer after resume", + ), + ); + expect( + thread.messages.find((message: ProviderRuntimeTestMessage) => message.id === resumedMessageId) + ?.text, + ).toBe("fresh answer after resume"); + }); + it("streams assistant deltas when thread.turn.start requests streaming mode", async () => { const harness = await createHarness({ serverSettings: { enableAssistantStreaming: true } }); 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 3e5978f4846..a94c3e1142d 100644 --- a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts +++ b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts @@ -202,6 +202,19 @@ function assistantSegmentMessageId(baseKey: string, segmentIndex: number): Messa segmentIndex === 0 ? `assistant:${baseKey}` : `assistant:${baseKey}:segment:${segmentIndex}`, ); } + +function assistantMessageIdFromEvent(event: ProviderRuntimeEvent): MessageId { + return assistantSegmentMessageId(assistantSegmentBaseKeyFromEvent(event), 0); +} + +function nextAssistantSegmentIndexFromItemId(itemId: string | undefined): number { + if (!itemId) { + return 1; + } + const match = itemId.match(/:segment:(\d+)$/); + return match ? Number(match[1]) + 1 : 1; +} + function buildContextWindowActivityPayload( event: ProviderRuntimeEvent, ): ThreadTokenUsageSnapshot | undefined { @@ -1365,6 +1378,36 @@ const make = Effect.gen(function* () { const proposedPlanDelta = event.type === "turn.proposed.delta" ? event.payload.delta : undefined; + if (event.type === "item.started" && event.payload.itemType === "assistant_message") { + const turnId = toTurnId(event.turnId); + if (turnId) { + const messageId = assistantMessageIdFromEvent(event); + const detailedThread = yield* getLoadedThreadDetail(); + const existingMessage = detailedThread + ? findMessageById(detailedThread.messages, messageId) + : undefined; + if (existingMessage && !existingMessage.streaming) { + yield* orchestrationEngine.dispatch({ + type: "thread.message.assistant.reset", + commandId: yield* providerCommandId(event, "assistant-reset-on-item-started"), + threadId: thread.id, + messageId, + turnId, + createdAt: now, + }); + } + yield* clearBufferedAssistantText(messageId); + yield* rememberAssistantMessageId(thread.id, turnId, messageId); + yield* setAssistantSegmentStateForTurn(thread.id, turnId, { + baseKey: assistantSegmentBaseKeyFromEvent(event), + nextSegmentIndex: nextAssistantSegmentIndexFromItemId( + event.itemId !== undefined ? String(event.itemId) : undefined, + ), + activeMessageId: messageId, + }); + } + } + if (assistantDelta && assistantDelta.length > 0) { const turnId = toTurnId(event.turnId); const assistantMessageId = yield* getOrCreateAssistantMessageId({ diff --git a/apps/server/src/orchestration/decider.ts b/apps/server/src/orchestration/decider.ts index 0d4af771ca8..e27ea53f25c 100644 --- a/apps/server/src/orchestration/decider.ts +++ b/apps/server/src/orchestration/decider.ts @@ -654,6 +654,34 @@ export const decideOrchestrationCommand = Effect.fn("decideOrchestrationCommand" }; } + case "thread.message.assistant.reset": { + yield* requireThread({ + readModel, + command, + threadId: command.threadId, + }); + return { + ...(yield* withEventBase({ + aggregateKind: "thread", + aggregateId: command.threadId, + occurredAt: command.createdAt, + commandId: command.commandId, + })), + type: "thread.message-sent", + payload: { + threadId: command.threadId, + messageId: command.messageId, + role: "assistant", + text: "", + turnId: command.turnId ?? null, + streaming: true, + reset: true, + createdAt: command.createdAt, + updatedAt: command.createdAt, + }, + }; + } + case "thread.proposed-plan.upsert": { yield* requireThread({ readModel, diff --git a/apps/server/src/orchestration/projector.ts b/apps/server/src/orchestration/projector.ts index fc6ab8f6fcf..fe56f2bed93 100644 --- a/apps/server/src/orchestration/projector.ts +++ b/apps/server/src/orchestration/projector.ts @@ -415,11 +415,14 @@ export function projectEvent( entry.id === message.id ? { ...entry, - text: message.streaming - ? `${entry.text}${message.text}` - : message.text.length > 0 + text: + payload.reset === true ? message.text - : entry.text, + : message.streaming + ? `${entry.text}${message.text}` + : message.text.length > 0 + ? message.text + : entry.text, streaming: message.streaming, updatedAt: message.updatedAt, turnId: message.turnId, diff --git a/apps/server/src/provider/Layers/ProviderService.ts b/apps/server/src/provider/Layers/ProviderService.ts index 4f9070c435a..b104fbec851 100644 --- a/apps/server/src/provider/Layers/ProviderService.ts +++ b/apps/server/src/provider/Layers/ProviderService.ts @@ -25,6 +25,7 @@ import { type ProviderSession, } from "@t3tools/contracts"; import { causeErrorTag } from "@t3tools/shared/observability"; +import * as Cause from "effect/Cause"; import * as DateTime from "effect/DateTime"; import * as Effect from "effect/Effect"; import * as Layer from "effect/Layer"; @@ -281,6 +282,80 @@ const makeProviderService = Effect.fn("makeProviderService")(function* ( }); }); + const syncPersistedRuntimeFromLifecycleEvent = ( + event: ProviderRuntimeEvent, + ): Effect.Effect => + Effect.gen(function* () { + const threadId = event.threadId; + if (threadId === undefined) { + return; + } + + const bindingOption = yield* directory.getBinding(threadId); + const binding = Option.getOrUndefined(bindingOption); + if (!binding || binding.providerInstanceId === undefined) { + return; + } + + const providerInstanceId = binding.providerInstanceId; + + const now = yield* nowIso; + switch (event.type) { + case "turn.started": + if (event.turnId === undefined) { + return; + } + yield* directory.upsert({ + threadId, + provider: binding.provider, + providerInstanceId, + status: "running", + runtimePayload: { + activeTurnId: event.turnId, + lastRuntimeEvent: "provider.turn.started", + lastRuntimeEventAt: now, + }, + }); + return; + case "turn.completed": + yield* directory.upsert({ + threadId, + provider: binding.provider, + providerInstanceId, + status: binding.status === "starting" ? "starting" : "running", + runtimePayload: { + activeTurnId: null, + lastRuntimeEvent: "provider.turn.completed", + lastRuntimeEventAt: now, + }, + }); + return; + case "session.exited": + yield* directory.upsert({ + threadId, + provider: binding.provider, + providerInstanceId, + status: "stopped", + runtimePayload: { + activeTurnId: null, + lastRuntimeEvent: "provider.session.exited", + lastRuntimeEventAt: now, + }, + }); + return; + default: + return; + } + }).pipe( + Effect.catchCause((cause) => + Effect.logWarning("provider runtime binding sync failed", { + eventType: event.type, + threadId: event.threadId, + cause: Cause.pretty(cause), + }), + ), + ); + const processRuntimeEvent = ( source: { readonly instanceId: ProviderInstanceId; @@ -290,10 +365,15 @@ const makeProviderService = Effect.fn("makeProviderService")(function* ( ): Effect.Effect => Effect.sync(() => correlateRuntimeEventWithInstance(source, event)).pipe( Effect.flatMap((canonicalEvent) => - increment(providerRuntimeEventsTotal, { - provider: canonicalEvent.provider, - eventType: canonicalEvent.type, - }).pipe(Effect.andThen(publishRuntimeEvent(canonicalEvent))), + syncPersistedRuntimeFromLifecycleEvent(canonicalEvent).pipe( + Effect.andThen( + increment(providerRuntimeEventsTotal, { + provider: canonicalEvent.provider, + eventType: canonicalEvent.type, + }), + ), + Effect.andThen(publishRuntimeEvent(canonicalEvent)), + ), ), ); @@ -683,15 +763,16 @@ const makeProviderService = Effect.fn("makeProviderService")(function* ( ...(input.modelSelection?.model ? { "provider.model": input.modelSelection.model } : {}), }); const turn = yield* routed.adapter.sendTurn(input); + // Persist resume cursor and bookkeeping only. Turn/session status and + // activeTurnId are owned by lifecycle events (turn.started/completed). yield* directory.upsert({ threadId: input.threadId, provider: routed.adapter.provider, providerInstanceId: routed.instanceId, - status: "running", ...(turn.resumeCursor !== undefined ? { resumeCursor: turn.resumeCursor } : {}), runtimePayload: { ...(input.modelSelection !== undefined ? { modelSelection: input.modelSelection } : {}), - activeTurnId: turn.turnId, + activeTurnId: null, lastRuntimeEvent: "provider.sendTurn", lastRuntimeEventAt: yield* nowIso, }, @@ -732,7 +813,9 @@ const makeProviderService = Effect.fn("makeProviderService")(function* ( const routed = yield* resolveRoutableSession({ threadId: input.threadId, operation: "ProviderService.interruptTurn", - allowRecovery: true, + // Stop must not resurrect a dead provider process; session settle + // happens in ProviderCommandReactor when interrupt is a no-op. + allowRecovery: false, }); metricProvider = routed.adapter.provider; yield* Effect.annotateCurrentSpan({ @@ -741,7 +824,9 @@ const makeProviderService = Effect.fn("makeProviderService")(function* ( "provider.thread_id": input.threadId, "provider.turn_id": input.turnId, }); - yield* routed.adapter.interruptTurn(routed.threadId, input.turnId); + if (routed.isActive) { + yield* routed.adapter.interruptTurn(routed.threadId, input.turnId); + } yield* analytics.record("provider.turn.interrupted", { provider: routed.adapter.provider, }); diff --git a/apps/server/src/provider/Layers/ProviderSessionDirectory.ts b/apps/server/src/provider/Layers/ProviderSessionDirectory.ts index 23075bd9a06..69c4750636f 100644 --- a/apps/server/src/provider/Layers/ProviderSessionDirectory.ts +++ b/apps/server/src/provider/Layers/ProviderSessionDirectory.ts @@ -125,6 +125,10 @@ const makeProviderSessionDirectory = Effect.gen(function* () { issue: "providerInstanceId is required for provider session runtime bindings.", }); } + const freshExisting = yield* repository + .getByThreadId({ threadId: resolvedThreadId }) + .pipe(Effect.mapError(toPersistenceError("ProviderSessionDirectory.upsert:getByThreadId"))); + const freshRuntime = Option.getOrUndefined(freshExisting); yield* repository .upsert({ threadId: resolvedThreadId, @@ -133,15 +137,19 @@ const makeProviderSessionDirectory = Effect.gen(function* () { adapterKey: binding.adapterKey ?? (providerChanged ? binding.provider : (existingRuntime?.adapterKey ?? binding.provider)), - runtimeMode: binding.runtimeMode ?? existingRuntime?.runtimeMode ?? "full-access", - status: binding.status ?? existingRuntime?.status ?? "running", + runtimeMode: + binding.runtimeMode ?? + freshRuntime?.runtimeMode ?? + existingRuntime?.runtimeMode ?? + "full-access", + status: binding.status ?? freshRuntime?.status ?? existingRuntime?.status ?? "running", lastSeenAt: now, resumeCursor: binding.resumeCursor !== undefined ? binding.resumeCursor - : (existingRuntime?.resumeCursor ?? null), + : (freshRuntime?.resumeCursor ?? existingRuntime?.resumeCursor ?? null), runtimePayload: mergeRuntimePayload( - existingRuntime?.runtimePayload ?? null, + freshRuntime?.runtimePayload ?? existingRuntime?.runtimePayload ?? null, binding.runtimePayload, ), }) diff --git a/apps/server/src/provider/Layers/ProviderSessionRecovery.test.ts b/apps/server/src/provider/Layers/ProviderSessionRecovery.test.ts new file mode 100644 index 00000000000..8e129babc07 --- /dev/null +++ b/apps/server/src/provider/Layers/ProviderSessionRecovery.test.ts @@ -0,0 +1,135 @@ +import { + ProviderDriverKind, + ProviderInstanceId, + ThreadId, + type OrchestrationCommand, +} from "@t3tools/contracts"; +import { describe, it, assert } from "@effect/vitest"; +import * as Crypto from "effect/Crypto"; +import * as Effect from "effect/Effect"; +import * as Layer from "effect/Layer"; +import * as Option from "effect/Option"; +import * as Ref from "effect/Ref"; + +import { OrchestrationEngineService } from "../../orchestration/Services/OrchestrationEngine.ts"; +import { SqlitePersistenceMemory } from "../../persistence/Layers/Sqlite.ts"; +import * as ProviderSessionRuntime from "../../persistence/ProviderSessionRuntime.ts"; +import { ProviderSessionDirectoryLive } from "./ProviderSessionDirectory.ts"; +import { recoverOrphanedProviderSessions } from "./ProviderSessionRecovery.ts"; + +const threadId = ThreadId.make("thread-recovery-1"); +const grokInstanceId = ProviderInstanceId.make("grok"); +const resumeCursor = { schemaVersion: 1, sessionId: "session-recovery-1" }; + +const makeTestLayer = (dispatchedRef: Ref.Ref>) => { + const runtimeRepositoryLayer = ProviderSessionRuntime.layer.pipe( + Layer.provide(SqlitePersistenceMemory), + ); + const directoryLayer = ProviderSessionDirectoryLive.pipe(Layer.provide(runtimeRepositoryLayer)); + const engineLayer = Layer.succeed(OrchestrationEngineService, { + dispatch: (command) => + Ref.update(dispatchedRef, (current) => [...current, command]).pipe( + Effect.as({ sequence: 1 }), + ), + readEvents: () => { + throw new Error("not implemented in test"); + }, + streamDomainEvents: Effect.die("not implemented in test") as never, + }); + const cryptoLayer = Layer.succeed( + Crypto.Crypto, + Crypto.make({ + randomBytes: (size) => new Uint8Array(size), + digest: () => Effect.succeed(new Uint8Array(32)), + }), + ); + + return Layer.mergeAll(directoryLayer, runtimeRepositoryLayer, engineLayer, cryptoLayer); +}; + +describe("ProviderSessionRecovery", () => { + it.effect("settles orphaned runtime bindings and marks them stopped on startup", () => + Effect.gen(function* () { + const dispatchedRef = yield* Ref.make>([]); + const testLayer = makeTestLayer(dispatchedRef); + + const dispatched = yield* Effect.gen(function* () { + const repository = yield* ProviderSessionRuntime.ProviderSessionRuntimeRepository; + yield* repository.upsert({ + threadId, + providerName: ProviderDriverKind.make("grok"), + providerInstanceId: grokInstanceId, + adapterKey: "grok", + runtimeMode: "full-access", + status: "running", + lastSeenAt: "2026-06-29T17:54:19.390Z", + resumeCursor, + runtimePayload: { + activeTurnId: "turn-orphaned", + lastRuntimeEvent: "provider.sendTurn", + }, + }); + + yield* recoverOrphanedProviderSessions; + + const persisted = yield* repository.getByThreadId({ threadId }); + assert.equal(Option.isSome(persisted), true); + if (Option.isSome(persisted)) { + assert.equal(persisted.value.status, "stopped"); + assert.deepEqual(persisted.value.resumeCursor, resumeCursor); + const payload = persisted.value.runtimePayload; + assert.equal( + payload !== null && typeof payload === "object" && !Array.isArray(payload), + true, + ); + if (payload !== null && typeof payload === "object" && !Array.isArray(payload)) { + assert.equal((payload as { activeTurnId: string | null }).activeTurnId, null); + assert.equal( + (payload as { lastRuntimeEvent: string }).lastRuntimeEvent, + "provider.session.recovered", + ); + } + } + + return yield* Ref.get(dispatchedRef); + }).pipe(Effect.provide(testLayer)); + + assert.equal(dispatched.length, 1); + assert.equal(dispatched[0]?.type, "thread.session.set"); + if (dispatched[0]?.type === "thread.session.set") { + assert.equal(dispatched[0].session.status, "ready"); + assert.equal(dispatched[0].session.activeTurnId, null); + assert.equal(dispatched[0].session.providerName, "grok"); + } + }), + ); + + it.effect("skips already-terminal runtime bindings", () => + Effect.gen(function* () { + const dispatchedRef = yield* Ref.make>([]); + const testLayer = makeTestLayer(dispatchedRef); + + const dispatched = yield* Effect.gen(function* () { + const repository = yield* ProviderSessionRuntime.ProviderSessionRuntimeRepository; + yield* repository.upsert({ + threadId, + providerName: ProviderDriverKind.make("cursor"), + providerInstanceId: ProviderInstanceId.make("cursor"), + adapterKey: "cursor", + runtimeMode: "full-access", + status: "stopped", + lastSeenAt: "2026-06-29T17:56:45.038Z", + resumeCursor, + runtimePayload: { + activeTurnId: null, + }, + }); + + yield* recoverOrphanedProviderSessions; + return yield* Ref.get(dispatchedRef); + }).pipe(Effect.provide(testLayer)); + + assert.equal(dispatched.length, 0); + }), + ); +}); diff --git a/apps/server/src/provider/Layers/ProviderSessionRecovery.ts b/apps/server/src/provider/Layers/ProviderSessionRecovery.ts new file mode 100644 index 00000000000..d2e1dff78cf --- /dev/null +++ b/apps/server/src/provider/Layers/ProviderSessionRecovery.ts @@ -0,0 +1,161 @@ +import { CommandId, type OrchestrationSession } from "@t3tools/contracts"; +import * as Cause from "effect/Cause"; +import * as Crypto from "effect/Crypto"; +import * as DateTime from "effect/DateTime"; +import * as Effect from "effect/Effect"; +import * as Layer from "effect/Layer"; + +import { OrchestrationEngineService } from "../../orchestration/Services/OrchestrationEngine.ts"; +import { ProviderSessionDirectory } from "../Services/ProviderSessionDirectory.ts"; + +/** + * ProviderSessionRecovery — settle provider sessions orphaned by a restart. + * + * Provider sessions are backed by child processes (`grok agent stdio`, + * `claude`, `codex app-server`, …) that do not survive a server restart, and + * can also be killed out from under a still-running server. The persisted + * `provider_session_runtime` row, however, keeps the last known status — so a + * session that was mid-turn when its provider went away stays marked + * `running`/`starting` with an unsettled turn. + * + * On the projection side a turn only settles when its session leaves the + * `running` status (see `settledTurnStateForSessionStatus` in + * `ProjectionPipeline`). With no live provider left to emit that transition, + * the thread is stuck showing a phantom "running" turn forever and the UI will + * not let the user continue talking. + * + * This recovery runs once on startup. For every persisted session still in a + * non-terminal status it: + * 1. dispatches a `thread.session.set` with status `ready` and no active + * turn, which settles the orphaned running turn (projected as + * `interrupted`) and unsticks the thread; and + * 2. marks the persisted runtime row `stopped` while preserving the resume + * cursor, so the next prompt resumes the existing provider session + * (continues the conversation) instead of resuming a dead one. + * + * It is best-effort and idempotent: terminal rows are skipped, re-running is a + * no-op, and a failure for one thread is logged without blocking the others or + * server startup. + * + * @module provider/Layers/ProviderSessionRecovery + */ + +/** + * Persisted runtime statuses that imply a live provider process. After a + * restart there is none, so these are the rows that need settling. `stopped` + * and `error` are already terminal and left untouched. + */ +const NON_TERMINAL_RUNTIME_STATUSES = new Set(["starting", "running"]); + +function readPersistedActiveTurnId(runtimePayload: unknown): string | null { + if ( + runtimePayload === null || + typeof runtimePayload !== "object" || + Array.isArray(runtimePayload) + ) { + return null; + } + const activeTurnId = (runtimePayload as { activeTurnId?: unknown }).activeTurnId; + return typeof activeTurnId === "string" && activeTurnId.length > 0 ? activeTurnId : null; +} + +export const recoverOrphanedProviderSessions = Effect.gen(function* () { + const directory = yield* ProviderSessionDirectory; + const engine = yield* OrchestrationEngineService; + const crypto = yield* Crypto.Crypto; + + const bindings = yield* directory.listBindings().pipe( + Effect.catchCause((cause) => + Effect.logWarning("provider session recovery failed to list bindings", { + cause: Cause.pretty(cause), + }).pipe(Effect.as([])), + ), + ); + + const orphaned = bindings.filter((binding) => { + if (binding.status !== undefined && NON_TERMINAL_RUNTIME_STATUSES.has(binding.status)) { + return true; + } + // Split-brain after restart: projection may already be ready while the + // persisted runtime row still carries a stale active turn id. + return readPersistedActiveTurnId(binding.runtimePayload) !== null; + }); + + if (orphaned.length === 0) { + return; + } + + yield* Effect.logInfo("recovering orphaned provider sessions on startup", { + count: orphaned.length, + }); + + yield* Effect.forEach( + orphaned, + (binding) => + Effect.gen(function* () { + const now = DateTime.formatIso(yield* DateTime.now); + const commandId = CommandId.make( + `server:provider-session-recovery:${yield* crypto.randomUUIDv4}`, + ); + const session: OrchestrationSession = { + threadId: binding.threadId, + status: "ready", + providerName: binding.provider, + ...(binding.providerInstanceId !== undefined + ? { providerInstanceId: binding.providerInstanceId } + : {}), + runtimeMode: binding.runtimeMode ?? "full-access", + // Provider turn ids are not orchestration turn ids; clearing the + // active turn is what settles the orphaned running turn. + activeTurnId: null, + lastError: null, + updatedAt: now, + }; + + // Settle the orphaned running turn (UI unstick). + yield* engine.dispatch({ + type: "thread.session.set", + commandId, + threadId: binding.threadId, + session, + createdAt: now, + }); + + // Mark the runtime row stopped, preserving the resume cursor so the + // next prompt resumes this session instead of starting a dead one. + yield* directory.upsert({ + threadId: binding.threadId, + provider: binding.provider, + ...(binding.providerInstanceId !== undefined + ? { providerInstanceId: binding.providerInstanceId } + : {}), + status: "stopped", + runtimePayload: { + activeTurnId: null, + lastRuntimeEvent: "provider.session.recovered", + lastRuntimeEventAt: now, + }, + }); + + yield* Effect.logInfo("recovered orphaned provider session", { + threadId: binding.threadId, + provider: binding.provider, + }); + }).pipe( + Effect.catchCause((cause) => + Effect.logWarning("provider session recovery failed for thread", { + threadId: binding.threadId, + cause: Cause.pretty(cause), + }), + ), + ), + { concurrency: 1, discard: true }, + ); +}); + +/** + * Startup layer that runs {@link recoverOrphanedProviderSessions} once when + * built. Requires `ProviderSessionDirectory` and `OrchestrationEngineService` + * (both provided by the server runtime services layer). + */ +export const ProviderSessionRecoveryLive = Layer.effectDiscard(recoverOrphanedProviderSessions); diff --git a/apps/server/src/provider/acp/AcpRuntimeModel.test.ts b/apps/server/src/provider/acp/AcpRuntimeModel.test.ts index 7682c5f5f9c..14b36a43501 100644 --- a/apps/server/src/provider/acp/AcpRuntimeModel.test.ts +++ b/apps/server/src/provider/acp/AcpRuntimeModel.test.ts @@ -8,6 +8,7 @@ import { parsePermissionRequest, parseSessionModeState, parseSessionUpdateEvent, + sessionLoadReplaySettleActivityMillis, sessionUpdateIsReplay, syntheticLoadSessionResponseFromInitialize, } from "./AcpRuntimeModel.ts"; @@ -336,6 +337,16 @@ describe("AcpRuntimeModel", () => { ]); }); + it("requires replay idle after load RPC completion even when replay stopped earlier", () => { + const idleGapMillis = 2_000; + const replayAt = 0; + const rpcAt = 2_000; + const settleFrom = sessionLoadReplaySettleActivityMillis(replayAt, rpcAt); + + expect(rpcAt - settleFrom >= idleGapMillis).toBe(false); + expect(rpcAt + idleGapMillis - settleFrom >= idleGapMillis).toBe(true); + }); + it("keeps permission request parsing compatible with loose extension payloads", () => { const request = parsePermissionRequest({ sessionId: "session-1", diff --git a/apps/server/src/provider/acp/AcpRuntimeModel.ts b/apps/server/src/provider/acp/AcpRuntimeModel.ts index e6bfc127e6e..df6b07563fc 100644 --- a/apps/server/src/provider/acp/AcpRuntimeModel.ts +++ b/apps/server/src/provider/acp/AcpRuntimeModel.ts @@ -465,6 +465,62 @@ export interface SessionLoadGate { readonly initializeResult: EffectAcpSchema.InitializeResponse; } +export function sessionLoadReplaySettleActivityMillis( + lastActivityAtMillis: number | undefined, + baselineMillis: number, +): number { + return Math.max(lastActivityAtMillis ?? 0, baselineMillis); +} + +/** + * Waits until session/load replay notifications have been quiet for `idleGap`. + * `baselineMillis` is recorded when the load RPC completes. The idle window must + * start after both the last replay activity and RPC completion, so a slow RPC + * cannot inherit replay quiet time that elapsed while the request was in flight. + */ +export const waitForSessionLoadReplayToSettle = (input: { + readonly gateRef: Ref.Ref>; + readonly baselineMillis: number; +}): Effect.Effect => + Effect.gen(function* () { + const pollInterval = Duration.millis(25); + while (true) { + const gate = yield* Ref.get(input.gateRef); + if (Option.isNone(gate) || !gate.value.active) { + return; + } + const idleGapMillis = Duration.toMillis(gate.value.idleGap); + const nowMillis = yield* Clock.currentTimeMillis; + const lastActivity = sessionLoadReplaySettleActivityMillis( + gate.value.lastActivityAtMillis, + input.baselineMillis, + ); + if (nowMillis - lastActivity >= idleGapMillis) { + return; + } + yield* Effect.sleep(pollInterval); + } + }); + +/** Records inbound protocol activity while the session/load replay gate is active. */ +export const touchSessionLoadReplayActivity = (input: { + readonly gateRef: Ref.Ref>; +}): Effect.Effect => + Effect.gen(function* () { + const gate = yield* Ref.get(input.gateRef); + if (Option.isNone(gate) || !gate.value.active) { + return; + } + const lastActivityAtMillis = yield* Clock.currentTimeMillis; + yield* Ref.set( + input.gateRef, + Option.some({ + ...gate.value, + lastActivityAtMillis, + }), + ); + }); + export const waitForSessionLoadReplayIdle = (input: { readonly gateRef: Ref.Ref>; }): Effect.Effect => diff --git a/apps/server/src/provider/acp/AcpSessionRuntime.ts b/apps/server/src/provider/acp/AcpSessionRuntime.ts index bc2df3aa8d4..120f156944b 100644 --- a/apps/server/src/provider/acp/AcpSessionRuntime.ts +++ b/apps/server/src/provider/acp/AcpSessionRuntime.ts @@ -4,6 +4,7 @@ import * as Context from "effect/Context"; import * as Deferred from "effect/Deferred"; import * as Duration from "effect/Duration"; import * as Effect from "effect/Effect"; +import * as Exit from "effect/Exit"; import * as Fiber from "effect/Fiber"; import * as Layer from "effect/Layer"; import * as Option from "effect/Option"; @@ -29,6 +30,8 @@ import { parseSessionUpdateEvent, sessionUpdateIsReplay, waitForSessionLoadReplayIdle, + touchSessionLoadReplayActivity, + waitForSessionLoadReplayToSettle, type SessionLoadGate, type AcpParsedSessionEvent, type AcpSessionModeState, @@ -242,6 +245,10 @@ export class AcpSessionRuntime extends Context.Service< method: string, payload: unknown, ) => Effect.Effect; + /** True while session/load replay is still settling and should suppress side effects. */ + readonly isSessionLoadReplayActive: Effect.Effect; + /** Extends the session/load replay idle window when non-session/update traffic arrives. */ + readonly touchSessionLoadReplayActivity: Effect.Effect; } >()("t3/provider/acp/AcpSessionRuntime") {} @@ -360,14 +367,7 @@ export const make = ( Effect.gen(function* () { const gate = yield* Ref.get(sessionLoadGateRef); if (Option.isSome(gate) && gate.value.active) { - const lastActivityAtMillis = yield* Clock.currentTimeMillis; - yield* Ref.set( - sessionLoadGateRef, - Option.some({ - ...gate.value, - lastActivityAtMillis, - }), - ); + yield* touchSessionLoadReplayActivity({ gateRef: sessionLoadGateRef }); return; } if (sessionUpdateIsReplay(notification)) { @@ -578,11 +578,35 @@ export const make = ( const idleFiber = yield* waitForSessionLoadReplayIdle({ gateRef: sessionLoadGateRef, }).pipe(Effect.forkIn(runtimeScope)); - const loaded = yield* Effect.raceFirst( - acp.agent.loadSession(loadPayload), - Fiber.join(idleFiber), - ).pipe( - Effect.ensuring(Fiber.interrupt(idleFiber).pipe(Effect.ignore)), + const rpcFiber = yield* acp.agent + .loadSession(loadPayload) + .pipe(Effect.forkIn(runtimeScope)); + const loaded = yield* Effect.gen(function* () { + const winner = yield* Effect.race( + Fiber.await(rpcFiber).pipe(Effect.map((exit) => ({ _tag: "Rpc" as const, exit }))), + Fiber.join(idleFiber).pipe( + Effect.map((synthetic) => ({ _tag: "Idle" as const, synthetic })), + ), + ); + + if (winner._tag === "Idle") { + yield* Fiber.interrupt(rpcFiber).pipe(Effect.ignore); + return winner.synthetic; + } + + if (Exit.isFailure(winner.exit)) { + yield* Fiber.interrupt(idleFiber).pipe(Effect.ignore); + return yield* Effect.failCause(winner.exit.cause); + } + + const rpcCompletedAtMillis = yield* Clock.currentTimeMillis; + yield* waitForSessionLoadReplayToSettle({ + gateRef: sessionLoadGateRef, + baselineMillis: rpcCompletedAtMillis, + }); + yield* Fiber.interrupt(idleFiber).pipe(Effect.ignore); + return winner.exit.value; + }).pipe( Effect.timeoutOption(sessionLoadTimeout), Effect.flatMap((result) => Option.match(result, { @@ -794,6 +818,12 @@ export const make = ( request: (method, payload) => runLoggedRequest(method, payload, acp.raw.request(method, payload)), notify: acp.raw.notify, + isSessionLoadReplayActive: Ref.get(sessionLoadGateRef).pipe( + Effect.map((gate) => Option.isSome(gate) && gate.value.active), + ), + touchSessionLoadReplayActivity: touchSessionLoadReplayActivity({ + gateRef: sessionLoadGateRef, + }), } satisfies AcpSessionRuntime["Service"]; }); diff --git a/apps/server/src/serverRuntimeStartup.ts b/apps/server/src/serverRuntimeStartup.ts index b52b577c5b5..c9eeebf233f 100644 --- a/apps/server/src/serverRuntimeStartup.ts +++ b/apps/server/src/serverRuntimeStartup.ts @@ -33,6 +33,7 @@ import * as ServerSettings from "./serverSettings.ts"; import * as AnalyticsService from "./telemetry/AnalyticsService.ts"; import * as ServerEnvironment from "./environment/ServerEnvironment.ts"; import * as EnvironmentAuth from "./auth/EnvironmentAuth.ts"; +import { recoverOrphanedProviderSessions } from "./provider/Layers/ProviderSessionRecovery.ts"; import * as ProviderSessionReaper from "./provider/Services/ProviderSessionReaper.ts"; import { formatHeadlessServeOutput, @@ -342,6 +343,13 @@ export const make = Effect.gen(function* () { "reactors.start", Effect.gen(function* () { yield* orchestrationReactor.start().pipe(Scope.provide(reactorScope)); + yield* recoverOrphanedProviderSessions.pipe( + Effect.catchCause((cause) => + Effect.logWarning("provider session recovery failed on startup", { + cause, + }), + ), + ); yield* providerSessionReaper.start().pipe(Scope.provide(reactorScope)); }), ); diff --git a/packages/contracts/src/orchestration.ts b/packages/contracts/src/orchestration.ts index 623fed0917b..266b8676325 100644 --- a/packages/contracts/src/orchestration.ts +++ b/packages/contracts/src/orchestration.ts @@ -725,6 +725,15 @@ const ThreadMessageAssistantCompleteCommand = Schema.Struct({ createdAt: IsoDateTime, }); +const ThreadMessageAssistantResetCommand = Schema.Struct({ + type: Schema.Literal("thread.message.assistant.reset"), + commandId: CommandId, + threadId: ThreadId, + messageId: MessageId, + turnId: Schema.optional(TurnId), + createdAt: IsoDateTime, +}); + const ThreadProposedPlanUpsertCommand = Schema.Struct({ type: Schema.Literal("thread.proposed-plan.upsert"), commandId: CommandId, @@ -767,6 +776,7 @@ const InternalOrchestrationCommand = Schema.Union([ ThreadSessionSetCommand, ThreadMessageAssistantDeltaCommand, ThreadMessageAssistantCompleteCommand, + ThreadMessageAssistantResetCommand, ThreadProposedPlanUpsertCommand, ThreadTurnDiffCompleteCommand, ThreadActivityAppendCommand, @@ -898,6 +908,7 @@ export const ThreadMessageSentPayload = Schema.Struct({ attachments: Schema.optional(Schema.Array(ChatAttachment)), turnId: Schema.NullOr(TurnId), streaming: Schema.Boolean, + reset: Schema.optional(Schema.Boolean), createdAt: IsoDateTime, updatedAt: IsoDateTime, });