diff --git a/apps/server/src/orchestration/Layers/OrphanSessionRecovery.ts b/apps/server/src/orchestration/Layers/OrphanSessionRecovery.ts index bddeac5bb1fa..66612d67b02e 100644 --- a/apps/server/src/orchestration/Layers/OrphanSessionRecovery.ts +++ b/apps/server/src/orchestration/Layers/OrphanSessionRecovery.ts @@ -21,7 +21,6 @@ import { type OrphanSessionRecoveryShape, } from "../Services/OrphanSessionRecovery.ts"; import { ProjectionSnapshotQuery } from "../Services/ProjectionSnapshotQuery.ts"; -import { readProviderRestartRecoveryMarker } from "../../provider/ProviderRestartRecovery.ts"; import { ProviderService } from "../../provider/Services/ProviderService.ts"; import { ProviderSessionDirectory } from "../../provider/Services/ProviderSessionDirectory.ts"; @@ -228,25 +227,14 @@ const make = Effect.gen(function* () { .pipe(Effect.orElseSucceed(() => ({ threads: [] as const }))); const bindings = yield* directory.listBindings().pipe(Effect.orElseSucceed(() => [])); const threadIds = new Set(); - const recoveryMarkedThreadIds = new Set( - bindings.flatMap((binding) => - readProviderRestartRecoveryMarker(binding.runtimePayload) === undefined - ? [] - : [String(binding.threadId)], - ), - ); let settledSessions = 0; let interruptedSessions = 0; for (const thread of snapshot.threads) { - if (recoveryMarkedThreadIds.has(String(thread.id))) { - continue; - } const claimsLive = isLiveClaimingSessionStatus(thread.session?.status); const processIsLive = claimsLive ? yield* hasLiveProcess(thread.id) : false; - // Provider restart reconciliation runs before this audit and may - // already have resumed the thread. Never classify that replacement - // process as an orphan merely because its projected session is live. + // A replacement process started after boot (for example #9167 + // continuation) is live — never classify that as an orphan. if ( !shouldSettleAfterServerRestart({ claimsLive, @@ -269,9 +257,6 @@ const make = Effect.gen(function* () { let settledRuntimes = 0; for (const binding of bindings) { - if (recoveryMarkedThreadIds.has(String(binding.threadId))) { - continue; - } const claimsLive = binding.status === "running" || binding.status === "starting"; if (!claimsLive) { continue; diff --git a/apps/server/src/orchestration/Layers/ProviderCommandReactor.test.ts b/apps/server/src/orchestration/Layers/ProviderCommandReactor.test.ts index 10effc900bc3..0b4e18fc1b09 100644 --- a/apps/server/src/orchestration/Layers/ProviderCommandReactor.test.ts +++ b/apps/server/src/orchestration/Layers/ProviderCommandReactor.test.ts @@ -17,9 +17,7 @@ import { DEFAULT_PROVIDER_INTERACTION_MODE, EnvironmentId, EventId, - IdentityUsername, MessageId, - PersonId, ProjectId, ThreadId, TurnId, @@ -62,10 +60,8 @@ import * as ThreadPlanProgress from "../ThreadPlanProgress.ts"; import { providerErrorLabel, providerErrorLabelFromInstanceHint, - RESTART_RECOVERY_CONTINUATION_INSTRUCTION, ProviderCommandReactorLive, } from "./ProviderCommandReactor.ts"; -import { makeProviderRestartRecoveryMarker } from "../../provider/ProviderRestartRecovery.ts"; import { OrchestrationEngineService } from "../Services/OrchestrationEngine.ts"; import { ProviderCommandReactor } from "../Services/ProviderCommandReactor.ts"; import { ProjectionSnapshotQuery } from "../Services/ProjectionSnapshotQuery.ts"; @@ -814,751 +810,6 @@ describe("ProviderCommandReactor", () => { expect(pendingAfter).toEqual([]); }); - it("continues an interrupted running turn without replaying its user message", async () => { - const modelSelection: ModelSelection = { - instanceId: ProviderInstanceId.make("codex"), - model: "gpt-5-recovery", - }; - const threadId = ThreadId.make("thread-1"); - const harness = await createHarness({ - deferReactorStart: true, - threadModelSelection: modelSelection, - providerBindings: [ - { - threadId, - provider: ProviderDriverKind.make("codex"), - providerInstanceId: ProviderInstanceId.make("codex"), - adapterKey: "codex", - runtimeMode: "full-access", - status: "stopped", - resumeCursor: { threadId: "provider-thread-resume" }, - runtimePayload: { - cwd: "/tmp/persisted-recovery-cwd", - modelSelection, - interactionMode: "plan", - activeTurnId: null, - restartRecovery: makeProviderRestartRecoveryMarker({ - interruptedProviderTurnId: asTurnId("provider-turn-before-restart"), - shutdownAt: "2026-01-01T00:00:01.000Z", - }), - }, - lastSeenAt: "2026-01-01T00:00:01.000Z", - }, - ], - }); - await Effect.runPromise( - harness.engine.dispatch({ - type: "thread.turn.start", - commandId: CommandId.make("cmd-running-before-restart"), - threadId, - message: { - messageId: asMessageId("running-user-message"), - role: "user", - text: "the original request must not be replayed", - attachments: [], - }, - modelSelection, - interactionMode: "plan", - runtimeMode: "full-access", - createdAt: "2026-01-01T00:00:00.000Z", - }), - ); - await Effect.runPromise( - harness.engine.dispatch({ - type: "thread.session.set", - commandId: CommandId.make("server:running-before-restart"), - threadId, - session: { - threadId, - status: "running", - providerName: "codex", - providerInstanceId: ProviderInstanceId.make("codex"), - runtimeMode: "full-access", - activeTurnId: asTurnId("orchestration-turn-before-restart"), - lastError: null, - updatedAt: "2026-01-01T00:00:00.500Z", - }, - createdAt: "2026-01-01T00:00:00.500Z", - }), - ); - - await harness.startReactor(); - await harness.drain(); - await waitFor(() => harness.sendTurn.mock.calls.length === 1); - - expect(harness.startSession.mock.calls[0]?.[1]).toMatchObject({ - cwd: "/tmp/persisted-recovery-cwd", - modelSelection, - resumeCursor: { threadId: "provider-thread-resume" }, - runtimeMode: "full-access", - providerInstanceId: ProviderInstanceId.make("codex"), - }); - expect(harness.sendTurn.mock.calls[0]?.[0]).toEqual({ - threadId, - input: RESTART_RECOVERY_CONTINUATION_INSTRUCTION, - attachments: [], - modelSelection, - interactionMode: "plan", - }); - const readModel = await harness.readModel(); - const thread = readModel.threads.find((entry) => entry.id === threadId); - const oldTurn = (await harness.readTurns(threadId)).find( - (turn) => turn.turnId === asTurnId("orchestration-turn-before-restart"), - ); - expect(oldTurn?.state).toBe("interrupted"); - expect(thread?.messages.map((message) => message.text)).toEqual([ - "the original request must not be replayed", - ]); - expect(harness.providerBindings.get(threadId)?.runtimePayload).toMatchObject({ - restartRecovery: null, - interactionMode: "plan", - }); - }); - - it("recovers a live shell turn when the adapter binding looks idle after prompt_complete", async () => { - const modelSelection: ModelSelection = { - instanceId: ProviderInstanceId.make("codex"), - model: "gpt-5-recovery", - }; - const threadId = ThreadId.make("thread-1"); - const harness = await createHarness({ - deferReactorStart: true, - threadModelSelection: modelSelection, - providerBindings: [ - { - threadId, - provider: ProviderDriverKind.make("codex"), - providerInstanceId: ProviderInstanceId.make("codex"), - adapterKey: "codex", - runtimeMode: "full-access", - status: "running", - resumeCursor: { threadId: "provider-thread-resume" }, - runtimePayload: { - cwd: "/tmp/persisted-recovery-cwd", - modelSelection, - interactionMode: "plan", - activeTurnId: null, - }, - lastSeenAt: "2026-01-01T00:00:01.000Z", - }, - ], - }); - await Effect.runPromise( - harness.engine.dispatch({ - type: "thread.turn.start", - commandId: CommandId.make("cmd-live-shell-after-prompt-complete"), - threadId, - message: { - messageId: asMessageId("live-shell-user-message"), - role: "user", - text: "keep going after the server restart", - attachments: [], - }, - modelSelection, - interactionMode: "plan", - runtimeMode: "full-access", - createdAt: "2026-01-01T00:00:00.000Z", - }), - ); - await Effect.runPromise( - harness.engine.dispatch({ - type: "thread.session.set", - commandId: CommandId.make("server:live-shell-before-restart"), - threadId, - session: { - threadId, - status: "running", - providerName: "codex", - providerInstanceId: ProviderInstanceId.make("codex"), - runtimeMode: "full-access", - activeTurnId: asTurnId("orchestration-turn-still-live"), - lastError: null, - updatedAt: "2026-01-01T00:00:00.500Z", - }, - createdAt: "2026-01-01T00:00:00.500Z", - }), - ); - - await harness.startReactor(); - await harness.drain(); - await waitFor(() => harness.sendTurn.mock.calls.length === 1); - - expect(harness.sendTurn.mock.calls[0]?.[0]).toMatchObject({ - threadId, - input: RESTART_RECOVERY_CONTINUATION_INSTRUCTION, - }); - }); - - it("sends desktop turns with identity-map commit attribution", async () => { - const harness = await createHarness({ - identityPeople: [ - { - personId: "patroza", - username: "patroza", - name: "Patrick Roza", - github: { login: "patroza", id: "42661" }, - }, - ], - }); - await harness.dispatch({ - type: "thread.turn.start", - commandId: CommandId.make("cmd-attributed-turn"), - threadId: ThreadId.make("thread-1"), - message: { - messageId: asMessageId("attributed-user-message"), - role: "user", - text: "ship this", - attachments: [], - }, - source: { - channel: "desktop", - personId: PersonId.make("patroza"), - username: IdentityUsername.make("patroza"), - }, - interactionMode: DEFAULT_PROVIDER_INTERACTION_MODE, - runtimeMode: "approval-required", - createdAt: "2026-01-01T00:00:00.000Z", - }); - - await waitFor(() => harness.sendTurn.mock.calls.length === 1); - expect(harness.sendTurn.mock.calls[0]?.[0]).toMatchObject({ - input: expect.stringContaining( - "Co-authored-by: Patrick Roza <42661+patroza@users.noreply.github.com>", - ), - }); - }); - - it("does not finish reactor startup before restart recovery claims complete", async () => { - const reconciliationGate = Effect.runSync(Deferred.make()); - const harness = await createHarness({ - deferReactorStart: true, - listBindingsEffect: () => Deferred.await(reconciliationGate).pipe(Effect.as([] as const)), - }); - - let startupFinished = false; - const startup = harness.startReactor().then(() => { - startupFinished = true; - }); - await Effect.runPromise(Effect.yieldNow); - expect(startupFinished).toBe(false); - - Effect.runSync(Deferred.succeed(reconciliationGate, undefined)); - await startup; - expect(startupFinished).toBe(true); - }); - - it("finishes reactor startup while recovery provider session start is still pending", async () => { - const modelSelection: ModelSelection = { - instanceId: ProviderInstanceId.make("codex"), - model: "gpt-5-recovery", - }; - const threadId = ThreadId.make("thread-1"); - const activation = Effect.runSync(Deferred.make()); - const releaseStart = Effect.runSync(Deferred.make()); - const harness = await createHarness({ - deferReactorStart: true, - threadModelSelection: modelSelection, - // Production path: recovery continuations park on ServerActivation. - serverActivation: Deferred.await(activation), - startSessionEffect: (session) => Deferred.await(releaseStart).pipe(Effect.as(session)), - providerBindings: [ - { - threadId, - provider: ProviderDriverKind.make("codex"), - providerInstanceId: ProviderInstanceId.make("codex"), - adapterKey: "codex", - runtimeMode: "full-access", - status: "stopped", - resumeCursor: { threadId: "provider-thread-resume" }, - runtimePayload: { - cwd: "/tmp/persisted-recovery-cwd", - modelSelection, - interactionMode: "plan", - activeTurnId: null, - restartRecovery: makeProviderRestartRecoveryMarker({ - interruptedProviderTurnId: asTurnId("provider-turn-before-restart"), - shutdownAt: "2026-01-01T00:00:01.000Z", - }), - }, - lastSeenAt: "2026-01-01T00:00:01.000Z", - }, - ], - }); - await Effect.runPromise( - harness.engine.dispatch({ - type: "thread.turn.start", - commandId: CommandId.make("cmd-running-before-restart-hang"), - threadId, - message: { - messageId: asMessageId("running-user-message-hang"), - role: "user", - text: "original request", - attachments: [], - }, - modelSelection, - interactionMode: "plan", - runtimeMode: "full-access", - createdAt: "2026-01-01T00:00:00.000Z", - }), - ); - - // Claim phase must complete without waiting on activation or provider start. - await harness.startReactor(); - expect(harness.startSession).not.toHaveBeenCalled(); - expect(harness.sendTurn).not.toHaveBeenCalled(); - expect(harness.providerBindings.get(threadId)?.runtimePayload).toMatchObject({ - restartRecovery: expect.objectContaining({ version: 1 }), - lastRuntimeEvent: "provider.restartRecovery.claimed", - }); - - Effect.runSync(Deferred.succeed(activation, undefined)); - await waitFor(() => harness.startSession.mock.calls.length === 1); - expect(harness.sendTurn).not.toHaveBeenCalled(); - - Effect.runSync(Deferred.succeed(releaseStart, undefined)); - await harness.drain(); - await waitFor(() => harness.sendTurn.mock.calls.length === 1); - expect(harness.sendTurn.mock.calls[0]?.[0]).toMatchObject({ - threadId, - input: RESTART_RECOVERY_CONTINUATION_INSTRUCTION, - }); - }); - - it("does not resume settled turns from stale recovery bindings", async () => { - const modelSelection: ModelSelection = { - instanceId: ProviderInstanceId.make("codex"), - model: "gpt-5-codex", - }; - const markerThreadId = ThreadId.make("thread-1"); - const legacyThreadId = ThreadId.make("thread-settled-legacy"); - const markerTurnId = asTurnId("turn-settled-marker"); - const legacyTurnId = asTurnId("turn-settled-legacy"); - const harness = await createHarness({ - deferReactorStart: true, - providerBindings: [ - { - threadId: markerThreadId, - provider: ProviderDriverKind.make("codex"), - providerInstanceId: ProviderInstanceId.make("codex"), - runtimeMode: "full-access", - status: "stopped", - resumeCursor: { threadId: "provider-thread-marker" }, - runtimePayload: { - modelSelection, - restartRecovery: makeProviderRestartRecoveryMarker({ - interruptedProviderTurnId: markerTurnId, - shutdownAt: "2026-01-01T00:00:02.000Z", - }), - }, - lastSeenAt: "2026-01-01T00:00:02.000Z", - }, - { - threadId: legacyThreadId, - provider: ProviderDriverKind.make("codex"), - providerInstanceId: ProviderInstanceId.make("codex"), - runtimeMode: "full-access", - status: "running", - resumeCursor: { threadId: "provider-thread-legacy" }, - runtimePayload: { - modelSelection, - activeTurnId: legacyTurnId, - }, - lastSeenAt: "2026-01-01T00:00:02.000Z", - }, - ], - }); - await Effect.runPromise( - harness.engine.dispatch({ - type: "thread.create", - commandId: CommandId.make("cmd-create-settled-legacy-thread"), - threadId: legacyThreadId, - projectId: asProjectId("project-1"), - title: "Settled legacy thread", - modelSelection, - interactionMode: "default", - runtimeMode: "full-access", - branch: null, - worktreePath: null, - createdAt: "2026-01-01T00:00:00.000Z", - }), - ); - - const settleTurn = async (threadId: ThreadId, turnId: TurnId, suffix: string) => { - await Effect.runPromise( - harness.engine.dispatch({ - type: "thread.turn.start", - commandId: CommandId.make(`cmd-settled-turn-${suffix}`), - threadId, - message: { - messageId: asMessageId(`message-settled-${suffix}`), - role: "user", - text: "already finished", - attachments: [], - }, - modelSelection, - interactionMode: "default", - runtimeMode: "full-access", - createdAt: "2026-01-01T00:00:00.500Z", - }), - ); - await Effect.runPromise( - harness.engine.dispatch({ - type: "thread.session.set", - commandId: CommandId.make(`server:settled-running-${suffix}`), - threadId, - session: { - threadId, - status: "running", - providerName: "codex", - providerInstanceId: ProviderInstanceId.make("codex"), - runtimeMode: "full-access", - activeTurnId: turnId, - lastError: null, - updatedAt: "2026-01-01T00:00:01.000Z", - }, - createdAt: "2026-01-01T00:00:01.000Z", - }), - ); - await Effect.runPromise( - harness.engine.dispatch({ - type: "thread.session.set", - commandId: CommandId.make(`server:settled-ready-${suffix}`), - threadId, - session: { - threadId, - status: "ready", - providerName: "codex", - providerInstanceId: ProviderInstanceId.make("codex"), - runtimeMode: "full-access", - activeTurnId: null, - lastError: null, - updatedAt: "2026-01-01T00:00:01.500Z", - }, - createdAt: "2026-01-01T00:00:01.500Z", - }), - ); - }; - - await settleTurn(markerThreadId, markerTurnId, "marker"); - await settleTurn(legacyThreadId, legacyTurnId, "legacy"); - expect( - (await harness.readTurns(markerThreadId)).find((turn) => turn.turnId === markerTurnId)?.state, - ).toBe("completed"); - expect( - (await harness.readTurns(legacyThreadId)).find((turn) => turn.turnId === legacyTurnId)?.state, - ).toBe("completed"); - - await harness.startReactor(); - await harness.drain(); - - expect(harness.startSession).not.toHaveBeenCalled(); - expect(harness.sendTurn).not.toHaveBeenCalled(); - for (const threadId of [markerThreadId, legacyThreadId]) { - expect(harness.providerBindings.get(threadId)).toMatchObject({ - status: "stopped", - runtimePayload: { - activeTurnId: null, - restartRecovery: null, - lastRuntimeEvent: "provider.restartRecovery.skipped", - }, - }); - } - }); - - it("recovers across temporary-SQLite runtimes and does not repeat a completed recovery", async () => { - const baseDir = NodeFS.mkdtempSync( - NodePath.join(NodeOS.tmpdir(), "t3code-reactor-two-runtime-"), - ); - const threadId = ThreadId.make("thread-1"); - const modelSelection: ModelSelection = { - instanceId: ProviderInstanceId.make("codex"), - model: "gpt-5-codex", - }; - const persistedBindings = new Map(); - const first = await createHarness({ baseDir, providerBindingsMap: persistedBindings }); - await Effect.runPromise( - first.engine.dispatch({ - type: "thread.turn.start", - commandId: CommandId.make("cmd-two-runtime-original-turn"), - threadId, - message: { - messageId: asMessageId("two-runtime-user-message"), - role: "user", - text: "finish this after the server restarts", - attachments: [], - }, - modelSelection, - interactionMode: "default", - runtimeMode: "approval-required", - createdAt: "2026-01-01T00:00:00.000Z", - }), - ); - await waitFor(() => first.sendTurn.mock.calls.length === 1); - await Effect.runPromise( - first.engine.dispatch({ - type: "thread.session.set", - commandId: CommandId.make("server:two-runtime-running"), - threadId, - session: { - threadId, - status: "running", - providerName: "codex", - providerInstanceId: ProviderInstanceId.make("codex"), - runtimeMode: "approval-required", - activeTurnId: asTurnId("orchestration-turn-two-runtime"), - lastError: null, - updatedAt: "2026-01-01T00:00:01.000Z", - }, - createdAt: "2026-01-01T00:00:01.000Z", - }), - ); - await Effect.runPromise( - first.directoryUpsert({ - threadId, - provider: ProviderDriverKind.make("codex"), - providerInstanceId: ProviderInstanceId.make("codex"), - runtimeMode: "approval-required", - status: "stopped", - resumeCursor: { threadId: "durable-provider-thread" }, - runtimePayload: { - cwd: "/tmp/provider-project", - modelSelection, - interactionMode: "default", - restartRecovery: makeProviderRestartRecoveryMarker({ - interruptedProviderTurnId: asTurnId("provider-turn-two-runtime"), - shutdownAt: "2026-01-01T00:00:02.000Z", - }), - }, - }), - ); - - await Effect.runPromise(Scope.close(scope!, Exit.void)); - scope = null; - await runtime!.dispose(); - runtime = null; - - const second = await createHarness({ - baseDir, - providerBindingsMap: persistedBindings, - skipDefaultSetup: true, - }); - await second.drain(); - await waitFor(() => second.sendTurn.mock.calls.length === 1); - expect(second.startSession.mock.calls[0]?.[1]).toMatchObject({ - resumeCursor: { threadId: "durable-provider-thread" }, - cwd: "/tmp/provider-project", - modelSelection, - runtimeMode: "approval-required", - }); - expect(second.sendTurn.mock.calls[0]?.[0]).toMatchObject({ - input: RESTART_RECOVERY_CONTINUATION_INSTRUCTION, - }); - const secondReadModel = await second.readModel(); - expect( - secondReadModel.threads - .find((thread) => thread.id === threadId) - ?.messages.map((message) => message.text), - ).toEqual(["finish this after the server restarts"]); - - await Effect.runPromise( - second.directoryUpsert({ - threadId, - provider: ProviderDriverKind.make("codex"), - providerInstanceId: ProviderInstanceId.make("codex"), - runtimeMode: "approval-required", - status: "running", - runtimePayload: { activeTurnId: null, restartRecovery: null }, - }), - ); - await Effect.runPromise( - second.engine.dispatch({ - type: "thread.session.set", - commandId: CommandId.make("server:two-runtime-recovery-complete"), - threadId, - session: { - threadId, - status: "ready", - providerName: "codex", - providerInstanceId: ProviderInstanceId.make("codex"), - runtimeMode: "approval-required", - activeTurnId: null, - lastError: null, - updatedAt: "2026-01-01T00:00:03.000Z", - }, - createdAt: "2026-01-01T00:00:03.000Z", - }), - ); - - await Effect.runPromise(Scope.close(scope!, Exit.void)); - scope = null; - await runtime!.dispose(); - runtime = null; - - const third = await createHarness({ - baseDir, - providerBindingsMap: persistedBindings, - skipDefaultSetup: true, - }); - await third.drain(); - expect(third.sendTurn).not.toHaveBeenCalled(); - }); - - it("isolates restart recovery failures and continues unrelated threads", async () => { - const modelSelection: ModelSelection = { - instanceId: ProviderInstanceId.make("codex"), - model: "gpt-5-codex", - }; - const failedThreadId = ThreadId.make("thread-1"); - const healthyThreadId = ThreadId.make("thread-2"); - const marker = makeProviderRestartRecoveryMarker({ - interruptedProviderTurnId: asTurnId("provider-turn-before-restart"), - shutdownAt: "2026-01-01T00:00:01.000Z", - }); - const makeBinding = ( - threadId: ThreadId, - resumeCursor: unknown | null, - ): ProviderRuntimeBindingWithMetadata => ({ - threadId, - provider: ProviderDriverKind.make("codex"), - providerInstanceId: ProviderInstanceId.make("codex"), - adapterKey: "codex", - runtimeMode: "approval-required", - status: "stopped", - resumeCursor, - runtimePayload: { - cwd: "/tmp/provider-project", - modelSelection, - interactionMode: "default", - restartRecovery: marker, - }, - lastSeenAt: "2026-01-01T00:00:01.000Z", - }); - const harness = await createHarness({ - deferReactorStart: true, - providerBindings: [ - makeBinding(failedThreadId, null), - makeBinding(healthyThreadId, { threadId: "healthy-provider-thread" }), - ], - }); - await Effect.runPromise( - harness.engine.dispatch({ - type: "thread.create", - commandId: CommandId.make("cmd-create-recovery-thread-2"), - threadId: healthyThreadId, - projectId: asProjectId("project-1"), - title: "Healthy recovery", - modelSelection, - interactionMode: "default", - runtimeMode: "approval-required", - branch: null, - worktreePath: null, - createdAt: "2026-01-01T00:00:00.000Z", - }), - ); - - await harness.startReactor(); - await harness.drain(); - expect(harness.sendTurn).toHaveBeenCalledTimes(1); - expect(harness.sendTurn.mock.calls[0]?.[0]).toMatchObject({ - threadId: healthyThreadId, - input: RESTART_RECOVERY_CONTINUATION_INSTRUCTION, - }); - - const readModel = await harness.readModel(); - const failedThread = readModel.threads.find((thread) => thread.id === failedThreadId); - expect(failedThread?.session?.status).toBe("error"); - expect( - failedThread?.activities.some( - (activity) => activity.kind === "provider.turn.recovery.failed", - ), - ).toBe(true); - }); - - it("surfaces a disabled provider instance without starting recovery work", async () => { - const threadId = ThreadId.make("thread-1"); - const modelSelection: ModelSelection = { - instanceId: ProviderInstanceId.make("codex"), - model: "gpt-5-codex", - }; - const harness = await createHarness({ - deferReactorStart: true, - providerInstanceEnabled: false, - providerBindings: [ - { - threadId, - provider: ProviderDriverKind.make("codex"), - providerInstanceId: ProviderInstanceId.make("codex"), - runtimeMode: "approval-required", - status: "stopped", - resumeCursor: { threadId: "disabled-provider-thread" }, - runtimePayload: { - modelSelection, - restartRecovery: makeProviderRestartRecoveryMarker({ - interruptedProviderTurnId: asTurnId("disabled-provider-turn"), - shutdownAt: "2026-01-01T00:00:01.000Z", - }), - }, - lastSeenAt: "2026-01-01T00:00:01.000Z", - }, - ], - }); - - await harness.startReactor(); - await harness.drain(); - expect(harness.startSession).not.toHaveBeenCalled(); - expect(harness.sendTurn).not.toHaveBeenCalled(); - const thread = (await harness.readModel()).threads.find((entry) => entry.id === threadId); - expect(thread?.session?.status).toBe("error"); - expect( - thread?.activities.find((activity) => activity.kind === "provider.turn.recovery.failed") - ?.payload, - ).toMatchObject({ detail: expect.stringContaining("disabled") }); - }); - - it("skips archived recovery candidates", async () => { - const threadId = ThreadId.make("thread-1"); - const modelSelection: ModelSelection = { - instanceId: ProviderInstanceId.make("codex"), - model: "gpt-5-codex", - }; - const harness = await createHarness({ - deferReactorStart: true, - providerBindings: [ - { - threadId, - provider: ProviderDriverKind.make("codex"), - providerInstanceId: ProviderInstanceId.make("codex"), - runtimeMode: "approval-required", - status: "stopped", - resumeCursor: { threadId: "archived-provider-thread" }, - runtimePayload: { - modelSelection, - restartRecovery: makeProviderRestartRecoveryMarker({ - interruptedProviderTurnId: asTurnId("archived-provider-turn"), - shutdownAt: "2026-01-01T00:00:01.000Z", - }), - }, - lastSeenAt: "2026-01-01T00:00:01.000Z", - }, - ], - }); - await Effect.runPromise( - harness.engine.dispatch({ - type: "thread.archive", - commandId: CommandId.make("cmd-archive-before-recovery"), - threadId, - }), - ); - - await harness.startReactor(); - await harness.drain(); - expect(harness.startSession).not.toHaveBeenCalled(); - expect(harness.sendTurn).not.toHaveBeenCalled(); - expect(harness.providerBindings.get(threadId)?.runtimePayload).toMatchObject({ - restartRecovery: expect.objectContaining({ version: 1 }), - }); - }); - effectIt.effect("retains a turn dispatched immediately after start until activation", () => Effect.gen(function* () { const activation = yield* Deferred.make(); diff --git a/apps/server/src/orchestration/Layers/ProviderCommandReactor.ts b/apps/server/src/orchestration/Layers/ProviderCommandReactor.ts index 381db13626a6..835a67be12ee 100644 --- a/apps/server/src/orchestration/Layers/ProviderCommandReactor.ts +++ b/apps/server/src/orchestration/Layers/ProviderCommandReactor.ts @@ -9,7 +9,6 @@ import { type OrchestrationSession, ThreadId, type ProviderSession, - type ProviderInteractionMode, type RuntimeMode, type TurnId, } from "@t3tools/contracts"; @@ -40,22 +39,9 @@ import { } from "../../observability/Metrics.ts"; import { ProviderAdapterRequestError } from "../../provider/Errors.ts"; import type { ProviderServiceError } from "../../provider/Errors.ts"; -import { - makeProviderRestartRecoveryMarker, - readLiveShellRestartRecoveryCandidate, - readPersistedProviderCwd, - readPersistedProviderInteractionMode, - readPersistedProviderModelSelection, - readProviderRestartRecoveryCandidate, - type ProviderRestartRecoveryCandidate, -} from "../../provider/ProviderRestartRecovery.ts"; import { TextGeneration } from "../../textGeneration/TextGeneration.ts"; import { ProviderService } from "../../provider/Services/ProviderService.ts"; import { ProviderRegistry } from "../../provider/Services/ProviderRegistry.ts"; -import { - ProviderSessionDirectory, - type ProviderRuntimeBindingWithMetadata, -} from "../../provider/Services/ProviderSessionDirectory.ts"; import { ProjectionTurnRepository } from "../../persistence/Services/ProjectionTurns.ts"; import { ProjectionTurnRepositoryLive } from "../../persistence/Layers/ProjectionTurns.ts"; import { OrchestrationEngineService } from "../Services/OrchestrationEngine.ts"; @@ -244,18 +230,6 @@ function formatThreadTitleContext(messages: ReadonlyArray): }; } const STARTUP_RECOVERY_CONCURRENCY = 4; -/** Bound recovery provider calls so a hung agent cannot pin reactor start forever. */ -const STARTUP_RECOVERY_PROVIDER_TIMEOUT = Duration.seconds(45); - -export const RESTART_RECOVERY_CONTINUATION_INSTRUCTION = - "The server restarted while you were working. Inspect the conversation and current workspace state, verify which side effects from the interrupted turn already happened, and continue the unfinished work safely. Do not repeat completed work or assume an earlier tool call failed merely because its response is absent."; - -function runtimePayloadRecord(value: unknown): Record { - if (value !== null && typeof value === "object" && !Array.isArray(value)) { - return { ...(value as Record) }; - } - return {}; -} export function providerErrorLabel(value: string | undefined): string { const normalized = value?.trim(); @@ -351,7 +325,6 @@ const make = Effect.gen(function* () { const projectionSnapshotQuery = yield* ProjectionSnapshotQuery; const projectionTurnRepository = yield* ProjectionTurnRepository; const providerService = yield* ProviderService; - const providerSessionDirectory = yield* ProviderSessionDirectory; const providerRegistry = yield* ProviderRegistry; const gitWorkflow = yield* GitWorkflowService; const fileSystem = yield* FileSystem.FileSystem; @@ -368,8 +341,6 @@ const make = Effect.gen(function* () { lookup: () => Effect.succeed(true), }); const startupReconciliationDone = yield* Deferred.make(); - const recoveredThreadIds = new Set(); - const interruptedRecoveryThreadIds = new Set(); const hasHandledTurnStartRecently = (key: string) => Cache.getOption(handledTurnStartKeys, key).pipe( @@ -1710,417 +1681,6 @@ const make = Effect.gen(function* () { } }); - const setRecoveryFailureState = Effect.fn("setRecoveryFailureState")(function* (input: { - readonly binding: ProviderRuntimeBindingWithMetadata; - readonly detail: string; - readonly createdAt: string; - }) { - const thread = yield* resolveThread(input.binding.threadId); - if (!thread) return; - yield* setThreadSession({ - threadId: thread.id, - session: { - threadId: thread.id, - status: "error", - providerName: input.binding.provider, - ...(input.binding.providerInstanceId !== undefined - ? { providerInstanceId: input.binding.providerInstanceId } - : {}), - runtimeMode: input.binding.runtimeMode ?? thread.runtimeMode, - activeTurnId: null, - lastError: input.detail, - updatedAt: input.createdAt, - }, - createdAt: input.createdAt, - }); - }); - - type ClaimedInterruptedRecovery = { - readonly binding: ProviderRuntimeBindingWithMetadata; - readonly candidate: ProviderRestartRecoveryCandidate; - readonly thread: { - readonly id: ThreadId; - readonly projectId: ProjectId; - readonly worktreePath: string | null; - readonly runtimeMode: RuntimeMode; - readonly modelSelection: ModelSelection; - readonly interactionMode: ProviderInteractionMode; - readonly latestTurn?: { readonly turnId?: TurnId | null } | null; - }; - readonly createdAt: string; - }; - - /** - * Sync claim only: durable marker + interrupted projection so orphan settle - * cannot wipe recovery intent, without waiting on provider start/sendTurn. - */ - const claimInterruptedTurn = Effect.fn("claimInterruptedTurn")(function* (input: { - readonly binding: ProviderRuntimeBindingWithMetadata; - readonly candidate: ProviderRestartRecoveryCandidate; - }) { - const { binding, candidate } = input; - if (recoveredThreadIds.has(binding.threadId)) { - yield* increment(providerTurnRecoveriesTotal, { - outcome: "skipped", - reason: "duplicate-in-boot", - provider: binding.provider, - }); - return undefined; - } - recoveredThreadIds.add(binding.threadId); - - const thread = yield* resolveThread(binding.threadId); - if (!thread) { - yield* Effect.logInfo("provider turn restart recovery skipped", { - threadId: binding.threadId, - provider: binding.provider, - reason: "thread-missing-archived-or-deleted", - }); - yield* increment(providerTurnRecoveriesTotal, { - outcome: "skipped", - reason: "inactive-thread", - provider: binding.provider, - }); - return undefined; - } - - const createdAt = DateTime.formatIso(yield* DateTime.now); - const projectedTurns = yield* projectionTurnRepository.listByThreadId({ - threadId: binding.threadId, - }); - const projectedCandidateTurn = - candidate.interruptedProviderTurnId === null - ? undefined - : projectedTurns.find((turn) => turn.turnId === candidate.interruptedProviderTurnId); - const latestProjectedTurn = projectedTurns.findLast((turn) => turn.turnId !== null); - const projectedRecoveryTurn = projectedCandidateTurn ?? latestProjectedTurn; - const shellStillLive = - thread.session?.status === "running" || - thread.session?.status === "starting" || - thread.session?.activeTurnId != null; - if ( - (projectedRecoveryTurn?.state === "completed" || projectedRecoveryTurn?.state === "error") && - !shellStillLive - ) { - yield* providerSessionDirectory - .upsert({ - threadId: binding.threadId, - provider: binding.provider, - ...(binding.providerInstanceId !== undefined - ? { providerInstanceId: binding.providerInstanceId } - : {}), - ...(binding.runtimeMode !== undefined ? { runtimeMode: binding.runtimeMode } : {}), - status: "stopped", - runtimePayload: { - activeTurnId: null, - restartRecovery: null, - lastRuntimeEvent: "provider.restartRecovery.skipped", - lastRuntimeEventAt: createdAt, - }, - }) - .pipe( - Effect.catchCause((cause) => - Effect.logWarning("stale provider restart recovery intent was not cleared", { - threadId: binding.threadId, - cause: Cause.pretty(cause), - }), - ), - ); - yield* Effect.logInfo("provider turn restart recovery skipped", { - threadId: binding.threadId, - provider: binding.provider, - reason: "projected-turn-settled", - projectedTurnId: projectedRecoveryTurn.turnId, - projectedTurnState: projectedRecoveryTurn.state, - }); - yield* increment(providerTurnRecoveriesTotal, { - outcome: "skipped", - reason: "projected-turn-settled", - provider: binding.provider, - }); - return undefined; - } - - interruptedRecoveryThreadIds.add(thread.id); - const runtimeMode = binding.runtimeMode ?? thread.runtimeMode; - const recoveryMarker = makeProviderRestartRecoveryMarker({ - interruptedProviderTurnId: candidate.interruptedProviderTurnId, - shutdownAt: candidate.shutdownAt, - }); - - // Persist a durable marker and demote live-claiming status before orphan - // audit. Orphan settle spreads payload fields, so the marker survives even - // if a late stop rewrites the binding; status "stopped" skips live claims. - yield* providerSessionDirectory - .upsert({ - threadId: binding.threadId, - provider: binding.provider, - ...(binding.providerInstanceId !== undefined - ? { providerInstanceId: binding.providerInstanceId } - : {}), - ...(binding.runtimeMode !== undefined ? { runtimeMode: binding.runtimeMode } : {}), - status: "stopped", - ...(binding.resumeCursor !== undefined && binding.resumeCursor !== null - ? { resumeCursor: binding.resumeCursor } - : {}), - runtimePayload: { - ...runtimePayloadRecord(binding.runtimePayload), - activeTurnId: null, - restartRecovery: recoveryMarker, - lastRuntimeEvent: "provider.restartRecovery.claimed", - lastRuntimeEventAt: createdAt, - }, - }) - .pipe( - Effect.catchCause((cause) => - Effect.logWarning("provider restart recovery claim was not persisted", { - threadId: binding.threadId, - cause: Cause.pretty(cause), - }), - ), - ); - - // Settle the concrete old projection row as interrupted before replacement. - yield* setThreadSession({ - threadId: thread.id, - session: { - threadId: thread.id, - status: "interrupted", - providerName: binding.provider, - ...(binding.providerInstanceId !== undefined - ? { providerInstanceId: binding.providerInstanceId } - : {}), - runtimeMode, - activeTurnId: null, - lastError: null, - updatedAt: createdAt, - }, - createdAt, - }); - - yield* Effect.logInfo("provider turn restart recovery claimed", { - threadId: thread.id, - provider: binding.provider, - recoverySource: candidate.source, - interruptedProviderTurnId: candidate.interruptedProviderTurnId, - }); - - return { - binding, - candidate, - thread, - createdAt, - } satisfies ClaimedInterruptedRecovery; - }); - - const continueInterruptedTurn = Effect.fn("continueInterruptedTurn")(function* ( - input: ClaimedInterruptedRecovery, - ) { - const { binding, candidate, thread, createdAt } = input; - - const recover = Effect.gen(function* () { - const runtimeMode = binding.runtimeMode ?? thread.runtimeMode; - - const providerInstanceId = binding.providerInstanceId; - if (providerInstanceId === undefined) { - return yield* new ProviderAdapterRequestError({ - provider: binding.provider, - method: "provider.turn.restart-recovery", - detail: `Persisted provider binding for thread '${binding.threadId}' has no provider instance id.`, - }); - } - if (binding.resumeCursor === null || binding.resumeCursor === undefined) { - return yield* new ProviderAdapterRequestError({ - provider: binding.provider, - method: "provider.turn.restart-recovery", - detail: `Cannot recover thread '${binding.threadId}' because no provider resume cursor is persisted.`, - }); - } - - const instanceInfo = yield* providerService.getInstanceInfo(providerInstanceId); - if (!instanceInfo.enabled) { - return yield* new ProviderAdapterRequestError({ - provider: binding.provider, - method: "provider.turn.restart-recovery", - detail: `Provider instance '${providerInstanceId}' is disabled in T3 Code settings.`, - }); - } - if (instanceInfo.driverKind !== binding.provider) { - return yield* new ProviderAdapterRequestError({ - provider: binding.provider, - method: "provider.turn.restart-recovery", - detail: `Persisted provider instance '${providerInstanceId}' now uses driver '${instanceInfo.driverKind}', not '${binding.provider}'.`, - }); - } - - const persistedModelSelection = readPersistedProviderModelSelection(binding.runtimePayload); - const modelSelection = persistedModelSelection ?? thread.modelSelection; - if (modelSelection.instanceId !== providerInstanceId) { - return yield* new ProviderAdapterRequestError({ - provider: binding.provider, - method: "provider.turn.restart-recovery", - detail: `Persisted model selection references provider instance '${modelSelection.instanceId}', but the recoverable session belongs to '${providerInstanceId}'.`, - }); - } - const interactionMode: ProviderInteractionMode = - readPersistedProviderInteractionMode(binding.runtimePayload) ?? thread.interactionMode; - const project = yield* resolveProject(thread.projectId); - const cwd = - readPersistedProviderCwd(binding.runtimePayload) ?? - resolveThreadWorkspaceCwd({ thread, projects: project ? [project] : [] }); - - const sessionResult = yield* providerService - .startSession(thread.id, { - threadId: thread.id, - provider: binding.provider, - providerInstanceId, - ...(cwd !== undefined ? { cwd } : {}), - modelSelection, - resumeCursor: binding.resumeCursor, - runtimeMode, - }) - .pipe(Effect.interruptible, Effect.timeoutOption(STARTUP_RECOVERY_PROVIDER_TIMEOUT)); - if (Option.isNone(sessionResult)) { - return yield* new ProviderAdapterRequestError({ - provider: binding.provider, - method: "provider.turn.restart-recovery", - detail: `Provider session start timed out after ${Duration.format(STARTUP_RECOVERY_PROVIDER_TIMEOUT)} during restart recovery.`, - }); - } - const session = sessionResult.value; - yield* setThreadSession({ - threadId: thread.id, - session: { - threadId: thread.id, - status: mapProviderSessionStatusToOrchestrationStatus(session.status), - providerName: session.provider, - providerInstanceId, - runtimeMode, - activeTurnId: null, - lastError: session.lastError ?? null, - updatedAt: session.updatedAt, - }, - createdAt, - }); - - // startSession stays bounded so a hung spawn cannot pin recovery. - // sendTurn is the live continuation: Grok's prompt RPC is the whole - // turn, so a 45s cap always fails after we claim ACP sessions. - // This fiber is already parked past HTTP readiness. - const replacement = yield* providerService.sendTurn({ - threadId: thread.id, - input: RESTART_RECOVERY_CONTINUATION_INSTRUCTION, - attachments: [], - modelSelection, - interactionMode, - }); - - // ProviderService clears this in the accepted sendTurn transaction. The - // explicit write keeps the reconciliation invariant local and obvious. - yield* providerSessionDirectory.upsert({ - threadId: thread.id, - provider: binding.provider, - providerInstanceId, - runtimeMode, - status: "running", - ...(replacement.resumeCursor !== undefined - ? { resumeCursor: replacement.resumeCursor } - : {}), - runtimePayload: { - activeTurnId: replacement.turnId, - modelSelection, - interactionMode, - restartRecovery: null, - lastRuntimeEvent: "provider.restartRecovery.accepted", - lastRuntimeEventAt: DateTime.formatIso(yield* DateTime.now), - }, - }); - - yield* Effect.logInfo("provider turn restart recovery accepted", { - threadId: thread.id, - provider: binding.provider, - providerInstanceId, - interruptedProviderTurnId: candidate.interruptedProviderTurnId, - replacementProviderTurnId: replacement.turnId, - recoverySource: candidate.source, - }); - yield* increment(providerTurnRecoveriesTotal, { - outcome: "continued", - source: candidate.source, - provider: binding.provider, - }); - }); - - yield* recover.pipe( - Effect.catchCause((cause) => { - const detail = formatFailureDetail(cause); - const persistFailureState = - binding.providerInstanceId === undefined - ? Effect.void - : providerSessionDirectory - .upsert({ - threadId: binding.threadId, - provider: binding.provider, - providerInstanceId: binding.providerInstanceId, - ...(binding.runtimeMode !== undefined - ? { runtimeMode: binding.runtimeMode } - : {}), - status: "error", - runtimePayload: { - activeTurnId: null, - lastError: detail, - lastRuntimeEvent: "provider.restartRecovery.failed", - lastRuntimeEventAt: createdAt, - }, - }) - .pipe( - Effect.catchCause((persistenceCause) => - Effect.logWarning("provider restart recovery failure state was not persisted", { - threadId: binding.threadId, - cause: Cause.pretty(persistenceCause), - }), - ), - ); - return persistFailureState.pipe( - Effect.andThen(setRecoveryFailureState({ binding, detail, createdAt })), - Effect.andThen( - appendProviderFailureActivity({ - threadId: binding.threadId, - kind: "provider.turn.recovery.failed", - summary: "Provider turn recovery failed", - detail, - turnId: thread.latestTurn?.turnId ?? null, - createdAt, - }), - ), - Effect.andThen( - Effect.logWarning("provider turn restart recovery failed", { - threadId: binding.threadId, - provider: binding.provider, - providerInstanceId: binding.providerInstanceId, - recoverySource: candidate.source, - cause: Cause.pretty(cause), - }), - ), - Effect.andThen( - increment(providerTurnRecoveriesTotal, { - outcome: "failed", - source: candidate.source, - provider: binding.provider, - }), - ), - Effect.catchCause((reportingCause) => - Effect.logWarning("provider turn restart recovery failure reporting failed", { - threadId: binding.threadId, - cause: Cause.pretty(reportingCause), - originalCause: Cause.pretty(cause), - }), - ), - ); - }), - ); - }); - const processDomainEvent = Effect.fn("processDomainEvent")(function* ( event: ProviderIntentEvent, ) { @@ -2207,71 +1767,19 @@ const make = Effect.gen(function* () { }); const reconcileStartup = Effect.fn("reconcileStartup")(function* () { - const bindings = yield* providerSessionDirectory.listBindings().pipe( - Effect.catchCause((cause) => - Effect.logWarning("provider restart recovery failed to list persisted bindings", { - cause: Cause.pretty(cause), - }).pipe(Effect.as([] as ReadonlyArray)), - ), - ); - const shells = yield* projectionSnapshotQuery.getShellSnapshot().pipe( - Effect.catchCause((cause) => - Effect.logWarning("provider restart recovery failed to read live thread shells", { - cause: Cause.pretty(cause), - }).pipe(Effect.as({ threads: [] as const })), - ), - ); - const liveShellByThreadId = new Map( - shells.threads.map((thread) => [String(thread.id), thread] as const), - ); - const recoveryCandidates = bindings.flatMap((binding) => { - const candidate = - readProviderRestartRecoveryCandidate({ - runtimePayload: binding.runtimePayload, - status: binding.status, - lastSeenAt: binding.lastSeenAt, - }) ?? - readLiveShellRestartRecoveryCandidate({ - sessionStatus: liveShellByThreadId.get(String(binding.threadId))?.session?.status, - activeTurnId: liveShellByThreadId.get(String(binding.threadId))?.session?.activeTurnId, - lastSeenAt: binding.lastSeenAt, - }); - return candidate === undefined ? [] : [{ binding, candidate }]; - }); const pendingTurnStarts = yield* projectionTurnRepository.listPendingTurnStarts().pipe( Effect.catchCause((cause) => - Effect.logWarning("provider restart recovery failed to list pending turn starts", { + Effect.logWarning("provider startup reconciliation failed to list pending turn starts", { cause: Cause.pretty(cause), }).pipe(Effect.as([])), ), ); - yield* Effect.logInfo("provider restart reconciliation candidates loaded", { - interruptedTurnCandidates: recoveryCandidates.length, + yield* Effect.logInfo("provider startup reconciliation candidates loaded", { pendingTurnStartCandidates: pendingTurnStarts.length, }); - if (recoveryCandidates.length > 0) { - yield* increment( - providerTurnRecoveriesTotal, - { outcome: "candidate", recoveryKind: "interrupted-turn" }, - recoveryCandidates.length, - ); - } - - // Claim phase is synchronous and must finish before orphan settle (next - // startup phase). Provider startSession/sendTurn runs after activation so a - // hung agent cannot block HTTP readiness / Discord oauth bootstrap. - const claimedRecoveries = yield* Effect.forEach(recoveryCandidates, claimInterruptedTurn, { - concurrency: STARTUP_RECOVERY_CONCURRENCY, - }).pipe( - Effect.map((claims) => claims.flatMap((claim) => (claim === undefined ? [] : [claim]))), - ); - - const pendingWithoutInterruptedRecovery = pendingTurnStarts.filter( - (pending) => !interruptedRecoveryThreadIds.has(pending.threadId), - ); - if (pendingWithoutInterruptedRecovery.length === 0) { - return claimedRecoveries; + if (pendingTurnStarts.length === 0) { + return; } const persistedEvents = yield* Stream.runCollect( @@ -2279,7 +1787,7 @@ const make = Effect.gen(function* () { ).pipe( Effect.map((events) => Array.from(events)), Effect.catchCause((cause) => - Effect.logWarning("provider restart recovery failed to read persisted turn starts", { + Effect.logWarning("provider startup reconciliation failed to read persisted turn starts", { cause: Cause.pretty(cause), }).pipe(Effect.as([] as ReadonlyArray)), ), @@ -2297,7 +1805,7 @@ const make = Effect.gen(function* () { } yield* Effect.forEach( - pendingWithoutInterruptedRecovery, + pendingTurnStarts, (pending) => Effect.gen(function* () { const thread = yield* resolveThread(pending.threadId); @@ -2362,8 +1870,6 @@ const make = Effect.gen(function* () { ), { concurrency: STARTUP_RECOVERY_CONCURRENCY, discard: true }, ); - - return claimedRecoveries; }); const start: ProviderCommandReactorShape["start"] = Effect.fn("start")(function* () { @@ -2426,37 +1932,14 @@ const make = Effect.gen(function* () { yield* forkParked(clearInterrupted); } - // Claim interrupted recoveries (and enqueue pending turn starts) before this - // reactor returns. Server startup then runs orphan settle against claimed - // durable markers, not against live-claiming zombie runtimes. - // - // Provider startSession/sendTurn is intentionally parked until activation so - // a hung agent cannot block command readiness (HTTP/oauth). drain still - // waits for those continuations via startupReconciliationDone. - const claimedRecoveries = yield* reconcileStartup().pipe( + yield* reconcileStartup().pipe( Effect.catchCause((cause) => - Effect.logWarning("provider restart reconciliation failed", { + Effect.logWarning("provider startup reconciliation failed", { cause: Cause.pretty(cause), - }).pipe(Effect.as([] as ReadonlyArray)), + }), ), - ); - - const runClaimedRecoveries = Effect.forEach(claimedRecoveries, continueInterruptedTurn, { - concurrency: STARTUP_RECOVERY_CONCURRENCY, - discard: true, - }).pipe( Effect.ensuring(Deferred.succeed(startupReconciliationDone, undefined).pipe(Effect.ignore)), ); - - // Mirror clearInterrupted: unit tests omit ServerActivation and run - // recovery inline for determinism. Production installs activation so - // startSession/sendTurn wait until after orphan settle + the readiness - // boundary, and cannot block HTTP/oauth on a hung agent. - if (activation === undefined) { - yield* runClaimedRecoveries; - } else { - yield* forkParked(runClaimedRecoveries); - } }); return { diff --git a/apps/server/src/provider/Layers/ProviderService.test.ts b/apps/server/src/provider/Layers/ProviderService.test.ts index e7cd87c1358e..0aa237405746 100644 --- a/apps/server/src/provider/Layers/ProviderService.test.ts +++ b/apps/server/src/provider/Layers/ProviderService.test.ts @@ -73,10 +73,6 @@ import * as ServerConfig from "../../config.ts"; import * as ServerSettings from "../../serverSettings.ts"; import * as AnalyticsService from "../../telemetry/AnalyticsService.ts"; import { makeAdapterRegistryMock } from "../testUtils/providerAdapterRegistryMock.ts"; -import { - makeProviderRestartRecoveryMarker, - readProviderRestartRecoveryMarker, -} from "../ProviderRestartRecovery.ts"; const defaultServerSettingsLayer = ServerSettings.ServerSettingsService.layerTest(); const serverConfigTestLayer = ServerConfig.layerTest(process.cwd(), process.cwd()).pipe( @@ -476,7 +472,7 @@ it.effect("ProviderServiceLive bounds a provider that wedges during shutdown", ( }), ); -it.effect("graceful shutdown preserves recovery intent only for working sessions", () => +it.effect("graceful shutdown interrupts working sessions before stopAll", () => Effect.gen(function* () { const tempDir = NodeFS.mkdtempSync(NodePath.join(NodeOS.tmpdir(), "t3-provider-recovery-")); const dbPath = NodePath.join(tempDir, "runtime.sqlite"); @@ -530,15 +526,6 @@ it.effect("graceful shutdown preserves recovery intent only for working sessions input: "keep working", interactionMode: "plan", }); - const armedRows = yield* Effect.gen(function* () { - const repository = yield* ProviderSessionRuntime.ProviderSessionRuntimeRepository; - return yield* repository.list(); - }).pipe(Effect.provide(runtimeRepositoryLayer)); - const armedRunning = armedRows.find((row) => row.threadId === runningThreadId); - assert.equal( - readProviderRestartRecoveryMarker(armedRunning?.runtimePayload)?.interruptedProviderTurnId, - asTurnId(`turn-${String(runningThreadId)}`), - ); codex.updateSession(runningThreadId, (session) => ({ ...session, status: "running", @@ -550,8 +537,8 @@ it.effect("graceful shutdown preserves recovery intent only for working sessions })); yield* provider.stopSession({ threadId: stoppedThreadId }); - // Model a provider protocol drain that never completes. Recovery intent - // must already be durable when the global shutdown deadline interrupts it. + // Model a provider protocol drain that never completes. Shutdown must still + // interrupt working sessions and bound the wedged stopAll. codex.stopAll.mockImplementation(() => Effect.never); const closeFiber = yield* Scope.close(scope, Exit.void).pipe( Effect.forkChild({ startImmediately: true }), @@ -565,20 +552,25 @@ it.effect("graceful shutdown preserves recovery intent only for working sessions }).pipe(Effect.provide(runtimeRepositoryLayer)); const byThreadId = new Map(rows.map((row) => [row.threadId, row])); - const runningMarker = readProviderRestartRecoveryMarker( - byThreadId.get(runningThreadId)?.runtimePayload, - ); - assert.equal(runningMarker?.interruptedProviderTurnId, asTurnId("provider-turn-running")); - assert.isDefined( - readProviderRestartRecoveryMarker(byThreadId.get(connectingThreadId)?.runtimePayload), - ); - assert.isUndefined( - readProviderRestartRecoveryMarker(byThreadId.get(readyThreadId)?.runtimePayload), - ); - assert.isUndefined( - readProviderRestartRecoveryMarker(byThreadId.get(stoppedThreadId)?.runtimePayload), - ); assert.equal(byThreadId.get(runningThreadId)?.status, "stopped"); + assert.equal(byThreadId.get(connectingThreadId)?.status, "stopped"); + assert.equal(byThreadId.get(readyThreadId)?.status, "stopped"); + assert.equal(byThreadId.get(stoppedThreadId)?.status, "stopped"); + const runningPayload = byThreadId.get(runningThreadId)?.runtimePayload as + | Record + | null + | undefined; + assert.equal(runningPayload?.continueAfterServerUpdate, "provider-turn-running"); + const readyPayload = byThreadId.get(readyThreadId)?.runtimePayload as + | Record + | null + | undefined; + assert.equal(readyPayload?.continueAfterServerUpdate, undefined); + const stoppedPayload = byThreadId.get(stoppedThreadId)?.runtimePayload as + | Record + | null + | undefined; + assert.equal(stoppedPayload?.continueAfterServerUpdate, undefined); // Working sessions must receive cooperative interrupt before hard stopAll. assert.isTrue(codex.interruptTurn.mock.calls.length >= 2); assert.deepEqual(codex.interruptTurn.mock.calls[0]?.[0], runningThreadId); @@ -587,7 +579,7 @@ it.effect("graceful shutdown preserves recovery intent only for working sessions }).pipe(Effect.provide(NodeServices.layer)), ); -it.effect("graceful shutdown recovers live bindings missing from adapter listSessions", () => +it.effect("graceful shutdown stops live bindings missing from adapter listSessions", () => Effect.gen(function* () { const tempDir = NodeFS.mkdtempSync( NodePath.join(NodeOS.tmpdir(), "t3-provider-recovery-orphan-binding-"), @@ -657,95 +649,9 @@ it.effect("graceful shutdown recovers live bindings missing from adapter listSes const orphan = rows.find((row) => row.threadId === orphanThreadId); assert.isDefined(orphan); assert.equal(orphan?.status, "stopped"); - const marker = readProviderRestartRecoveryMarker(orphan?.runtimePayload); - assert.equal(marker?.interruptedProviderTurnId, asTurnId("provider-turn-orphan")); - - NodeFS.rmSync(tempDir, { recursive: true, force: true }); - }).pipe(Effect.provide(NodeServices.layer)), -); - -it.effect("keeps restart recovery markers when session.exited arrives after stopAll", () => - Effect.gen(function* () { - const tempDir = NodeFS.mkdtempSync( - NodePath.join(NodeOS.tmpdir(), "t3-provider-recovery-exited-"), - ); - const dbPath = NodePath.join(tempDir, "runtime.sqlite"); - const persistenceLayer = makeSqlitePersistenceLive(dbPath); - const runtimeRepositoryLayer = ProviderSessionRuntime.layer.pipe( - Layer.provide(persistenceLayer), - ); - const directoryLayer = ProviderSessionDirectoryLive.pipe(Layer.provide(runtimeRepositoryLayer)); - const codex = makeFakeCodexAdapter(); - const providerLayer = makeProviderServiceLive().pipe( - Layer.provide( - Layer.succeed( - ProviderAdapterRegistry.ProviderAdapterRegistry, - makeAdapterRegistryMock({ [CODEX_DRIVER]: codex.adapter }), - ), - ), - Layer.provide(directoryLayer), - Layer.provide(defaultServerSettingsLayer), - Layer.provide(AnalyticsService.layerTest), - Layer.provide(serverConfigTestLayer), - Layer.provide( - Layer.succeed( - ProviderEventLoggers.ProviderEventLoggers, - ProviderEventLoggers.NoOpProviderEventLoggers, - ), - ), - ); - const scope = yield* Scope.make(); - const services = yield* Layer.build( - Layer.mergeAll(providerLayer, runtimeRepositoryLayer, directoryLayer), - ).pipe(Scope.provide(scope)); - const provider = yield* ProviderService.ProviderService.pipe(Effect.provide(services)); - const threadId = asThreadId("thread-exited-after-marker"); - yield* provider.startSession(threadId, { - provider: CODEX_DRIVER, - providerInstanceId: codexInstanceId, - threadId, - runtimeMode: "full-access", - }); - const directory = yield* ProviderSessionDirectory.ProviderSessionDirectory.pipe( - Effect.provide(services), - ); - const marker = makeProviderRestartRecoveryMarker({ - interruptedProviderTurnId: asTurnId("provider-turn-exited"), - shutdownAt: "2026-01-01T00:00:02.000Z", - }); - yield* directory.upsert({ - threadId, - provider: CODEX_DRIVER, - providerInstanceId: codexInstanceId, - status: "stopped", - resumeCursor: { threadId: "provider-exited" }, - runtimePayload: { - restartRecovery: marker, - lastRuntimeEvent: "provider.stopAll", - lastRuntimeEventAt: "2026-01-01T00:00:02.000Z", - }, - }); - - codex.emit({ - type: "session.exited", - eventId: asEventId("evt-session-exited"), - provider: CODEX_DRIVER, - createdAt: "2026-01-01T00:00:03.000Z", - threadId, - }); - yield* advanceTestClock(50); + const orphanPayload = orphan?.runtimePayload as Record | null | undefined; + assert.equal(orphanPayload?.continueAfterServerUpdate, "provider-turn-orphan"); - const rows = yield* Effect.gen(function* () { - const repository = yield* ProviderSessionRuntime.ProviderSessionRuntimeRepository; - return yield* repository.list(); - }).pipe(Effect.provide(runtimeRepositoryLayer)); - const persisted = rows.find((row) => row.threadId === threadId); - assert.equal( - readProviderRestartRecoveryMarker(persisted?.runtimePayload)?.interruptedProviderTurnId, - asTurnId("provider-turn-exited"), - ); - - yield* Scope.close(scope, Exit.void); NodeFS.rmSync(tempDir, { recursive: true, force: true }); }).pipe(Effect.provide(NodeServices.layer)), ); @@ -1939,10 +1845,6 @@ routing.layer("ProviderServiceLive routing", (it) => { assert.equal(runtimePayload.activeTurnId, `turn-${String(session.threadId)}`); assert.equal(runtimePayload.lastError, null); assert.equal(runtimePayload.lastRuntimeEvent, "provider.sendTurn"); - assert.equal( - readProviderRestartRecoveryMarker(runtimePayload)?.interruptedProviderTurnId, - asTurnId(`turn-${String(session.threadId)}`), - ); } } @@ -1963,9 +1865,7 @@ routing.layer("ProviderServiceLive routing", (it) => { }); assert.equal(Option.isSome(completedRuntime), true); if (Option.isSome(completedRuntime)) { - assert.isUndefined( - readProviderRestartRecoveryMarker(completedRuntime.value.runtimePayload), - ); + assert.equal(completedRuntime.value.status, "running"); } }), ); diff --git a/apps/server/src/provider/Layers/ProviderService.ts b/apps/server/src/provider/Layers/ProviderService.ts index 66b3b9670045..191e18633508 100644 --- a/apps/server/src/provider/Layers/ProviderService.ts +++ b/apps/server/src/provider/Layers/ProviderService.ts @@ -12,6 +12,7 @@ import { NonNegativeInt, ThreadId, + TurnId, ProviderCompactSessionInput, ProviderInterruptTurnInput, ProviderRespondToRequestInput, @@ -66,13 +67,14 @@ import * as AnalyticsService from "../../telemetry/AnalyticsService.ts"; import * as McpProviderSession from "../../mcp/McpProviderSession.ts"; import * as McpSessionRegistry from "../../mcp/McpSessionRegistry.ts"; import { - makeProviderRestartRecoveryMarker, readPersistedProviderActiveTurnId, readPersistedProviderCwd, readPersistedProviderModelSelection, - readProviderRestartRecoveryMarker, - type ProviderRestartRecoveryMarker, } from "../ProviderRestartRecovery.ts"; +import { + readRuntimePayload, + SERVER_UPDATE_CONTINUATION_KEY, +} from "../providerSessionContinuation.ts"; import * as ServerSettings from "../../serverSettings.ts"; /** @@ -83,17 +85,16 @@ import * as ServerSettings from "../../serverSettings.ts"; export interface ProviderServiceLiveOptions { readonly canonicalEventLogger?: EventNdjsonLogger; /** - * After writing restartRecovery markers, wait this long after cooperative - * `interruptTurn` on in-flight sessions before hard `adapter.stopAll()`. - * Gives providers time to cancel tools / flush before process teardown. - * Default `30 seconds` — well under ops' ~150s main-pid SIGTERM reap window - * (interrupt grace + stopAll grace must fit inside that). + * Wait this long after cooperative `interruptTurn` on in-flight sessions + * before hard `adapter.stopAll()`. Gives providers time to cancel tools / + * flush before process teardown. Default `30 seconds` — well under ops' + * ~150s main-pid SIGTERM reap window (interrupt grace + stopAll grace must + * fit inside that). */ readonly shutdownInterruptGracePeriod?: Duration.Input; /** * Maximum time the server gives all provider adapters, collectively, to - * stop during process shutdown (after the interrupt grace). Recovery intent - * is persisted before either clock starts. Default `1 minute`. + * stop during process shutdown (after the interrupt grace). Default `1 minute`. */ readonly shutdownGracePeriod?: Duration.Input; /** @@ -167,7 +168,6 @@ function toRuntimePayloadFromSession( readonly interactionMode?: unknown; readonly lastRuntimeEvent?: string; readonly lastRuntimeEventAt?: string; - readonly restartRecovery?: unknown; }, ): Record { return { @@ -181,7 +181,6 @@ function toRuntimePayloadFromSession( ...(extra?.lastRuntimeEventAt !== undefined ? { lastRuntimeEventAt: extra.lastRuntimeEventAt } : {}), - ...(extra?.restartRecovery !== undefined ? { restartRecovery: extra.restartRecovery } : {}), }; } @@ -343,20 +342,20 @@ const makeProviderService = Effect.fn("makeProviderService")(function* ( ) { return undefined; } - return { status: "running" as const, activeTurnId: null, restartRecovery: null }; + return { status: "running" as const, activeTurnId: null }; case "session.exited": - return { status: "stopped" as const, activeTurnId: null, restartRecovery: null }; + return { status: "stopped" as const, activeTurnId: null }; case "session.state.changed": switch (event.payload.state) { case "starting": return { status: "starting" as const }; case "error": - return { status: "error" as const, activeTurnId: null, restartRecovery: null }; + return { status: "error" as const, activeTurnId: null }; case "stopped": - return { status: "stopped" as const, activeTurnId: null, restartRecovery: null }; + return { status: "stopped" as const, activeTurnId: null }; case "ready": case "waiting": - return { status: "running" as const, activeTurnId: null, restartRecovery: null }; + return { status: "running" as const, activeTurnId: null }; case "running": return { status: "running" as const }; } @@ -366,18 +365,6 @@ const makeProviderService = Effect.fn("makeProviderService")(function* ( })(); if (lifecycle === undefined) return; - const payload = - binding.runtimePayload !== null && - typeof binding.runtimePayload === "object" && - !Array.isArray(binding.runtimePayload) - ? (binding.runtimePayload as Record) - : undefined; - // stopAll already persisted recovery intent. A following Grok/ACP - // turn.completed (cancelled) or session.exited must not delete it. - const preserveShutdownRecoveryMarker = - payload?.lastRuntimeEvent === "provider.stopAll" && - readProviderRestartRecoveryMarker(binding.runtimePayload) !== undefined; - yield* directory.upsert({ threadId: event.threadId, provider: binding.provider, @@ -386,9 +373,6 @@ const makeProviderService = Effect.fn("makeProviderService")(function* ( status: lifecycle.status, runtimePayload: { ...(lifecycle.activeTurnId !== undefined ? { activeTurnId: lifecycle.activeTurnId } : {}), - ...(lifecycle.restartRecovery !== undefined && !preserveShutdownRecoveryMarker - ? { restartRecovery: lifecycle.restartRecovery } - : {}), lastRuntimeEvent: event.type, lastRuntimeEventAt: event.createdAt, }, @@ -421,7 +405,6 @@ const makeProviderService = Effect.fn("makeProviderService")(function* ( readonly interactionMode?: unknown; readonly lastRuntimeEvent?: string; readonly lastRuntimeEventAt?: string; - readonly restartRecovery?: unknown; }, ) => Effect.gen(function* () { @@ -975,14 +958,6 @@ const makeProviderService = Effect.fn("makeProviderService")(function* ( ? { interactionMode: input.interactionMode } : {}), activeTurnId: turn.turnId, - // Arm recovery while the turn is live, rather than waiting for a - // process-shutdown finalizer. Supervisors and desktop updaters can - // still exhaust their graceful-stop budget or crash after TERM; the - // next process must already have durable intent to resume the turn. - restartRecovery: makeProviderRestartRecoveryMarker({ - interruptedProviderTurnId: turn.turnId, - shutdownAt: turnStartedAt, - }), lastRuntimeEvent: "provider.sendTurn", lastRuntimeEventAt: turnStartedAt, }, @@ -1192,7 +1167,6 @@ const makeProviderService = Effect.fn("makeProviderService")(function* ( status: "stopped", runtimePayload: { activeTurnId: null, - restartRecovery: null, }, }); yield* analytics.record("provider.session.stopped", { @@ -1389,12 +1363,6 @@ const makeProviderService = Effect.fn("makeProviderService")(function* ( const threadIds = yield* directory.listThreadIds(); const currentAdapters = yield* getAdapterEntries; const lastRuntimeEventAt = yield* nowIso; - // Durable recovery intent for the next boot. Must not depend solely on - // in-memory adapter.listSessions(): during SIGTERM teardown adapters can - // already be empty while SQLite still has starting/running bindings. The - // final "stopped" write also clears activeTurnId, so without a marker the - // next process cannot recover interrupted turns. - const recoveryByThreadId = new Map(); const activeSessions = yield* Effect.forEach(currentAdapters, ([instanceId, adapter]) => adapter.listSessions().pipe( @@ -1407,101 +1375,84 @@ const makeProviderService = Effect.fn("makeProviderService")(function* ( ), ).pipe(Effect.map((sessionsByAdapter) => sessionsByAdapter.flatMap((sessions) => sessions))); + // Same payload key as in-app #9167. Startup reconcileProviderSessions + // continues marked threads after this process is replaced (systemd + // restart, deploy). Crash/hard-kill never reaches stopAll, so those + // still settle as interrupted. + const persistedBindings = yield* directory.listBindings().pipe(Effect.orElseSucceed(() => [])); + const bindingByThreadId = new Map( + persistedBindings.map((binding) => [binding.threadId, binding] as const), + ); + const continuationByThreadId = new Map(); + const considerContinuation = ( + threadId: ThreadId, + turnId: TurnId | null | undefined, + resumeCursor: unknown, + ) => { + if (turnId === null || turnId === undefined) return; + if (resumeCursor === null || resumeCursor === undefined) return; + if (continuationByThreadId.has(threadId)) return; + continuationByThreadId.set(threadId, turnId); + }; + for (const session of activeSessions) { - // Adapter-level "ready" is idle — only connecting/running need continuation. - const wasWorking = session.status === "connecting" || session.status === "running"; - if (wasWorking) { - recoveryByThreadId.set( - session.threadId, - makeProviderRestartRecoveryMarker({ - interruptedProviderTurnId: session.activeTurnId, - shutdownAt: lastRuntimeEventAt, - }), - ); - } + const working = session.status === "connecting" || session.status === "running"; + if (!working) continue; + const binding = bindingByThreadId.get(session.threadId); + considerContinuation( + session.threadId, + session.activeTurnId ?? readPersistedProviderActiveTurnId(binding?.runtimePayload), + binding?.resumeCursor ?? session.resumeCursor, + ); } - - const adapterThreadIds = new Set(activeSessions.map((session) => session.threadId)); - const persistedBindings = yield* directory.listBindings().pipe(Effect.orElseSucceed(() => [])); for (const binding of persistedBindings) { - // Persistence maps ready→running, so require an active turn id to avoid - // marking idle ready sessions as recovery candidates. - if (binding.status !== "starting" && binding.status !== "running") { - continue; - } - if (recoveryByThreadId.has(binding.threadId) || adapterThreadIds.has(binding.threadId)) { - continue; - } - const activeTurnId = readPersistedProviderActiveTurnId(binding.runtimePayload); - if (activeTurnId === undefined) { - continue; - } - recoveryByThreadId.set( + if (binding.status !== "starting" && binding.status !== "running") continue; + considerContinuation( binding.threadId, - makeProviderRestartRecoveryMarker({ - interruptedProviderTurnId: activeTurnId, - shutdownAt: lastRuntimeEventAt, - }), + readPersistedProviderActiveTurnId(binding.runtimePayload), + binding.resumeCursor, ); } - yield* Effect.forEach(activeSessions, (session) => { - const marker = recoveryByThreadId.get(session.threadId); - return upsertSessionBinding(session, session.threadId, { - lastRuntimeEvent: "provider.stopAll", - lastRuntimeEventAt, - // Omit (undefined) for idle sessions so we do not clobber a marker that - // was recorded from a persisted live binding for the same thread. - ...(marker !== undefined ? { restartRecovery: marker } : {}), - }); - }).pipe(Effect.asVoid); - - // Persist markers only for bindings that never appeared in adapter.listSessions - // (adapter path already wrote via upsertSessionBinding and must keep its - // resumeCursor / payload fields intact). yield* Effect.forEach( - persistedBindings.filter( - (binding) => - recoveryByThreadId.has(binding.threadId) && !adapterThreadIds.has(binding.threadId), - ), - (binding) => + [...continuationByThreadId.entries()], + ([threadId, turnId]) => Effect.gen(function* () { - const providerInstanceId = dieOnMissingBindingInstanceId( - "ProviderService.stopAll", - binding, - ); - const marker = recoveryByThreadId.get(binding.threadId); - if (marker === undefined) { - return; - } + const binding = bindingByThreadId.get(threadId); + if (binding === undefined) return; yield* directory.upsert({ - threadId: binding.threadId, - provider: binding.provider, - providerInstanceId, - ...(binding.runtimeMode !== undefined ? { runtimeMode: binding.runtimeMode } : {}), - ...(binding.resumeCursor !== undefined ? { resumeCursor: binding.resumeCursor } : {}), - status: binding.status ?? "running", + ...binding, runtimePayload: { - restartRecovery: marker, - lastRuntimeEvent: "provider.stopAll", - lastRuntimeEventAt, + ...readRuntimePayload(binding.runtimePayload), + [SERVER_UPDATE_CONTINUATION_KEY]: turnId, }, }); - }), + }).pipe( + Effect.catchCause((cause) => + Effect.logWarning("failed to mark provider session for restart continuation", { + threadId, + cause: Cause.pretty(cause), + }), + ), + ), { discard: true }, ); - - if (recoveryByThreadId.size > 0) { - yield* Effect.logInfo("persisted provider restart recovery markers before stopAll", { - recoveryCount: recoveryByThreadId.size, - adapterSessionCount: activeSessions.length, + if (continuationByThreadId.size > 0) { + yield* Effect.logInfo("marked provider sessions to continue after graceful restart", { + continuationCount: continuationByThreadId.size, }); } + yield* Effect.forEach(activeSessions, (session) => + upsertSessionBinding(session, session.threadId, { + lastRuntimeEvent: "provider.stopAll", + lastRuntimeEventAt, + }), + ).pipe(Effect.asVoid); + // Cooperative interrupt before hard session teardown so running provider - // turns receive cancel (tools stop, agents can settle) while markers are - // already durable. Sequence under ops SIGTERM reap (~150s): - // markers (fast) → interrupt → interruptGrace (~30s) → stopAll (~60s). + // turns receive cancel (tools stop, agents can settle). Sequence under ops + // SIGTERM reap (~150s): interrupt → interruptGrace (~30s) → stopAll (~60s). const workingSessions = activeSessions.filter( (session) => session.status === "connecting" || session.status === "running", ); @@ -1558,7 +1509,6 @@ const makeProviderService = Effect.fn("makeProviderService")(function* ( yield* Effect.logWarning("provider shutdown grace period elapsed", { timeout: String(stopAllGrace), sessionCount: activeSessions.length, - recoveryCount: recoveryByThreadId.size, }); } else { yield* Effect.forEach( @@ -1579,11 +1529,6 @@ const makeProviderService = Effect.fn("makeProviderService")(function* ( "ProviderService.stopAll", binding, ); - // Prefer the marker recorded at the start of this stopAll; fall back to - // anything already durable (partial prior shutdown / race). - const marker = - recoveryByThreadId.get(binding.threadId) ?? - readProviderRestartRecoveryMarker(binding.runtimePayload); return yield* directory.upsert({ threadId: binding.threadId, provider: binding.provider, @@ -1593,9 +1538,6 @@ const makeProviderService = Effect.fn("makeProviderService")(function* ( activeTurnId: null, lastRuntimeEvent: "provider.stopAll", lastRuntimeEventAt: yield* nowIso, - // Explicitly re-write so merge cannot leave us marker-less after we - // clear activeTurnId (legacy recovery requires starting/running + id). - restartRecovery: marker ?? null, }, }); }), diff --git a/apps/server/src/provider/ProviderRestartRecovery.test.ts b/apps/server/src/provider/ProviderRestartRecovery.test.ts index 9f646c9b17f7..27eb34000572 100644 --- a/apps/server/src/provider/ProviderRestartRecovery.test.ts +++ b/apps/server/src/provider/ProviderRestartRecovery.test.ts @@ -2,95 +2,37 @@ import { ModelSelection, ProviderInstanceId, TurnId } from "@t3tools/contracts"; import { describe, expect, it } from "vite-plus/test"; import { - makeProviderRestartRecoveryMarker, + readPersistedProviderActiveTurnId, readPersistedProviderCwd, readPersistedProviderInteractionMode, readPersistedProviderModelSelection, - readLiveShellRestartRecoveryCandidate, - readProviderRestartRecoveryCandidate, } from "./ProviderRestartRecovery.ts"; describe("ProviderRestartRecovery", () => { - it("reads typed recovery metadata and persisted restart settings", () => { + it("reads persisted session fields from the runtime payload", () => { const modelSelection: ModelSelection = { instanceId: ProviderInstanceId.make("codex-work"), model: "gpt-5.4", options: [{ id: "reasoningEffort", value: "high" }], }; - const marker = makeProviderRestartRecoveryMarker({ - interruptedProviderTurnId: TurnId.make("turn-interrupted"), - shutdownAt: "2026-07-22T00:00:00.000Z", - }); const runtimePayload = { cwd: " /tmp/project ", modelSelection, interactionMode: "plan", - restartRecovery: marker, + activeTurnId: TurnId.make("turn-live"), }; - expect( - readProviderRestartRecoveryCandidate({ - runtimePayload, - status: "stopped", - lastSeenAt: "2026-07-22T00:00:01.000Z", - }), - ).toEqual({ ...marker, source: "marker" }); expect(readPersistedProviderCwd(runtimePayload)).toBe("/tmp/project"); expect(readPersistedProviderModelSelection(runtimePayload)).toEqual(modelSelection); expect(readPersistedProviderInteractionMode(runtimePayload)).toBe("plan"); + expect(readPersistedProviderActiveTurnId(runtimePayload)).toBe(TurnId.make("turn-live")); }); - it("recognizes crash-style legacy running rows with an active turn", () => { - expect( - readProviderRestartRecoveryCandidate({ - runtimePayload: { activeTurnId: TurnId.make("turn-before-crash") }, - status: "running", - lastSeenAt: "2026-07-22T00:00:00.000Z", - }), - ).toEqual({ - version: 1, - interruptedProviderTurnId: TurnId.make("turn-before-crash"), - shutdownAt: "2026-07-22T00:00:00.000Z", - source: "legacy-active-turn", - }); - }); - - it("recovers a live orchestration turn when the adapter binding looks idle", () => { - expect( - readLiveShellRestartRecoveryCandidate({ - sessionStatus: "running", - activeTurnId: TurnId.make("turn-still-working"), - lastSeenAt: "2026-07-22T00:00:00.000Z", - }), - ).toEqual({ - version: 1, - interruptedProviderTurnId: TurnId.make("turn-still-working"), - shutdownAt: "2026-07-22T00:00:00.000Z", - source: "legacy-active-turn", - }); - expect( - readLiveShellRestartRecoveryCandidate({ - sessionStatus: "ready", - activeTurnId: TurnId.make("turn-still-working"), - lastSeenAt: "2026-07-22T00:00:00.000Z", - }), - ).toBeUndefined(); - }); - - it("does not recover idle, stopped, or malformed legacy rows", () => { - expect( - readProviderRestartRecoveryCandidate({ - runtimePayload: { activeTurnId: TurnId.make("turn-ready") }, - status: "stopped", - lastSeenAt: "2026-07-22T00:00:00.000Z", - }), - ).toBeUndefined(); - expect( - readProviderRestartRecoveryCandidate({ - runtimePayload: { activeTurnId: null }, - status: "running", - lastSeenAt: "2026-07-22T00:00:00.000Z", - }), - ).toBeUndefined(); + it("ignores missing or malformed payload fields", () => { + expect(readPersistedProviderCwd(null)).toBeUndefined(); + expect(readPersistedProviderCwd({ cwd: " " })).toBeUndefined(); + expect(readPersistedProviderModelSelection({})).toBeUndefined(); + expect(readPersistedProviderInteractionMode({ interactionMode: "nope" })).toBeUndefined(); + expect(readPersistedProviderActiveTurnId({ activeTurnId: null })).toBeUndefined(); }); }); diff --git a/apps/server/src/provider/ProviderRestartRecovery.ts b/apps/server/src/provider/ProviderRestartRecovery.ts index 02922255e1be..087bd60cae44 100644 --- a/apps/server/src/provider/ProviderRestartRecovery.ts +++ b/apps/server/src/provider/ProviderRestartRecovery.ts @@ -1,26 +1,6 @@ -import { - IsoDateTime, - ModelSelection, - ProviderInteractionMode, - TurnId, - type ProviderSessionRuntimeStatus, -} from "@t3tools/contracts"; +import { ModelSelection, ProviderInteractionMode, TurnId } from "@t3tools/contracts"; import * as Schema from "effect/Schema"; -export const PROVIDER_RESTART_RECOVERY_PAYLOAD_KEY = "restartRecovery"; - -export const ProviderRestartRecoveryMarker = Schema.Struct({ - version: Schema.Literal(1), - interruptedProviderTurnId: Schema.NullOr(TurnId), - shutdownAt: IsoDateTime, -}); -export type ProviderRestartRecoveryMarker = typeof ProviderRestartRecoveryMarker.Type; - -export interface ProviderRestartRecoveryCandidate extends ProviderRestartRecoveryMarker { - readonly source: "marker" | "legacy-active-turn"; -} - -const isProviderRestartRecoveryMarker = Schema.is(ProviderRestartRecoveryMarker); const isModelSelection = Schema.is(ModelSelection); const isProviderInteractionMode = Schema.is(ProviderInteractionMode); const isTurnId = Schema.is(TurnId); @@ -29,71 +9,6 @@ function isRecord(value: unknown): value is Record { return value !== null && typeof value === "object" && !Array.isArray(value); } -export function makeProviderRestartRecoveryMarker(input: { - readonly interruptedProviderTurnId: TurnId | null | undefined; - readonly shutdownAt: string; -}): ProviderRestartRecoveryMarker { - return { - version: 1, - interruptedProviderTurnId: input.interruptedProviderTurnId ?? null, - shutdownAt: IsoDateTime.make(input.shutdownAt), - }; -} - -export function readProviderRestartRecoveryMarker( - runtimePayload: unknown, -): ProviderRestartRecoveryMarker | undefined { - if (!isRecord(runtimePayload)) return undefined; - const marker = runtimePayload[PROVIDER_RESTART_RECOVERY_PAYLOAD_KEY]; - return isProviderRestartRecoveryMarker(marker) ? marker : undefined; -} - -export function readProviderRestartRecoveryCandidate(input: { - readonly runtimePayload: unknown; - readonly status: ProviderSessionRuntimeStatus | undefined; - readonly lastSeenAt: string; -}): ProviderRestartRecoveryCandidate | undefined { - const marker = readProviderRestartRecoveryMarker(input.runtimePayload); - if (marker !== undefined) { - return { ...marker, source: "marker" }; - } - if (input.status !== "starting" && input.status !== "running") { - return undefined; - } - if (!isRecord(input.runtimePayload)) return undefined; - const activeTurnId = input.runtimePayload.activeTurnId; - if (!isTurnId(activeTurnId)) return undefined; - return { - version: 1, - interruptedProviderTurnId: activeTurnId, - shutdownAt: IsoDateTime.make(input.lastSeenAt), - source: "legacy-active-turn", - }; -} - -/** - * Grok (and other ACP adapters) map idle-between-prompts to session - * `ready`, which persistence stores as a running binding with no - * `activeTurnId`. The orchestration shell still claims the T3 turn. - * Recover from that live shell instead of orphan-settling. - */ -export function readLiveShellRestartRecoveryCandidate(input: { - readonly sessionStatus: string | undefined | null; - readonly activeTurnId: unknown; - readonly lastSeenAt: string; -}): ProviderRestartRecoveryCandidate | undefined { - if (input.sessionStatus !== "starting" && input.sessionStatus !== "running") { - return undefined; - } - if (!isTurnId(input.activeTurnId)) return undefined; - return { - version: 1, - interruptedProviderTurnId: input.activeTurnId, - shutdownAt: IsoDateTime.make(input.lastSeenAt), - source: "legacy-active-turn", - }; -} - export function readPersistedProviderCwd(runtimePayload: unknown): string | undefined { if (!isRecord(runtimePayload)) return undefined; const cwd = runtimePayload.cwd; diff --git a/apps/server/src/provider/providerSessionContinuation.ts b/apps/server/src/provider/providerSessionContinuation.ts new file mode 100644 index 000000000000..a622dfb6c404 --- /dev/null +++ b/apps/server/src/provider/providerSessionContinuation.ts @@ -0,0 +1,32 @@ +import { TurnId } from "@t3tools/contracts"; + +/** Runtime payload key shared by in-app server updates and graceful process stop. */ +export const SERVER_UPDATE_CONTINUATION_KEY = "continueAfterServerUpdate"; +export const SERVER_UPDATE_CONTINUATION_PROMPT = "Continue where you left off."; + +export function readRuntimePayload(runtimePayload: unknown): Record { + return runtimePayload !== null && + typeof runtimePayload === "object" && + !Array.isArray(runtimePayload) + ? (runtimePayload as Record) + : {}; +} + +export function hasServerUpdateContinuationMarker( + runtimePayload: unknown, +): runtimePayload is Record { + return ( + runtimePayload !== null && + typeof runtimePayload === "object" && + !Array.isArray(runtimePayload) && + SERVER_UPDATE_CONTINUATION_KEY in runtimePayload + ); +} + +export function readServerUpdateContinuationTurnId(runtimePayload: unknown): TurnId | null { + if (!hasServerUpdateContinuationMarker(runtimePayload)) { + return null; + } + const value = runtimePayload[SERVER_UPDATE_CONTINUATION_KEY]; + return typeof value === "string" && value.length > 0 ? TurnId.make(value) : null; +} diff --git a/apps/server/src/serverRuntimeStartup.reconcile.test.ts b/apps/server/src/serverRuntimeStartup.reconcile.test.ts index 0553afb43a11..df42ff59afb1 100644 --- a/apps/server/src/serverRuntimeStartup.reconcile.test.ts +++ b/apps/server/src/serverRuntimeStartup.reconcile.test.ts @@ -19,7 +19,6 @@ import * as ProjectionSnapshotQuery from "./orchestration/Services/ProjectionSna import { ProviderSessionDirectoryPersistenceError } from "./provider/Errors.ts"; import * as ProviderService from "./provider/Services/ProviderService.ts"; import * as ProviderSessionDirectory from "./provider/Services/ProviderSessionDirectory.ts"; -import { makeProviderRestartRecoveryMarker } from "./provider/ProviderRestartRecovery.ts"; import * as ServerRuntimeStartup from "./serverRuntimeStartup.ts"; const providerInstanceId = ProviderInstanceId.make("codex"); @@ -438,49 +437,6 @@ it.effect("retries continuation preparation before settling a persistent failure ); }); -it.effect("does not settle Tim Smart restart-recovery markers as update orphans", () => { - const recovering = makeThread("thread-restart-recovery", "running", TurnId.make("turn-live")); - const dispatched: OrchestrationCommand[] = []; - const upserts: ProviderSessionDirectory.ProviderRuntimeBinding[] = []; - const marker = makeProviderRestartRecoveryMarker({ - interruptedProviderTurnId: TurnId.make("turn-live"), - shutdownAt: updatedAt, - }); - - return runReconciliation({ - threads: [recovering], - directory: { - getBinding: (candidate) => - Effect.succeed( - Option.some({ - threadId: candidate, - provider: ProviderDriverKind.make("codex"), - providerInstanceId, - status: "stopped" as const, - resumeCursor: { cursor: candidate }, - runtimePayload: { - activeTurnId: null, - restartRecovery: marker, - }, - }), - ), - upsert: (binding) => Effect.sync(() => upserts.push(binding)), - getProvider: () => Effect.die("unused"), - listThreadIds: () => Effect.die("unused"), - listBindings: () => Effect.die("unused"), - }, - dispatch: (command) => - Effect.sync(() => dispatched.push(command)).pipe(Effect.as({ sequence: dispatched.length })), - }).pipe( - Effect.tap(() => - Effect.sync(() => { - assert.deepStrictEqual(dispatched, []); - assert.deepStrictEqual(upserts, []); - }), - ), - ); -}); - it.effect("reconciles multiple active and archived orphans but skips live sessions", () => { const starting = makeThread("thread-starting", "starting"); const running = makeThread("thread-running", "running", TurnId.make("turn-running")); diff --git a/apps/server/src/serverRuntimeStartup.ts b/apps/server/src/serverRuntimeStartup.ts index 3e3885df0da6..47a36d68efb7 100644 --- a/apps/server/src/serverRuntimeStartup.ts +++ b/apps/server/src/serverRuntimeStartup.ts @@ -7,7 +7,6 @@ import { ProjectId, ProviderInstanceId, ThreadId, - TurnId, } from "@t3tools/contracts"; import * as Cause from "effect/Cause"; import * as Console from "effect/Console"; @@ -43,8 +42,14 @@ import * as AnalyticsService from "./telemetry/AnalyticsService.ts"; import * as ServerEnvironment from "./environment/ServerEnvironment.ts"; import * as EnvironmentAuth from "./auth/EnvironmentAuth.ts"; import * as ProviderService from "./provider/Services/ProviderService.ts"; -import { readProviderRestartRecoveryMarker } from "./provider/ProviderRestartRecovery.ts"; import * as ProviderSessionDirectory from "./provider/Services/ProviderSessionDirectory.ts"; +import { + hasServerUpdateContinuationMarker, + readRuntimePayload, + readServerUpdateContinuationTurnId, + SERVER_UPDATE_CONTINUATION_KEY, + SERVER_UPDATE_CONTINUATION_PROMPT, +} from "./provider/providerSessionContinuation.ts"; import * as ProviderSessionReaper from "./provider/Services/ProviderSessionReaper.ts"; import * as OrphanSessionRecovery from "./orchestration/Services/OrphanSessionRecovery.ts"; import { forkParked } from "./serverActivation.ts"; @@ -342,8 +347,6 @@ export function interruptSessionAfterServerRestart( const ORPHANED_PROVIDER_SESSION_ERROR = "Provider session did not survive a server restart. Send a new message to continue."; -const SERVER_UPDATE_CONTINUATION_KEY = "continueAfterServerUpdate"; -const SERVER_UPDATE_CONTINUATION_PROMPT = "Continue where you left off."; class ProviderSessionContinuationError extends Schema.TaggedErrorClass()( "ProviderSessionContinuationError", @@ -367,35 +370,8 @@ export class ServerUpdateThreadContinuationError extends Schema.TaggedErrorClass } } -function hasServerUpdateContinuationMarker( - runtimePayload: unknown, -): runtimePayload is Record { - return ( - runtimePayload !== null && - typeof runtimePayload === "object" && - !Array.isArray(runtimePayload) && - SERVER_UPDATE_CONTINUATION_KEY in runtimePayload - ); -} - -function readRuntimePayload(runtimePayload: unknown): Record { - return runtimePayload !== null && - typeof runtimePayload === "object" && - !Array.isArray(runtimePayload) - ? (runtimePayload as Record) - : {}; -} - const isServerUpdateThreadContinuationError = Schema.is(ServerUpdateThreadContinuationError); -function readServerUpdateContinuationTurnId(runtimePayload: unknown): TurnId | null { - if (!hasServerUpdateContinuationMarker(runtimePayload)) { - return null; - } - const value = runtimePayload[SERVER_UPDATE_CONTINUATION_KEY]; - return typeof value === "string" && value.length > 0 ? TurnId.make(value) : null; -} - const toServerUpdateThreadContinuationError = (cause: unknown) => isServerUpdateThreadContinuationError(cause) ? cause @@ -518,12 +494,6 @@ export const reconcileProviderSessions = Effect.gen(function* () { const continuationMarked = continuationTurnId !== null && (session.activeTurnId === null || continuationTurnId === session.activeTurnId); - const restartRecoveryMarked = - Option.isSome(binding) && - readProviderRestartRecoveryMarker(binding.value.runtimePayload) !== undefined; - if (restartRecoveryMarked) { - continue; - } const settleAsError = (lastError: string) => Effect.gen(function* () { yield* Effect.gen(function* () { diff --git a/apps/web/src/components/settings/SettingsPanels.tsx b/apps/web/src/components/settings/SettingsPanels.tsx index cb725bfc8931..fade18ac7a7d 100644 --- a/apps/web/src/components/settings/SettingsPanels.tsx +++ b/apps/web/src/components/settings/SettingsPanels.tsx @@ -2205,7 +2205,7 @@ export function GeneralSettingsPanel() { { "new-threads", "start-from-origin", ]); + expect(searchSettings("deploy").map((item) => item.id)).toContain( + "continue-threads-after-server-update", + ); expect(searchSettings("glass").map((item) => item.id)).toEqual(["setting-glass-opacity"]); expect(searchSettings("thè\u{1ab0}mes")[0]?.id).toBe("theme"); const localeLowerCase = vi.spyOn(String.prototype, "toLocaleLowerCase").mockReturnValue("gıt"); diff --git a/apps/web/src/components/settings/settingsSearch.ts b/apps/web/src/components/settings/settingsSearch.ts index 48d23583b7a9..51a54bf62d54 100644 --- a/apps/web/src/components/settings/settingsSearch.ts +++ b/apps/web/src/components/settings/settingsSearch.ts @@ -195,7 +195,7 @@ export const SETTINGS_SEARCH_ITEMS = [ id: "continue-threads-after-server-update", title: "Continue threads after server updates", to: "/settings/general", - searchTerms: ["resume running active work restart desktop update automatically"], + searchTerms: ["resume running active work restart desktop update automatically deploy"], }, { id: "background-activity", diff --git a/docs/user/updating.md b/docs/user/updating.md index 18fc82ad202d..ef08692d098a 100644 --- a/docs/user/updating.md +++ b/docs/user/updating.md @@ -15,12 +15,13 @@ update the server, and the version difference remains visible in Connections. ## Before You Update -Updating restarts the server, so the connection will disappear briefly. **Settings** → **General** -has a **Continue threads after server updates** preference. It is off by default. When enabled, the -update buttons automatically resume supported provider threads after the replacement server is -ready. Providers with native promptless continuation use it; other providers receive a short -instruction to continue where they left off. Terminal commands and other running work may still be -interrupted during the update. +Updating restarts the server, so the connection will disappear briefly. A graceful server restart +(including a deploy that stops and starts the service) resumes supported provider threads that were +actively working. **Settings** → **General** has a **Continue threads after server updates** +preference. It is off by default. When enabled, the in-app update buttons also resume those threads +after the replacement server is ready. Providers with native promptless continuation use it; other +providers receive a short instruction to continue where they left off. Terminal commands and other +running work may still be interrupted during the update. A crash or forced kill does not resume. The update does not remove saved threads, settings, or project files.