diff --git a/apps/server/src/provider/Layers/ProviderService.test.ts b/apps/server/src/provider/Layers/ProviderService.test.ts index 5738c86d228c..e7cd87c1358e 100644 --- a/apps/server/src/provider/Layers/ProviderService.test.ts +++ b/apps/server/src/provider/Layers/ProviderService.test.ts @@ -73,7 +73,10 @@ 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 { readProviderRestartRecoveryMarker } from "../ProviderRestartRecovery.ts"; +import { + makeProviderRestartRecoveryMarker, + readProviderRestartRecoveryMarker, +} from "../ProviderRestartRecovery.ts"; const defaultServerSettingsLayer = ServerSettings.ServerSettingsService.layerTest(); const serverConfigTestLayer = ServerConfig.layerTest(process.cwd(), process.cwd()).pipe( @@ -661,6 +664,92 @@ it.effect("graceful shutdown recovers live bindings missing from adapter listSes }).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 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)), +); + it.effect("ProviderServiceLive rejects new sessions for disabled providers", () => Effect.gen(function* () { const codex = makeFakeCodexAdapter(); diff --git a/apps/server/src/provider/Layers/ProviderService.ts b/apps/server/src/provider/Layers/ProviderService.ts index 4407935d2e52..66b3b9670045 100644 --- a/apps/server/src/provider/Layers/ProviderService.ts +++ b/apps/server/src/provider/Layers/ProviderService.ts @@ -366,6 +366,18 @@ 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, @@ -374,7 +386,7 @@ const makeProviderService = Effect.fn("makeProviderService")(function* ( status: lifecycle.status, runtimePayload: { ...(lifecycle.activeTurnId !== undefined ? { activeTurnId: lifecycle.activeTurnId } : {}), - ...(lifecycle.restartRecovery !== undefined + ...(lifecycle.restartRecovery !== undefined && !preserveShutdownRecoveryMarker ? { restartRecovery: lifecycle.restartRecovery } : {}), lastRuntimeEvent: event.type, diff --git a/apps/server/src/serverRuntimeStartup.reconcile.test.ts b/apps/server/src/serverRuntimeStartup.reconcile.test.ts index df42ff59afb1..0553afb43a11 100644 --- a/apps/server/src/serverRuntimeStartup.reconcile.test.ts +++ b/apps/server/src/serverRuntimeStartup.reconcile.test.ts @@ -19,6 +19,7 @@ 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"); @@ -437,6 +438,49 @@ 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 6a98eaffb7d4..3e3885df0da6 100644 --- a/apps/server/src/serverRuntimeStartup.ts +++ b/apps/server/src/serverRuntimeStartup.ts @@ -43,6 +43,7 @@ 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 * as ProviderSessionReaper from "./provider/Services/ProviderSessionReaper.ts"; import * as OrphanSessionRecovery from "./orchestration/Services/OrphanSessionRecovery.ts"; @@ -517,6 +518,12 @@ 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* () {