diff --git a/apps/server/integration/OrchestrationEngineHarness.integration.ts b/apps/server/integration/OrchestrationEngineHarness.integration.ts index 6f9f4c6f44fd..e1c6c4848a0e 100644 --- a/apps/server/integration/OrchestrationEngineHarness.integration.ts +++ b/apps/server/integration/OrchestrationEngineHarness.integration.ts @@ -39,6 +39,7 @@ import { ProjectionPendingApprovalRepository } from "../src/persistence/Services import { ProviderUnsupportedError } from "../src/provider/Errors.ts"; import { ProviderAdapterRegistry } from "../src/provider/Services/ProviderAdapterRegistry.ts"; import { ProviderSessionDirectoryLive } from "../src/provider/Layers/ProviderSessionDirectory.ts"; +import { ProviderSessionDirectoryEventsLive } from "../src/provider/Layers/ProviderSessionDirectoryEvents.ts"; import { ServerSettingsService } from "../src/serverSettings.ts"; import { makeProviderServiceLive } from "../src/provider/Layers/ProviderService.ts"; import { makeCodexAdapterLive } from "../src/provider/Layers/CodexAdapter.ts"; @@ -258,8 +259,10 @@ export const makeOrchestrationIntegrationHarness = ( Layer.provide(OrchestrationEventStoreLive), Layer.provide(OrchestrationCommandReceiptRepositoryLive), ); + const providerSessionDirectoryEventsLayer = ProviderSessionDirectoryEventsLive; const providerSessionDirectoryLayer = ProviderSessionDirectoryLive.pipe( Layer.provide(ProviderSessionRuntimeRepositoryLive), + Layer.provide(providerSessionDirectoryEventsLayer), ); const realCodexRegistry = Layer.effect( ProviderAdapterRegistry, @@ -277,16 +280,19 @@ export const makeOrchestrationIntegrationHarness = ( Layer.provide(makeCodexAdapterLive()), Layer.provideMerge(ServerConfig.layerTest(workspaceDir, rootDir)), Layer.provideMerge(NodeServices.layer), + Layer.provideMerge(providerSessionDirectoryEventsLayer), Layer.provideMerge(providerSessionDirectoryLayer), ); const providerLayer = useRealCodex ? makeProviderServiceLive().pipe( Layer.provide(providerSessionDirectoryLayer), + Layer.provide(providerSessionDirectoryEventsLayer), Layer.provide(realCodexRegistry), Layer.provide(AnalyticsService.layerTest), ) : makeProviderServiceLive().pipe( Layer.provide(providerSessionDirectoryLayer), + Layer.provide(providerSessionDirectoryEventsLayer), Layer.provide(fakeRegistry!), Layer.provide(AnalyticsService.layerTest), ); diff --git a/apps/server/integration/providerService.integration.test.ts b/apps/server/integration/providerService.integration.test.ts index 89cf6ac153d9..b76944b0417b 100644 --- a/apps/server/integration/providerService.integration.test.ts +++ b/apps/server/integration/providerService.integration.test.ts @@ -8,6 +8,7 @@ import { Effect, FileSystem, Layer, Path, Queue, Stream } from "effect"; import { ProviderUnsupportedError } from "../src/provider/Errors.ts"; import { ProviderAdapterRegistry } from "../src/provider/Services/ProviderAdapterRegistry.ts"; import { ProviderSessionDirectoryLive } from "../src/provider/Layers/ProviderSessionDirectory.ts"; +import { ProviderSessionDirectoryEventsLive } from "../src/provider/Layers/ProviderSessionDirectoryEvents.ts"; import { makeProviderServiceLive } from "../src/provider/Layers/ProviderService.ts"; import { ProviderService, @@ -55,12 +56,15 @@ const makeIntegrationFixture = Effect.gen(function* () { listProviders: () => Effect.succeed(["codex"]), }; + const directoryEventsLayer = ProviderSessionDirectoryEventsLive; const directoryLayer = ProviderSessionDirectoryLive.pipe( Layer.provide(ProviderSessionRuntimeRepositoryLive), + Layer.provide(directoryEventsLayer), ); const shared = Layer.mergeAll( directoryLayer, + directoryEventsLayer, Layer.succeed(ProviderAdapterRegistry, registry), ServerSettingsService.layerTest(DEFAULT_SERVER_SETTINGS), AnalyticsService.layerTest, diff --git a/apps/server/src/cli-config.test.ts b/apps/server/src/cli-config.test.ts index 5adece730201..b4b7782f3c30 100644 --- a/apps/server/src/cli-config.test.ts +++ b/apps/server/src/cli-config.test.ts @@ -7,6 +7,10 @@ import { NetService } from "@t3tools/shared/Net"; import * as NodeServices from "@effect/platform-node/NodeServices"; import { deriveServerPaths } from "./config.ts"; import { resolveServerConfig } from "./cli.ts"; +import { + DEFAULT_PROVIDER_SESSION_REAPER_FALLBACK_RECONCILE_INTERVAL_MS, + DEFAULT_PROVIDER_SESSION_REAPER_INACTIVITY_THRESHOLD_MS, +} from "./provider/Services/ProviderSessionReaper.ts"; it.layer(NodeServices.layer)("cli config resolution", (it) => { const defaultObservabilityConfig = { @@ -19,6 +23,10 @@ it.layer(NodeServices.layer)("cli config resolution", (it) => { otlpMetricsUrl: undefined, otlpExportIntervalMs: 10_000, otlpServiceName: "t3-server", + providerSessionReaperInactivityThresholdMs: + DEFAULT_PROVIDER_SESSION_REAPER_INACTIVITY_THRESHOLD_MS, + providerSessionReaperFallbackReconcileIntervalMs: + DEFAULT_PROVIDER_SESSION_REAPER_FALLBACK_RECONCILE_INTERVAL_MS, } as const; const openBootstrapFd = Effect.fn(function* (payload: Record) { @@ -91,6 +99,63 @@ it.layer(NodeServices.layer)("cli config resolution", (it) => { }), ); + it.effect("reads provider session reaper tuning env vars", () => + Effect.gen(function* () { + const fs = yield* FileSystem.FileSystem; + const baseDir = yield* fs.makeTempDirectoryScoped({ prefix: "t3-cli-config-reaper-" }); + const derivedPaths = yield* deriveServerPaths(baseDir, undefined); + const resolved = yield* resolveServerConfig( + { + mode: Option.some("desktop"), + port: Option.some(4888), + host: Option.none(), + baseDir: Option.some(baseDir), + cwd: Option.none(), + devUrl: Option.none(), + noBrowser: Option.none(), + bootstrapFd: Option.none(), + autoBootstrapProjectFromCwd: Option.none(), + logWebSocketEvents: Option.none(), + }, + Option.none(), + ).pipe( + Effect.provide( + Layer.mergeAll( + ConfigProvider.layer( + ConfigProvider.fromEnv({ + env: { + T3CODE_PROVIDER_SESSION_REAPER_INACTIVITY_THRESHOLD_MS: "1500", + T3CODE_PROVIDER_SESSION_REAPER_FALLBACK_RECONCILE_INTERVAL_MS: "2500", + }, + }), + ), + NetService.layer, + ), + ), + ); + + expect(resolved).toEqual({ + logLevel: "Info", + ...defaultObservabilityConfig, + providerSessionReaperInactivityThresholdMs: 1500, + providerSessionReaperFallbackReconcileIntervalMs: 2500, + mode: "desktop", + port: 4888, + cwd: process.cwd(), + baseDir, + ...derivedPaths, + host: "127.0.0.1", + staticDir: resolved.staticDir, + devUrl: undefined, + noBrowser: true, + startupPresentation: "browser", + desktopBootstrapToken: undefined, + autoBootstrapProjectFromCwd: false, + logWebSocketEvents: false, + }); + }), + ); + it.effect("uses CLI flags when provided", () => Effect.gen(function* () { const { join } = yield* Path.Path; diff --git a/apps/server/src/cli.test.ts b/apps/server/src/cli.test.ts index 7ebde01067a7..2c4d4648786f 100644 --- a/apps/server/src/cli.test.ts +++ b/apps/server/src/cli.test.ts @@ -63,6 +63,8 @@ const makeCliTestServerConfig = (baseDir: string) => otlpMetricsUrl: undefined, otlpExportIntervalMs: 10_000, otlpServiceName: "t3-server", + providerSessionReaperInactivityThresholdMs: 30 * 60 * 1000, + providerSessionReaperFallbackReconcileIntervalMs: 30 * 60 * 1000, mode: "web", port: 0, host: "127.0.0.1", diff --git a/apps/server/src/cli.ts b/apps/server/src/cli.ts index 4fc23a1ded09..2edf1e2f6e8a 100644 --- a/apps/server/src/cli.ts +++ b/apps/server/src/cli.ts @@ -57,6 +57,10 @@ import { OrchestrationEngineService } from "./orchestration/Services/Orchestrati import { ProjectionSnapshotQuery } from "./orchestration/Services/ProjectionSnapshotQuery.ts"; import { OrchestrationLayerLive } from "./orchestration/runtimeLayer.ts"; import { layerConfig as SqlitePersistenceLayerLive } from "./persistence/Layers/Sqlite.ts"; +import { + DEFAULT_PROVIDER_SESSION_REAPER_FALLBACK_RECONCILE_INTERVAL_MS, + DEFAULT_PROVIDER_SESSION_REAPER_INACTIVITY_THRESHOLD_MS, +} from "./provider/Services/ProviderSessionReaper.ts"; import { RepositoryIdentityResolverLive } from "./project/Layers/RepositoryIdentityResolver.ts"; import { getAutoBootstrapDefaultModelSelection } from "./serverRuntimeStartup.ts"; import { @@ -150,6 +154,12 @@ const EnvServerConfig = Config.all({ Config.withDefault(10_000), ), otlpServiceName: Config.string("T3CODE_OTLP_SERVICE_NAME").pipe(Config.withDefault("t3-server")), + providerSessionReaperInactivityThresholdMs: Config.int( + "T3CODE_PROVIDER_SESSION_REAPER_INACTIVITY_THRESHOLD_MS", + ).pipe(Config.withDefault(DEFAULT_PROVIDER_SESSION_REAPER_INACTIVITY_THRESHOLD_MS)), + providerSessionReaperFallbackReconcileIntervalMs: Config.int( + "T3CODE_PROVIDER_SESSION_REAPER_FALLBACK_RECONCILE_INTERVAL_MS", + ).pipe(Config.withDefault(DEFAULT_PROVIDER_SESSION_REAPER_FALLBACK_RECONCILE_INTERVAL_MS)), mode: Config.schema(RuntimeMode, "T3CODE_MODE").pipe( Config.option, Config.map(Option.getOrUndefined), @@ -351,6 +361,9 @@ export const resolveServerConfig = ( persistedObservabilitySettings.otlpMetricsUrl, otlpExportIntervalMs: env.otlpExportIntervalMs, otlpServiceName: env.otlpServiceName, + providerSessionReaperInactivityThresholdMs: env.providerSessionReaperInactivityThresholdMs, + providerSessionReaperFallbackReconcileIntervalMs: + env.providerSessionReaperFallbackReconcileIntervalMs, mode, port, cwd, diff --git a/apps/server/src/config.ts b/apps/server/src/config.ts index 7840c7611512..2be01ce9ae8d 100644 --- a/apps/server/src/config.ts +++ b/apps/server/src/config.ts @@ -8,6 +8,11 @@ */ import { Effect, FileSystem, Layer, LogLevel, Path, Schema, Context } from "effect"; +import { + DEFAULT_PROVIDER_SESSION_REAPER_FALLBACK_RECONCILE_INTERVAL_MS, + DEFAULT_PROVIDER_SESSION_REAPER_INACTIVITY_THRESHOLD_MS, +} from "./provider/Services/ProviderSessionReaper.ts"; + export const DEFAULT_PORT = 3773; export const RuntimeMode = Schema.Literals(["web", "desktop"]); @@ -53,6 +58,8 @@ export interface ServerConfigShape extends ServerDerivedPaths { readonly otlpMetricsUrl: string | undefined; readonly otlpExportIntervalMs: number; readonly otlpServiceName: string; + readonly providerSessionReaperInactivityThresholdMs: number; + readonly providerSessionReaperFallbackReconcileIntervalMs: number; readonly mode: RuntimeMode; readonly port: number; readonly host: string | undefined; @@ -152,6 +159,10 @@ export class ServerConfig extends Context.Service EventId.make(value); const asThreadId = (value: string): ThreadId => ThreadId.make(value); const asTurnId = (value: string): TurnId => TurnId.make(value); +function makeDirectoryLayer( + runtimeRepositoryLayer: Layer.Layer, +) { + const directoryEventsLayer = ProviderSessionDirectoryEventsLive; + return Layer.mergeAll( + directoryEventsLayer, + ProviderSessionDirectoryLive.pipe( + Layer.provide(runtimeRepositoryLayer), + Layer.provide(directoryEventsLayer), + ), + ); +} + type LegacyProviderRuntimeEvent = { readonly type: string; readonly eventId: EventId; @@ -263,7 +277,7 @@ function makeProviderServiceLayer() { const runtimeRepositoryLayer = ProviderSessionRuntimeRepositoryLive.pipe( Layer.provide(SqlitePersistenceMemory), ); - const directoryLayer = ProviderSessionDirectoryLive.pipe(Layer.provide(runtimeRepositoryLayer)); + const directoryLayer = makeDirectoryLayer(runtimeRepositoryLayer); const layer = it.layer( Layer.mergeAll( @@ -312,7 +326,7 @@ it.effect("ProviderServiceLive rejects new sessions for disabled providers", () const runtimeRepositoryLayer = ProviderSessionRuntimeRepositoryLive.pipe( Layer.provide(SqlitePersistenceMemory), ); - const directoryLayer = ProviderSessionDirectoryLive.pipe(Layer.provide(runtimeRepositoryLayer)); + const directoryLayer = makeDirectoryLayer(runtimeRepositoryLayer); const providerLayer = makeProviderServiceLive().pipe( Layer.provide(providerAdapterLayer), Layer.provide(directoryLayer), @@ -354,7 +368,7 @@ it.effect("ProviderServiceLive writes canonical events to the emitting thread se const runtimeRepositoryLayer = ProviderSessionRuntimeRepositoryLive.pipe( Layer.provide(SqlitePersistenceMemory), ); - const directoryLayer = ProviderSessionDirectoryLive.pipe(Layer.provide(runtimeRepositoryLayer)); + const directoryLayer = makeDirectoryLayer(runtimeRepositoryLayer); const providerLayer = makeProviderServiceLive({ canonicalEventLogger: { filePath: "memory://provider-canonical-events", @@ -412,7 +426,7 @@ it.effect("ProviderServiceLive keeps persisted resumable sessions on startup", ( const runtimeRepositoryLayer = ProviderSessionRuntimeRepositoryLive.pipe( Layer.provide(persistenceLayer), ); - const directoryLayer = ProviderSessionDirectoryLive.pipe(Layer.provide(runtimeRepositoryLayer)); + const directoryLayer = makeDirectoryLayer(runtimeRepositoryLayer); yield* Effect.gen(function* () { const directory = yield* ProviderSessionDirectory; @@ -479,9 +493,7 @@ it.effect( listProviders: () => Effect.succeed(["codex"]), }; - const firstDirectoryLayer = ProviderSessionDirectoryLive.pipe( - Layer.provide(runtimeRepositoryLayer), - ); + const firstDirectoryLayer = makeDirectoryLayer(runtimeRepositoryLayer); const firstProviderLayer = makeProviderServiceLive().pipe( Layer.provide(Layer.succeed(ProviderAdapterRegistry, firstRegistry)), Layer.provide(firstDirectoryLayer), @@ -531,9 +543,7 @@ it.effect( : Effect.fail(new ProviderUnsupportedError({ provider })), listProviders: () => Effect.succeed(["codex"]), }; - const secondDirectoryLayer = ProviderSessionDirectoryLive.pipe( - Layer.provide(runtimeRepositoryLayer), - ); + const secondDirectoryLayer = makeDirectoryLayer(runtimeRepositoryLayer); const secondProviderLayer = makeProviderServiceLive().pipe( Layer.provide(Layer.succeed(ProviderAdapterRegistry, secondRegistry)), Layer.provide(secondDirectoryLayer), @@ -985,9 +995,7 @@ routing.layer("ProviderServiceLive routing", (it) => { : Effect.fail(new ProviderUnsupportedError({ provider })), listProviders: () => Effect.succeed(["claudeAgent"]), }; - const firstDirectoryLayer = ProviderSessionDirectoryLive.pipe( - Layer.provide(runtimeRepositoryLayer), - ); + const firstDirectoryLayer = makeDirectoryLayer(runtimeRepositoryLayer); const firstProviderLayer = makeProviderServiceLive().pipe( Layer.provide(Layer.succeed(ProviderAdapterRegistry, firstRegistry)), Layer.provide(firstDirectoryLayer), @@ -1018,9 +1026,7 @@ routing.layer("ProviderServiceLive routing", (it) => { : Effect.fail(new ProviderUnsupportedError({ provider })), listProviders: () => Effect.succeed(["claudeAgent"]), }; - const secondDirectoryLayer = ProviderSessionDirectoryLive.pipe( - Layer.provide(runtimeRepositoryLayer), - ); + const secondDirectoryLayer = makeDirectoryLayer(runtimeRepositoryLayer); const secondProviderLayer = makeProviderServiceLive().pipe( Layer.provide(Layer.succeed(ProviderAdapterRegistry, secondRegistry)), Layer.provide(secondDirectoryLayer), diff --git a/apps/server/src/provider/Layers/ProviderSessionDirectory.test.ts b/apps/server/src/provider/Layers/ProviderSessionDirectory.test.ts index 35bdec1e37d2..028756cd1282 100644 --- a/apps/server/src/provider/Layers/ProviderSessionDirectory.test.ts +++ b/apps/server/src/provider/Layers/ProviderSessionDirectory.test.ts @@ -6,7 +6,7 @@ import * as NodeServices from "@effect/platform-node/NodeServices"; import { ThreadId } from "@t3tools/contracts"; import { it, assert } from "@effect/vitest"; import { assertSome } from "@effect/vitest/utils"; -import { Effect, Layer, Option } from "effect"; +import { Cause, Effect, Exit, Fiber, Layer, Option, Stream } from "effect"; import * as SqlClient from "effect/unstable/sql/SqlClient"; import { @@ -16,15 +16,22 @@ import { import { ProviderSessionRuntimeRepositoryLive } from "../../persistence/Layers/ProviderSessionRuntime.ts"; import { ProviderSessionRuntimeRepository } from "../../persistence/Services/ProviderSessionRuntime.ts"; import { ProviderSessionDirectory } from "../Services/ProviderSessionDirectory.ts"; +import { ProviderSessionDirectoryEvents } from "../Services/ProviderSessionDirectoryEvents.ts"; +import { ProviderSessionDirectoryEventsLive } from "./ProviderSessionDirectoryEvents.ts"; import { ProviderSessionDirectoryLive } from "./ProviderSessionDirectory.ts"; function makeDirectoryLayer(persistenceLayer: Layer.Layer) { const runtimeRepositoryLayer = ProviderSessionRuntimeRepositoryLive.pipe( Layer.provide(persistenceLayer), ); + const directoryEventsLayer = ProviderSessionDirectoryEventsLive; return Layer.mergeAll( runtimeRepositoryLayer, - ProviderSessionDirectoryLive.pipe(Layer.provide(runtimeRepositoryLayer)), + directoryEventsLayer, + ProviderSessionDirectoryLive.pipe( + Layer.provide(runtimeRepositoryLayer), + Layer.provide(directoryEventsLayer), + ), NodeServices.layer, ); } @@ -77,6 +84,26 @@ it.layer(makeDirectoryLayer(SqlitePersistenceMemory))("ProviderSessionDirectoryL assert.deepEqual(threadIds, [nextThreadId]); })); + it("publishes change notifications after successful upserts", () => + Effect.gen(function* () { + const directory = yield* ProviderSessionDirectory; + const directoryEvents = yield* ProviderSessionDirectoryEvents; + const threadId = ThreadId.make("thread-change-notification"); + + const changeFiber = yield* Stream.take(directoryEvents.changes, 1).pipe( + Stream.runCollect, + Effect.forkScoped, + ); + yield* Effect.yieldNow; + yield* directory.upsert({ + provider: "codex", + threadId, + }); + + const changes = yield* Fiber.join(changeFiber); + assert.deepEqual(Array.from(changes), [{ threadId }]); + })); + it("persists runtime fields and merges payload updates", () => Effect.gen(function* () { const directory = yield* ProviderSessionDirectory; @@ -265,4 +292,80 @@ it.layer(makeDirectoryLayer(SqlitePersistenceMemory))("ProviderSessionDirectoryL fs.rmSync(tempDir, { recursive: true, force: true }); })); + + it("persists bindings even when change notification publishing fails", () => + Effect.gen(function* () { + const runtimeRepositoryLayer = ProviderSessionRuntimeRepositoryLive.pipe( + Layer.provide(SqlitePersistenceMemory), + ); + const failingDirectoryEventsLayer = Layer.succeed(ProviderSessionDirectoryEvents, { + publishChanged: () => Effect.die(new Error("simulated publish failure")), + changes: Stream.empty, + }); + const directoryLayer = Layer.mergeAll( + runtimeRepositoryLayer, + failingDirectoryEventsLayer, + ProviderSessionDirectoryLive.pipe( + Layer.provide(runtimeRepositoryLayer), + Layer.provide(failingDirectoryEventsLayer), + ), + NodeServices.layer, + ); + + const threadId = ThreadId.make("thread-failing-directory-events"); + const runtime = yield* Effect.gen(function* () { + const directory = yield* ProviderSessionDirectory; + const repository = yield* ProviderSessionRuntimeRepository; + + yield* directory.upsert({ + provider: "codex", + threadId, + status: "running", + }); + + return yield* repository.getByThreadId({ threadId }); + }).pipe(Effect.provide(directoryLayer)); + + assert.equal(Option.isSome(runtime), true); + if (Option.isSome(runtime)) { + assert.equal(runtime.value.threadId, threadId); + assert.equal(runtime.value.providerName, "codex"); + } + })); + + it("preserves interruption when change notification publishing is interrupted", () => + Effect.gen(function* () { + const runtimeRepositoryLayer = ProviderSessionRuntimeRepositoryLive.pipe( + Layer.provide(SqlitePersistenceMemory), + ); + const interruptedDirectoryEventsLayer = Layer.succeed(ProviderSessionDirectoryEvents, { + publishChanged: () => Effect.interrupt, + changes: Stream.empty, + }); + const directoryLayer = Layer.mergeAll( + runtimeRepositoryLayer, + interruptedDirectoryEventsLayer, + ProviderSessionDirectoryLive.pipe( + Layer.provide(runtimeRepositoryLayer), + Layer.provide(interruptedDirectoryEventsLayer), + ), + NodeServices.layer, + ); + + const exit = yield* Effect.exit( + Effect.gen(function* () { + const directory = yield* ProviderSessionDirectory; + yield* directory.upsert({ + provider: "codex", + threadId: ThreadId.make("thread-interrupted-directory-events"), + status: "running", + }); + }).pipe(Effect.provide(directoryLayer)), + ); + + assert.equal(Exit.isFailure(exit), true); + if (Exit.isFailure(exit)) { + assert.equal(Cause.hasInterruptsOnly(exit.cause), true); + } + })); }); diff --git a/apps/server/src/provider/Layers/ProviderSessionDirectory.ts b/apps/server/src/provider/Layers/ProviderSessionDirectory.ts index 4cb7147180ad..1ea06d29d6a8 100644 --- a/apps/server/src/provider/Layers/ProviderSessionDirectory.ts +++ b/apps/server/src/provider/Layers/ProviderSessionDirectory.ts @@ -1,5 +1,5 @@ import { ProviderKind, type ThreadId } from "@t3tools/contracts"; -import { Effect, Layer, Option, Schema } from "effect"; +import { Cause, Effect, Layer, Option, Schema } from "effect"; import type { ProviderSessionRuntime } from "../../persistence/Services/ProviderSessionRuntime.ts"; import { ProviderSessionRuntimeRepository } from "../../persistence/Services/ProviderSessionRuntime.ts"; @@ -10,6 +10,7 @@ import { type ProviderRuntimeBindingWithMetadata, type ProviderSessionDirectoryShape, } from "../Services/ProviderSessionDirectory.ts"; +import { ProviderSessionDirectoryEvents } from "../Services/ProviderSessionDirectoryEvents.ts"; function toPersistenceError(operation: string) { return (cause: unknown) => @@ -76,6 +77,7 @@ function toRuntimeBinding( const makeProviderSessionDirectory = Effect.gen(function* () { const repository = yield* ProviderSessionRuntimeRepository; + const directoryEvents = yield* ProviderSessionDirectoryEvents; const getBinding = (threadId: ThreadId) => repository.getByThreadId({ threadId }).pipe( @@ -128,6 +130,17 @@ const makeProviderSessionDirectory = Effect.gen(function* () { ), }) .pipe(Effect.mapError(toPersistenceError("ProviderSessionDirectory.upsert:upsert"))); + yield* directoryEvents.publishChanged(resolvedThreadId).pipe( + Effect.catchCause((cause) => { + if (Cause.hasInterruptsOnly(cause)) { + return Effect.failCause(cause); + } + return Effect.logDebug("provider.session.directory.change-signal-failed", { + threadId: resolvedThreadId, + cause, + }); + }), + ); }); const getProvider: ProviderSessionDirectoryShape["getProvider"] = (threadId) => diff --git a/apps/server/src/provider/Layers/ProviderSessionDirectoryEvents.test.ts b/apps/server/src/provider/Layers/ProviderSessionDirectoryEvents.test.ts new file mode 100644 index 000000000000..c32c80b18e2b --- /dev/null +++ b/apps/server/src/provider/Layers/ProviderSessionDirectoryEvents.test.ts @@ -0,0 +1,30 @@ +import { assert, it } from "@effect/vitest"; +import { ThreadId } from "@t3tools/contracts"; +import { Effect, Fiber, Stream } from "effect"; + +import { ProviderSessionDirectoryEvents } from "../Services/ProviderSessionDirectoryEvents.ts"; +import { ProviderSessionDirectoryEventsLive } from "./ProviderSessionDirectoryEvents.ts"; + +it.effect("ProviderSessionDirectoryEventsLive fans out changes to each subscriber", () => + Effect.gen(function* () { + const directoryEvents = yield* ProviderSessionDirectoryEvents; + const threadId = ThreadId.make("thread-directory-events-fanout"); + + const firstFiber = yield* Stream.take(directoryEvents.changes, 1).pipe( + Stream.runCollect, + Effect.forkScoped, + ); + const secondFiber = yield* Stream.take(directoryEvents.changes, 1).pipe( + Stream.runCollect, + Effect.forkScoped, + ); + yield* Effect.yieldNow; + + yield* directoryEvents.publishChanged(threadId); + + const first = yield* Fiber.join(firstFiber); + const second = yield* Fiber.join(secondFiber); + assert.deepEqual(Array.from(first), [{ threadId }]); + assert.deepEqual(Array.from(second), [{ threadId }]); + }).pipe(Effect.provide(ProviderSessionDirectoryEventsLive)), +); diff --git a/apps/server/src/provider/Layers/ProviderSessionDirectoryEvents.ts b/apps/server/src/provider/Layers/ProviderSessionDirectoryEvents.ts new file mode 100644 index 000000000000..1e41d6074688 --- /dev/null +++ b/apps/server/src/provider/Layers/ProviderSessionDirectoryEvents.ts @@ -0,0 +1,26 @@ +import { type ThreadId } from "@t3tools/contracts"; +import { Effect, Layer, PubSub, Stream } from "effect"; + +import { + ProviderSessionDirectoryEvents, + type ProviderSessionDirectoryEventsShape, +} from "../Services/ProviderSessionDirectoryEvents.ts"; + +const makeProviderSessionDirectoryEvents = Effect.gen(function* () { + const pubSub = yield* Effect.acquireRelease( + PubSub.unbounded<{ readonly threadId: ThreadId }>(), + PubSub.shutdown, + ); + + return { + publishChanged: (threadId) => PubSub.publish(pubSub, { threadId }).pipe(Effect.asVoid), + get changes(): ProviderSessionDirectoryEventsShape["changes"] { + return Stream.fromPubSub(pubSub); + }, + } satisfies ProviderSessionDirectoryEventsShape; +}); + +export const ProviderSessionDirectoryEventsLive = Layer.effect( + ProviderSessionDirectoryEvents, + makeProviderSessionDirectoryEvents, +); diff --git a/apps/server/src/provider/Layers/ProviderSessionReaper.test.ts b/apps/server/src/provider/Layers/ProviderSessionReaper.test.ts index 45199a02b2af..1a663200cc0f 100644 --- a/apps/server/src/provider/Layers/ProviderSessionReaper.test.ts +++ b/apps/server/src/provider/Layers/ProviderSessionReaper.test.ts @@ -1,8 +1,30 @@ import * as NodeServices from "@effect/platform-node/NodeServices"; -import { ProjectId, ThreadId, TurnId } from "@t3tools/contracts"; -import { Effect, Exit, Layer, ManagedRuntime, Option, Scope, Stream } from "effect"; +import { + MessageId, + ProjectId, + ThreadId, + TurnId, + type OrchestrationEvent, +} from "@t3tools/contracts"; +import { + Effect, + Exit, + Layer, + Logger, + ManagedRuntime, + Metric, + Option, + PubSub, + References, + Scope, + Stream, + Tracer, +} from "effect"; import { afterEach, describe, expect, it, vi } from "vitest"; +import { makeLocalFileTracer } from "../../observability/LocalFileTracer.ts"; +import type { EffectTraceRecord } from "../../observability/TraceRecord.ts"; +import type { TraceSink } from "../../observability/TraceSink.ts"; import { OrchestrationEngineService, type OrchestrationEngineShape, @@ -13,6 +35,7 @@ import { ProviderSessionRuntimeRepository } from "../../persistence/Services/Pro import { ProviderValidationError } from "../Errors.ts"; import { ProviderSessionReaper } from "../Services/ProviderSessionReaper.ts"; import { ProviderService, type ProviderServiceShape } from "../Services/ProviderService.ts"; +import { ProviderSessionDirectoryEventsLive } from "./ProviderSessionDirectoryEvents.ts"; import { ProviderSessionDirectoryLive } from "./ProviderSessionDirectory.ts"; import { makeProviderSessionReaperLive } from "./ProviderSessionReaper.ts"; @@ -21,6 +44,46 @@ const defaultModelSelection = { model: "gpt-5-codex", } as const; +const hasMetricSnapshot = ( + snapshots: ReadonlyArray, + id: string, + attributes: Readonly>, +) => + snapshots.some( + (snapshot) => + snapshot.id === id && + Object.entries(attributes).every(([key, value]) => snapshot.attributes?.[key] === value), + ); + +const makeTraceTestLayer = (records: Array) => { + const sink = { + filePath: "provider-session-reaper-test", + push(record) { + if (record.type === "effect-span") { + records.push(record); + } + }, + flush: Effect.void, + close: () => Effect.void, + } satisfies TraceSink; + + return Layer.mergeAll( + Layer.effect( + Tracer.Tracer, + makeLocalFileTracer({ + filePath: sink.filePath, + maxBytes: 1024 * 1024, + maxFiles: 1, + batchWindowMs: 10_000, + sink, + }), + ), + Logger.layer([Logger.tracerLogger], { mergeWithExisting: false }), + Layer.succeed(References.MinimumLogLevel, "Info"), + Layer.succeed(References.TracerTimingEnabled, true), + ); +}; + async function waitFor( predicate: () => boolean | Promise, timeoutMs = 2_000, @@ -45,10 +108,18 @@ const unsupported = () => Effect.die(new Error("Unsupported provider call in tes function makeReadModel( threads: ReadonlyArray<{ readonly id: ThreadId; + readonly latestTurn?: { + readonly turnId: TurnId; + readonly state: "running" | "interrupted" | "completed" | "error"; + readonly requestedAt: string; + readonly startedAt: string | null; + readonly completedAt: string | null; + readonly assistantMessageId: MessageId | null; + } | null; readonly session: { readonly threadId: ThreadId; readonly status: "starting" | "running" | "ready" | "interrupted" | "stopped" | "error"; - readonly providerName: "codex" | "claudeAgent"; + readonly providerName: "codex" | "claudeAgent" | "cursor" | "opencode"; readonly runtimeMode: "approval-required" | "full-access" | "auto-accept-edits"; readonly activeTurnId: TurnId | null; readonly lastError: string | null; @@ -86,7 +157,7 @@ function makeReadModel( createdAt: now, updatedAt: now, archivedAt: null, - latestTurn: null, + latestTurn: thread.latestTurn ?? null, messages: [], session: thread.session, activities: [], @@ -117,6 +188,8 @@ describe("ProviderSessionReaper", () => { async function createHarness(input: { readonly readModel: ReturnType; + readonly streamDomainEvents?: Stream.Stream; + readonly traceRecords?: Array; readonly stopSessionImplementation?: (input: { readonly threadId: ThreadId; }) => ReturnType; @@ -144,34 +217,114 @@ describe("ProviderSessionReaper", () => { streamEvents: Stream.empty, }; + const domainEventPubSub = Effect.runSync(PubSub.unbounded()); + const currentReadModel = input.readModel; const orchestrationEngine: OrchestrationEngineShape = { - getReadModel: () => Effect.succeed(input.readModel), + getReadModel: () => Effect.succeed(currentReadModel), readEvents: () => Stream.empty, dispatch: () => unsupported(), - streamDomainEvents: Stream.empty, + get streamDomainEvents() { + return input.streamDomainEvents ?? Stream.fromPubSub(domainEventPubSub); + }, }; const runtimeRepositoryLayer = ProviderSessionRuntimeRepositoryLive.pipe( Layer.provide(SqlitePersistenceMemory), ); + const directoryEventsLayer = ProviderSessionDirectoryEventsLive; const providerSessionDirectoryLayer = ProviderSessionDirectoryLive.pipe( Layer.provide(runtimeRepositoryLayer), + Layer.provide(directoryEventsLayer), ); - const layer = makeProviderSessionReaperLive({ + const baseLayer = makeProviderSessionReaperLive({ inactivityThresholdMs: 1_000, - sweepIntervalMs: 60_000, + fallbackReconcileIntervalMs: 60_000, }).pipe( Layer.provideMerge(providerSessionDirectoryLayer), + Layer.provideMerge(directoryEventsLayer), Layer.provideMerge(runtimeRepositoryLayer), Layer.provideMerge(Layer.succeed(ProviderService, providerService)), Layer.provideMerge(Layer.succeed(OrchestrationEngineService, orchestrationEngine)), Layer.provideMerge(NodeServices.layer), ); + const layer = + input.traceRecords === undefined + ? baseLayer + : baseLayer.pipe(Layer.provideMerge(makeTraceTestLayer(input.traceRecords))); runtime = ManagedRuntime.make(layer); - return { stopSession, stoppedThreadIds }; + return { + stopSession, + stoppedThreadIds, + }; } + it("emits root start and iteration spans without signal feed spans", async () => { + const traceRecords: Array = []; + await createHarness({ + readModel: makeReadModel([]), + traceRecords, + }); + + const reaper = await runtime!.runPromise(Effect.service(ProviderSessionReaper)); + scope = await Effect.runPromise(Scope.make("sequential")); + await runtime!.runPromise( + reaper.start().pipe(Scope.provide(scope), Effect.withSpan("test.parent")), + ); + + await waitFor(() => + traceRecords.some((record) => record.name === "provider.session.reaper.iteration"), + ); + + await Effect.runPromise(Scope.close(scope, Exit.void)); + scope = null; + + const start = traceRecords.find((record) => record.name === "provider.session.reaper.start"); + const iteration = traceRecords.find( + (record) => record.name === "provider.session.reaper.iteration", + ); + const reconcile = traceRecords.find( + (record) => record.name === "provider.session.reaper.reconcile", + ); + + expect(start).toBeDefined(); + expect(start?.parentSpanId).toBeUndefined(); + expect(iteration).toBeDefined(); + expect(iteration?.parentSpanId).toBeUndefined(); + expect(reconcile?.parentSpanId).toBe(iteration?.spanId); + expect( + traceRecords.some((record) => record.name === "provider.session.reaper.signal_feed"), + ).toBe(false); + }); + + it("records signal feed restart metrics without signal feed spans", async () => { + const traceRecords: Array = []; + await createHarness({ + readModel: makeReadModel([]), + streamDomainEvents: Stream.die(new Error("simulated orchestration stream failure")), + traceRecords, + }); + + const reaper = await runtime!.runPromise(Effect.service(ProviderSessionReaper)); + scope = await Effect.runPromise(Scope.make("sequential")); + await runtime!.runPromise(reaper.start().pipe(Scope.provide(scope))); + + await waitFor(async () => { + const snapshots = await runtime!.runPromise(Metric.snapshot); + return hasMetricSnapshot(snapshots, "t3_provider_session_reaper_signal_feed_restarts_total", { + mode: "deadline", + feed: "orchestration-domain-events", + }); + }); + + await Effect.runPromise(Scope.close(scope, Exit.void)); + scope = null; + + expect( + traceRecords.some((record) => record.name === "provider.session.reaper.signal_feed"), + ).toBe(false); + }); + it("reaps stale persisted sessions without active turns", async () => { const threadId = ThreadId.make("thread-reaper-stale"); const now = new Date().toISOString(); @@ -218,6 +371,116 @@ describe("ProviderSessionReaper", () => { expect(harness.stoppedThreadIds.has(threadId)).toBe(true); }); + it("does not reap when the latest turn is still within the inactivity threshold", async () => { + const threadId = ThreadId.make("thread-reaper-latest-turn-fresh"); + const turnId = TurnId.make("turn-reaper-latest-fresh"); + const now = new Date().toISOString(); + const harness = await createHarness({ + readModel: makeReadModel([ + { + id: threadId, + latestTurn: { + turnId, + state: "completed", + requestedAt: now, + startedAt: now, + completedAt: now, + assistantMessageId: null, + }, + session: { + threadId, + status: "ready", + providerName: "claudeAgent", + runtimeMode: "full-access", + activeTurnId: null, + lastError: null, + updatedAt: now, + }, + }, + ]), + }); + const repository = await runtime!.runPromise(Effect.service(ProviderSessionRuntimeRepository)); + + await runtime!.runPromise( + repository.upsert({ + threadId, + providerName: "claudeAgent", + adapterKey: "claudeAgent", + runtimeMode: "full-access", + status: "running", + lastSeenAt: "2026-04-14T00:00:00.000Z", + resumeCursor: { + opaque: "resume-latest-turn-fresh", + }, + runtimePayload: null, + }), + ); + + const reaper = await runtime!.runPromise(Effect.service(ProviderSessionReaper)); + scope = await Effect.runPromise(Scope.make("sequential")); + await Effect.runPromise(reaper.start().pipe(Scope.provide(scope))); + await new Promise((resolve) => setTimeout(resolve, 50)); + + expect(harness.stopSession).not.toHaveBeenCalled(); + const remaining = await runtime!.runPromise(repository.getByThreadId({ threadId })); + expect(Option.isSome(remaining)).toBe(true); + }); + + it("reaps based on the latest turn even when lastSeenAt is fresher", async () => { + const threadId = ThreadId.make("thread-reaper-latest-turn-stale"); + const turnId = TurnId.make("turn-reaper-latest-stale"); + const now = new Date().toISOString(); + const harness = await createHarness({ + readModel: makeReadModel([ + { + id: threadId, + latestTurn: { + turnId, + state: "completed", + requestedAt: "2026-04-14T00:00:00.000Z", + startedAt: "2026-04-14T00:00:00.000Z", + completedAt: "2026-04-14T00:00:00.000Z", + assistantMessageId: null, + }, + session: { + threadId, + status: "ready", + providerName: "claudeAgent", + runtimeMode: "full-access", + activeTurnId: null, + lastError: null, + updatedAt: now, + }, + }, + ]), + }); + const repository = await runtime!.runPromise(Effect.service(ProviderSessionRuntimeRepository)); + + await runtime!.runPromise( + repository.upsert({ + threadId, + providerName: "claudeAgent", + adapterKey: "claudeAgent", + runtimeMode: "full-access", + status: "running", + lastSeenAt: now, + resumeCursor: { + opaque: "resume-latest-turn-stale", + }, + runtimePayload: null, + }), + ); + + const reaper = await runtime!.runPromise(Effect.service(ProviderSessionReaper)); + scope = await Effect.runPromise(Scope.make("sequential")); + await Effect.runPromise(reaper.start().pipe(Scope.provide(scope))); + + await waitFor(() => harness.stopSession.mock.calls.length === 1); + + expect(harness.stopSession.mock.calls[0]?.[0]).toEqual({ threadId }); + expect(harness.stoppedThreadIds.has(threadId)).toBe(true); + }); + it("skips stale sessions when the thread still has an active turn", async () => { const threadId = ThreadId.make("thread-reaper-active-turn"); const turnId = TurnId.make("turn-reaper-active"); diff --git a/apps/server/src/provider/Layers/ProviderSessionReaper.timing.test.ts b/apps/server/src/provider/Layers/ProviderSessionReaper.timing.test.ts new file mode 100644 index 000000000000..2f0995d476bc --- /dev/null +++ b/apps/server/src/provider/Layers/ProviderSessionReaper.timing.test.ts @@ -0,0 +1,835 @@ +import { assert, it, vi } from "@effect/vitest"; +import { + CommandId, + EventId, + MessageId, + ProjectId, + ThreadId, + TurnId, + type OrchestrationEvent, + type OrchestrationSession, +} from "@t3tools/contracts"; +import { Duration, Effect, Layer, PubSub, Ref, Stream } from "effect"; +import { TestClock } from "effect/testing"; + +import { + OrchestrationEngineService, + type OrchestrationEngineShape, +} from "../../orchestration/Services/OrchestrationEngine.ts"; +import { SqlitePersistenceMemory } from "../../persistence/Layers/Sqlite.ts"; +import { ProviderSessionRuntimeRepositoryLive } from "../../persistence/Layers/ProviderSessionRuntime.ts"; +import { ProviderSessionRuntimeRepository } from "../../persistence/Services/ProviderSessionRuntime.ts"; +import { ProviderValidationError } from "../Errors.ts"; +import { ProviderSessionReaper } from "../Services/ProviderSessionReaper.ts"; +import { ProviderService, type ProviderServiceShape } from "../Services/ProviderService.ts"; +import { ProviderSessionDirectory } from "../Services/ProviderSessionDirectory.ts"; +import { ProviderSessionDirectoryEventsLive } from "./ProviderSessionDirectoryEvents.ts"; +import { ProviderSessionDirectoryLive } from "./ProviderSessionDirectory.ts"; +import { makeProviderSessionReaperLive } from "./ProviderSessionReaper.ts"; + +const defaultModelSelection = { + provider: "codex", + model: "gpt-5-codex", +} as const; + +const unsupported = () => Effect.die(new Error("Unsupported provider call in test")) as never; + +function makeReadModel( + threads: ReadonlyArray<{ + readonly id: ThreadId; + readonly latestTurn?: { + readonly turnId: TurnId; + readonly state: "running" | "interrupted" | "completed" | "error"; + readonly requestedAt: string; + readonly startedAt: string | null; + readonly completedAt: string | null; + readonly assistantMessageId: MessageId | null; + } | null; + readonly session: { + readonly threadId: ThreadId; + readonly status: "starting" | "running" | "ready" | "interrupted" | "stopped" | "error"; + readonly providerName: "codex" | "claudeAgent" | "cursor" | "opencode"; + readonly runtimeMode: "approval-required" | "full-access" | "auto-accept-edits"; + readonly activeTurnId: TurnId | null; + readonly lastError: string | null; + readonly updatedAt: string; + } | null; + }>, +) { + const now = new Date(0).toISOString(); + const projectId = ProjectId.make("project-provider-session-reaper-timing"); + + return { + snapshotSequence: 0, + updatedAt: now, + projects: [ + { + id: projectId, + title: "Provider Reaper Timing Project", + workspaceRoot: "/tmp/provider-reaper-timing", + defaultModelSelection, + scripts: [], + createdAt: now, + updatedAt: now, + deletedAt: null, + }, + ], + threads: threads.map((thread) => ({ + id: thread.id, + projectId, + title: `Thread ${thread.id}`, + modelSelection: defaultModelSelection, + interactionMode: "default" as const, + runtimeMode: "full-access" as const, + branch: null, + worktreePath: null, + createdAt: now, + updatedAt: now, + archivedAt: null, + latestTurn: thread.latestTurn ?? null, + messages: [], + session: thread.session, + activities: [], + proposedPlans: [], + checkpoints: [], + deletedAt: null, + })), + }; +} + +type ReadModelSession = NonNullable[0][number]["session"]>; + +function makeThreadSessionSetEvent( + threadId: ThreadId, + session: OrchestrationSession, +): OrchestrationEvent { + return { + sequence: 0, + eventId: EventId.make(`evt-${String(threadId)}-session-set`), + aggregateKind: "thread", + aggregateId: threadId, + type: "thread.session-set", + occurredAt: session.updatedAt, + commandId: CommandId.make(`cmd-${String(threadId)}-session-set`), + causationEventId: null, + correlationId: null, + metadata: {}, + payload: { + threadId, + session, + }, + }; +} + +function makeHarness(input: { + readonly initialReadModel: ReturnType; + readonly inactivityThresholdMs: number; + readonly fallbackReconcileIntervalMs: number; + readonly stopFailureRetryIntervalMs?: number; + readonly streamDomainEvents?: ( + defaultStream: Stream.Stream, + ) => Stream.Stream; + readonly stopSessionImplementation?: (request: { + readonly threadId: ThreadId; + }) => ReturnType; +}) { + return Effect.gen(function* () { + const readModelRef = yield* Ref.make(input.initialReadModel); + const domainEventPubSub = yield* PubSub.unbounded(); + const stopSession = vi.fn((request) => + input.stopSessionImplementation ? input.stopSessionImplementation(request) : Effect.void, + ); + + const providerService: ProviderServiceShape = { + startSession: () => unsupported(), + sendTurn: () => unsupported(), + interruptTurn: () => unsupported(), + respondToRequest: () => unsupported(), + respondToUserInput: () => unsupported(), + stopSession, + listSessions: () => Effect.succeed([]), + getCapabilities: () => Effect.succeed({ sessionModelSwitch: "in-session" }), + rollbackConversation: () => unsupported(), + streamEvents: Stream.empty, + }; + + const orchestrationEngine: OrchestrationEngineShape = { + getReadModel: () => Ref.get(readModelRef), + readEvents: () => Stream.empty, + dispatch: () => unsupported(), + get streamDomainEvents() { + const defaultStream = Stream.fromPubSub(domainEventPubSub); + return input.streamDomainEvents?.(defaultStream) ?? defaultStream; + }, + }; + + const runtimeRepositoryLayer = ProviderSessionRuntimeRepositoryLive.pipe( + Layer.provide(SqlitePersistenceMemory), + ); + const directoryEventsLayer = ProviderSessionDirectoryEventsLive; + const directoryLayer = ProviderSessionDirectoryLive.pipe( + Layer.provide(runtimeRepositoryLayer), + Layer.provide(directoryEventsLayer), + ); + const layer = makeProviderSessionReaperLive({ + inactivityThresholdMs: input.inactivityThresholdMs, + fallbackReconcileIntervalMs: input.fallbackReconcileIntervalMs, + ...(input.stopFailureRetryIntervalMs !== undefined + ? { stopFailureRetryIntervalMs: input.stopFailureRetryIntervalMs } + : {}), + }).pipe( + Layer.provideMerge(directoryLayer), + Layer.provideMerge(directoryEventsLayer), + Layer.provideMerge(runtimeRepositoryLayer), + Layer.provideMerge(Layer.succeed(ProviderService, providerService)), + Layer.provideMerge(Layer.succeed(OrchestrationEngineService, orchestrationEngine)), + ); + + return { + layer, + stopSession, + setReadModel: (readModel: ReturnType) => + Ref.set(readModelRef, readModel), + publishDomainEvent: (event: OrchestrationEvent) => + PubSub.publish(domainEventPubSub, event).pipe(Effect.asVoid), + }; + }); +} + +it.effect("reaps at the exact inactivity deadline", () => + Effect.scoped( + Effect.gen(function* () { + const threadId = ThreadId.make("thread-deadline-exact"); + const harness = yield* makeHarness({ + initialReadModel: makeReadModel([ + { + id: threadId, + session: { + threadId, + status: "ready", + providerName: "codex", + runtimeMode: "full-access", + activeTurnId: null, + lastError: null, + updatedAt: new Date(0).toISOString(), + }, + }, + ]), + inactivityThresholdMs: 1_000, + fallbackReconcileIntervalMs: 60_000, + }); + + yield* Effect.gen(function* () { + const repository = yield* ProviderSessionRuntimeRepository; + const reaper = yield* ProviderSessionReaper; + + yield* repository.upsert({ + threadId, + providerName: "codex", + adapterKey: "codex", + runtimeMode: "full-access", + status: "running", + lastSeenAt: new Date(0).toISOString(), + resumeCursor: null, + runtimePayload: null, + }); + + yield* reaper.start(); + yield* Effect.yieldNow; + assert.equal(harness.stopSession.mock.calls.length, 0); + + yield* TestClock.adjust(Duration.millis(999)); + yield* Effect.yieldNow; + assert.equal(harness.stopSession.mock.calls.length, 0); + + yield* TestClock.adjust(Duration.millis(2)); + yield* Effect.yieldNow; + assert.equal(harness.stopSession.mock.calls.length, 1); + }).pipe(Effect.provide(harness.layer)); + }).pipe(Effect.provide(TestClock.layer())), + ), +); + +it.effect("reschedules future deadlines after a reap without relying on a directory wake", () => + Effect.scoped( + Effect.gen(function* () { + const dueThreadId = ThreadId.make("thread-post-stop-due"); + const futureThreadId = ThreadId.make("thread-post-stop-future"); + const harness = yield* makeHarness({ + initialReadModel: makeReadModel([ + { + id: dueThreadId, + session: { + threadId: dueThreadId, + status: "ready", + providerName: "codex", + runtimeMode: "full-access", + activeTurnId: null, + lastError: null, + updatedAt: new Date(0).toISOString(), + }, + }, + { + id: futureThreadId, + session: { + threadId: futureThreadId, + status: "ready", + providerName: "claudeAgent", + runtimeMode: "full-access", + activeTurnId: null, + lastError: null, + updatedAt: new Date(500).toISOString(), + }, + }, + ]), + inactivityThresholdMs: 1_000, + fallbackReconcileIntervalMs: 60_000, + stopSessionImplementation: (request) => + Effect.gen(function* () { + const repository = yield* ProviderSessionRuntimeRepository; + const providerName = request.threadId === dueThreadId ? "codex" : "claudeAgent"; + yield* repository.upsert({ + threadId: request.threadId, + providerName, + adapterKey: providerName, + runtimeMode: "full-access", + status: "stopped", + lastSeenAt: new Date(0).toISOString(), + resumeCursor: null, + runtimePayload: { + activeTurnId: null, + }, + }); + }) as ReturnType, + }); + + yield* Effect.gen(function* () { + const repository = yield* ProviderSessionRuntimeRepository; + const reaper = yield* ProviderSessionReaper; + + yield* repository.upsert({ + threadId: dueThreadId, + providerName: "codex", + adapterKey: "codex", + runtimeMode: "full-access", + status: "running", + lastSeenAt: new Date(0).toISOString(), + resumeCursor: null, + runtimePayload: null, + }); + yield* repository.upsert({ + threadId: futureThreadId, + providerName: "claudeAgent", + adapterKey: "claudeAgent", + runtimeMode: "full-access", + status: "running", + lastSeenAt: new Date(500).toISOString(), + resumeCursor: null, + runtimePayload: null, + }); + + yield* reaper.start(); + yield* Effect.yieldNow; + + yield* TestClock.adjust(Duration.millis(1_001)); + yield* Effect.yieldNow; + assert.deepEqual( + harness.stopSession.mock.calls.map(([request]) => request.threadId), + [dueThreadId], + ); + + yield* TestClock.adjust(Duration.millis(498)); + yield* Effect.yieldNow; + assert.deepEqual( + harness.stopSession.mock.calls.map(([request]) => request.threadId), + [dueThreadId], + ); + + yield* TestClock.adjust(Duration.millis(2)); + yield* Effect.yieldNow; + assert.deepEqual( + harness.stopSession.mock.calls.map(([request]) => request.threadId), + [dueThreadId, futureThreadId], + ); + }).pipe(Effect.provide(harness.layer)); + }).pipe(Effect.provide(TestClock.layer())), + ), +); + +it.effect("keeps future deadlines scheduled while retrying failed stops", () => + Effect.scoped( + Effect.gen(function* () { + const failedThreadId = ThreadId.make("thread-stop-failure-retry"); + const futureThreadId = ThreadId.make("thread-stop-failure-future"); + const harness = yield* makeHarness({ + initialReadModel: makeReadModel([ + { + id: failedThreadId, + session: { + threadId: failedThreadId, + status: "ready", + providerName: "codex", + runtimeMode: "full-access", + activeTurnId: null, + lastError: null, + updatedAt: new Date(0).toISOString(), + }, + }, + { + id: futureThreadId, + session: { + threadId: futureThreadId, + status: "ready", + providerName: "claudeAgent", + runtimeMode: "full-access", + activeTurnId: null, + lastError: null, + updatedAt: new Date(300).toISOString(), + }, + }, + ]), + inactivityThresholdMs: 1_000, + fallbackReconcileIntervalMs: 60_000, + stopFailureRetryIntervalMs: 1_000, + stopSessionImplementation: (request) => + request.threadId === failedThreadId + ? Effect.fail( + new ProviderValidationError({ + operation: "ProviderSessionReaper.timing.test", + issue: "simulated stop failure", + }), + ) + : (Effect.gen(function* () { + const repository = yield* ProviderSessionRuntimeRepository; + yield* repository.upsert({ + threadId: request.threadId, + providerName: "claudeAgent", + adapterKey: "claudeAgent", + runtimeMode: "full-access", + status: "stopped", + lastSeenAt: new Date(300).toISOString(), + resumeCursor: null, + runtimePayload: { + activeTurnId: null, + }, + }); + }) as ReturnType), + }); + + yield* Effect.gen(function* () { + const repository = yield* ProviderSessionRuntimeRepository; + const reaper = yield* ProviderSessionReaper; + + yield* repository.upsert({ + threadId: failedThreadId, + providerName: "codex", + adapterKey: "codex", + runtimeMode: "full-access", + status: "running", + lastSeenAt: new Date(0).toISOString(), + resumeCursor: null, + runtimePayload: null, + }); + yield* repository.upsert({ + threadId: futureThreadId, + providerName: "claudeAgent", + adapterKey: "claudeAgent", + runtimeMode: "full-access", + status: "running", + lastSeenAt: new Date(300).toISOString(), + resumeCursor: null, + runtimePayload: null, + }); + + yield* reaper.start(); + yield* Effect.yieldNow; + + yield* TestClock.adjust(Duration.millis(1_001)); + yield* Effect.yieldNow; + assert.deepEqual( + harness.stopSession.mock.calls.map(([request]) => request.threadId), + [failedThreadId], + ); + + yield* TestClock.adjust(Duration.millis(298)); + yield* Effect.yieldNow; + assert.deepEqual( + harness.stopSession.mock.calls.map(([request]) => request.threadId), + [failedThreadId], + ); + + yield* TestClock.adjust(Duration.millis(2)); + yield* Effect.yieldNow; + assert.deepEqual( + harness.stopSession.mock.calls.map(([request]) => request.threadId), + [failedThreadId, failedThreadId, futureThreadId], + ); + + yield* TestClock.adjust(Duration.millis(500)); + yield* Effect.yieldNow; + assert.deepEqual( + harness.stopSession.mock.calls.map(([request]) => request.threadId), + [failedThreadId, failedThreadId, futureThreadId], + ); + + yield* TestClock.adjust(Duration.millis(600)); + yield* Effect.yieldNow; + assert.deepEqual( + harness.stopSession.mock.calls.map(([request]) => request.threadId), + [failedThreadId, failedThreadId, futureThreadId, failedThreadId], + ); + }).pipe(Effect.provide(harness.layer)); + }).pipe(Effect.provide(TestClock.layer())), + ), +); + +it.effect("cancels a pending reap when an active turn starts just before the deadline", () => + Effect.scoped( + Effect.gen(function* () { + const threadId = ThreadId.make("thread-active-mid-wait"); + const activeTurnId = TurnId.make("turn-active-mid-wait"); + const harness = yield* makeHarness({ + initialReadModel: makeReadModel([ + { + id: threadId, + session: { + threadId, + status: "ready", + providerName: "codex", + runtimeMode: "full-access", + activeTurnId: null, + lastError: null, + updatedAt: new Date(0).toISOString(), + }, + }, + ]), + inactivityThresholdMs: 1_000, + fallbackReconcileIntervalMs: 60_000, + }); + + yield* Effect.gen(function* () { + const repository = yield* ProviderSessionRuntimeRepository; + const reaper = yield* ProviderSessionReaper; + + yield* repository.upsert({ + threadId, + providerName: "codex", + adapterKey: "codex", + runtimeMode: "full-access", + status: "running", + lastSeenAt: new Date(0).toISOString(), + resumeCursor: null, + runtimePayload: null, + }); + + yield* reaper.start(); + yield* Effect.yieldNow; + yield* TestClock.adjust(Duration.millis(999)); + yield* Effect.yieldNow; + + const runningSession: ReadModelSession = { + threadId, + status: "running", + providerName: "codex", + runtimeMode: "full-access", + activeTurnId, + lastError: null, + updatedAt: new Date(999).toISOString(), + }; + yield* harness.setReadModel( + makeReadModel([ + { + id: threadId, + latestTurn: { + turnId: activeTurnId, + state: "running", + requestedAt: new Date(999).toISOString(), + startedAt: new Date(999).toISOString(), + completedAt: null, + assistantMessageId: null, + }, + session: runningSession, + }, + ]), + ); + yield* harness.publishDomainEvent(makeThreadSessionSetEvent(threadId, runningSession)); + yield* Effect.yieldNow; + + yield* TestClock.adjust(Duration.millis(2)); + yield* Effect.yieldNow; + assert.equal(harness.stopSession.mock.calls.length, 0); + }).pipe(Effect.provide(harness.layer)); + }).pipe(Effect.provide(TestClock.layer())), + ), +); + +it.effect("reconciles overdue bindings after a long suspended sleep", () => + Effect.scoped( + Effect.gen(function* () { + const firstThreadId = ThreadId.make("thread-sleep-first"); + const secondThreadId = ThreadId.make("thread-sleep-second"); + const harness = yield* makeHarness({ + initialReadModel: makeReadModel([ + { + id: firstThreadId, + session: { + threadId: firstThreadId, + status: "ready", + providerName: "codex", + runtimeMode: "full-access", + activeTurnId: null, + lastError: null, + updatedAt: new Date(0).toISOString(), + }, + }, + { + id: secondThreadId, + session: { + threadId: secondThreadId, + status: "ready", + providerName: "claudeAgent", + runtimeMode: "full-access", + activeTurnId: null, + lastError: null, + updatedAt: new Date(0).toISOString(), + }, + }, + ]), + inactivityThresholdMs: 1_000, + fallbackReconcileIntervalMs: 60_000, + }); + + yield* Effect.gen(function* () { + const repository = yield* ProviderSessionRuntimeRepository; + const reaper = yield* ProviderSessionReaper; + + yield* repository.upsert({ + threadId: firstThreadId, + providerName: "codex", + adapterKey: "codex", + runtimeMode: "full-access", + status: "running", + lastSeenAt: new Date(0).toISOString(), + resumeCursor: null, + runtimePayload: null, + }); + yield* repository.upsert({ + threadId: secondThreadId, + providerName: "claudeAgent", + adapterKey: "claudeAgent", + runtimeMode: "full-access", + status: "running", + lastSeenAt: new Date(500).toISOString(), + resumeCursor: null, + runtimePayload: null, + }); + + yield* reaper.start(); + yield* Effect.yieldNow; + yield* TestClock.adjust(Duration.hours(3)); + yield* Effect.yieldNow; + + const stoppedThreadIds = new Set( + harness.stopSession.mock.calls.map(([request]) => request.threadId), + ); + assert.equal(stoppedThreadIds.has(firstThreadId), true); + assert.equal(stoppedThreadIds.has(secondThreadId), true); + }).pipe(Effect.provide(harness.layer)); + }).pipe(Effect.provide(TestClock.layer())), + ), +); + +it.effect("restarts a failed orchestration feed and wakes on later domain events", () => + Effect.scoped( + Effect.gen(function* () { + const threadId = ThreadId.make("thread-feed-restart"); + let streamSubscriptionCount = 0; + const harness = yield* makeHarness({ + initialReadModel: makeReadModel([ + { + id: threadId, + session: { + threadId, + status: "ready", + providerName: "codex", + runtimeMode: "full-access", + activeTurnId: null, + lastError: null, + updatedAt: new Date(0).toISOString(), + }, + }, + ]), + inactivityThresholdMs: 1_000, + fallbackReconcileIntervalMs: 60_000, + streamDomainEvents: (defaultStream) => { + streamSubscriptionCount += 1; + return streamSubscriptionCount === 1 + ? Stream.die(new Error("simulated transient orchestration feed failure")) + : defaultStream; + }, + }); + + yield* Effect.gen(function* () { + const repository = yield* ProviderSessionRuntimeRepository; + const reaper = yield* ProviderSessionReaper; + + yield* reaper.start(); + yield* Effect.yieldNow; + + yield* TestClock.adjust(Duration.millis(1_000)); + yield* Effect.yieldNow; + assert.equal(streamSubscriptionCount, 2); + + yield* repository.upsert({ + threadId, + providerName: "codex", + adapterKey: "codex", + runtimeMode: "full-access", + status: "running", + lastSeenAt: new Date(0).toISOString(), + resumeCursor: null, + runtimePayload: null, + }); + yield* harness.publishDomainEvent( + makeThreadSessionSetEvent(threadId, { + threadId, + status: "ready", + providerName: "codex", + runtimeMode: "full-access", + activeTurnId: null, + lastError: null, + updatedAt: new Date(1_000).toISOString(), + }), + ); + yield* Effect.yieldNow; + + assert.equal(harness.stopSession.mock.calls.length, 1); + }).pipe(Effect.provide(harness.layer)); + }).pipe(Effect.provide(TestClock.layer())), + ), +); + +it.effect("uses the fallback reconcile to notice bindings that changed without a wake signal", () => + Effect.scoped( + Effect.gen(function* () { + const threadId = ThreadId.make("thread-fallback-reconcile"); + const harness = yield* makeHarness({ + initialReadModel: makeReadModel([ + { + id: threadId, + session: { + threadId, + status: "ready", + providerName: "codex", + runtimeMode: "full-access", + activeTurnId: null, + lastError: null, + updatedAt: new Date(0).toISOString(), + }, + }, + ]), + inactivityThresholdMs: 50, + fallbackReconcileIntervalMs: 100, + }); + + yield* Effect.gen(function* () { + const repository = yield* ProviderSessionRuntimeRepository; + const reaper = yield* ProviderSessionReaper; + + yield* reaper.start(); + yield* Effect.yieldNow; + + yield* repository.upsert({ + threadId, + providerName: "codex", + adapterKey: "codex", + runtimeMode: "full-access", + status: "running", + lastSeenAt: new Date(0).toISOString(), + resumeCursor: null, + runtimePayload: null, + }); + + yield* TestClock.adjust(Duration.millis(100)); + yield* Effect.yieldNow; + assert.equal(harness.stopSession.mock.calls.length, 1); + }).pipe(Effect.provide(harness.layer)); + }).pipe(Effect.provide(TestClock.layer())), + ), +); + +it.effect( + "applies the recent provider.sendTurn deadline floor before orchestration catches up", + () => + Effect.scoped( + Effect.gen(function* () { + const threadId = ThreadId.make("thread-send-turn-floor"); + const staleCompletedAt = new Date(0).toISOString(); + const sendTurnAt = new Date(999).toISOString(); + const harness = yield* makeHarness({ + initialReadModel: makeReadModel([ + { + id: threadId, + latestTurn: { + turnId: TurnId.make("turn-send-turn-stale"), + state: "completed", + requestedAt: staleCompletedAt, + startedAt: staleCompletedAt, + completedAt: staleCompletedAt, + assistantMessageId: null, + }, + session: { + threadId, + status: "ready", + providerName: "codex", + runtimeMode: "full-access", + activeTurnId: null, + lastError: null, + updatedAt: staleCompletedAt, + }, + }, + ]), + inactivityThresholdMs: 1_000, + fallbackReconcileIntervalMs: 60_000, + }); + + yield* Effect.gen(function* () { + const directory = yield* ProviderSessionDirectory; + const repository = yield* ProviderSessionRuntimeRepository; + const reaper = yield* ProviderSessionReaper; + + yield* repository.upsert({ + threadId, + providerName: "codex", + adapterKey: "codex", + runtimeMode: "full-access", + status: "running", + lastSeenAt: staleCompletedAt, + resumeCursor: null, + runtimePayload: null, + }); + + yield* reaper.start(); + yield* Effect.yieldNow; + yield* TestClock.adjust(Duration.millis(999)); + yield* Effect.yieldNow; + + yield* directory.upsert({ + threadId, + provider: "codex", + status: "running", + runtimePayload: { + activeTurnId: TurnId.make("turn-send-turn-new"), + lastRuntimeEvent: "provider.sendTurn", + lastRuntimeEventAt: sendTurnAt, + }, + }); + yield* Effect.yieldNow; + + yield* TestClock.adjust(Duration.millis(2)); + yield* Effect.yieldNow; + assert.equal(harness.stopSession.mock.calls.length, 0); + }).pipe(Effect.provide(harness.layer)); + }).pipe(Effect.provide(TestClock.layer())), + ), +); diff --git a/apps/server/src/provider/Layers/ProviderSessionReaper.ts b/apps/server/src/provider/Layers/ProviderSessionReaper.ts index aa31c8c7d7a9..6445425a5eb8 100644 --- a/apps/server/src/provider/Layers/ProviderSessionReaper.ts +++ b/apps/server/src/provider/Layers/ProviderSessionReaper.ts @@ -1,126 +1,742 @@ -import { Duration, Effect, Layer, Schedule } from "effect"; +import type { ProviderKind } from "@t3tools/contracts"; +import { + Cause, + Clock, + Context, + Duration, + Effect, + Layer, + Metric, + Option, + Queue, + Scope, + Stream, + Tracer, +} from "effect"; +import { + increment, + metricAttributes, + providerSessionReaperDueCandidatesTotal, + providerSessionReaperReapLag, + providerSessionReaperReapedTotal, + providerSessionReaperReconcileDuration, + providerSessionReaperScheduleSize, + providerSessionReaperSignalFeedRestartsTotal, + providerSessionReaperWakeCoalescedTotal, + providerSessionReaperWakeupsTotal, +} from "../../observability/Metrics.ts"; import { OrchestrationEngineService } from "../../orchestration/Services/OrchestrationEngine.ts"; import { ProviderSessionDirectory } from "../Services/ProviderSessionDirectory.ts"; +import { ProviderSessionDirectoryEvents } from "../Services/ProviderSessionDirectoryEvents.ts"; import { + DEFAULT_PROVIDER_SESSION_REAPER_FALLBACK_RECONCILE_INTERVAL_MS, + DEFAULT_PROVIDER_SESSION_REAPER_INACTIVITY_THRESHOLD_MS, + DEFAULT_PROVIDER_SESSION_REAPER_STOP_FAILURE_RETRY_INTERVAL_MS, ProviderSessionReaper, type ProviderSessionReaperShape, } from "../Services/ProviderSessionReaper.ts"; import { ProviderService } from "../Services/ProviderService.ts"; +import type { InvalidAnchorEntry, ReapScheduleEntry } from "./reaperDeadlines.ts"; +import { deriveReapEntries } from "./reaperDeadlines.ts"; -const DEFAULT_INACTIVITY_THRESHOLD_MS = 30 * 60 * 1000; -const DEFAULT_SWEEP_INTERVAL_MS = 5 * 60 * 1000; +const REAPER_MODE = "deadline"; +const SIGNAL_FEED_RESTART_DELAY_MS = 1_000; export interface ProviderSessionReaperLiveOptions { readonly inactivityThresholdMs?: number; - readonly sweepIntervalMs?: number; + readonly fallbackReconcileIntervalMs?: number; + readonly stopFailureRetryIntervalMs?: number; +} + +type ReaperSignal = + | { readonly type: "startup" } + | { readonly type: "reconcile-all"; readonly reason: string } + | { readonly type: "runtime-binding-changed"; readonly threadId: string } + | { readonly type: "orchestration-thread-changed"; readonly threadId: string } + | { readonly type: "thread-deleted"; readonly threadId: string }; + +interface ReconcileSnapshot { + readonly bindingCount: number; + readonly readModelThreadCount: number; + readonly entries: ReadonlyArray; + readonly skippedStopped: number; + readonly skippedActiveTurn: number; + readonly invalidAnchors: ReadonlyArray; + readonly reconciledAtMs: number; +} + +interface CoalescedWake { + readonly signal: (signal: ReaperSignal) => Effect.Effect; + readonly await: Effect.Effect; +} + +interface SchedulerIterationResult { + readonly nextDeadlineAtMs: number | undefined; + readonly observedStartupWake: boolean; +} + +function buildReaperLogContext(input: { + readonly now: number; + readonly inactivityThresholdMs: number; + readonly entry: ReapScheduleEntry; +}) { + const inactivityAnchorMs = Date.parse(input.entry.anchorAt); + const idleDurationMs = input.now - inactivityAnchorMs; + + return { + threadId: input.entry.threadId, + provider: input.entry.provider, + bindingStatus: input.entry.bindingStatus, + readModelThreadPresent: input.entry.readModelThreadPresent, + sessionStatus: input.entry.sessionStatus, + sessionUpdatedAt: input.entry.sessionUpdatedAt, + activeTurnId: input.entry.activeTurnId, + lastSeenAt: input.entry.lastSeenAt, + latestTurnId: input.entry.latestTurnId, + latestTurnState: input.entry.latestTurnState, + latestTurnRequestedAt: input.entry.latestTurnRequestedAt, + latestTurnStartedAt: input.entry.latestTurnStartedAt, + latestTurnCompletedAt: input.entry.latestTurnCompletedAt, + inactivityAnchorAt: input.entry.anchorAt, + inactivityAnchorSource: input.entry.anchorSource, + inactivityAnchorMs, + deadlineBasisAt: input.entry.deadlineBasisAt, + deadlineBasisSource: input.entry.deadlineBasisSource, + deadlineAt: new Date(input.entry.deadlineAtMs).toISOString(), + deadlineAtMs: input.entry.deadlineAtMs, + inactivityThresholdMs: input.inactivityThresholdMs, + idleDurationMs, + remainingUntilReapMs: Math.max(0, input.entry.deadlineAtMs - input.now), + reapLagMs: Math.max(0, input.now - input.entry.deadlineAtMs), + }; +} + +function buildInvalidAnchorLogContext(input: { + readonly now: number; + readonly inactivityThresholdMs: number; + readonly entry: InvalidAnchorEntry; +}) { + return { + threadId: input.entry.threadId, + provider: input.entry.provider, + bindingStatus: input.entry.bindingStatus, + readModelThreadPresent: input.entry.readModelThreadPresent, + sessionStatus: input.entry.sessionStatus, + sessionUpdatedAt: input.entry.sessionUpdatedAt, + activeTurnId: input.entry.activeTurnId, + lastSeenAt: input.entry.lastSeenAt, + latestTurnId: input.entry.latestTurnId, + latestTurnState: input.entry.latestTurnState, + latestTurnRequestedAt: input.entry.latestTurnRequestedAt, + latestTurnStartedAt: input.entry.latestTurnStartedAt, + latestTurnCompletedAt: input.entry.latestTurnCompletedAt, + inactivityAnchorAt: input.entry.inactivityAnchorAt, + inactivityAnchorSource: input.entry.inactivityAnchorSource, + inactivityAnchorMs: null, + reconciledAt: new Date(input.now).toISOString(), + inactivityThresholdMs: input.inactivityThresholdMs, + idleDurationMs: null, + remainingUntilReapMs: null, + }; +} + +function wakeReason(signal: ReaperSignal | "timeout"): string { + if (signal === "timeout") { + return "timeout"; + } + if (signal.type === "startup") { + return "startup"; + } + if (signal.type === "reconcile-all" && signal.reason === "fallback-tick") { + return "fallback"; + } + return `signal:${signal.type}`; } +function wakeLogContext(signal: ReaperSignal | "timeout"): Record { + if (signal === "timeout" || signal.type === "startup") { + return {}; + } + if (signal.type === "reconcile-all") { + return { + reconcileReason: signal.reason, + }; + } + return { + threadId: signal.threadId, + }; +} + +function providerBreakdown( + entries: ReadonlyArray, +): Record { + const counts: Record = { + codex: 0, + claudeAgent: 0, + cursor: 0, + opencode: 0, + }; + for (const entry of entries) { + counts[entry.provider] += 1; + } + return counts; +} + +function earliestDeadline(...deadlines: ReadonlyArray): number | undefined { + let earliest: number | undefined = undefined; + for (const deadline of deadlines) { + if (deadline === undefined) { + continue; + } + earliest = earliest === undefined ? deadline : Math.min(earliest, deadline); + } + return earliest; +} + +const makeCoalescedWake = () => + Effect.gen(function* () { + const queue = yield* Effect.acquireRelease(Queue.dropping(1), Queue.shutdown); + + return { + signal: (signal) => + Queue.offer(queue, signal).pipe( + Effect.flatMap((enqueued) => + enqueued + ? Effect.void + : increment(providerSessionReaperWakeCoalescedTotal, { mode: REAPER_MODE }), + ), + ), + await: Queue.take(queue), + } satisfies CoalescedWake; + }); + const makeProviderSessionReaper = (options?: ProviderSessionReaperLiveOptions) => Effect.gen(function* () { const providerService = yield* ProviderService; const directory = yield* ProviderSessionDirectory; + const directoryEvents = yield* ProviderSessionDirectoryEvents; const orchestrationEngine = yield* OrchestrationEngineService; const inactivityThresholdMs = Math.max( 1, - options?.inactivityThresholdMs ?? DEFAULT_INACTIVITY_THRESHOLD_MS, + options?.inactivityThresholdMs ?? DEFAULT_PROVIDER_SESSION_REAPER_INACTIVITY_THRESHOLD_MS, ); - const sweepIntervalMs = Math.max(1, options?.sweepIntervalMs ?? DEFAULT_SWEEP_INTERVAL_MS); - - const sweep = Effect.gen(function* () { - const readModel = yield* orchestrationEngine.getReadModel(); - const threadsById = new Map(readModel.threads.map((thread) => [thread.id, thread] as const)); - const bindings = yield* directory.listBindings(); - const now = Date.now(); - let reapedCount = 0; - - for (const binding of bindings) { - if (binding.status === "stopped") { - continue; - } + const fallbackReconcileIntervalMs = Math.max( + 1, + options?.fallbackReconcileIntervalMs ?? + DEFAULT_PROVIDER_SESSION_REAPER_FALLBACK_RECONCILE_INTERVAL_MS, + ); + const stopFailureRetryIntervalMs = Math.max( + 1, + options?.stopFailureRetryIntervalMs ?? + DEFAULT_PROVIDER_SESSION_REAPER_STOP_FAILURE_RETRY_INTERVAL_MS, + ); + const mode = REAPER_MODE; - const lastSeenMs = Date.parse(binding.lastSeenAt); - if (Number.isNaN(lastSeenMs)) { - yield* Effect.logWarning("provider.session.reaper.invalid-last-seen", { - threadId: binding.threadId, - provider: binding.provider, - lastSeenAt: binding.lastSeenAt, - }); - continue; - } + const recordWake = (signal: ReaperSignal | "timeout", extra?: Record) => + increment(providerSessionReaperWakeupsTotal, { + mode, + reason: wakeReason(signal), + }).pipe( + Effect.andThen( + Effect.logDebug("provider.session.reaper.wake", { + mode, + reason: wakeReason(signal), + ...wakeLogContext(signal), + ...extra, + }), + ), + ); - const idleDurationMs = now - lastSeenMs; - if (idleDurationMs < inactivityThresholdMs) { - continue; - } + const reconcileAuthoritativeState = () => + Effect.gen(function* () { + yield* Effect.annotateCurrentSpan({ + "provider.session_reaper.mode": mode, + "provider.session_reaper.inactivity_threshold_ms": inactivityThresholdMs, + }); + const startedAtMs = yield* Clock.currentTimeMillis; + yield* Effect.logDebug("provider.session.reaper.reconcile-started", { + mode, + inactivityThresholdMs, + }); - const thread = threadsById.get(binding.threadId); - if (thread?.session?.activeTurnId != null) { - yield* Effect.logDebug("provider.session.reaper.skipped-active-turn", { - threadId: binding.threadId, - activeTurnId: thread.session.activeTurnId, - idleDurationMs, - }); - continue; - } + const readModel = yield* orchestrationEngine.getReadModel(); + const bindings = yield* directory.listBindings(); + const derived = deriveReapEntries({ + bindings, + readModel, + inactivityThresholdMs, + }); + const reconciledAtMs = yield* Clock.currentTimeMillis; + const reconcileDurationMs = Math.max(0, reconciledAtMs - startedAtMs); + yield* Effect.annotateCurrentSpan({ + "provider.session_reaper.binding_count": bindings.length, + "provider.session_reaper.read_model_thread_count": readModel.threads.length, + "provider.session_reaper.schedule_size": derived.entries.length, + "provider.session_reaper.skipped_stopped_count": derived.skippedStopped, + "provider.session_reaper.skipped_active_turn_count": derived.skippedActiveTurn, + "provider.session_reaper.invalid_anchor_count": derived.invalidAnchors.length, + "provider.session_reaper.reconcile_duration_ms": reconcileDurationMs, + }); + + yield* Metric.update( + Metric.withAttributes(providerSessionReaperReconcileDuration, metricAttributes({ mode })), + Duration.millis(reconcileDurationMs), + ); - const reaped = yield* providerService.stopSession({ threadId: binding.threadId }).pipe( - Effect.tap(() => - Effect.logInfo("provider.session.reaped", { - threadId: binding.threadId, - provider: binding.provider, - idleDurationMs, - reason: "inactivity_threshold", + for (const invalidAnchor of derived.invalidAnchors) { + yield* Effect.logWarning( + "provider.session.reaper.invalid-inactivity-anchor", + buildInvalidAnchorLogContext({ + now: reconciledAtMs, + inactivityThresholdMs, + entry: invalidAnchor, }), - ), - Effect.as(true), - Effect.catchCause((cause) => - Effect.logWarning("provider.session.reaper.stop-failed", { - threadId: binding.threadId, - provider: binding.provider, - idleDurationMs, - cause, - }).pipe(Effect.as(false)), - ), + ); + } + + yield* Effect.logDebug("provider.session.reaper.reconcile-completed", { + mode, + bindingCount: bindings.length, + readModelThreadCount: readModel.threads.length, + scheduleSize: derived.entries.length, + skippedStoppedCount: derived.skippedStopped, + skippedActiveTurnCount: derived.skippedActiveTurn, + invalidAnchorCount: derived.invalidAnchors.length, + reconcileDurationMs, + }); + + return { + bindingCount: bindings.length, + readModelThreadCount: readModel.threads.length, + entries: derived.entries, + skippedStopped: derived.skippedStopped, + skippedActiveTurn: derived.skippedActiveTurn, + invalidAnchors: derived.invalidAnchors, + reconciledAtMs, + } satisfies ReconcileSnapshot; + }).pipe( + Effect.withSpan("provider.session.reaper.reconcile", { + attributes: { + "provider.session_reaper.mode": mode, + }, + }), + ); + + const stopDueEntries = (input: { + readonly entries: ReadonlyArray; + readonly nowMs: number; + }) => + Effect.gen(function* () { + yield* Effect.annotateCurrentSpan({ + "provider.session_reaper.mode": mode, + "provider.session_reaper.due_count": input.entries.length, + }); + yield* increment( + providerSessionReaperDueCandidatesTotal, + { + mode, + }, + input.entries.length, ); - if (reaped) { - reapedCount += 1; + let reapedCount = 0; + let stopFailedCount = 0; + + for (const entry of input.entries) { + const reaped = yield* Effect.gen(function* () { + const logContext = buildReaperLogContext({ + now: input.nowMs, + inactivityThresholdMs, + entry, + }); + + yield* Effect.annotateCurrentSpan({ + "provider.session_reaper.mode": mode, + "provider.thread_id": entry.threadId, + "provider.kind": entry.provider, + "provider.session_reaper.deadline_basis_source": entry.deadlineBasisSource, + "provider.session_reaper.deadline_at_ms": entry.deadlineAtMs, + "provider.session_reaper.reap_lag_ms": logContext.reapLagMs, + }); + + yield* Metric.update( + Metric.withAttributes( + providerSessionReaperReapLag, + metricAttributes({ + mode, + provider: entry.provider, + }), + ), + Duration.millis(Math.max(0, input.nowMs - entry.deadlineAtMs)), + ); + + yield* Effect.logDebug("provider.session.reaper.reap-candidate", { + ...logContext, + decision: "attempt_stop_session", + }); + + const reaped = yield* providerService + .stopSession({ + threadId: entry.threadId, + }) + .pipe( + Effect.tap(() => + Effect.logInfo("provider.session.reaped", { + ...logContext, + reason: "inactivity_threshold", + }).pipe( + Effect.andThen( + increment(providerSessionReaperReapedTotal, { + mode, + provider: entry.provider, + }), + ), + ), + ), + Effect.as(true), + Effect.catchCause((cause) => { + if (Cause.hasInterruptsOnly(cause)) { + return Effect.failCause(cause); + } + stopFailedCount += 1; + return Effect.logWarning("provider.session.reaper.stop-failed", { + ...logContext, + reason: "inactivity_threshold", + cause, + }).pipe(Effect.as(false)); + }), + ); + + yield* Effect.annotateCurrentSpan({ + "provider.session_reaper.stop_outcome": reaped ? "reaped" : "failed", + }); + return reaped; + }).pipe( + Effect.withSpan("provider.session.reaper.stop_session", { + kind: "client", + attributes: { + "provider.session_reaper.mode": mode, + "provider.thread_id": entry.threadId, + "provider.kind": entry.provider, + "provider.session_reaper.deadline_basis_source": entry.deadlineBasisSource, + }, + }), + ); + + if (reaped) { + reapedCount += 1; + } } - } - if (reapedCount > 0) { - yield* Effect.logInfo("provider.session.reaper.sweep-complete", { - reapedCount, - totalBindings: bindings.length, + yield* Effect.annotateCurrentSpan({ + "provider.session_reaper.reaped_count": reapedCount, + "provider.session_reaper.stop_failed_count": stopFailedCount, }); + + return { + reapedCount, + stopFailedCount, + } as const; + }).pipe( + Effect.withSpan("provider.session.reaper.stop_due", { + kind: "internal", + attributes: { + "provider.session_reaper.mode": mode, + }, + }), + ); + + const runDeadlineScheduler = Effect.gen(function* () { + const wake = yield* makeCoalescedWake(); + + const signalWake = (signal: ReaperSignal) => wake.signal(signal); + + const forkFeed = (name: string, effect: () => Effect.Effect) => { + const run = (): Effect.Effect => + Effect.suspend(effect).pipe( + Effect.catchCause((cause) => { + if (Cause.hasInterruptsOnly(cause)) { + return Effect.failCause(cause); + } + return increment(providerSessionReaperSignalFeedRestartsTotal, { + mode, + feed: name, + }).pipe( + Effect.andThen( + Effect.logWarning("provider.session.reaper.signal-stream-failed", { + mode, + feed: name, + cause, + }), + ), + Effect.andThen( + signalWake({ type: "reconcile-all", reason: `feed-failed:${name}` }), + ), + Effect.andThen(Effect.sleep(Duration.millis(SIGNAL_FEED_RESTART_DELAY_MS))), + Effect.flatMap(() => run()), + ); + }), + ); + + return Effect.forkScoped(run()); + }; + + yield* forkFeed("provider-session-directory-events", () => + Stream.runForEach(directoryEvents.changes, (change) => + signalWake({ + type: "runtime-binding-changed", + threadId: change.threadId, + }), + ), + ); + + yield* forkFeed("orchestration-domain-events", () => + Stream.runForEach(orchestrationEngine.streamDomainEvents, (event) => { + switch (event.type) { + case "thread.deleted": + return signalWake({ + type: "thread-deleted", + threadId: event.payload.threadId, + }); + case "thread.reverted": + case "thread.session-set": + case "thread.turn-diff-completed": + return signalWake({ + type: "orchestration-thread-changed", + threadId: event.payload.threadId, + }); + default: + return Effect.void; + } + }), + ); + + yield* Effect.forkScoped( + signalWake({ type: "reconcile-all", reason: "fallback-tick" }).pipe( + Effect.delay(Duration.millis(fallbackReconcileIntervalMs)), + Effect.forever, + ), + ); + + yield* signalWake({ type: "startup" }); + let nextDeadlineAtMs: number | undefined = undefined; + let observedStartupWake = false; + + while (true) { + const iterationResult: SchedulerIterationResult = yield* Effect.gen(function* () { + yield* Effect.annotateCurrentSpan({ + "provider.session_reaper.mode": mode, + ...(nextDeadlineAtMs !== undefined + ? { + "provider.session_reaper.next_deadline_at_ms": nextDeadlineAtMs, + "provider.session_reaper.next_deadline_at": new Date( + nextDeadlineAtMs, + ).toISOString(), + } + : {}), + }); + + let wakeSignal: ReaperSignal | "timeout"; + if (nextDeadlineAtMs === undefined) { + const signal = yield* wake.await; + observedStartupWake ||= signal.type === "startup"; + wakeSignal = signal; + yield* recordWake(signal); + } else { + const waitMs = Math.max(0, nextDeadlineAtMs - (yield* Clock.currentTimeMillis)); + const signal = yield* wake.await.pipe(Effect.timeoutOption(Duration.millis(waitMs))); + if (Option.isSome(signal)) { + observedStartupWake ||= signal.value.type === "startup"; + wakeSignal = signal.value; + yield* recordWake(signal.value, { + deadlineAt: new Date(nextDeadlineAtMs).toISOString(), + waitMs, + }); + } else { + wakeSignal = "timeout"; + yield* recordWake("timeout", { + deadlineAt: new Date(nextDeadlineAtMs).toISOString(), + waitMs, + }); + } + } + + yield* Effect.annotateCurrentSpan({ + "provider.session_reaper.wake_reason": wakeReason(wakeSignal), + ...(wakeSignal !== "timeout" && "threadId" in wakeSignal + ? { "provider.thread_id": wakeSignal.threadId } + : {}), + ...(wakeSignal !== "timeout" && wakeSignal.type === "reconcile-all" + ? { "provider.session_reaper.reconcile_reason": wakeSignal.reason } + : {}), + }); + + const snapshotOption = yield* reconcileAuthoritativeState().pipe( + Effect.map(Option.some), + Effect.catchCause((cause) => { + if (Cause.hasInterruptsOnly(cause)) { + return Effect.failCause(cause); + } + return Effect.logWarning("provider.session.reaper.reconcile-failed", { + mode, + cause, + }).pipe(Effect.as(Option.none())); + }), + ); + + if (Option.isNone(snapshotOption)) { + yield* Effect.annotateCurrentSpan({ + "provider.session_reaper.reconcile_outcome": "failed", + }); + return { + nextDeadlineAtMs: undefined, + observedStartupWake, + }; + } + + const snapshot = snapshotOption.value; + const now = snapshot.reconciledAtMs; + const dueEntries = snapshot.entries.filter((entry) => entry.deadlineAtMs <= now); + const futureEntries = snapshot.entries.filter((entry) => entry.deadlineAtMs > now); + + yield* Metric.update( + Metric.withAttributes(providerSessionReaperScheduleSize, metricAttributes({ mode })), + futureEntries.length, + ); + + yield* Effect.annotateCurrentSpan({ + "provider.session_reaper.reconcile_outcome": "success", + "provider.session_reaper.schedule_size": futureEntries.length, + "provider.session_reaper.due_count": dueEntries.length, + "provider.session_reaper.skipped_stopped_count": snapshot.skippedStopped, + "provider.session_reaper.skipped_active_turn_count": snapshot.skippedActiveTurn, + "provider.session_reaper.invalid_anchor_count": snapshot.invalidAnchors.length, + }); + + if (observedStartupWake && dueEntries.length > 0) { + observedStartupWake = false; + yield* Effect.logInfo("provider.session.reaper.overdue-on-startup", { + mode, + overdueCount: dueEntries.length, + providers: providerBreakdown(dueEntries), + }); + } + + if (dueEntries.length > 0) { + const stopResult = yield* stopDueEntries({ + entries: dueEntries, + nowMs: now, + }); + yield* Effect.annotateCurrentSpan({ + "provider.session_reaper.reaped_count": stopResult.reapedCount, + "provider.session_reaper.stop_failed_count": stopResult.stopFailedCount, + }); + if (stopResult.stopFailedCount > 0) { + const nextFutureDeadlineAtMs = futureEntries[0]?.deadlineAtMs; + const retryDeadlineAtMs = now + stopFailureRetryIntervalMs; + const nextDeadlineAtMs = earliestDeadline(retryDeadlineAtMs, nextFutureDeadlineAtMs); + yield* Effect.logDebug("provider.session.reaper.stop-failure-retry-scheduled", { + mode, + stopFailedCount: stopResult.stopFailedCount, + reapedCount: stopResult.reapedCount, + stopFailureRetryIntervalMs, + retryDeadlineAtMs, + retryDeadlineAt: new Date(retryDeadlineAtMs).toISOString(), + ...(nextFutureDeadlineAtMs !== undefined + ? { + nextFutureDeadlineAtMs, + nextFutureDeadlineAt: new Date(nextFutureDeadlineAtMs).toISOString(), + } + : {}), + ...(nextDeadlineAtMs !== undefined + ? { + nextDeadlineAtMs, + nextDeadlineAt: new Date(nextDeadlineAtMs).toISOString(), + waitMs: Math.max(0, nextDeadlineAtMs - now), + } + : {}), + }); + return { + nextDeadlineAtMs, + observedStartupWake: false, + }; + } + if (stopResult.reapedCount > 0 && futureEntries.length > 0) { + yield* signalWake({ type: "reconcile-all", reason: "post-stop" }); + } + return { + nextDeadlineAtMs: undefined, + observedStartupWake: false, + }; + } + + const nextEntry = futureEntries[0]; + if (nextEntry === undefined) { + return { + nextDeadlineAtMs: undefined, + observedStartupWake: false, + }; + } + + yield* Effect.logDebug("provider.session.reaper.next-deadline-selected", { + ...buildReaperLogContext({ + now, + inactivityThresholdMs, + entry: nextEntry, + }), + mode, + waitMs: Math.max(0, nextEntry.deadlineAtMs - now), + }); + + return { + nextDeadlineAtMs: nextEntry.deadlineAtMs, + observedStartupWake: false, + }; + }).pipe( + Effect.withSpan("provider.session.reaper.iteration", { + kind: "internal", + root: true, + attributes: { + "provider.session_reaper.mode": mode, + }, + }), + ); + + nextDeadlineAtMs = iterationResult.nextDeadlineAtMs; + observedStartupWake = iterationResult.observedStartupWake; } }); const start: ProviderSessionReaperShape["start"] = () => Effect.gen(function* () { yield* Effect.forkScoped( - sweep.pipe( - Effect.catch((error: unknown) => - Effect.logWarning("provider.session.reaper.sweep-failed", { - error, - }), - ), - Effect.catchDefect((defect: unknown) => - Effect.logWarning("provider.session.reaper.sweep-defect", { - defect, - }), + runDeadlineScheduler.pipe( + Effect.updateContext((context: Context.Context) => + Context.omit(Tracer.ParentSpan)(context), ), - Effect.repeat(Schedule.spaced(Duration.millis(sweepIntervalMs))), ), ); - yield* Effect.logInfo("provider.session.reaper.started", { + yield* Effect.logInfo("provider.session.reaper.scheduler-started", { + mode, inactivityThresholdMs, - sweepIntervalMs, + fallbackReconcileIntervalMs, + stopFailureRetryIntervalMs, }); - }); + }).pipe( + Effect.withSpan("provider.session.reaper.start", { + kind: "internal", + root: true, + attributes: { + "provider.session_reaper.mode": mode, + "provider.session_reaper.inactivity_threshold_ms": inactivityThresholdMs, + "provider.session_reaper.fallback_reconcile_interval_ms": fallbackReconcileIntervalMs, + "provider.session_reaper.stop_failure_retry_interval_ms": stopFailureRetryIntervalMs, + }, + }), + ); return { start, @@ -129,5 +745,3 @@ const makeProviderSessionReaper = (options?: ProviderSessionReaperLiveOptions) = export const makeProviderSessionReaperLive = (options?: ProviderSessionReaperLiveOptions) => Layer.effect(ProviderSessionReaper, makeProviderSessionReaper(options)); - -export const ProviderSessionReaperLive = makeProviderSessionReaperLive(); diff --git a/apps/server/src/provider/Layers/reaperDeadlines.test.ts b/apps/server/src/provider/Layers/reaperDeadlines.test.ts new file mode 100644 index 000000000000..eefd8d6bc450 --- /dev/null +++ b/apps/server/src/provider/Layers/reaperDeadlines.test.ts @@ -0,0 +1,347 @@ +import { MessageId, ProjectId, ThreadId, TurnId } from "@t3tools/contracts"; +import { describe, expect, it } from "vitest"; + +import { deriveReapEntries } from "./reaperDeadlines.ts"; + +const defaultModelSelection = { + provider: "codex", + model: "gpt-5-codex", +} as const; + +function makeReadModel( + threads: ReadonlyArray<{ + readonly id: ThreadId; + readonly latestTurn?: { + readonly turnId: TurnId; + readonly state: "running" | "interrupted" | "completed" | "error"; + readonly requestedAt: string; + readonly startedAt: string | null; + readonly completedAt: string | null; + readonly assistantMessageId: MessageId | null; + } | null; + readonly session?: { + readonly threadId: ThreadId; + readonly status: "starting" | "running" | "ready" | "interrupted" | "stopped" | "error"; + readonly providerName: "codex" | "claudeAgent" | "cursor" | "opencode"; + readonly runtimeMode: "approval-required" | "full-access" | "auto-accept-edits"; + readonly activeTurnId: TurnId | null; + readonly lastError: string | null; + readonly updatedAt: string; + } | null; + }>, +) { + const now = "2026-04-20T12:00:00.000Z"; + const projectId = ProjectId.make("project-provider-reaper-deadlines"); + + return { + snapshotSequence: 0, + updatedAt: now, + projects: [ + { + id: projectId, + title: "Provider Reaper Deadlines Project", + workspaceRoot: "/tmp/provider-reaper-deadlines", + defaultModelSelection, + scripts: [], + createdAt: now, + updatedAt: now, + deletedAt: null, + }, + ], + threads: threads.map((thread) => ({ + id: thread.id, + projectId, + title: `Thread ${thread.id}`, + modelSelection: defaultModelSelection, + interactionMode: "default" as const, + runtimeMode: "full-access" as const, + branch: null, + worktreePath: null, + createdAt: now, + updatedAt: now, + archivedAt: null, + latestTurn: thread.latestTurn ?? null, + messages: [], + session: thread.session ?? null, + activities: [], + proposedPlans: [], + checkpoints: [], + deletedAt: null, + })), + }; +} + +describe("deriveReapEntries", () => { + it("prefers latestTurn.completedAt over newer runtime metadata", () => { + const threadId = ThreadId.make("thread-anchor-completed"); + const turnId = TurnId.make("turn-anchor-completed"); + const result = deriveReapEntries({ + inactivityThresholdMs: 1_000, + readModel: makeReadModel([ + { + id: threadId, + latestTurn: { + turnId, + state: "completed", + requestedAt: "2026-04-20T10:00:00.000Z", + startedAt: "2026-04-20T10:01:00.000Z", + completedAt: "2026-04-20T10:02:00.000Z", + assistantMessageId: null, + }, + }, + ]), + bindings: [ + { + threadId, + provider: "codex", + status: "running", + lastSeenAt: "2026-04-20T11:00:00.000Z", + runtimeMode: "full-access", + }, + ], + }); + + expect(result.entries).toHaveLength(1); + expect(result.entries[0]).toMatchObject({ + anchorAt: "2026-04-20T10:02:00.000Z", + anchorSource: "latest_turn_completed_at", + deadlineBasisAt: "2026-04-20T10:02:00.000Z", + deadlineBasisSource: "anchor", + }); + }); + + it("falls back to startedAt, then requestedAt, then lastSeenAt", () => { + const startedThreadId = ThreadId.make("thread-anchor-started"); + const requestedThreadId = ThreadId.make("thread-anchor-requested"); + const lastSeenThreadId = ThreadId.make("thread-anchor-last-seen"); + const result = deriveReapEntries({ + inactivityThresholdMs: 1_000, + readModel: makeReadModel([ + { + id: startedThreadId, + latestTurn: { + turnId: TurnId.make("turn-anchor-started"), + state: "running", + requestedAt: "2026-04-20T10:00:00.000Z", + startedAt: "2026-04-20T10:01:00.000Z", + completedAt: null, + assistantMessageId: null, + }, + }, + { + id: requestedThreadId, + latestTurn: { + turnId: TurnId.make("turn-anchor-requested"), + state: "running", + requestedAt: "2026-04-20T10:05:00.000Z", + startedAt: null, + completedAt: null, + assistantMessageId: null, + }, + }, + ]), + bindings: [ + { + threadId: startedThreadId, + provider: "codex", + status: "running", + lastSeenAt: "2026-04-20T11:00:00.000Z", + runtimeMode: "full-access", + }, + { + threadId: requestedThreadId, + provider: "codex", + status: "running", + lastSeenAt: "2026-04-20T11:05:00.000Z", + runtimeMode: "full-access", + }, + { + threadId: lastSeenThreadId, + provider: "claudeAgent", + status: "running", + lastSeenAt: "2026-04-20T11:10:00.000Z", + runtimeMode: "full-access", + }, + ], + }); + + expect(result.entries).toMatchObject([ + { + threadId: startedThreadId, + anchorAt: "2026-04-20T10:01:00.000Z", + anchorSource: "latest_turn_started_at", + }, + { + threadId: requestedThreadId, + anchorAt: "2026-04-20T10:05:00.000Z", + anchorSource: "latest_turn_requested_at", + }, + { + threadId: lastSeenThreadId, + anchorAt: "2026-04-20T11:10:00.000Z", + anchorSource: "session_last_seen_at", + }, + ]); + }); + + it("filters stopped bindings and active turns", () => { + const stoppedThreadId = ThreadId.make("thread-stopped"); + const activeThreadId = ThreadId.make("thread-active"); + const result = deriveReapEntries({ + inactivityThresholdMs: 1_000, + readModel: makeReadModel([ + { + id: activeThreadId, + session: { + threadId: activeThreadId, + status: "running", + providerName: "codex", + runtimeMode: "full-access", + activeTurnId: TurnId.make("turn-active"), + lastError: null, + updatedAt: "2026-04-20T10:00:00.000Z", + }, + }, + ]), + bindings: [ + { + threadId: stoppedThreadId, + provider: "codex", + status: "stopped", + lastSeenAt: "2026-04-20T09:00:00.000Z", + runtimeMode: "full-access", + }, + { + threadId: activeThreadId, + provider: "codex", + status: "running", + lastSeenAt: "2026-04-20T09:05:00.000Z", + runtimeMode: "full-access", + }, + ], + }); + + expect(result.entries).toHaveLength(0); + expect(result.skippedStopped).toBe(1); + expect(result.skippedActiveTurn).toBe(1); + }); + + it("reports invalid anchors and skips scheduling them", () => { + const threadId = ThreadId.make("thread-invalid-anchor"); + const result = deriveReapEntries({ + inactivityThresholdMs: 1_000, + readModel: makeReadModel([ + { + id: threadId, + latestTurn: { + turnId: TurnId.make("turn-invalid-anchor"), + state: "completed", + requestedAt: "2026-04-20T10:00:00.000Z", + startedAt: "2026-04-20T10:01:00.000Z", + completedAt: "not-a-date", + assistantMessageId: null, + }, + }, + ]), + bindings: [ + { + threadId, + provider: "codex", + status: "running", + lastSeenAt: "2026-04-20T11:00:00.000Z", + runtimeMode: "full-access", + }, + ], + }); + + expect(result.entries).toHaveLength(0); + expect(result.invalidAnchors).toMatchObject([ + { + threadId, + inactivityAnchorAt: "not-a-date", + inactivityAnchorSource: "latest_turn_completed_at", + }, + ]); + }); + + it("uses a recent provider.sendTurn timestamp as a deadline floor", () => { + const threadId = ThreadId.make("thread-send-turn-floor"); + const result = deriveReapEntries({ + inactivityThresholdMs: 1_000, + readModel: makeReadModel([ + { + id: threadId, + latestTurn: { + turnId: TurnId.make("turn-send-turn-floor"), + state: "completed", + requestedAt: "2026-04-20T09:00:00.000Z", + startedAt: "2026-04-20T09:01:00.000Z", + completedAt: "2026-04-20T09:02:00.000Z", + assistantMessageId: null, + }, + }, + ]), + bindings: [ + { + threadId, + provider: "codex", + status: "running", + lastSeenAt: "2026-04-20T11:00:00.000Z", + runtimeMode: "full-access", + runtimePayload: { + lastRuntimeEvent: "provider.sendTurn", + lastRuntimeEventAt: "2026-04-20T11:30:00.000Z", + }, + }, + ], + }); + + expect(result.entries).toHaveLength(1); + expect(result.entries[0]).toMatchObject({ + anchorAt: "2026-04-20T09:02:00.000Z", + anchorSource: "latest_turn_completed_at", + deadlineBasisAt: "2026-04-20T11:30:00.000Z", + deadlineBasisSource: "recent_provider_send_turn", + deadlineAtMs: Date.parse("2026-04-20T11:30:00.000Z") + 1_000, + }); + }); + + it("ignores invalid or unrelated runtime metadata for deadline floors", () => { + const threadId = ThreadId.make("thread-ignored-floor"); + const result = deriveReapEntries({ + inactivityThresholdMs: 1_000, + readModel: makeReadModel([ + { + id: threadId, + latestTurn: { + turnId: TurnId.make("turn-ignored-floor"), + state: "completed", + requestedAt: "2026-04-20T09:00:00.000Z", + startedAt: "2026-04-20T09:01:00.000Z", + completedAt: "2026-04-20T09:02:00.000Z", + assistantMessageId: null, + }, + }, + ]), + bindings: [ + { + threadId, + provider: "codex", + status: "running", + lastSeenAt: "2026-04-20T11:00:00.000Z", + runtimeMode: "full-access", + runtimePayload: { + lastRuntimeEvent: "provider.stopAll", + lastRuntimeEventAt: "not-a-date", + }, + }, + ], + }); + + expect(result.entries).toHaveLength(1); + expect(result.entries[0]).toMatchObject({ + deadlineBasisAt: "2026-04-20T09:02:00.000Z", + deadlineBasisSource: "anchor", + }); + }); +}); diff --git a/apps/server/src/provider/Layers/reaperDeadlines.ts b/apps/server/src/provider/Layers/reaperDeadlines.ts new file mode 100644 index 000000000000..a620ec52e494 --- /dev/null +++ b/apps/server/src/provider/Layers/reaperDeadlines.ts @@ -0,0 +1,217 @@ +import type { + OrchestrationLatestTurn, + OrchestrationReadModel, + OrchestrationSession, + ProviderKind, + ProviderSessionRuntimeStatus, + ThreadId, + TurnId, +} from "@t3tools/contracts"; + +import type { ProviderRuntimeBindingWithMetadata } from "../Services/ProviderSessionDirectory.ts"; + +export type InactivityAnchorSource = + | "latest_turn_completed_at" + | "latest_turn_started_at" + | "latest_turn_requested_at" + | "session_last_seen_at"; + +export type DeadlineBasisSource = "anchor" | "recent_provider_send_turn"; + +export interface ReapScheduleEntry { + readonly threadId: ThreadId; + readonly provider: ProviderKind; + readonly deadlineAtMs: number; + readonly anchorAt: string; + readonly anchorSource: InactivityAnchorSource; + readonly deadlineBasisAt: string; + readonly deadlineBasisSource: DeadlineBasisSource; + readonly bindingStatus: ProviderSessionRuntimeStatus | undefined; + readonly sessionStatus: OrchestrationSession["status"] | null; + readonly activeTurnId: TurnId | null; + readonly readModelThreadPresent: boolean; + readonly sessionUpdatedAt: string | null; + readonly lastSeenAt: string; + readonly latestTurnId: TurnId | null; + readonly latestTurnState: OrchestrationLatestTurn["state"] | null; + readonly latestTurnRequestedAt: string | null; + readonly latestTurnStartedAt: string | null; + readonly latestTurnCompletedAt: string | null; +} + +export interface InvalidAnchorEntry { + readonly threadId: ThreadId; + readonly provider: ProviderKind; + readonly bindingStatus: ProviderSessionRuntimeStatus | undefined; + readonly readModelThreadPresent: boolean; + readonly sessionStatus: OrchestrationSession["status"] | null; + readonly sessionUpdatedAt: string | null; + readonly activeTurnId: TurnId | null; + readonly lastSeenAt: string; + readonly latestTurnId: TurnId | null; + readonly latestTurnState: OrchestrationLatestTurn["state"] | null; + readonly latestTurnRequestedAt: string | null; + readonly latestTurnStartedAt: string | null; + readonly latestTurnCompletedAt: string | null; + readonly inactivityAnchorAt: string; + readonly inactivityAnchorSource: InactivityAnchorSource; +} + +export interface DeriveReapEntriesResult { + readonly entries: ReadonlyArray; + readonly skippedStopped: number; + readonly skippedActiveTurn: number; + readonly invalidAnchors: ReadonlyArray; +} + +function isRecord(value: unknown): value is Record { + return value !== null && typeof value === "object" && !Array.isArray(value); +} + +function readRecentProviderSendTurnAt(runtimePayload: unknown): string | undefined { + if (!isRecord(runtimePayload) || runtimePayload.lastRuntimeEvent !== "provider.sendTurn") { + return undefined; + } + + const lastRuntimeEventAt = runtimePayload.lastRuntimeEventAt; + return typeof lastRuntimeEventAt === "string" && lastRuntimeEventAt.length > 0 + ? lastRuntimeEventAt + : undefined; +} + +function resolveInactivityAnchor(input: { + readonly latestTurn: { + readonly requestedAt: string; + readonly startedAt: string | null; + readonly completedAt: string | null; + } | null; + readonly lastSeenAt: string; +}): { readonly at: string; readonly source: InactivityAnchorSource } { + const latestTurn = input.latestTurn; + if (latestTurn?.completedAt !== null && latestTurn?.completedAt !== undefined) { + return { + at: latestTurn.completedAt, + source: "latest_turn_completed_at", + }; + } + if (latestTurn?.startedAt !== null && latestTurn?.startedAt !== undefined) { + return { + at: latestTurn.startedAt, + source: "latest_turn_started_at", + }; + } + if (latestTurn !== null && latestTurn !== undefined) { + return { + at: latestTurn.requestedAt, + source: "latest_turn_requested_at", + }; + } + return { + at: input.lastSeenAt, + source: "session_last_seen_at", + }; +} + +export function deriveReapEntries(input: { + readonly bindings: ReadonlyArray; + readonly readModel: OrchestrationReadModel; + readonly inactivityThresholdMs: number; +}): DeriveReapEntriesResult { + const threadsById = new Map( + input.readModel.threads.map((thread) => [thread.id, thread] as const), + ); + const entries: ReapScheduleEntry[] = []; + const invalidAnchors: InvalidAnchorEntry[] = []; + let skippedStopped = 0; + let skippedActiveTurn = 0; + + for (const binding of input.bindings) { + if (binding.status === "stopped") { + skippedStopped += 1; + continue; + } + + const thread = threadsById.get(binding.threadId); + const inactivityAnchor = resolveInactivityAnchor({ + latestTurn: thread?.latestTurn ?? null, + lastSeenAt: binding.lastSeenAt, + }); + const inactivityAnchorMs = Date.parse(inactivityAnchor.at); + + if (Number.isNaN(inactivityAnchorMs)) { + invalidAnchors.push({ + threadId: binding.threadId, + provider: binding.provider, + bindingStatus: binding.status, + readModelThreadPresent: thread !== undefined, + sessionStatus: thread?.session?.status ?? null, + sessionUpdatedAt: thread?.session?.updatedAt ?? null, + activeTurnId: thread?.session?.activeTurnId ?? null, + lastSeenAt: binding.lastSeenAt, + latestTurnId: thread?.latestTurn?.turnId ?? null, + latestTurnState: thread?.latestTurn?.state ?? null, + latestTurnRequestedAt: thread?.latestTurn?.requestedAt ?? null, + latestTurnStartedAt: thread?.latestTurn?.startedAt ?? null, + latestTurnCompletedAt: thread?.latestTurn?.completedAt ?? null, + inactivityAnchorAt: inactivityAnchor.at, + inactivityAnchorSource: inactivityAnchor.source, + }); + continue; + } + + if (thread?.session?.activeTurnId !== null && thread?.session?.activeTurnId !== undefined) { + skippedActiveTurn += 1; + continue; + } + + const recentSendTurnAt = readRecentProviderSendTurnAt(binding.runtimePayload); + const recentSendTurnAtMs = + recentSendTurnAt === undefined ? Number.NaN : Date.parse(recentSendTurnAt); + let deadlineBasisAt = inactivityAnchor.at; + let deadlineBasisSource: DeadlineBasisSource = "anchor"; + let deadlineBasisMs = inactivityAnchorMs; + if ( + recentSendTurnAt !== undefined && + !Number.isNaN(recentSendTurnAtMs) && + recentSendTurnAtMs > inactivityAnchorMs + ) { + deadlineBasisAt = recentSendTurnAt; + deadlineBasisSource = "recent_provider_send_turn"; + deadlineBasisMs = recentSendTurnAtMs; + } + + entries.push({ + threadId: binding.threadId, + provider: binding.provider, + deadlineAtMs: deadlineBasisMs + input.inactivityThresholdMs, + anchorAt: inactivityAnchor.at, + anchorSource: inactivityAnchor.source, + deadlineBasisAt, + deadlineBasisSource, + bindingStatus: binding.status, + sessionStatus: thread?.session?.status ?? null, + activeTurnId: thread?.session?.activeTurnId ?? null, + readModelThreadPresent: thread !== undefined, + sessionUpdatedAt: thread?.session?.updatedAt ?? null, + lastSeenAt: binding.lastSeenAt, + latestTurnId: thread?.latestTurn?.turnId ?? null, + latestTurnState: thread?.latestTurn?.state ?? null, + latestTurnRequestedAt: thread?.latestTurn?.requestedAt ?? null, + latestTurnStartedAt: thread?.latestTurn?.startedAt ?? null, + latestTurnCompletedAt: thread?.latestTurn?.completedAt ?? null, + }); + } + + entries.sort( + (left, right) => + left.deadlineAtMs - right.deadlineAtMs || + String(left.threadId).localeCompare(String(right.threadId)), + ); + + return { + entries, + skippedStopped, + skippedActiveTurn, + invalidAnchors, + }; +} diff --git a/apps/server/src/provider/Services/ProviderSessionDirectoryEvents.ts b/apps/server/src/provider/Services/ProviderSessionDirectoryEvents.ts new file mode 100644 index 000000000000..824574030fac --- /dev/null +++ b/apps/server/src/provider/Services/ProviderSessionDirectoryEvents.ts @@ -0,0 +1,13 @@ +import { type ThreadId } from "@t3tools/contracts"; +import { Context } from "effect"; +import type { Effect, Stream } from "effect"; + +export interface ProviderSessionDirectoryEventsShape { + readonly publishChanged: (threadId: ThreadId) => Effect.Effect; + readonly changes: Stream.Stream<{ readonly threadId: ThreadId }>; +} + +export class ProviderSessionDirectoryEvents extends Context.Service< + ProviderSessionDirectoryEvents, + ProviderSessionDirectoryEventsShape +>()("t3/provider/Services/ProviderSessionDirectoryEvents") {} diff --git a/apps/server/src/provider/Services/ProviderSessionReaper.ts b/apps/server/src/provider/Services/ProviderSessionReaper.ts index b13b6f7e0c7b..c73f5547b443 100644 --- a/apps/server/src/provider/Services/ProviderSessionReaper.ts +++ b/apps/server/src/provider/Services/ProviderSessionReaper.ts @@ -1,6 +1,10 @@ import { Context } from "effect"; import type { Effect, Scope } from "effect"; +export const DEFAULT_PROVIDER_SESSION_REAPER_INACTIVITY_THRESHOLD_MS = 30 * 60 * 1000; +export const DEFAULT_PROVIDER_SESSION_REAPER_FALLBACK_RECONCILE_INTERVAL_MS = 30 * 60 * 1000; +export const DEFAULT_PROVIDER_SESSION_REAPER_STOP_FAILURE_RETRY_INTERVAL_MS = 5 * 1000; + export interface ProviderSessionReaperShape { /** * Start the background provider session reaper within the provided scope. diff --git a/apps/server/src/server.test.ts b/apps/server/src/server.test.ts index 47e159d3036e..10be14638146 100644 --- a/apps/server/src/server.test.ts +++ b/apps/server/src/server.test.ts @@ -354,6 +354,8 @@ const buildAppUnderTest = (options?: { otlpMetricsUrl: undefined, otlpExportIntervalMs: 10_000, otlpServiceName: "t3-server", + providerSessionReaperInactivityThresholdMs: 30 * 60 * 1000, + providerSessionReaperFallbackReconcileIntervalMs: 30 * 60 * 1000, mode: "desktop", port: 0, host: "127.0.0.1", diff --git a/apps/server/src/server.ts b/apps/server/src/server.ts index f94bbb34b5bf..afaa70de8c7c 100644 --- a/apps/server/src/server.ts +++ b/apps/server/src/server.ts @@ -18,6 +18,7 @@ import { ServerLifecycleEventsLive } from "./serverLifecycleEvents.ts"; import { AnalyticsServiceLayerLive } from "./telemetry/Layers/AnalyticsService.ts"; import { makeEventNdjsonLogger } from "./provider/Layers/EventNdjsonLogger.ts"; import { ProviderSessionDirectoryLive } from "./provider/Layers/ProviderSessionDirectory.ts"; +import { ProviderSessionDirectoryEventsLive } from "./provider/Layers/ProviderSessionDirectoryEvents.ts"; import { ProviderSessionRuntimeRepositoryLive } from "./persistence/Layers/ProviderSessionRuntime.ts"; import { makeCodexAdapterLive } from "./provider/Layers/CodexAdapter.ts"; import { makeClaudeAdapterLive } from "./provider/Layers/ClaudeAdapter.ts"; @@ -25,7 +26,7 @@ import { makeCursorAdapterLive } from "./provider/Layers/CursorAdapter.ts"; import { makeOpenCodeAdapterLive } from "./provider/Layers/OpenCodeAdapter.ts"; import { ProviderAdapterRegistryLive } from "./provider/Layers/ProviderAdapterRegistry.ts"; import { makeProviderServiceLive } from "./provider/Layers/ProviderService.ts"; -import { ProviderSessionReaperLive } from "./provider/Layers/ProviderSessionReaper.ts"; +import { makeProviderSessionReaperLive } from "./provider/Layers/ProviderSessionReaper.ts"; import { CheckpointDiffQueryLive } from "./checkpointing/Layers/CheckpointDiffQuery.ts"; import { CheckpointStoreLive } from "./checkpointing/Layers/CheckpointStore.ts"; import { GitCoreLive } from "./git/Layers/GitCore.ts"; @@ -140,8 +141,11 @@ const CheckpointingLayerLive = Layer.empty.pipe( Layer.provideMerge(CheckpointStoreLive), ); +const ProviderSessionDirectoryEventsLayerLive = ProviderSessionDirectoryEventsLive; + const ProviderSessionDirectoryLayerLive = ProviderSessionDirectoryLive.pipe( Layer.provide(ProviderSessionRuntimeRepositoryLive), + Layer.provide(ProviderSessionDirectoryEventsLayerLive), ); const ProviderLayerLive = Layer.unwrap( @@ -219,8 +223,19 @@ const AuthLayerLive = ServerAuthLive.pipe( Layer.provide(ServerSecretStoreLive), ); -const ProviderRuntimeLayerLive = ProviderSessionReaperLive.pipe( +const ProviderSessionReaperLayerLive = Layer.unwrap( + Effect.gen(function* () { + const config = yield* ServerConfig; + return makeProviderSessionReaperLive({ + inactivityThresholdMs: config.providerSessionReaperInactivityThresholdMs, + fallbackReconcileIntervalMs: config.providerSessionReaperFallbackReconcileIntervalMs, + }); + }), +); + +const ProviderRuntimeLayerLive = ProviderSessionReaperLayerLive.pipe( Layer.provideMerge(ProviderLayerLive), + Layer.provideMerge(ProviderSessionDirectoryEventsLayerLive), Layer.provideMerge(OrchestrationLayerLive), );