From b1796d4e5440cd80fc5e2352a62d09d20615a1f4 Mon Sep 17 00:00:00 2001 From: Julius Marminge Date: Sun, 20 Sep 2026 13:57:52 -0700 Subject: [PATCH 1/6] refactor(server): share scheduling for tasks and limit recovery --- .../ThreadLaunchService.test.ts | 3 +- .../UsageLimitRecoveryWorker.ts | 15 +- .../src/orchestration-v2/runtimeLayer.ts | 3 +- .../ScheduledTaskService.schedule.test.ts | 3 + .../ScheduledTaskService.test.ts | 3 + .../scheduledTasks/ScheduledTaskService.ts | 10 +- .../scheduling/Scheduler.integration.test.ts | 157 ++++++++++++++++++ apps/server/src/scheduling/Scheduler.test.ts | 84 ++++++++++ apps/server/src/scheduling/Scheduler.ts | 63 +++++++ 9 files changed, 322 insertions(+), 19 deletions(-) create mode 100644 apps/server/src/scheduling/Scheduler.integration.test.ts create mode 100644 apps/server/src/scheduling/Scheduler.test.ts create mode 100644 apps/server/src/scheduling/Scheduler.ts diff --git a/apps/server/src/orchestration-v2/ThreadLaunchService.test.ts b/apps/server/src/orchestration-v2/ThreadLaunchService.test.ts index de2e1b1a59bd..e3eff03b6697 100644 --- a/apps/server/src/orchestration-v2/ThreadLaunchService.test.ts +++ b/apps/server/src/orchestration-v2/ThreadLaunchService.test.ts @@ -1,3 +1,4 @@ +import { layer as schedulerLayer } from "../scheduling/Scheduler.ts"; import * as WorktreeSetupTracker from "../project/WorktreeSetupTracker.ts"; import * as ProjectCloneTracker from "../project/ProjectCloneTracker.ts"; import * as TerminalManager from "../terminal/Manager.ts"; @@ -267,7 +268,7 @@ for (const target of ["new", "existing"] as const) { () => { const harness = makeHarness(); const scheduledTasks = ScheduledTasks.layer.pipe( - Layer.provide(Layer.mergeAll(harness.layer, NodeCrypto.layer)), + Layer.provide(Layer.mergeAll(harness.layer, NodeCrypto.layer, schedulerLayer)), ); return Effect.gen(function* () { const tasks = yield* ScheduledTasks.ScheduledTaskService; diff --git a/apps/server/src/orchestration-v2/UsageLimitRecoveryWorker.ts b/apps/server/src/orchestration-v2/UsageLimitRecoveryWorker.ts index 8d7a573ac023..f268b5242905 100644 --- a/apps/server/src/orchestration-v2/UsageLimitRecoveryWorker.ts +++ b/apps/server/src/orchestration-v2/UsageLimitRecoveryWorker.ts @@ -7,7 +7,7 @@ import { import * as DateTime from "effect/DateTime"; import * as Effect from "effect/Effect"; import * as Layer from "effect/Layer"; -import * as Schedule from "effect/Schedule"; +import { Scheduler } from "../scheduling/Scheduler.ts"; import * as ServerSettings from "../serverSettings.ts"; import * as ProjectionStore from "./ProjectionStore.ts"; import * as ThreadManagement from "./ThreadManagementService.ts"; @@ -104,17 +104,12 @@ const makeSweep = Effect.gen(function* () { }); }); -// The schedule is derived from persisted failures and thread recovery choices, -// so restarts need no timer restoration and disconnected clients need not run it. +// The shared scheduler derives due work from persisted failures and recovery +// choices, so restarts need no timer restoration or connected client. export const workerLive = Layer.effectDiscard( Effect.gen(function* () { const sweep = yield* makeSweep; - yield* sweep().pipe( - Effect.catchCause((cause) => - Effect.logWarning("orchestration-v2.limit-recovery.sweep-failed", { cause }), - ), - Effect.repeat(Schedule.spaced("30 seconds")), - Effect.forkScoped, - ); + const scheduler = yield* Scheduler; + yield* scheduler.register("usage-limit-recovery", sweep()); }), ); diff --git a/apps/server/src/orchestration-v2/runtimeLayer.ts b/apps/server/src/orchestration-v2/runtimeLayer.ts index 83b176495a75..39a9f99c5bb0 100644 --- a/apps/server/src/orchestration-v2/runtimeLayer.ts +++ b/apps/server/src/orchestration-v2/runtimeLayer.ts @@ -1,4 +1,5 @@ import * as UsageLimitRecoveryWorker from "./UsageLimitRecoveryWorker.ts"; +import { layer as schedulerLayer } from "../scheduling/Scheduler.ts"; import * as Layer from "effect/Layer"; import { OrchestrationEventInfrastructureLayerLive, @@ -301,4 +302,4 @@ export const OrchestrationV2ProductionLayerLive = Layer.mergeAll( ), providerContinuationWorkerProvided, agentSessionImporterProvided, -).pipe(Layer.provideMerge(OrchestrationLayerLive)); +).pipe(Layer.provide(schedulerLayer), Layer.provideMerge(OrchestrationLayerLive)); diff --git a/apps/server/src/scheduledTasks/ScheduledTaskService.schedule.test.ts b/apps/server/src/scheduledTasks/ScheduledTaskService.schedule.test.ts index 8e2d67d9fe7e..ebc18e216fd9 100644 --- a/apps/server/src/scheduledTasks/ScheduledTaskService.schedule.test.ts +++ b/apps/server/src/scheduledTasks/ScheduledTaskService.schedule.test.ts @@ -1,3 +1,4 @@ +import { layer as schedulerLayer } from "../scheduling/Scheduler.ts"; import * as NodeCrypto from "@effect/platform-node/NodeCrypto"; import { expect, it } from "@effect/vitest"; import { ScheduledTaskUpsertInput } from "@t3tools/contracts"; @@ -18,6 +19,7 @@ it.effect("rejects a stale form save after deletion while preserving explicit-id Effect.gen(function* () { const dependencies = Layer.mergeAll( NodeCrypto.layer, + schedulerLayer, Layer.mock(ThreadLaunchService)({}), Layer.mock(ThreadManagementService)({}), ); @@ -59,6 +61,7 @@ it.effect("preserves a due run when a save only pads the scheduled hour", () => const dependencies = Layer.mergeAll( NodeCrypto.layer, + schedulerLayer, Layer.mock(ThreadLaunchService)({}), Layer.mock(ThreadManagementService)({}), ); diff --git a/apps/server/src/scheduledTasks/ScheduledTaskService.test.ts b/apps/server/src/scheduledTasks/ScheduledTaskService.test.ts index 700a32a97a22..719f08e14f9d 100644 --- a/apps/server/src/scheduledTasks/ScheduledTaskService.test.ts +++ b/apps/server/src/scheduledTasks/ScheduledTaskService.test.ts @@ -1,3 +1,4 @@ +import { layer as schedulerLayer } from "../scheduling/Scheduler.ts"; import * as NodeUtil from "node:util"; import * as NodeCrypto from "@effect/platform-node/NodeCrypto"; @@ -211,6 +212,7 @@ it.effect( }), Layer.mock(ThreadManagementService.ThreadManagementService)({}), NodeCrypto.layer, + schedulerLayer, ), ), ); @@ -315,6 +317,7 @@ it.effect( }), Layer.mock(ThreadManagementService.ThreadManagementService)({}), NodeCrypto.layer, + schedulerLayer, ), ), ); diff --git a/apps/server/src/scheduledTasks/ScheduledTaskService.ts b/apps/server/src/scheduledTasks/ScheduledTaskService.ts index 107493fee0e7..54482d243ba1 100644 --- a/apps/server/src/scheduledTasks/ScheduledTaskService.ts +++ b/apps/server/src/scheduledTasks/ScheduledTaskService.ts @@ -18,7 +18,6 @@ import * as Cause from "effect/Cause"; import * as Context from "effect/Context"; import * as Crypto from "effect/Crypto"; import * as DateTime from "effect/DateTime"; -import * as Duration from "effect/Duration"; import * as Effect from "effect/Effect"; import * as Layer from "effect/Layer"; import * as Option from "effect/Option"; @@ -31,6 +30,7 @@ import * as SqlClient from "effect/unstable/sql/SqlClient"; import * as ThreadLaunchService from "../orchestration-v2/ThreadLaunchService.ts"; import * as ThreadManagementService from "../orchestration-v2/ThreadManagementService.ts"; +import { Scheduler } from "../scheduling/Scheduler.ts"; import { isMissedFixedTimeRun, isSameSchedule, nextScheduledRunAt } from "./Schedule.ts"; const decodeTask = Schema.decodeUnknownEffect(ScheduledTask); @@ -209,6 +209,7 @@ export const layer = Layer.effect( const crypto = yield* Crypto.Crypto; const threadLaunch = yield* ThreadLaunchService.ThreadLaunchService; const threadManagement = yield* ThreadManagementService.ThreadManagementService; + const scheduler = yield* Scheduler; const activeRuns = yield* Ref.make>(new Set()); // Sliding(1) coalesces the dirty-signal: every notification triggers a // full list() re-emit anyway, so a slow subscriber only ever needs the @@ -696,12 +697,7 @@ export const layer = Layer.effect( ), ); - yield* runDueTasks().pipe( - Effect.catch((cause) => Effect.logWarning("Scheduled task polling failed", { cause })), - Effect.delay(Duration.seconds(5)), - Effect.forever, - Effect.forkScoped, - ); + yield* scheduler.register("scheduled-tasks", runDueTasks()); const list: ScheduledTaskService["Service"]["list"] = () => listRows().pipe( diff --git a/apps/server/src/scheduling/Scheduler.integration.test.ts b/apps/server/src/scheduling/Scheduler.integration.test.ts new file mode 100644 index 000000000000..5809a4773a13 --- /dev/null +++ b/apps/server/src/scheduling/Scheduler.integration.test.ts @@ -0,0 +1,157 @@ +import * as NodeCrypto from "@effect/platform-node/NodeCrypto"; +import { assert, expect, it } from "@effect/vitest"; +import { + CommandId, + DEFAULT_SERVER_SETTINGS, + ProjectId, + ProviderInstanceId, + RunId, + ThreadId, + type OrchestrationV2Command, + type OrchestrationV2ThreadShell, +} from "@t3tools/contracts"; +import * as DateTime from "effect/DateTime"; +import * as Effect from "effect/Effect"; +import * as Layer from "effect/Layer"; +import * as Queue from "effect/Queue"; +import * as Ref from "effect/Ref"; +import * as SqlClient from "effect/unstable/sql/SqlClient"; +import * as TestClock from "effect/testing/TestClock"; + +import { ProjectionStoreV2 } from "../orchestration-v2/ProjectionStore.ts"; +import { ThreadLaunchService } from "../orchestration-v2/ThreadLaunchService.ts"; +import { ThreadManagementService } from "../orchestration-v2/ThreadManagementService.ts"; +import { workerLive as recoveryWorker } from "../orchestration-v2/UsageLimitRecoveryWorker.ts"; +import { SqlitePersistenceMemory } from "../persistence/Layers/Sqlite.ts"; +import * as ScheduledTasks from "../scheduledTasks/ScheduledTaskService.ts"; +import { ServerSettingsService } from "../serverSettings.ts"; +import { layer as schedulerLayer } from "./Scheduler.ts"; + +it.effect.each(["on time", "after restart"])( + "runs Scheduled Tasks and a persisted limit retry through the same scheduler %s", + (scenario) => + Effect.gen(function* () { + const now = yield* DateTime.now; + const resetAt = DateTime.formatIso(DateTime.add(now, { seconds: 60 })); + const thread: OrchestrationV2ThreadShell = { + id: ThreadId.make("thread:limited"), + projectId: ProjectId.make("project:test"), + title: "Limited", + providerInstanceId: ProviderInstanceId.make("codex"), + modelSelection: { instanceId: ProviderInstanceId.make("codex"), model: "gpt-5.4" }, + runtimeMode: "full-access", + interactionMode: "default", + createdBy: "user", + creationSource: "web", + branch: null, + worktreePath: null, + lineage: { + rootThreadId: ThreadId.make("thread:limited"), + parentThreadId: null, + relationshipToParent: null, + }, + forkedFrom: null, + activeProviderThreadId: null, + latestRunId: RunId.make("run:limited"), + latestRunCompletedAt: now, + activeRunId: null, + status: "failed", + lastErrorClass: "usage_limit", + usageLimitResetAt: resetAt, + pendingRuntimeRequest: null, + latestVisibleMessage: null, + latestUserMessageAt: now, + hasActionableProposedPlan: false, + itemCount: 1, + visibleItemCount: 1, + createdAt: now, + updatedAt: now, + archivedAt: null, + settledOverride: null, + settledAt: null, + lastVisitedAt: null, + pinnedAt: null, + deletedAt: null, + limitRecovery: { + runId: RunId.make("run:limited"), + resetAt, + autoResume: true, + requestId: CommandId.make("recovery:choice"), + }, + }; + if (scenario === "after restart") { + yield* TestClock.adjust("65 seconds"); + } + const current = yield* Ref.make(thread); + const commands = yield* Ref.make>([]); + const receipts = yield* Queue.unbounded<"task" | "retry">(); + const dependencies = Layer.mergeAll( + NodeCrypto.layer, + Layer.mock(ThreadLaunchService)({ + launch: () => + Queue.offer(receipts, "task").pipe( + Effect.andThen(Effect.die("fixture dispatch failure")), + ), + }), + Layer.mock(ThreadManagementService)({ + dispatch: (command) => + Ref.update(commands, (all) => [...all, command]).pipe( + Effect.andThen( + Ref.update(current, (shell) => ({ ...shell, status: "running" as const })), + ), + Effect.andThen(Queue.offer(receipts, "retry")), + Effect.as({ sequence: 1, storedEvents: [] }), + ), + }), + Layer.mock(ProjectionStoreV2)({ + getShellSnapshot: () => + Ref.get(current).pipe( + Effect.map((shell) => ({ + schemaVersion: 2, + snapshotSequence: 0, + threads: [shell], + archivedThreads: [], + })), + ), + }), + Layer.mock(ServerSettingsService)({ getSettings: Effect.succeed(DEFAULT_SERVER_SETTINGS) }), + ); + const workers = Layer.mergeAll(ScheduledTasks.layer, recoveryWorker).pipe( + Layer.provide(dependencies), + Layer.provide(schedulerLayer), + ); + yield* Effect.gen(function* () { + const tasks = yield* ScheduledTasks.ScheduledTaskService; + const { task } = yield* tasks.upsert({ + title: "Scheduled work", + prompt: "Run scheduled work", + enabled: true, + schedule: { type: "interval", everyMs: 60_000 }, + projectId: thread.projectId, + workspaceStrategy: { type: "root" }, + modelSelection: thread.modelSelection, + runtimeMode: thread.runtimeMode, + interactionMode: thread.interactionMode, + }); + const sql = yield* SqlClient.SqlClient; + yield* sql`UPDATE scheduled_tasks SET next_run_at = ${resetAt} WHERE task_id = ${task.id}`; + if (scenario === "on time") { + yield* TestClock.adjust("55 seconds"); + assert.deepEqual(yield* Ref.get(commands), []); + } + yield* TestClock.adjust("5 seconds"); + assert.deepEqual([yield* Queue.take(receipts), yield* Queue.take(receipts)].sort(), [ + "retry", + "task", + ]); + const delivered = yield* Ref.get(commands); + assert.equal(delivered.length, 1); + expect(delivered[0]).toMatchObject({ + type: "message.dispatch", + usageLimitContinuationOfRunId: thread.latestRunId, + usageLimitRecoveryRequestId: thread.limitRecovery!.requestId, + text: "Continue where you left off.", + }); + }).pipe(Effect.provide(workers)); + }).pipe(Effect.provide(SqlitePersistenceMemory)), +); diff --git a/apps/server/src/scheduling/Scheduler.test.ts b/apps/server/src/scheduling/Scheduler.test.ts new file mode 100644 index 000000000000..649ca5536515 --- /dev/null +++ b/apps/server/src/scheduling/Scheduler.test.ts @@ -0,0 +1,84 @@ +import { assert, it } from "@effect/vitest"; +import * as Deferred from "effect/Deferred"; +import * as Effect from "effect/Effect"; +import * as Fiber from "effect/Fiber"; +import * as Queue from "effect/Queue"; +import * as Ref from "effect/Ref"; +import * as TestClock from "effect/testing/TestClock"; + +import { Scheduler, layer } from "./Scheduler.ts"; + +it.effect("keeps other sources running after a source defects", () => + Effect.gen(function* () { + const scheduler = yield* Scheduler; + const receipts = yield* Queue.unbounded(); + yield* scheduler.register( + "broken", + Queue.offer(receipts, "broken").pipe(Effect.andThen(Effect.die("fixture failure"))), + ); + yield* scheduler.register("healthy", Queue.offer(receipts, "healthy").pipe(Effect.asVoid)); + for (let tick = 0; tick < 2; tick += 1) { + yield* TestClock.adjust("5 seconds"); + assert.deepEqual([yield* Queue.take(receipts), yield* Queue.take(receipts)].sort(), [ + "broken", + "healthy", + ]); + } + }).pipe(Effect.provide(layer)), +); + +it.effect("does not overlap slow work or hold up another source", () => + Effect.gen(function* () { + const scheduler = yield* Scheduler; + const started = yield* Deferred.make(); + const release = yield* Deferred.make(); + const runs = yield* Ref.make(0); + const receipts = yield* Queue.unbounded(); + yield* scheduler.register( + "slow", + Ref.update(runs, (n) => n + 1).pipe( + Effect.andThen(Deferred.succeed(started, undefined)), + Effect.andThen(Deferred.await(release)), + ), + ); + yield* scheduler.register("healthy", Queue.offer(receipts, undefined).pipe(Effect.asVoid)); + yield* TestClock.adjust("5 seconds"); + yield* Deferred.await(started); + yield* Queue.take(receipts); + yield* TestClock.adjust("10 seconds"); + yield* Queue.take(receipts); + yield* Queue.take(receipts); + assert.equal(yield* Ref.get(runs), 1); + yield* Deferred.succeed(release, undefined); + }).pipe(Effect.provide(layer)), +); + +it.effect("unregisters closed sources and interrupts their in-flight work", () => + Effect.gen(function* () { + const scheduler = yield* Scheduler; + const started = yield* Deferred.make(); + const stopped = yield* Deferred.make(); + const runs = yield* Ref.make(0); + const registration = yield* Effect.scoped( + scheduler + .register( + "scoped", + Ref.update(runs, (n) => n + 1).pipe( + Effect.andThen(Deferred.succeed(started, undefined)), + Effect.andThen(Effect.never), + Effect.ensuring(Deferred.succeed(stopped, undefined)), + ), + ) + .pipe(Effect.andThen(Effect.never)), + ).pipe(Effect.forkChild); + yield* TestClock.adjust("5 seconds"); + yield* Deferred.await(started); + yield* Fiber.interrupt(registration); + yield* Deferred.await(stopped); + const receipts = yield* Queue.unbounded(); + yield* scheduler.register("remaining", Queue.offer(receipts, undefined).pipe(Effect.asVoid)); + yield* TestClock.adjust("5 seconds"); + yield* Queue.take(receipts); + assert.equal(yield* Ref.get(runs), 1); + }).pipe(Effect.provide(layer)), +); diff --git a/apps/server/src/scheduling/Scheduler.ts b/apps/server/src/scheduling/Scheduler.ts new file mode 100644 index 000000000000..2f21be9602ed --- /dev/null +++ b/apps/server/src/scheduling/Scheduler.ts @@ -0,0 +1,63 @@ +import * as Cause from "effect/Cause"; +import * as Context from "effect/Context"; +import * as Effect from "effect/Effect"; +import * as Layer from "effect/Layer"; +import * as Ref from "effect/Ref"; +import * as Scope from "effect/Scope"; +import * as Semaphore from "effect/Semaphore"; + +/** Sources load due work from their durable state; registration owns its execution lifetime. */ +export class Scheduler extends Context.Service< + Scheduler, + { + readonly register: ( + name: string, + runDueWork: Effect.Effect, + ) => Effect.Effect; + } +>()("t3/scheduling/Scheduler") {} + +export const layer = Layer.effect( + Scheduler, + Effect.gen(function* () { + const sources = yield* Ref.make(new Map>()); + const register: Scheduler["Service"]["register"] = Effect.fn("Scheduler.register")(function* < + E, + R, + >(name: string, runDueWork: Effect.Effect) { + const scope = yield* Effect.scope; + const context = yield* Effect.context(); + const permit = yield* Semaphore.make(1); + const run = runDueWork.pipe( + Effect.provideContext(context), + Effect.catchCause((cause) => + Cause.hasInterruptsOnly(cause) + ? Effect.void + : Effect.logWarning("Scheduler source failed", { source: name, cause }), + ), + permit.withPermitsIfAvailable(1), + Effect.forkIn(scope), + Effect.asVoid, + ); + const id = Symbol(name); + yield* Effect.acquireRelease( + Ref.update(sources, (current) => new Map(current).set(id, run)), + () => + Ref.update(sources, (current) => { + const next = new Map(current); + next.delete(id); + return next; + }), + ); + }); + const tick = Ref.get(sources).pipe( + Effect.flatMap((current) => + Effect.forEach(current.values(), (run) => run, { discard: true }), + ), + ); + // One clock for all due-work sources. A slow source cannot block another + // source or overlap itself, and no extra missed-tick backlog is queued. + yield* Effect.sleep("5 seconds").pipe(Effect.andThen(tick), Effect.forever, Effect.forkScoped); + return Scheduler.of({ register }); + }), +); From 74072e7ae1dda9004a9307341818c603f3479678 Mon Sep 17 00:00:00 2001 From: Julius Marminge Date: Sun, 20 Sep 2026 13:59:54 -0700 Subject: [PATCH 2/6] test(server): ensure scheduler skips overlapping ticks --- apps/server/src/scheduling/Scheduler.test.ts | 3 +++ 1 file changed, 3 insertions(+) diff --git a/apps/server/src/scheduling/Scheduler.test.ts b/apps/server/src/scheduling/Scheduler.test.ts index 649ca5536515..f9958dba8d7e 100644 --- a/apps/server/src/scheduling/Scheduler.test.ts +++ b/apps/server/src/scheduling/Scheduler.test.ts @@ -50,6 +50,9 @@ it.effect("does not overlap slow work or hold up another source", () => yield* Queue.take(receipts); assert.equal(yield* Ref.get(runs), 1); yield* Deferred.succeed(release, undefined); + yield* TestClock.adjust("5 seconds"); + yield* Queue.take(receipts); + assert.equal(yield* Ref.get(runs), 2); }).pipe(Effect.provide(layer)), ); From b126b04841ab684c06d8a266605f67ebe41370ae Mon Sep 17 00:00:00 2001 From: Julius Marminge Date: Sun, 20 Sep 2026 14:01:08 -0700 Subject: [PATCH 3/6] fix(server): check due work when scheduler sources register --- apps/server/src/scheduling/Scheduler.test.ts | 7 +++---- apps/server/src/scheduling/Scheduler.ts | 1 + 2 files changed, 4 insertions(+), 4 deletions(-) diff --git a/apps/server/src/scheduling/Scheduler.test.ts b/apps/server/src/scheduling/Scheduler.test.ts index f9958dba8d7e..af80c4438807 100644 --- a/apps/server/src/scheduling/Scheduler.test.ts +++ b/apps/server/src/scheduling/Scheduler.test.ts @@ -17,8 +17,8 @@ it.effect("keeps other sources running after a source defects", () => Queue.offer(receipts, "broken").pipe(Effect.andThen(Effect.die("fixture failure"))), ); yield* scheduler.register("healthy", Queue.offer(receipts, "healthy").pipe(Effect.asVoid)); - for (let tick = 0; tick < 2; tick += 1) { - yield* TestClock.adjust("5 seconds"); + for (let tick = 0; tick < 3; tick += 1) { + if (tick > 0) yield* TestClock.adjust("5 seconds"); assert.deepEqual([yield* Queue.take(receipts), yield* Queue.take(receipts)].sort(), [ "broken", "healthy", @@ -42,7 +42,6 @@ it.effect("does not overlap slow work or hold up another source", () => ), ); yield* scheduler.register("healthy", Queue.offer(receipts, undefined).pipe(Effect.asVoid)); - yield* TestClock.adjust("5 seconds"); yield* Deferred.await(started); yield* Queue.take(receipts); yield* TestClock.adjust("10 seconds"); @@ -74,12 +73,12 @@ it.effect("unregisters closed sources and interrupts their in-flight work", () = ) .pipe(Effect.andThen(Effect.never)), ).pipe(Effect.forkChild); - yield* TestClock.adjust("5 seconds"); yield* Deferred.await(started); yield* Fiber.interrupt(registration); yield* Deferred.await(stopped); const receipts = yield* Queue.unbounded(); yield* scheduler.register("remaining", Queue.offer(receipts, undefined).pipe(Effect.asVoid)); + yield* Queue.take(receipts); yield* TestClock.adjust("5 seconds"); yield* Queue.take(receipts); assert.equal(yield* Ref.get(runs), 1); diff --git a/apps/server/src/scheduling/Scheduler.ts b/apps/server/src/scheduling/Scheduler.ts index 2f21be9602ed..4a8ad3bae44f 100644 --- a/apps/server/src/scheduling/Scheduler.ts +++ b/apps/server/src/scheduling/Scheduler.ts @@ -49,6 +49,7 @@ export const layer = Layer.effect( return next; }), ); + yield* run; }); const tick = Ref.get(sources).pipe( Effect.flatMap((current) => From 2d944361b08aa22bcaeb416cfb5715be3c900fd3 Mon Sep 17 00:00:00 2001 From: Julius Marminge Date: Sun, 20 Sep 2026 14:05:13 -0700 Subject: [PATCH 4/6] refactor(server): follow scheduler service conventions --- .../ThreadLaunchService.test.ts | 4 +- .../UsageLimitRecoveryWorker.ts | 4 +- .../src/orchestration-v2/runtimeLayer.ts | 4 +- .../ScheduledTaskService.schedule.test.ts | 6 +- .../ScheduledTaskService.test.ts | 6 +- .../scheduledTasks/ScheduledTaskService.ts | 4 +- .../scheduling/Scheduler.integration.test.ts | 4 +- apps/server/src/scheduling/Scheduler.test.ts | 14 ++-- apps/server/src/scheduling/Scheduler.ts | 83 +++++++++---------- 9 files changed, 63 insertions(+), 66 deletions(-) diff --git a/apps/server/src/orchestration-v2/ThreadLaunchService.test.ts b/apps/server/src/orchestration-v2/ThreadLaunchService.test.ts index e3eff03b6697..cb6f66d4e5b2 100644 --- a/apps/server/src/orchestration-v2/ThreadLaunchService.test.ts +++ b/apps/server/src/orchestration-v2/ThreadLaunchService.test.ts @@ -1,4 +1,4 @@ -import { layer as schedulerLayer } from "../scheduling/Scheduler.ts"; +import * as Scheduler from "../scheduling/Scheduler.ts"; import * as WorktreeSetupTracker from "../project/WorktreeSetupTracker.ts"; import * as ProjectCloneTracker from "../project/ProjectCloneTracker.ts"; import * as TerminalManager from "../terminal/Manager.ts"; @@ -268,7 +268,7 @@ for (const target of ["new", "existing"] as const) { () => { const harness = makeHarness(); const scheduledTasks = ScheduledTasks.layer.pipe( - Layer.provide(Layer.mergeAll(harness.layer, NodeCrypto.layer, schedulerLayer)), + Layer.provide(Layer.mergeAll(harness.layer, NodeCrypto.layer, Scheduler.layer)), ); return Effect.gen(function* () { const tasks = yield* ScheduledTasks.ScheduledTaskService; diff --git a/apps/server/src/orchestration-v2/UsageLimitRecoveryWorker.ts b/apps/server/src/orchestration-v2/UsageLimitRecoveryWorker.ts index f268b5242905..ae03b6fa45b0 100644 --- a/apps/server/src/orchestration-v2/UsageLimitRecoveryWorker.ts +++ b/apps/server/src/orchestration-v2/UsageLimitRecoveryWorker.ts @@ -7,7 +7,7 @@ import { import * as DateTime from "effect/DateTime"; import * as Effect from "effect/Effect"; import * as Layer from "effect/Layer"; -import { Scheduler } from "../scheduling/Scheduler.ts"; +import * as Scheduler from "../scheduling/Scheduler.ts"; import * as ServerSettings from "../serverSettings.ts"; import * as ProjectionStore from "./ProjectionStore.ts"; import * as ThreadManagement from "./ThreadManagementService.ts"; @@ -109,7 +109,7 @@ const makeSweep = Effect.gen(function* () { export const workerLive = Layer.effectDiscard( Effect.gen(function* () { const sweep = yield* makeSweep; - const scheduler = yield* Scheduler; + const scheduler = yield* Scheduler.Scheduler; yield* scheduler.register("usage-limit-recovery", sweep()); }), ); diff --git a/apps/server/src/orchestration-v2/runtimeLayer.ts b/apps/server/src/orchestration-v2/runtimeLayer.ts index 39a9f99c5bb0..b4de5bb5d633 100644 --- a/apps/server/src/orchestration-v2/runtimeLayer.ts +++ b/apps/server/src/orchestration-v2/runtimeLayer.ts @@ -1,5 +1,5 @@ import * as UsageLimitRecoveryWorker from "./UsageLimitRecoveryWorker.ts"; -import { layer as schedulerLayer } from "../scheduling/Scheduler.ts"; +import * as Scheduler from "../scheduling/Scheduler.ts"; import * as Layer from "effect/Layer"; import { OrchestrationEventInfrastructureLayerLive, @@ -302,4 +302,4 @@ export const OrchestrationV2ProductionLayerLive = Layer.mergeAll( ), providerContinuationWorkerProvided, agentSessionImporterProvided, -).pipe(Layer.provide(schedulerLayer), Layer.provideMerge(OrchestrationLayerLive)); +).pipe(Layer.provide(Scheduler.layer), Layer.provideMerge(OrchestrationLayerLive)); diff --git a/apps/server/src/scheduledTasks/ScheduledTaskService.schedule.test.ts b/apps/server/src/scheduledTasks/ScheduledTaskService.schedule.test.ts index ebc18e216fd9..ae5abcc86f7c 100644 --- a/apps/server/src/scheduledTasks/ScheduledTaskService.schedule.test.ts +++ b/apps/server/src/scheduledTasks/ScheduledTaskService.schedule.test.ts @@ -1,4 +1,4 @@ -import { layer as schedulerLayer } from "../scheduling/Scheduler.ts"; +import * as Scheduler from "../scheduling/Scheduler.ts"; import * as NodeCrypto from "@effect/platform-node/NodeCrypto"; import { expect, it } from "@effect/vitest"; import { ScheduledTaskUpsertInput } from "@t3tools/contracts"; @@ -19,7 +19,7 @@ it.effect("rejects a stale form save after deletion while preserving explicit-id Effect.gen(function* () { const dependencies = Layer.mergeAll( NodeCrypto.layer, - schedulerLayer, + Scheduler.layer, Layer.mock(ThreadLaunchService)({}), Layer.mock(ThreadManagementService)({}), ); @@ -61,7 +61,7 @@ it.effect("preserves a due run when a save only pads the scheduled hour", () => const dependencies = Layer.mergeAll( NodeCrypto.layer, - schedulerLayer, + Scheduler.layer, Layer.mock(ThreadLaunchService)({}), Layer.mock(ThreadManagementService)({}), ); diff --git a/apps/server/src/scheduledTasks/ScheduledTaskService.test.ts b/apps/server/src/scheduledTasks/ScheduledTaskService.test.ts index 719f08e14f9d..d571013a1605 100644 --- a/apps/server/src/scheduledTasks/ScheduledTaskService.test.ts +++ b/apps/server/src/scheduledTasks/ScheduledTaskService.test.ts @@ -1,4 +1,4 @@ -import { layer as schedulerLayer } from "../scheduling/Scheduler.ts"; +import * as Scheduler from "../scheduling/Scheduler.ts"; import * as NodeUtil from "node:util"; import * as NodeCrypto from "@effect/platform-node/NodeCrypto"; @@ -212,7 +212,7 @@ it.effect( }), Layer.mock(ThreadManagementService.ThreadManagementService)({}), NodeCrypto.layer, - schedulerLayer, + Scheduler.layer, ), ), ); @@ -317,7 +317,7 @@ it.effect( }), Layer.mock(ThreadManagementService.ThreadManagementService)({}), NodeCrypto.layer, - schedulerLayer, + Scheduler.layer, ), ), ); diff --git a/apps/server/src/scheduledTasks/ScheduledTaskService.ts b/apps/server/src/scheduledTasks/ScheduledTaskService.ts index 54482d243ba1..400ee5d0b111 100644 --- a/apps/server/src/scheduledTasks/ScheduledTaskService.ts +++ b/apps/server/src/scheduledTasks/ScheduledTaskService.ts @@ -30,7 +30,7 @@ import * as SqlClient from "effect/unstable/sql/SqlClient"; import * as ThreadLaunchService from "../orchestration-v2/ThreadLaunchService.ts"; import * as ThreadManagementService from "../orchestration-v2/ThreadManagementService.ts"; -import { Scheduler } from "../scheduling/Scheduler.ts"; +import * as Scheduler from "../scheduling/Scheduler.ts"; import { isMissedFixedTimeRun, isSameSchedule, nextScheduledRunAt } from "./Schedule.ts"; const decodeTask = Schema.decodeUnknownEffect(ScheduledTask); @@ -209,7 +209,7 @@ export const layer = Layer.effect( const crypto = yield* Crypto.Crypto; const threadLaunch = yield* ThreadLaunchService.ThreadLaunchService; const threadManagement = yield* ThreadManagementService.ThreadManagementService; - const scheduler = yield* Scheduler; + const scheduler = yield* Scheduler.Scheduler; const activeRuns = yield* Ref.make>(new Set()); // Sliding(1) coalesces the dirty-signal: every notification triggers a // full list() re-emit anyway, so a slow subscriber only ever needs the diff --git a/apps/server/src/scheduling/Scheduler.integration.test.ts b/apps/server/src/scheduling/Scheduler.integration.test.ts index 5809a4773a13..6f7b06cd7ad0 100644 --- a/apps/server/src/scheduling/Scheduler.integration.test.ts +++ b/apps/server/src/scheduling/Scheduler.integration.test.ts @@ -25,7 +25,7 @@ import { workerLive as recoveryWorker } from "../orchestration-v2/UsageLimitReco import { SqlitePersistenceMemory } from "../persistence/Layers/Sqlite.ts"; import * as ScheduledTasks from "../scheduledTasks/ScheduledTaskService.ts"; import { ServerSettingsService } from "../serverSettings.ts"; -import { layer as schedulerLayer } from "./Scheduler.ts"; +import * as Scheduler from "./Scheduler.ts"; it.effect.each(["on time", "after restart"])( "runs Scheduled Tasks and a persisted limit retry through the same scheduler %s", @@ -118,7 +118,7 @@ it.effect.each(["on time", "after restart"])( ); const workers = Layer.mergeAll(ScheduledTasks.layer, recoveryWorker).pipe( Layer.provide(dependencies), - Layer.provide(schedulerLayer), + Layer.provide(Scheduler.layer), ); yield* Effect.gen(function* () { const tasks = yield* ScheduledTasks.ScheduledTaskService; diff --git a/apps/server/src/scheduling/Scheduler.test.ts b/apps/server/src/scheduling/Scheduler.test.ts index af80c4438807..dd3ac3948a12 100644 --- a/apps/server/src/scheduling/Scheduler.test.ts +++ b/apps/server/src/scheduling/Scheduler.test.ts @@ -6,11 +6,11 @@ import * as Queue from "effect/Queue"; import * as Ref from "effect/Ref"; import * as TestClock from "effect/testing/TestClock"; -import { Scheduler, layer } from "./Scheduler.ts"; +import * as Scheduler from "./Scheduler.ts"; it.effect("keeps other sources running after a source defects", () => Effect.gen(function* () { - const scheduler = yield* Scheduler; + const scheduler = yield* Scheduler.Scheduler; const receipts = yield* Queue.unbounded(); yield* scheduler.register( "broken", @@ -24,12 +24,12 @@ it.effect("keeps other sources running after a source defects", () => "healthy", ]); } - }).pipe(Effect.provide(layer)), + }).pipe(Effect.provide(Scheduler.layer)), ); it.effect("does not overlap slow work or hold up another source", () => Effect.gen(function* () { - const scheduler = yield* Scheduler; + const scheduler = yield* Scheduler.Scheduler; const started = yield* Deferred.make(); const release = yield* Deferred.make(); const runs = yield* Ref.make(0); @@ -52,12 +52,12 @@ it.effect("does not overlap slow work or hold up another source", () => yield* TestClock.adjust("5 seconds"); yield* Queue.take(receipts); assert.equal(yield* Ref.get(runs), 2); - }).pipe(Effect.provide(layer)), + }).pipe(Effect.provide(Scheduler.layer)), ); it.effect("unregisters closed sources and interrupts their in-flight work", () => Effect.gen(function* () { - const scheduler = yield* Scheduler; + const scheduler = yield* Scheduler.Scheduler; const started = yield* Deferred.make(); const stopped = yield* Deferred.make(); const runs = yield* Ref.make(0); @@ -82,5 +82,5 @@ it.effect("unregisters closed sources and interrupts their in-flight work", () = yield* TestClock.adjust("5 seconds"); yield* Queue.take(receipts); assert.equal(yield* Ref.get(runs), 1); - }).pipe(Effect.provide(layer)), + }).pipe(Effect.provide(Scheduler.layer)), ); diff --git a/apps/server/src/scheduling/Scheduler.ts b/apps/server/src/scheduling/Scheduler.ts index 4a8ad3bae44f..076a18dee20a 100644 --- a/apps/server/src/scheduling/Scheduler.ts +++ b/apps/server/src/scheduling/Scheduler.ts @@ -17,48 +17,45 @@ export class Scheduler extends Context.Service< } >()("t3/scheduling/Scheduler") {} -export const layer = Layer.effect( - Scheduler, - Effect.gen(function* () { - const sources = yield* Ref.make(new Map>()); - const register: Scheduler["Service"]["register"] = Effect.fn("Scheduler.register")(function* < - E, - R, - >(name: string, runDueWork: Effect.Effect) { - const scope = yield* Effect.scope; - const context = yield* Effect.context(); - const permit = yield* Semaphore.make(1); - const run = runDueWork.pipe( - Effect.provideContext(context), - Effect.catchCause((cause) => - Cause.hasInterruptsOnly(cause) - ? Effect.void - : Effect.logWarning("Scheduler source failed", { source: name, cause }), - ), - permit.withPermitsIfAvailable(1), - Effect.forkIn(scope), - Effect.asVoid, - ); - const id = Symbol(name); - yield* Effect.acquireRelease( - Ref.update(sources, (current) => new Map(current).set(id, run)), - () => - Ref.update(sources, (current) => { - const next = new Map(current); - next.delete(id); - return next; - }), - ); - yield* run; - }); - const tick = Ref.get(sources).pipe( - Effect.flatMap((current) => - Effect.forEach(current.values(), (run) => run, { discard: true }), +export const make = Effect.gen(function* () { + const sources = yield* Ref.make(new Map>()); + const register: Scheduler["Service"]["register"] = Effect.fn("Scheduler.register")(function* < + E, + R, + >(name: string, runDueWork: Effect.Effect) { + const scope = yield* Effect.scope; + const context = yield* Effect.context(); + const permit = yield* Semaphore.make(1); + const run = runDueWork.pipe( + Effect.provideContext(context), + Effect.catchCause((cause) => + Cause.hasInterruptsOnly(cause) + ? Effect.void + : Effect.logWarning("Scheduler source failed", { source: name, cause }), ), + permit.withPermitsIfAvailable(1), + Effect.forkIn(scope), + Effect.asVoid, + ); + const id = Symbol(name); + yield* Effect.acquireRelease( + Ref.update(sources, (current) => new Map(current).set(id, run)), + () => + Ref.update(sources, (current) => { + const next = new Map(current); + next.delete(id); + return next; + }), ); - // One clock for all due-work sources. A slow source cannot block another - // source or overlap itself, and no extra missed-tick backlog is queued. - yield* Effect.sleep("5 seconds").pipe(Effect.andThen(tick), Effect.forever, Effect.forkScoped); - return Scheduler.of({ register }); - }), -); + yield* run; + }); + const tick = Ref.get(sources).pipe( + Effect.flatMap((current) => Effect.forEach(current.values(), (run) => run, { discard: true })), + ); + // One clock for all due-work sources. A slow source cannot block another + // source or overlap itself, and no extra missed-tick backlog is queued. + yield* Effect.sleep("5 seconds").pipe(Effect.andThen(tick), Effect.forever, Effect.forkScoped); + return Scheduler.of({ register }); +}); + +export const layer = Layer.effect(Scheduler, make); From 19602ca51c49206fcf3f633d964ef1ae2230bbbf Mon Sep 17 00:00:00 2001 From: Julius Marminge Date: Sun, 20 Sep 2026 16:16:56 -0700 Subject: [PATCH 5/6] perf(server): query only actionable usage-limit recovery --- .../orchestration-v2/ProjectionStore.test.ts | 105 +++++++++++ .../src/orchestration-v2/ProjectionStore.ts | 166 ++++++++++++++++++ .../ProviderTurnControlService.test.ts | 1 + .../UsageLimitRecoveryWorker.ts | 20 +-- .../scheduling/Scheduler.integration.test.ts | 9 +- 5 files changed, 284 insertions(+), 17 deletions(-) diff --git a/apps/server/src/orchestration-v2/ProjectionStore.test.ts b/apps/server/src/orchestration-v2/ProjectionStore.test.ts index c44cbc55d708..acdb7320bb03 100644 --- a/apps/server/src/orchestration-v2/ProjectionStore.test.ts +++ b/apps/server/src/orchestration-v2/ProjectionStore.test.ts @@ -1,6 +1,7 @@ import { assert, it, vi } from "@effect/vitest"; import { EventId, + CommandId, CheckpointId, CheckpointRef, CheckpointScopeId, @@ -2003,6 +2004,29 @@ it.layer(TestLayer)("ProjectionStoreV2", (it) => { ); } assert.isNull(shells.threads.find((row) => row.id === otherThreadId)!.lastErrorClass); + const candidates = yield* store.getLimitRecoveryCandidates({ + now, + autoResume: true, + snooze: false, + }); + const candidate = candidates.find((row) => row.id === threadId); + if (lastErrorClass === "usage_limit") { + assert.deepEqual(candidate, { + id: sqlShell.id, + status: sqlShell.status, + lastErrorClass: sqlShell.lastErrorClass, + usageLimitResetAt: sqlShell.usageLimitResetAt, + latestRunId: sqlShell.latestRunId, + latestRunCompletedAt: sqlShell.latestRunCompletedAt, + updatedAt: sqlShell.updatedAt, + archivedAt: sqlShell.archivedAt, + settledOverride: sqlShell.settledOverride, + pendingRuntimeRequest: null, + limitRecovery: sqlShell.limitRecovery, + snoozedUntil: sqlShell.snoozedUntil, + }); + } else assert.isUndefined(candidate); + assert.isUndefined(candidates.find((row) => row.id === otherThreadId)); }); yield* store.apply({ id: EventId.make("event:limit-shell:error"), @@ -2013,6 +2037,87 @@ it.layer(TestLayer)("ProjectionStoreV2", (it) => { }); yield* applyRun("failed"); yield* assertSummary("Plan limit reached.", "usage_limit"); + const sql = yield* SqlClient.SqlClient; + const [originalRow] = yield* sql<{ + payload_json: string; + }>`SELECT payload_json FROM orchestration_v2_projection_threads WHERE thread_id = ${threadId}`; + for (const [field, value] of [ + ["archivedAt", DateTime.formatIso(now)], + ["settledOverride", "settled"], + ]) { + yield* sql`UPDATE orchestration_v2_projection_threads + SET payload_json = json_set(payload_json, ${`$.${field}`}, ${value}) + WHERE thread_id = ${threadId}`; + assert.isUndefined( + (yield* store.getLimitRecoveryCandidates({ now, autoResume: true, snooze: false })).find( + (row) => row.id === threadId, + ), + ); + yield* sql`UPDATE orchestration_v2_projection_threads + SET payload_json = ${originalRow!.payload_json} WHERE thread_id = ${threadId}`; + } + yield* sql`UPDATE orchestration_v2_projection_threads SET deleted_at = ${DateTime.formatIso(now)} WHERE thread_id = ${threadId}`; + assert.isUndefined( + (yield* store.getLimitRecoveryCandidates({ now, autoResume: true, snooze: false })).find( + (row) => row.id === threadId, + ), + ); + yield* sql`UPDATE orchestration_v2_projection_threads SET deleted_at = NULL WHERE thread_id = ${threadId}`; + const recoveryOptions = { now, autoResume: false, snooze: false }; + assert.isUndefined( + (yield* store.getLimitRecoveryCandidates(recoveryOptions)).find( + (row) => row.id === threadId, + ), + ); + const reset = DateTime.makeUnsafe(limitItem.failure.resetAt); + const recovery = { + runId: original.id, + resetAt: limitItem.failure.resetAt, + autoResume: true, + requestId: CommandId.make("recovery:choice"), + }; + yield* sql`UPDATE orchestration_v2_projection_threads + SET payload_json = json_set(payload_json, '$.limitRecovery', json(${encodeUnknownJsonString(recovery)})) + WHERE thread_id = ${threadId}`; + // Armed future retries need no state decoding until they become due. + assert.isUndefined( + (yield* store.getLimitRecoveryCandidates({ ...recoveryOptions, autoResume: true })).find( + (row) => row.id === threadId, + ), + ); + const due = (yield* store.getLimitRecoveryCandidates({ + ...recoveryOptions, + now: reset, + })).find((row) => row.id === threadId)!; + assert.deepEqual(due.limitRecovery, recovery); + yield* sql`INSERT INTO orchestration_v2_projection_runtime_requests + (runtime_request_id, thread_id, node_id, kind, status, created_at, payload_json) + VALUES ('limit-shell:pending-request', ${threadId}, ${original.rootNodeId}, 'approval', 'pending', ${DateTime.formatIso(now)}, '{}')`; + assert.isUndefined( + (yield* store.getLimitRecoveryCandidates({ ...recoveryOptions, now: reset })).find( + (row) => row.id === threadId, + ), + ); + yield* sql`DELETE FROM orchestration_v2_projection_runtime_requests WHERE runtime_request_id = 'limit-shell:pending-request'`; + yield* sql`UPDATE orchestration_v2_projection_threads + SET payload_json = json_set(payload_json, '$.snoozedUntil', ${DateTime.formatIso(DateTime.add(reset, { minutes: 1 }))}) + WHERE thread_id = ${threadId}`; + assert.isUndefined( + (yield* store.getLimitRecoveryCandidates({ ...recoveryOptions, now: reset })).find( + (row) => row.id === threadId, + ), + ); + yield* sql`UPDATE orchestration_v2_projection_threads + SET payload_json = json_set(payload_json, '$.snoozedUntil', NULL, '$.limitRecovery.autoResume', json('false')) + WHERE thread_id = ${threadId}`; + assert.isUndefined( + (yield* store.getLimitRecoveryCandidates({ + ...recoveryOptions, + now: reset, + autoResume: true, + })).find((row) => row.id === threadId), + ); + yield* sql`UPDATE orchestration_v2_projection_threads SET payload_json = ${originalRow!.payload_json} WHERE thread_id = ${threadId}`; const session = { id: ProviderSessionId.make("session:limit-shell:shared"), driver, diff --git a/apps/server/src/orchestration-v2/ProjectionStore.ts b/apps/server/src/orchestration-v2/ProjectionStore.ts index a8e90be26d38..191ca76103c8 100644 --- a/apps/server/src/orchestration-v2/ProjectionStore.ts +++ b/apps/server/src/orchestration-v2/ProjectionStore.ts @@ -130,6 +130,23 @@ export type ProjectionRecoveryKind = | "subagent-results" | "delegated-completions"; +/** Persisted state needed for limit recovery, without transcript or fork history. */ +export type ProjectionLimitRecoveryCandidate = Pick< + OrchestrationV2ThreadShell, + | "id" + | "status" + | "lastErrorClass" + | "latestRunId" + | "usageLimitResetAt" + | "archivedAt" + | "settledOverride" + | "pendingRuntimeRequest" + | "latestRunCompletedAt" + | "updatedAt" + | "limitRecovery" + | "snoozedUntil" +>; + /** Thread activity needed by settlement, without transcript or fork history. */ export type ProjectionSettlementCandidate = Pick< OrchestrationV2ThreadShell, @@ -255,6 +272,11 @@ export interface ProjectionStoreV2Shape { readonly getThread: ( threadId: ThreadId, ) => Effect.Effect; + readonly getLimitRecoveryCandidates: (options: { + readonly now: DateTime.Utc; + readonly autoResume: boolean; + readonly snooze: boolean; + }) => Effect.Effect, ProjectionStoreV2Error>; readonly getSettlementCandidates: () => Effect.Effect< ReadonlyArray, ProjectionStoreV2Error @@ -3054,6 +3076,107 @@ export const layer: Layer.Layer = ), ); + const getLimitRecoveryCandidates = Effect.fn("ProjectionStore.getLimitRecoveryCandidates")( + function* (options: Parameters[0]) { + // Indexed latest-run and root-error lookups avoid reading run histories, + // counting transcript items, or walking fork ancestors on scheduler ticks. + const rows = yield* sql<{ + readonly payload_json: string; + readonly run_id: string; + readonly completed_at: string | null; + readonly failure_payload_json: string; + readonly last_error: string | null; + }>` + SELECT t.payload_json, r.run_id, r.completed_at, + item.payload_json AS failure_payload_json, + ( + SELECT json_extract(session.payload_json, '$.lastError') + FROM orchestration_v2_projection_provider_sessions session + INNER JOIN orchestration_v2_projection_provider_session_bindings binding + ON binding.provider_session_id = session.provider_session_id + WHERE binding.thread_id = t.thread_id + AND session.provider_instance_id = t.provider_instance_id + ORDER BY session.updated_at DESC, session.provider_session_id DESC + LIMIT 1 + ) AS last_error + FROM orchestration_v2_projection_threads t + INNER JOIN orchestration_v2_projection_runs r ON r.run_id = ( + SELECT latest.run_id FROM orchestration_v2_projection_runs latest + WHERE latest.thread_id = t.thread_id + ORDER BY latest.ordinal DESC, latest.run_id DESC LIMIT 1 + ) AND r.status = 'failed' + INNER JOIN orchestration_v2_projection_turn_items item ON item.turn_item_id = ( + SELECT error.turn_item_id FROM orchestration_v2_projection_turn_items error + WHERE error.thread_id = t.thread_id AND error.run_id = r.run_id + AND error.type = 'error' AND error.status = 'failed' + AND error.node_id IS json_extract(r.payload_json, '$.rootNodeId') + ORDER BY error.updated_at DESC, error.ordinal DESC, error.turn_item_id DESC + LIMIT 1 + ) + WHERE t.deleted_at IS NULL + AND json_extract(t.payload_json, '$.archivedAt') IS NULL + AND json_extract(t.payload_json, '$.settledOverride') IS NOT 'settled' + AND json_extract(item.payload_json, '$.failure.class') = 'usage_limit' + AND json_extract(item.payload_json, '$.failure.resetAt') IS NOT NULL + AND julianday(json_extract(item.payload_json, '$.failure.resetAt')) > julianday(COALESCE(r.completed_at, json_extract(t.payload_json, '$.updatedAt'))) + AND ( + ( + json_extract(t.payload_json, '$.limitRecovery.runId') IS r.run_id + AND json_extract(t.payload_json, '$.limitRecovery.resetAt') IS json_extract(item.payload_json, '$.failure.resetAt') + AND json_extract(t.payload_json, '$.limitRecovery.autoResume') = 1 + AND julianday(json_extract(item.payload_json, '$.failure.resetAt')) <= julianday(${DateTime.formatIso(options.now)}) + AND ( + json_extract(t.payload_json, '$.snoozedUntil') IS NULL + OR julianday(json_extract(t.payload_json, '$.snoozedUntil')) <= julianday(${DateTime.formatIso(options.now)}) + ) + ) + OR ( + ( + json_extract(t.payload_json, '$.limitRecovery.runId') IS NOT r.run_id + OR json_extract(t.payload_json, '$.limitRecovery.resetAt') IS NOT json_extract(item.payload_json, '$.failure.resetAt') + ) + AND ( + ${options.autoResume} + OR (${options.snooze} AND julianday(json_extract(item.payload_json, '$.failure.resetAt')) > julianday(${DateTime.formatIso(options.now)})) + ) + ) + ) + AND NOT EXISTS ( + SELECT 1 FROM orchestration_v2_projection_runtime_requests request + WHERE request.thread_id = t.thread_id AND request.status = 'pending' + ) + ORDER BY t.thread_id + `; + const candidates: Array = []; + for (const row of rows) { + const thread = yield* decodeThreadPayload(row.payload_json); + const item = yield* decodeTurnItemPayload(row.failure_payload_json); + const summary = threadErrorSummary( + item.type === "error" ? item.failure : null, + row.last_error, + ); + if (summary.lastErrorClass !== "usage_limit") continue; + candidates.push({ + id: thread.id, + status: "failed", + lastErrorClass: summary.lastErrorClass, + usageLimitResetAt: summary.usageLimitResetAt, + latestRunId: RunId.make(row.run_id), + latestRunCompletedAt: + row.completed_at === null ? null : DateTime.makeUnsafe(row.completed_at), + updatedAt: thread.updatedAt, + archivedAt: thread.archivedAt, + settledOverride: thread.settledOverride, + pendingRuntimeRequest: null, + limitRecovery: thread.limitRecovery ?? null, + snoozedUntil: thread.snoozedUntil ?? null, + }); + } + return candidates; + }, + Effect.mapError((cause) => new ProjectionStoreSetupError({ cause })), + ); + const getRecoveryThreadIds = Effect.fn("ProjectionStore.getRecoveryThreadIds")( function* (kind: ProjectionRecoveryKind) { const candidates = (() => { @@ -4813,6 +4936,7 @@ export const layer: Layer.Layer = getRuntimeRequest, getPlan, getProviderControlContext, + getLimitRecoveryCandidates, getRecoveryThreadIds, getUnreadableThreadIds, getThreadSnapshot, @@ -4911,6 +5035,48 @@ export const layerMemory: Layer.Layer = Layer.effect( left.id.localeCompare(right.id), ); }), + getLimitRecoveryCandidates: (options) => + Ref.get(replayState).pipe( + Effect.map((state) => + [...state.projections.values()] + .filter( + ({ thread }) => + thread.deletedAt === null && + thread.archivedAt === null && + thread.settledOverride !== "settled", + ) + .map(threadShellFromProjection) + .filter( + (thread) => + thread.status === "failed" && + thread.lastErrorClass === "usage_limit" && + thread.usageLimitResetAt !== null && + thread.pendingRuntimeRequest === null, + ) + .filter((thread) => { + const resetMs = Date.parse(thread.usageLimitResetAt!); + const nowMs = DateTime.toEpochMillis(options.now); + if ( + !Number.isFinite(resetMs) || + resetMs <= DateTime.toEpochMillis(thread.latestRunCompletedAt ?? thread.updatedAt) + ) + return false; + if ( + thread.limitRecovery?.runId !== thread.latestRunId || + thread.limitRecovery.resetAt !== thread.usageLimitResetAt + ) { + return options.autoResume || (options.snooze && resetMs > nowMs); + } + return ( + thread.limitRecovery.autoResume && + resetMs <= nowMs && + (thread.snoozedUntil == null || + DateTime.toEpochMillis(thread.snoozedUntil) <= nowMs) + ); + }) + .toSorted((left, right) => left.id.localeCompare(right.id)), + ), + ), getRecoveryThreadIds: (kind) => Effect.gen(function* () { const projections = (yield* Ref.get(replayState)).projections; diff --git a/apps/server/src/orchestration-v2/ProviderTurnControlService.test.ts b/apps/server/src/orchestration-v2/ProviderTurnControlService.test.ts index 13751fb16ac7..c9f005110e27 100644 --- a/apps/server/src/orchestration-v2/ProviderTurnControlService.test.ts +++ b/apps/server/src/orchestration-v2/ProviderTurnControlService.test.ts @@ -214,6 +214,7 @@ it.effect( ProjectionStoreV2, ProjectionStoreV2.of({ apply: () => Effect.void, + getLimitRecoveryCandidates: () => Effect.die("unused getLimitRecoveryCandidates"), getShellSnapshot: () => Effect.die("unused getShellSnapshot"), getThreadShell: () => Effect.die("unused getThreadShell"), getThread: () => Ref.get(projection).pipe(Effect.map((state) => state.thread)), diff --git a/apps/server/src/orchestration-v2/UsageLimitRecoveryWorker.ts b/apps/server/src/orchestration-v2/UsageLimitRecoveryWorker.ts index ae03b6fa45b0..c8d3e539ee28 100644 --- a/apps/server/src/orchestration-v2/UsageLimitRecoveryWorker.ts +++ b/apps/server/src/orchestration-v2/UsageLimitRecoveryWorker.ts @@ -1,9 +1,4 @@ -import { - CommandId, - MessageId, - type OrchestrationV2ThreadShell, - type OrchestrationV2Command, -} from "@t3tools/contracts"; +import { CommandId, MessageId, type OrchestrationV2Command } from "@t3tools/contracts"; import * as DateTime from "effect/DateTime"; import * as Effect from "effect/Effect"; import * as Layer from "effect/Layer"; @@ -14,7 +9,7 @@ import * as ThreadManagement from "./ThreadManagementService.ts"; /** The persisted run and reset form the identity of one recovery opportunity. */ export function limitRecoveryCommand( - thread: OrchestrationV2ThreadShell, + thread: ProjectionStore.ProjectionLimitRecoveryCandidate, autoResume: boolean, nowMs: number, snooze = false, @@ -82,9 +77,14 @@ const makeSweep = Effect.gen(function* () { const settings = yield* ServerSettings.ServerSettingsService; return Effect.fn("UsageLimitRecoveryWorker.sweep")(function* () { const preferences = yield* settings.getSettings; - const snapshot = yield* projections.getShellSnapshot(); - const nowMs = DateTime.toEpochMillis(yield* DateTime.now); - for (const thread of snapshot.threads) { + const now = yield* DateTime.now; + const candidates = yield* projections.getLimitRecoveryCandidates({ + now, + autoResume: preferences.autoResumeLimitedThreads, + snooze: preferences.snoozeLimitedThreads, + }); + const nowMs = DateTime.toEpochMillis(now); + for (const thread of candidates) { const command = limitRecoveryCommand( thread, preferences.autoResumeLimitedThreads, diff --git a/apps/server/src/scheduling/Scheduler.integration.test.ts b/apps/server/src/scheduling/Scheduler.integration.test.ts index 6f7b06cd7ad0..e63e4565e8ba 100644 --- a/apps/server/src/scheduling/Scheduler.integration.test.ts +++ b/apps/server/src/scheduling/Scheduler.integration.test.ts @@ -104,14 +104,9 @@ it.effect.each(["on time", "after restart"])( ), }), Layer.mock(ProjectionStoreV2)({ - getShellSnapshot: () => + getLimitRecoveryCandidates: () => Ref.get(current).pipe( - Effect.map((shell) => ({ - schemaVersion: 2, - snapshotSequence: 0, - threads: [shell], - archivedThreads: [], - })), + Effect.map((shell) => (shell.status === "failed" ? [shell] : [])), ), }), Layer.mock(ServerSettingsService)({ getSettings: Effect.succeed(DEFAULT_SERVER_SETTINGS) }), From 4dc66ff1cc7df33711ecdd0201d4bc3a837d7cef Mon Sep 17 00:00:00 2001 From: Julius Marminge Date: Sun, 20 Sep 2026 18:14:14 -0700 Subject: [PATCH 6/6] fix(server): keep scheduler construction private --- apps/server/src/scheduling/Scheduler.ts | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/apps/server/src/scheduling/Scheduler.ts b/apps/server/src/scheduling/Scheduler.ts index 076a18dee20a..957ba11e2ea0 100644 --- a/apps/server/src/scheduling/Scheduler.ts +++ b/apps/server/src/scheduling/Scheduler.ts @@ -17,7 +17,7 @@ export class Scheduler extends Context.Service< } >()("t3/scheduling/Scheduler") {} -export const make = Effect.gen(function* () { +const make = Effect.gen(function* () { const sources = yield* Ref.make(new Map>()); const register: Scheduler["Service"]["register"] = Effect.fn("Scheduler.register")(function* < E,