From 9c6668648d259c749e6b0c0cba3748361dca9db2 Mon Sep 17 00:00:00 2001 From: Guillermo Casanova Date: Mon, 14 Sep 2026 22:57:40 -0300 Subject: [PATCH 1/2] fix(server): preserve provider sessions after host sleep --- .../Layers/ProviderSessionReaper.test.ts | 115 ++++++++++++++++-- .../provider/Layers/ProviderSessionReaper.ts | 25 +++- 2 files changed, 130 insertions(+), 10 deletions(-) diff --git a/apps/server/src/provider/Layers/ProviderSessionReaper.test.ts b/apps/server/src/provider/Layers/ProviderSessionReaper.test.ts index 7b1fec90f867..b34dd940be84 100644 --- a/apps/server/src/provider/Layers/ProviderSessionReaper.test.ts +++ b/apps/server/src/provider/Layers/ProviderSessionReaper.test.ts @@ -150,20 +150,23 @@ describe("ProviderSessionReaper", () => { ); } - async function sweepAt(nowMs: number) { + async function sweepAt( + nowMs: number | { wallMs: number; monotonicMs: number }, + monotonicMs = typeof nowMs === "number" ? nowMs : nowMs.monotonicMs, + ) { + const time = typeof nowMs === "number" ? { wallMs: nowMs, monotonicMs } : nowMs; await runtime!.runPromise( Effect.gen(function* () { const reaper = yield* ProviderSessionReaper; - const clock = yield* Clock.Clock; const swept = yield* Deferred.make(); yield* reaper.start().pipe( Effect.provideService(Clock.Clock, { - currentTimeMillis: Effect.succeed(nowMs), - currentTimeMillisUnsafe: () => nowMs, - currentTimeNanos: Effect.succeed(BigInt(nowMs) * 1_000_000n), - currentTimeNanosUnsafe: () => BigInt(nowMs) * 1_000_000n, - monotonicTimeNanos: clock.monotonicTimeNanos, - monotonicTimeNanosUnsafe: () => clock.monotonicTimeNanosUnsafe(), + currentTimeMillis: Effect.sync(() => time.wallMs), + currentTimeMillisUnsafe: () => time.wallMs, + currentTimeNanos: Effect.sync(() => BigInt(time.wallMs) * 1_000_000n), + currentTimeNanosUnsafe: () => BigInt(time.wallMs) * 1_000_000n, + monotonicTimeNanos: Effect.sync(() => BigInt(time.monotonicMs) * 1_000_000n), + monotonicTimeNanosUnsafe: () => BigInt(time.monotonicMs) * 1_000_000n, // Reaching the next scheduled sleep proves this sweep has finished. sleep: () => Deferred.succeed(swept, undefined).pipe(Effect.andThen(Effect.never)), }), @@ -521,6 +524,102 @@ describe("ProviderSessionReaper", () => { }, ); + it.each([ + { sleepMs: 8 * 60 * 60 * 1_000, monotonicAdvances: false }, + { sleepMs: 8 * 60 * 60 * 1_000, monotonicAdvances: true }, + { sleepMs: 2_000, monotonicAdvances: false }, + ])( + "gives sessions a fresh idle window after $sleepMs ms sleep, monotonic clock advances=$monotonicAdvances", + async ({ sleepMs, monotonicAdvances }) => { + const threadId = ThreadId.make("thread-reaper-sleep"); + const now = "2026-04-14T01:00:00.000Z"; + const nowMs = Date.parse(now); + const harness = await createHarness({ + readModel: makeReadModel([{ id: threadId, session: null }]), + }); + const repository = await runtime!.runPromise( + Effect.service(ProviderSessionRuntime.ProviderSessionRuntimeRepository), + ); + await runtime!.runPromise( + repository.upsert({ + threadId, + providerName: "claudeAgent", + providerInstanceId: null, + adapterKey: "claudeAgent", + runtimeMode: "full-access", + status: "running", + lastSeenAt: now, + resumeCursor: { opaque: "resume-sleep" }, + runtimePayload: null, + }), + ); + + await sweepAt(nowMs, 0); + expect(harness.stopSession).not.toHaveBeenCalled(); + + const resumedAt = nowMs + sleepMs; + const monotonicMs = monotonicAdvances ? sleepMs : 0; + await sweepAt(resumedAt, monotonicMs); + expect(harness.stopSession).not.toHaveBeenCalled(); + await sweepAt(resumedAt + 999, monotonicMs + 999); + expect(harness.stopSession).not.toHaveBeenCalled(); + const secondResumeAt = resumedAt + 999 + sleepMs; + const secondMonotonicMs = monotonicMs + 999 + (monotonicAdvances ? sleepMs : 0); + await sweepAt(secondResumeAt, secondMonotonicMs); + expect(harness.stopSession).not.toHaveBeenCalled(); + await sweepAt(secondResumeAt + 999, secondMonotonicMs + 999); + expect(harness.stopSession).not.toHaveBeenCalled(); + await sweepAt(secondResumeAt + 1_000, secondMonotonicMs + 1_000); + expect(harness.stopSession).toHaveBeenCalledExactlyOnceWith({ threadId }); + }, + ); + + it("does not mistake a slow sweep for host sleep", async () => { + const threadId = ThreadId.make("thread-reaper-slow-stop"); + const lastSeenAt = "2026-04-14T01:00:00.000Z"; + const time = { wallMs: Date.parse(lastSeenAt) + 1_000, monotonicMs: 0 }; + const harness = await createHarness({ + readModel: makeReadModel([{ id: threadId, session: null }]), + stopSessionImplementation: () => + Effect.sync(() => { + time.wallMs += 2_000; + time.monotonicMs += 2_000; + }).pipe( + Effect.andThen( + Effect.fail( + new ProviderValidationError({ + operation: "ProviderSessionReaper.test", + issue: "slow stop failed; retry on the next sweep", + }), + ), + ), + ), + }); + const repository = await runtime!.runPromise( + Effect.service(ProviderSessionRuntime.ProviderSessionRuntimeRepository), + ); + await runtime!.runPromise( + repository.upsert({ + threadId, + providerName: "claudeAgent", + providerInstanceId: null, + adapterKey: "claudeAgent", + runtimeMode: "full-access", + status: "running", + lastSeenAt, + resumeCursor: { opaque: "resume-slow-stop" }, + runtimePayload: null, + }), + ); + await sweepAt(time); + expect(harness.stopSession).toHaveBeenCalledTimes(1); + + time.wallMs += 60_000; + time.monotonicMs += 60_000; + await sweepAt(time); + expect(harness.stopSession).toHaveBeenCalledTimes(2); + }); + it("skips persisted sessions that are already marked stopped", async () => { const threadId = ThreadId.make("thread-reaper-stopped"); const now = "2026-01-01T00:00:00.000Z"; diff --git a/apps/server/src/provider/Layers/ProviderSessionReaper.ts b/apps/server/src/provider/Layers/ProviderSessionReaper.ts index bf8199f80eac..fb461c4e47f8 100644 --- a/apps/server/src/provider/Layers/ProviderSessionReaper.ts +++ b/apps/server/src/provider/Layers/ProviderSessionReaper.ts @@ -16,6 +16,7 @@ import { ProviderService } from "../Services/ProviderService.ts"; const DEFAULT_INACTIVITY_THRESHOLD_MS = 30 * 60 * 1000; const DEFAULT_SWEEP_INTERVAL_MS = 5 * 60 * 1000; +const CLOCK_JUMP_TOLERANCE_MS = 1_000; export interface ProviderSessionReaperLiveOptions { readonly inactivityThresholdMs?: number; @@ -33,12 +34,26 @@ const makeProviderSessionReaper = (options?: ProviderSessionReaperLiveOptions) = options?.inactivityThresholdMs ?? DEFAULT_INACTIVITY_THRESHOLD_MS, ); const sweepIntervalMs = Math.max(1, options?.sweepIntervalMs ?? DEFAULT_SWEEP_INTERVAL_MS); + let previousSweep: { wallMs: number; monotonicMs: number } | undefined; + let lastResumeMs = 0; const sweep = Effect.gen(function* () { + const now = yield* Clock.currentTimeMillis; + const monotonicMs = Number(yield* Clock.monotonicTimeNanos) / 1_000_000; + if ( + previousSweep && + now - + previousSweep.wallMs - + Math.min(monotonicMs - previousSweep.monotonicMs, sweepIntervalMs) > + CLOCK_JUMP_TOLERANCE_MS + ) { + // Monotonic clocks may include sleep, so also detect overdue sweeps. + // ponytail: scheduler stalls also grant grace; use power events for exact accounting. + lastResumeMs = now; + } // Stopped rows stay for their resume cursors and far outnumber live // ones, so the query skips them. const bindings = yield* directory.listBindings({ excludeStopped: true }); - const now = yield* Clock.currentTimeMillis; let reapedCount = 0; for (const binding of bindings) { @@ -52,7 +67,7 @@ const makeProviderSessionReaper = (options?: ProviderSessionReaperLiveOptions) = continue; } - if (now - lastSeenMs < inactivityThresholdMs) { + if (now - Math.max(lastSeenMs, lastResumeMs) < inactivityThresholdMs) { continue; } @@ -64,6 +79,7 @@ const makeProviderSessionReaper = (options?: ProviderSessionReaperLiveOptions) = // even though the binding was last touched when the turn was sent. const lastActivityMs = Math.max( lastSeenMs, + lastResumeMs, Date.parse(thread?.session?.updatedAt ?? binding.lastSeenAt), ); const idleDurationMs = now - lastActivityMs; @@ -123,6 +139,11 @@ const makeProviderSessionReaper = (options?: ProviderSessionReaperLiveOptions) = liveBindings: bindings.length, }); } + // Measure only the scheduled wait, excluding time spent reaping sessions. + previousSweep = { + wallMs: yield* Clock.currentTimeMillis, + monotonicMs: Number(yield* Clock.monotonicTimeNanos) / 1_000_000, + }; }); const start: ProviderSessionReaperShape["start"] = () => From 5cd92799f012432018f409e21ccc939ffc1be435 Mon Sep 17 00:00:00 2001 From: Guillermo Casanova Date: Mon, 14 Sep 2026 23:30:21 -0300 Subject: [PATCH 2/2] fix(server): retain resume detection across session sweeps --- .../Layers/ProviderSessionReaper.test.ts | 99 +++++++++++++++++-- .../provider/Layers/ProviderSessionReaper.ts | 33 ++++--- 2 files changed, 111 insertions(+), 21 deletions(-) diff --git a/apps/server/src/provider/Layers/ProviderSessionReaper.test.ts b/apps/server/src/provider/Layers/ProviderSessionReaper.test.ts index b34dd940be84..1c177392736b 100644 --- a/apps/server/src/provider/Layers/ProviderSessionReaper.test.ts +++ b/apps/server/src/provider/Layers/ProviderSessionReaper.test.ts @@ -22,6 +22,7 @@ import { SqlitePersistenceMemory } from "../../persistence/Layers/Sqlite.ts"; import * as ProviderSessionRuntime from "../../persistence/ProviderSessionRuntime.ts"; import { ProviderValidationError } from "../Errors.ts"; import { ProviderSessionReaper } from "../Services/ProviderSessionReaper.ts"; +import { ProviderSessionDirectory } from "../Services/ProviderSessionDirectory.ts"; import { ProviderService, type ProviderServiceShape } from "../Services/ProviderService.ts"; import { ProviderSessionDirectoryLive } from "./ProviderSessionDirectory.ts"; import { makeProviderSessionReaperLive } from "./ProviderSessionReaper.ts"; @@ -122,7 +123,10 @@ function makeReadModel( describe("ProviderSessionReaper", () => { let runtime: ManagedRuntime.ManagedRuntime< - ProviderSessionReaper | ProviderSessionRuntime.ProviderSessionRuntimeRepository, + | ProviderSessionReaper + | ProviderSessionRuntime.ProviderSessionRuntimeRepository + | ProviderSessionDirectory + | ProjectionSnapshotQuery, unknown > | null = null; let scope: Scope.Closeable | null = null; @@ -153,6 +157,7 @@ describe("ProviderSessionReaper", () => { async function sweepAt( nowMs: number | { wallMs: number; monotonicMs: number }, monotonicMs = typeof nowMs === "number" ? nowMs : nowMs.monotonicMs, + afterWallRead?: () => void, ) { const time = typeof nowMs === "number" ? { wallMs: nowMs, monotonicMs } : nowMs; await runtime!.runPromise( @@ -161,7 +166,12 @@ describe("ProviderSessionReaper", () => { const swept = yield* Deferred.make(); yield* reaper.start().pipe( Effect.provideService(Clock.Clock, { - currentTimeMillis: Effect.sync(() => time.wallMs), + currentTimeMillis: Effect.sync(() => { + const wallMs = time.wallMs; + afterWallRead?.(); + afterWallRead = undefined; + return wallMs; + }), currentTimeMillisUnsafe: () => time.wallMs, currentTimeNanos: Effect.sync(() => BigInt(time.wallMs) * 1_000_000n), currentTimeNanosUnsafe: () => BigInt(time.wallMs) * 1_000_000n, @@ -614,12 +624,89 @@ describe("ProviderSessionReaper", () => { await sweepAt(time); expect(harness.stopSession).toHaveBeenCalledTimes(1); - time.wallMs += 60_000; - time.monotonicMs += 60_000; - await sweepAt(time); - expect(harness.stopSession).toHaveBeenCalledTimes(2); + for (let attempt = 2; attempt <= 4; attempt++) { + time.wallMs += 62_000; + time.monotonicMs += 62_000; + await sweepAt(time); + expect(harness.stopSession).toHaveBeenCalledTimes(attempt); + } }); + it.each( + [false, true].flatMap((monotonicAdvances) => + (["list", "query", "stop", "clock"] as const).flatMap((stage) => + (stage === "stop" ? [true] : stage === "clock" ? [false] : [false, true]).map((fails) => ({ + stage, + fails, + monotonicAdvances, + })), + ), + ), + )( + "preserves resume grace across $stage, fails=$fails, monotonic advances=$monotonicAdvances", + async ({ stage, fails, monotonicAdvances }) => { + const threadId = ThreadId.make("thread-reaper-mid-sweep-sleep"); + const lastSeenAt = "2026-04-14T01:00:00.000Z"; + const time = { wallMs: Date.parse(lastSeenAt) + 1_000, monotonicMs: 0 }; + const harness = await createHarness({ + readModel: makeReadModel([{ id: threadId, session: null }]), + }); + const repository = await runtime!.runPromise( + Effect.service(ProviderSessionRuntime.ProviderSessionRuntimeRepository), + ); + await runtime!.runPromise( + repository.upsert({ + threadId, + providerName: "claudeAgent", + providerInstanceId: null, + adapterKey: "claudeAgent", + runtimeMode: "full-access", + status: "running", + lastSeenAt, + resumeCursor: { opaque: "resume-mid-sweep" }, + runtimePayload: null, + }), + ); + const advanceTime = () => { + const sleepMs = 8 * 60 * 60 * 1_000; + time.wallMs += sleepMs; + if (monotonicAdvances) time.monotonicMs += sleepMs; + }; + const suspend = Effect.sync(advanceTime).pipe( + Effect.andThen(fails ? Effect.die("failed after resume") : Effect.void), + ); + if (stage === "list") { + const directory = await runtime!.runPromise(Effect.service(ProviderSessionDirectory)); + const listBindings = directory.listBindings; + vi.spyOn(directory, "listBindings").mockImplementationOnce(() => + suspend.pipe(Effect.andThen(listBindings())), + ); + } else if (stage === "query") { + const query = await runtime!.runPromise(Effect.service(ProjectionSnapshotQuery)); + const getThread = query.getThreadShellById; + vi.spyOn(query, "getThreadShellById").mockImplementationOnce((id) => + suspend.pipe(Effect.andThen(getThread(id))), + ); + } else if (stage === "stop") { + harness.stopSession.mockImplementationOnce(() => suspend); + } + + const callsBeforeResume = stage === "stop" ? 1 : 0; + await sweepAt(time, time.monotonicMs, stage === "clock" ? advanceTime : undefined); + expect(harness.stopSession).toHaveBeenCalledTimes(callsBeforeResume); + await sweepAt(time); + expect(harness.stopSession).toHaveBeenCalledTimes(callsBeforeResume); + time.wallMs += 999; + time.monotonicMs += 999; + await sweepAt(time); + expect(harness.stopSession).toHaveBeenCalledTimes(callsBeforeResume); + time.wallMs += 1; + time.monotonicMs += 1; + await sweepAt(time); + expect(harness.stopSession).toHaveBeenCalledTimes(callsBeforeResume + 1); + }, + ); + it("skips persisted sessions that are already marked stopped", async () => { const threadId = ThreadId.make("thread-reaper-stopped"); const now = "2026-01-01T00:00:00.000Z"; diff --git a/apps/server/src/provider/Layers/ProviderSessionReaper.ts b/apps/server/src/provider/Layers/ProviderSessionReaper.ts index fb461c4e47f8..ba651656b3ff 100644 --- a/apps/server/src/provider/Layers/ProviderSessionReaper.ts +++ b/apps/server/src/provider/Layers/ProviderSessionReaper.ts @@ -34,26 +34,32 @@ const makeProviderSessionReaper = (options?: ProviderSessionReaperLiveOptions) = options?.inactivityThresholdMs ?? DEFAULT_INACTIVITY_THRESHOLD_MS, ); const sweepIntervalMs = Math.max(1, options?.sweepIntervalMs ?? DEFAULT_SWEEP_INTERVAL_MS); - let previousSweep: { wallMs: number; monotonicMs: number } | undefined; + let previousSample: { wallMs: number; monotonicMs: number } | undefined; let lastResumeMs = 0; - const sweep = Effect.gen(function* () { - const now = yield* Clock.currentTimeMillis; + const checkForResume = Effect.gen(function* () { const monotonicMs = Number(yield* Clock.monotonicTimeNanos) / 1_000_000; + const now = yield* Clock.currentTimeMillis; if ( - previousSweep && - now - - previousSweep.wallMs - - Math.min(monotonicMs - previousSweep.monotonicMs, sweepIntervalMs) > - CLOCK_JUMP_TOLERANCE_MS + previousSample && + (now - previousSample.wallMs - (monotonicMs - previousSample.monotonicMs) > + CLOCK_JUMP_TOLERANCE_MS || + now - previousSample.wallMs > sweepIntervalMs * 2) ) { - // Monotonic clocks may include sleep, so also detect overdue sweeps. - // ponytail: scheduler stalls also grant grace; use power events for exact accounting. + // Some clocks count suspend. Require a missed sweep, not ordinary timer lateness. + // ponytail: multi-minute stalls also grant grace; use power events for exact accounting. lastResumeMs = now; } + previousSample = { wallMs: now, monotonicMs }; + return now; + }); + + const sweep = Effect.gen(function* () { + yield* checkForResume; // Stopped rows stay for their resume cursors and far outnumber live // ones, so the query skips them. const bindings = yield* directory.listBindings({ excludeStopped: true }); + let now = yield* checkForResume; let reapedCount = 0; for (const binding of bindings) { @@ -74,6 +80,7 @@ const makeProviderSessionReaper = (options?: ProviderSessionReaperLiveOptions) = const thread = yield* projectionSnapshotQuery .getThreadShellById(binding.threadId) .pipe(Effect.map(Option.getOrUndefined)); + now = yield* checkForResume; // Ingestion updates this timestamp alongside activeTurnId when a turn // settles. Long turns must get a full idle window after that transition, // even though the binding was last touched when the turn was sent. @@ -139,11 +146,7 @@ const makeProviderSessionReaper = (options?: ProviderSessionReaperLiveOptions) = liveBindings: bindings.length, }); } - // Measure only the scheduled wait, excluding time spent reaping sessions. - previousSweep = { - wallMs: yield* Clock.currentTimeMillis, - monotonicMs: Number(yield* Clock.monotonicTimeNanos) / 1_000_000, - }; + yield* checkForResume; }); const start: ProviderSessionReaperShape["start"] = () =>