Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
91 changes: 90 additions & 1 deletion apps/server/src/provider/Layers/ProviderService.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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(
Expand Down Expand Up @@ -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();
Expand Down
14 changes: 13 additions & 1 deletion apps/server/src/provider/Layers/ProviderService.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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<string, unknown>)
: 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,
Expand All @@ -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,
Expand Down
44 changes: 44 additions & 0 deletions apps/server/src/serverRuntimeStartup.reconcile.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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");
Expand Down Expand Up @@ -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"));
Expand Down
7 changes: 7 additions & 0 deletions apps/server/src/serverRuntimeStartup.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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";
Expand Down Expand Up @@ -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* () {
Expand Down
Loading