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
72 changes: 72 additions & 0 deletions apps/server/src/provider/Layers/ProviderService.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -648,6 +648,78 @@ it.effect("ProviderServiceLive catches stopAll failures during shutdown", () =>
}),
);

it.effect("ProviderServiceLive shutdown leaves settled session rows untouched", () =>
Effect.gen(function* () {
const recordedAnalytics = makeRecordingAnalytics();
const codex = makeFakeCodexAdapter();
const persistence = yield* Layer.build(
ProviderSessionDirectoryLive.pipe(
Layer.provide(ProviderSessionRuntime.layer.pipe(Layer.provide(SqlitePersistenceMemory))),
),
);
const directory = yield* ProviderSessionDirectory.ProviderSessionDirectory.pipe(
Effect.provide(persistence),
);
const seed = (threadId: ThreadId, status: "running" | "stopped", activeTurnId: TurnId | null) =>
directory.upsert({
threadId,
provider: CODEX_DRIVER,
providerInstanceId: codexInstanceId,
status,
runtimePayload: { cwd: "/repo", activeTurnId },
});
const readBindings = directory
.listBindings()
.pipe(
Effect.map((bindings) => new Map(bindings.map((binding) => [binding.threadId, binding]))),
);
const settledId = asThreadId("shutdown-settled");
const runningId = asThreadId("shutdown-running");
const stoppedWithTurnId = asThreadId("shutdown-stopped-with-turn");
yield* seed(settledId, "stopped", null);
yield* seed(runningId, "running", asTurnId("running-turn"));
yield* seed(stoppedWithTurnId, "stopped", asTurnId("stale-turn"));
const settledBefore = (yield* readBindings).get(settledId);
assert(settledBefore !== undefined);

const scope = yield* Scope.make();
yield* Layer.build(
makeProviderServiceLive().pipe(
Layer.provide(NodeServices.layer),
Layer.provide(Layer.succeed(ProviderSessionDirectory.ProviderSessionDirectory, directory)),
Layer.provide(
Layer.succeed(
ProviderAdapterRegistry.ProviderAdapterRegistry,
makeStaticInstanceRegistry([[codexInstanceId, codex.adapter]]),
),
),
Layer.provide(defaultServerSettingsLayer),
Layer.provide(serverConfigTestLayer),
Layer.provide(recordedAnalytics.layer),
Layer.provide(
Layer.succeed(
ProviderEventLoggers.ProviderEventLoggers,
ProviderEventLoggers.NoOpProviderEventLoggers,
),
),
),
).pipe(Scope.provide(scope));
yield* TestClock.adjust("1 minute");
yield* Scope.close(scope, Exit.void);

const byThread = yield* readBindings;
assert.deepStrictEqual(byThread.get(settledId), settledBefore);
for (const threadId of [runningId, stoppedWithTurnId]) {
const binding = byThread.get(threadId);
assert.equal(binding?.status, "stopped");
assert.propertyVal(binding?.runtimePayload, "activeTurnId", null);
assert.propertyVal(binding?.runtimePayload, "lastRuntimeEvent", "provider.stopAll");
}
const [stoppedAll] = recordedAnalytics.eventsByName("provider.sessions.stopped_all");
assert.equal(stoppedAll?.properties?.stoppedSessionCount, 2);
}).pipe(Effect.provide(NodeServices.layer)),
);

it.effect("ProviderServiceLive flushes deferred completions during shutdown", () =>
Effect.gen(function* () {
const recordedAnalytics = makeRecordingAnalytics();
Expand Down
20 changes: 17 additions & 3 deletions apps/server/src/provider/Layers/ProviderService.ts
Original file line number Diff line number Diff line change
Expand Up @@ -427,6 +427,14 @@ function readPersistedCwd(
return trimmed.length > 0 ? trimmed : undefined;
}

/** Stopped rows with no active turn are settled; shutdown leaves them untouched. */
function isSettledBinding(binding: ProviderSessionDirectory.ProviderRuntimeBinding): boolean {
if (binding.status !== "stopped") return false;
const payload = binding.runtimePayload;
if (!payload || typeof payload !== "object" || Array.isArray(payload)) return true;
return !("activeTurnId" in payload) || payload.activeTurnId == null;
}

const dieOnMissingBindingInstanceId = (
operation: string,
payload: {
Expand Down Expand Up @@ -2331,7 +2339,6 @@ const makeProviderService = Effect.fn("makeProviderService")(function* (
return [completed, state] as const;
});
yield* recordCompletedTurnProperties(properties);
const threadIds = yield* directory.listThreadIds();
const currentAdapters = yield* getAdapterEntries;
const activeSessions = yield* Effect.forEach(currentAdapters, ([instanceId, adapter]) =>
adapter.listSessions().pipe(
Expand Down Expand Up @@ -2362,7 +2369,12 @@ const makeProviderService = Effect.fn("makeProviderService")(function* (
yield* Effect.forEach(currentAdapters, ([, adapter]) => adapter.stopAll()).pipe(Effect.asVoid);
yield* McpSessionRegistry.revokeAllActiveMcpCredentials();
McpProviderSession.clearAllMcpProviderSessions();
const bindings = yield* directory.listBindings().pipe(Effect.orElseSucceed(() => []));
// Stopped rows stay for their resume cursors, so long-lived installs hold
// thousands. Only rewrite the ones this shutdown actually stops.
const bindings = yield* directory.listBindings().pipe(
Effect.map((all) => all.filter((binding) => !isSettledBinding(binding))),
Effect.orElseSucceed(() => []),
);
yield* Effect.forEach(bindings, (binding) =>
Effect.gen(function* () {
const providerInstanceId = dieOnMissingBindingInstanceId(
Expand All @@ -2382,8 +2394,10 @@ const makeProviderService = Effect.fn("makeProviderService")(function* (
});
}),
).pipe(Effect.asVoid);
// Not `sessionCount`: that older property counted every row, so a new name
// keeps the two meanings in separate series.
yield* analytics.record("provider.sessions.stopped_all", {
sessionCount: threadIds.length,
stoppedSessionCount: bindings.length,
});
yield* analytics.flush;
});
Expand Down
Loading