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/ThreadLaunchService.test.ts b/apps/server/src/orchestration-v2/ThreadLaunchService.test.ts index de2e1b1a59bd..cb6f66d4e5b2 100644 --- a/apps/server/src/orchestration-v2/ThreadLaunchService.test.ts +++ b/apps/server/src/orchestration-v2/ThreadLaunchService.test.ts @@ -1,3 +1,4 @@ +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"; @@ -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, 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 8d7a573ac023..c8d3e539ee28 100644 --- a/apps/server/src/orchestration-v2/UsageLimitRecoveryWorker.ts +++ b/apps/server/src/orchestration-v2/UsageLimitRecoveryWorker.ts @@ -1,20 +1,15 @@ -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"; -import * as Schedule from "effect/Schedule"; +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"; /** 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, @@ -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.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..b4de5bb5d633 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 * as Scheduler 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(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 8e2d67d9fe7e..ae5abcc86f7c 100644 --- a/apps/server/src/scheduledTasks/ScheduledTaskService.schedule.test.ts +++ b/apps/server/src/scheduledTasks/ScheduledTaskService.schedule.test.ts @@ -1,3 +1,4 @@ +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"; @@ -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, + Scheduler.layer, 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, + 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 700a32a97a22..d571013a1605 100644 --- a/apps/server/src/scheduledTasks/ScheduledTaskService.test.ts +++ b/apps/server/src/scheduledTasks/ScheduledTaskService.test.ts @@ -1,3 +1,4 @@ +import * as Scheduler 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, + Scheduler.layer, ), ), ); @@ -315,6 +317,7 @@ it.effect( }), Layer.mock(ThreadManagementService.ThreadManagementService)({}), NodeCrypto.layer, + Scheduler.layer, ), ), ); diff --git a/apps/server/src/scheduledTasks/ScheduledTaskService.ts b/apps/server/src/scheduledTasks/ScheduledTaskService.ts index 107493fee0e7..400ee5d0b111 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 * as 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.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..e63e4565e8ba --- /dev/null +++ b/apps/server/src/scheduling/Scheduler.integration.test.ts @@ -0,0 +1,152 @@ +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 * 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", + (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)({ + getLimitRecoveryCandidates: () => + Ref.get(current).pipe( + Effect.map((shell) => (shell.status === "failed" ? [shell] : [])), + ), + }), + Layer.mock(ServerSettingsService)({ getSettings: Effect.succeed(DEFAULT_SERVER_SETTINGS) }), + ); + const workers = Layer.mergeAll(ScheduledTasks.layer, recoveryWorker).pipe( + Layer.provide(dependencies), + Layer.provide(Scheduler.layer), + ); + 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..dd3ac3948a12 --- /dev/null +++ b/apps/server/src/scheduling/Scheduler.test.ts @@ -0,0 +1,86 @@ +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 * as Scheduler from "./Scheduler.ts"; + +it.effect("keeps other sources running after a source defects", () => + Effect.gen(function* () { + const scheduler = yield* Scheduler.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 < 3; tick += 1) { + if (tick > 0) yield* TestClock.adjust("5 seconds"); + assert.deepEqual([yield* Queue.take(receipts), yield* Queue.take(receipts)].sort(), [ + "broken", + "healthy", + ]); + } + }).pipe(Effect.provide(Scheduler.layer)), +); + +it.effect("does not overlap slow work or hold up another source", () => + Effect.gen(function* () { + const scheduler = yield* Scheduler.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* 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); + yield* TestClock.adjust("5 seconds"); + yield* Queue.take(receipts); + assert.equal(yield* Ref.get(runs), 2); + }).pipe(Effect.provide(Scheduler.layer)), +); + +it.effect("unregisters closed sources and interrupts their in-flight work", () => + Effect.gen(function* () { + const scheduler = yield* Scheduler.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* 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); + }).pipe(Effect.provide(Scheduler.layer)), +); diff --git a/apps/server/src/scheduling/Scheduler.ts b/apps/server/src/scheduling/Scheduler.ts new file mode 100644 index 000000000000..957ba11e2ea0 --- /dev/null +++ b/apps/server/src/scheduling/Scheduler.ts @@ -0,0 +1,61 @@ +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") {} + +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; + }), + ); + 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);