diff --git a/apps/mobile/src/features/settings/SettingsServerControlsRouteScreen.tsx b/apps/mobile/src/features/settings/SettingsServerControlsRouteScreen.tsx index 090467a8cd90..3445d8d0f3b4 100644 --- a/apps/mobile/src/features/settings/SettingsServerControlsRouteScreen.tsx +++ b/apps/mobile/src/features/settings/SettingsServerControlsRouteScreen.tsx @@ -1,7 +1,9 @@ import { ScreenScrollView as ScrollView } from "../../components/ScreenScrollView"; import { SymbolView } from "../../components/AppSymbol"; -import { AppText as Text } from "../../components/AppText"; +import { AppText as Text, AppTextInput } from "../../components/AppText"; import { + PROVIDER_SEND_TURN_MAX_INPUT_CHARS, + UsageLimitContinuationPrompt, type ResponseStreamingMode, type ServerSettings, type ServerSettingsPatch, @@ -10,6 +12,7 @@ import { type ProjectScopedServerSettingKey, } from "@t3tools/contracts"; import { useRef, useState, type ComponentProps } from "react"; +import * as Schema from "effect/Schema"; import { Alert, Platform, Pressable, View } from "react-native"; import { useSafeAreaInsets } from "react-native-safe-area-context"; @@ -33,6 +36,8 @@ import { type ScopedMobileSettingsTarget, } from "./settings-scoped-server"; +const isUsageLimitContinuationPrompt = Schema.is(UsageLimitContinuationPrompt); + type SettingsPage = "new-threads" | "source-control" | "agent-behavior" | "maintenance"; const PAGE_TITLES: Record = { @@ -365,6 +370,30 @@ function ServerSettingsDetail(props: { readonly page: SettingsPage }) { onValueChange={(value) => write({ continueThreadsAfterServerUpdate: value })} /> + + write({ continueThreadsAfterUsageLimit: value })} + /> + + {uniform("continueThreadsAfterUsageLimit") === true ? ( + target.environment.environmentId).join(",")}:${projectSelected}:${uniform("usageLimitContinuationPrompt")}`} + value={uniform("usageLimitContinuationPrompt")} + disabled={disabledFor("usageLimitContinuationPrompt")} + onSave={(usageLimitContinuationPrompt) => + write({ usageLimitContinuationPrompt }) + } + /> + ) : null} ) : null} @@ -375,6 +404,52 @@ function ServerSettingsDetail(props: { readonly page: SettingsPage }) { ); } +function UsageLimitPromptEditor(props: { + readonly value: string | null; + readonly disabled: boolean; + readonly onSave: (value: string) => void; +}) { + const [draft, setDraft] = useState(null); + const [error, setError] = useState(null); + return ( + + Continuation prompt + { + setDraft(text); + setError(null); + }} + onBlur={() => { + if (draft === null || props.disabled) return; + const prompt = draft.trim(); + if (!isUsageLimitContinuationPrompt(prompt)) { + setError("Enter a continuation prompt."); + return; + } + setDraft(prompt); + if (prompt !== props.value) props.onSave(prompt); + }} + /> + {error ? ( + + {error} + + ) : null} + + The server must remain running to continue. + + + ); +} + function ChoiceRow(props: { readonly label: string; readonly description: string; diff --git a/apps/server/integration/OrchestrationEngineHarness.integration.ts b/apps/server/integration/OrchestrationEngineHarness.integration.ts index 0df54be2f701..8998178e9ae1 100644 --- a/apps/server/integration/OrchestrationEngineHarness.integration.ts +++ b/apps/server/integration/OrchestrationEngineHarness.integration.ts @@ -58,6 +58,7 @@ import { RuntimeReceiptBusTest } from "../src/orchestration/Layers/RuntimeReceip import { OrchestrationReactorLive } from "../src/orchestration/Layers/OrchestrationReactor.ts"; import { ProviderCommandReactorLive } from "../src/orchestration/Layers/ProviderCommandReactor.ts"; import { ProviderRuntimeIngestionLive } from "../src/orchestration/Layers/ProviderRuntimeIngestion.ts"; +import { UsageLimitContinuation } from "../src/orchestration/UsageLimitContinuation.ts"; import { CheckpointReactor } from "../src/orchestration/Services/CheckpointReactor.ts"; import { ProviderRuntimeIngestionService } from "../src/orchestration/Services/ProviderRuntimeIngestion.ts"; import { @@ -268,6 +269,7 @@ export const makeOrchestrationIntegrationHarness = ( const persistenceLayer = makeSqlitePersistenceLive(dbPath); const orchestrationLayer = OrchestrationEngineLive.pipe( + Layer.provide(ServerSettingsService.layerTest()), Layer.provide(OrchestrationProjectionPipelineLive), Layer.provide(OrchestrationEventStoreLive), Layer.provide(OrchestrationCommandReceiptRepositoryLive), @@ -388,6 +390,12 @@ export const makeOrchestrationIntegrationHarness = ( }), ), Layer.provideMerge(runtimeIngestionLayer), + Layer.provideMerge( + Layer.mock(UsageLimitContinuation)({ + start: () => Effect.void, + recordFailure: () => Effect.void, + }), + ), Layer.provideMerge(providerCommandReactorLayer), Layer.provideMerge(checkpointReactorLayer), Layer.provideMerge( diff --git a/apps/server/integration/orphanedProviderSessionStartup.integration.test.ts b/apps/server/integration/orphanedProviderSessionStartup.integration.test.ts index 77173af1b377..22fbeef2b28d 100644 --- a/apps/server/integration/orphanedProviderSessionStartup.integration.test.ts +++ b/apps/server/integration/orphanedProviderSessionStartup.integration.test.ts @@ -59,6 +59,7 @@ const stoppedBindingResumeCursor = { const makePersistedRuntimeLayer = (dbPath: string) => { const persistence = makeSqlitePersistenceLive(dbPath); const orchestration = OrchestrationLayerLive.pipe( + Layer.provide(ServerSettings.layerTest()), Layer.provideMerge(RepositoryIdentityResolver.layer), Layer.provideMerge(persistence), ); diff --git a/apps/server/src/bin.test.ts b/apps/server/src/bin.test.ts index 9bc20fff84e1..f25d4bdbeb0a 100644 --- a/apps/server/src/bin.test.ts +++ b/apps/server/src/bin.test.ts @@ -39,6 +39,7 @@ import * as ServerEnvironment from "./environment/ServerEnvironment.ts"; import * as ProjectionSnapshotQuery from "./orchestration/Services/ProjectionSnapshotQuery.ts"; import * as OrchestrationEngine from "./orchestration/Services/OrchestrationEngine.ts"; import { OrchestrationLayerLive } from "./orchestration/runtimeLayer.ts"; +import { ServerSettingsService } from "./serverSettings.ts"; import { orchestrationHttpApiLayer } from "./orchestration/http.ts"; import * as ProjectCloneTracker from "./project/ProjectCloneTracker.ts"; import { layerConfig as SqlitePersistenceLayerLive } from "./persistence/Layers/Sqlite.ts"; @@ -128,6 +129,7 @@ const makeCliTestServerConfig = (baseDir: string) => const makeProjectPersistenceLayer = (config: ServerConfig.ServerConfig["Service"]) => Layer.mergeAll( OrchestrationLayerLive.pipe( + Layer.provide(ServerSettingsService.layerTest()), Layer.provideMerge(RepositoryIdentityResolver.layer), Layer.provideMerge(SqlitePersistenceLayerLive), ), diff --git a/apps/server/src/cli/project.ts b/apps/server/src/cli/project.ts index cef4dc879073..e7c27b915f84 100644 --- a/apps/server/src/cli/project.ts +++ b/apps/server/src/cli/project.ts @@ -23,6 +23,8 @@ import { FetchHttpClient, HttpClient, HttpClientError } from "effect/unstable/ht import * as HttpApiClient from "effect/unstable/httpapi/HttpApiClient"; import * as EnvironmentAuth from "../auth/EnvironmentAuth.ts"; +import * as ServerSecretStore from "../auth/ServerSecretStore.ts"; +import * as ServerSettings from "../serverSettings.ts"; import * as ServerConfig from "../config.ts"; import * as OrchestrationEngine from "../orchestration/Services/OrchestrationEngine.ts"; @@ -200,6 +202,7 @@ const projectCommandUuid = Crypto.Crypto.pipe( const ProjectCliRuntimeLive = Layer.mergeAll( WorkspacePaths.layer, OrchestrationLayerLive.pipe( + Layer.provide(ServerSettings.layer.pipe(Layer.provide(ServerSecretStore.layer))), Layer.provideMerge(RepositoryIdentityResolver.layer), Layer.provideMerge(SqlitePersistenceLayerLive), ), diff --git a/apps/server/src/orchestration/Layers/CheckpointReactor.test.ts b/apps/server/src/orchestration/Layers/CheckpointReactor.test.ts index 1d0b3da1bdbc..9e243ea84180 100644 --- a/apps/server/src/orchestration/Layers/CheckpointReactor.test.ts +++ b/apps/server/src/orchestration/Layers/CheckpointReactor.test.ts @@ -1,3 +1,4 @@ +import { ServerSettingsService } from "../../serverSettings.ts"; // @effect-diagnostics nodeBuiltinImport:off import * as NodeFS from "node:fs"; import * as NodeOS from "node:os"; @@ -323,6 +324,7 @@ describe("CheckpointReactor", () => { options?.providerName ?? ProviderDriverKind.make("codex"), ); const orchestrationLayer = OrchestrationEngineLive.pipe( + Layer.provide(ServerSettingsService.layerTest()), Layer.provide(OrchestrationProjectionSnapshotQueryLive), Layer.provide(ThreadBackgroundLiveness.layer), Layer.provide(ThreadPlanProgress.layer), diff --git a/apps/server/src/orchestration/Layers/OrchestrationEngine.test.ts b/apps/server/src/orchestration/Layers/OrchestrationEngine.test.ts index 078750967471..c764d098b44d 100644 --- a/apps/server/src/orchestration/Layers/OrchestrationEngine.test.ts +++ b/apps/server/src/orchestration/Layers/OrchestrationEngine.test.ts @@ -54,6 +54,7 @@ import { } from "../Services/ProjectionPipeline.ts"; import { ProjectionSnapshotQuery } from "../Services/ProjectionSnapshotQuery.ts"; import { ServerConfig } from "../../config.ts"; +import { ServerSettingsService } from "../../serverSettings.ts"; const asProjectId = (value: string): ProjectId => ProjectId.make(value); const asMessageId = (value: string): MessageId => MessageId.make(value); @@ -72,6 +73,7 @@ function makeOrchestrationLayer( }); return Layer.mergeAll( OrchestrationEngineLive.pipe( + Layer.provideMerge(ServerSettingsService.layerTest({ continueThreadsAfterUsageLimit: true })), Layer.provide(OrchestrationProjectionSnapshotQueryLive), Layer.provide(OrchestrationProjectionPipelineLive), ), @@ -417,6 +419,7 @@ describe("OrchestrationEngine", () => { let fullSnapshotReadCount = 0; const layer = OrchestrationEngineLive.pipe( + Layer.provideMerge(ServerSettingsService.layerTest({ continueThreadsAfterUsageLimit: true })), Layer.provide( Layer.succeed(ProjectionSnapshotQuery, { getUserInputActivity: () => Effect.die("unused"), @@ -1511,6 +1514,9 @@ describe("OrchestrationEngine", () => { const runtime = ManagedRuntime.make( OrchestrationEngineLive.pipe( + Layer.provideMerge( + ServerSettingsService.layerTest({ continueThreadsAfterUsageLimit: true }), + ), Layer.provide(OrchestrationProjectionSnapshotQueryLive), Layer.provide(ThreadBackgroundLiveness.layer), Layer.provide(ThreadPlanProgress.layer), @@ -1619,6 +1625,9 @@ describe("OrchestrationEngine", () => { const runtime = ManagedRuntime.make( OrchestrationEngineLive.pipe( + Layer.provideMerge( + ServerSettingsService.layerTest({ continueThreadsAfterUsageLimit: true }), + ), Layer.provide(OrchestrationProjectionSnapshotQueryLive), Layer.provide(ThreadBackgroundLiveness.layer), Layer.provide(ThreadPlanProgress.layer), @@ -1768,6 +1777,9 @@ describe("OrchestrationEngine", () => { const runtime = ManagedRuntime.make( OrchestrationEngineLive.pipe( + Layer.provideMerge( + ServerSettingsService.layerTest({ continueThreadsAfterUsageLimit: true }), + ), Layer.provide(OrchestrationProjectionSnapshotQueryLive), Layer.provide(ThreadBackgroundLiveness.layer), Layer.provide(ThreadPlanProgress.layer), diff --git a/apps/server/src/orchestration/Layers/OrchestrationEngine.ts b/apps/server/src/orchestration/Layers/OrchestrationEngine.ts index fb2fadde5e63..cd2eed1404af 100644 --- a/apps/server/src/orchestration/Layers/OrchestrationEngine.ts +++ b/apps/server/src/orchestration/Layers/OrchestrationEngine.ts @@ -40,7 +40,8 @@ import { type OrchestrationDispatchError, type OrchestrationProjectorDecodeError, } from "../Errors.ts"; -import { decideOrchestrationCommand } from "../decider.ts"; +import { cancelsUsageLimitContinuation, decideOrchestrationCommand } from "../decider.ts"; +import { ServerSettingsService } from "../../serverSettings.ts"; import { createEmptyReadModel, projectEvent } from "../projector.ts"; import { OrchestrationProjectionPipeline } from "../Services/ProjectionPipeline.ts"; import { ProjectionSnapshotQuery } from "../Services/ProjectionSnapshotQuery.ts"; @@ -83,6 +84,7 @@ function commandToAggregateRef(command: OrchestrationCommand): { const makeOrchestrationEngine = Effect.gen(function* () { const sql = yield* SqlClient.SqlClient; + const settingsService = yield* ServerSettingsService; const eventStore = yield* OrchestrationEventStore; const commandReceiptRepository = yield* OrchestrationCommandReceiptRepository; const projectionPipeline = yield* OrchestrationProjectionPipeline; @@ -185,6 +187,48 @@ const makeOrchestrationEngine = Effect.gen(function* () { }); } + if ( + (envelope.command.type === "thread.turn.start" || + envelope.command.type === "thread.session.set") && + envelope.command.expectedUsageLimit !== undefined + ) { + const command = envelope.command; + const canceled = yield* eventStore + .readAggregateRange({ + aggregateKind: "thread", + aggregateId: command.threadId, + fromSequenceExclusive: envelope.command.expectedUsageLimit.snapshotSequence, + toSequenceInclusive: commandReadModel.snapshotSequence, + limit: + commandReadModel.snapshotSequence - + envelope.command.expectedUsageLimit.snapshotSequence, + }) + .pipe(Stream.filter(cancelsUsageLimitContinuation), Stream.runHead); + if (Option.isSome(canceled)) { + return yield* new OrchestrationCommandInvariantError({ + commandType: command.type, + detail: `thread ${command.threadId} changed before usage-limit continuation`, + }); + } + if (command.type === "thread.turn.start") { + const settings = yield* settingsService.getSettings.pipe( + Effect.mapError( + () => + new OrchestrationCommandInvariantError({ + commandType: command.type, + detail: "Could not read usage-limit continuation settings.", + }), + ), + ); + if (!settings.continueThreadsAfterUsageLimit) { + return yield* new OrchestrationCommandInvariantError({ + commandType: command.type, + detail: "Usage-limit continuation is disabled.", + }); + } + } + } + // The decider compares the lookup inputs. Only recreation needs an // event check, since it can reset a thread to the same field values. if ( diff --git a/apps/server/src/orchestration/Layers/OrchestrationReactor.test.ts b/apps/server/src/orchestration/Layers/OrchestrationReactor.test.ts index d25c442ac8b2..0148faa8abfd 100644 --- a/apps/server/src/orchestration/Layers/OrchestrationReactor.test.ts +++ b/apps/server/src/orchestration/Layers/OrchestrationReactor.test.ts @@ -15,6 +15,7 @@ import * as ThreadPullRequestReactor from "../ThreadPullRequestReactor.ts"; import { OrchestrationReactor } from "../Services/OrchestrationReactor.ts"; import { makeOrchestrationReactor } from "./OrchestrationReactor.ts"; import * as AgentAwarenessRelay from "../../relay/AgentAwarenessRelay.ts"; +import { UsageLimitContinuation } from "../UsageLimitContinuation.ts"; import { StorageCleanup } from "../../storageCleanup.ts"; describe("OrchestrationReactor", () => { @@ -32,6 +33,7 @@ describe("OrchestrationReactor", () => { runtime = ManagedRuntime.make( Layer.effect(OrchestrationReactor, makeOrchestrationReactor).pipe( + Layer.provide(Layer.mock(UsageLimitContinuation)({ start: () => Effect.void })), Layer.provideMerge( Layer.succeed(StorageCleanup, { start: () => { diff --git a/apps/server/src/orchestration/Layers/OrchestrationReactor.ts b/apps/server/src/orchestration/Layers/OrchestrationReactor.ts index da5cc337e420..4ef466c03e76 100644 --- a/apps/server/src/orchestration/Layers/OrchestrationReactor.ts +++ b/apps/server/src/orchestration/Layers/OrchestrationReactor.ts @@ -10,6 +10,7 @@ import { ProviderCommandReactor } from "../Services/ProviderCommandReactor.ts"; import { ProviderRuntimeIngestionService } from "../Services/ProviderRuntimeIngestion.ts"; import { ThreadDeletionReactor } from "../Services/ThreadDeletionReactor.ts"; import * as ThreadSettlementReactor from "../ThreadSettlementReactor.ts"; +import * as UsageLimitContinuation from "../UsageLimitContinuation.ts"; import * as PullRequestSyncReactor from "../PullRequestSyncReactor.ts"; import * as ThreadPullRequestReactor from "../ThreadPullRequestReactor.ts"; import * as AgentAwarenessRelay from "../../relay/AgentAwarenessRelay.ts"; @@ -21,12 +22,14 @@ export const makeOrchestrationReactor = Effect.gen(function* () { const checkpointReactor = yield* CheckpointReactor; const threadDeletionReactor = yield* ThreadDeletionReactor; const threadSettlementReactor = yield* ThreadSettlementReactor.ThreadSettlementReactor; + const usageLimitContinuation = yield* UsageLimitContinuation.UsageLimitContinuation; const pullRequestSyncReactor = yield* PullRequestSyncReactor.PullRequestSyncReactor; const threadPullRequestReactor = yield* ThreadPullRequestReactor.ThreadPullRequestReactor; const agentAwarenessRelay = yield* AgentAwarenessRelay.AgentAwarenessRelay; const storageCleanup = yield* StorageCleanup.StorageCleanup; const start: OrchestrationReactorShape["start"] = Effect.fn("start")(function* () { + yield* usageLimitContinuation.start(); yield* providerRuntimeIngestion.start(); yield* providerCommandReactor.start(); yield* checkpointReactor.start(); diff --git a/apps/server/src/orchestration/Layers/ProjectionPipeline.test.ts b/apps/server/src/orchestration/Layers/ProjectionPipeline.test.ts index 179d04843c7e..80e4523d3765 100644 --- a/apps/server/src/orchestration/Layers/ProjectionPipeline.test.ts +++ b/apps/server/src/orchestration/Layers/ProjectionPipeline.test.ts @@ -1,3 +1,4 @@ +import { ServerSettingsService } from "../../serverSettings.ts"; import { ApprovalRequestId, CheckpointRef, @@ -4355,6 +4356,7 @@ it.effect("restores pending turn-start metadata across projection pipeline resta const engineLayer = it.layer( OrchestrationEngineLive.pipe( + Layer.provide(ServerSettingsService.layerTest()), Layer.provideMerge(OrchestrationProjectionSnapshotQueryLive), Layer.provide(ThreadBackgroundLiveness.layer), Layer.provide(ThreadPlanProgress.layer), diff --git a/apps/server/src/orchestration/Layers/ProjectionSnapshotQuery.test.ts b/apps/server/src/orchestration/Layers/ProjectionSnapshotQuery.test.ts index 843eb8343d84..2bb0908b5a32 100644 --- a/apps/server/src/orchestration/Layers/ProjectionSnapshotQuery.test.ts +++ b/apps/server/src/orchestration/Layers/ProjectionSnapshotQuery.test.ts @@ -1052,6 +1052,15 @@ projectionSnapshotLayer("ProjectionSnapshotQuery", (it) => { ); assert.equal(archivedShellSnapshot.threads[0]?.archivedAt, "2026-04-06T00:00:06.000Z"); assert.deepEqual(archivedShellSnapshot.threads[0]?.branchPullRequest, branchPullRequest); + assert.equal( + (yield* snapshotQuery.getThreadShellById(ThreadId.make("thread-archived")))._tag, + "None", + ); + const archivedThread = yield* snapshotQuery.getThreadShellById( + ThreadId.make("thread-archived"), + { includeArchived: true }, + ); + assert.equal(Option.getOrThrow(archivedThread).archivedAt, "2026-04-06T00:00:06.000Z"); const activeContext = yield* snapshotQuery.getThreadRuntimeContext( ThreadId.make("thread-active"), ); @@ -1086,6 +1095,12 @@ projectionSnapshotLayer("ProjectionSnapshotQuery", (it) => { }, ]); assert.deepEqual((yield* snapshotQuery.getArchivedShellSnapshot()).threads, []); + assert.equal( + (yield* snapshotQuery.getThreadShellById(ThreadId.make("thread-archived"), { + includeArchived: true, + }))._tag, + "None", + ); yield* sql` UPDATE projection_projects SET deleted_at = '2026-04-06T00:00:10.000Z', updated_at = '2026-04-06T00:00:10.000Z' diff --git a/apps/server/src/orchestration/Layers/ProjectionSnapshotQuery.ts b/apps/server/src/orchestration/Layers/ProjectionSnapshotQuery.ts index 1e7058742e25..fb72adf4fc53 100644 --- a/apps/server/src/orchestration/Layers/ProjectionSnapshotQuery.ts +++ b/apps/server/src/orchestration/Layers/ProjectionSnapshotQuery.ts @@ -1237,10 +1237,10 @@ const makeProjectionSnapshotQuery = Effect.gen(function* () { `, }); - const getActiveThreadRowById = SqlSchema.findOneOption({ - Request: ThreadIdLookupInput, + const getThreadRowById = SqlSchema.findOneOption({ + Request: Schema.Struct({ threadId: ThreadId, includeArchived: Schema.Boolean }), Result: ProjectionThreadDbRowSchema, - execute: ({ threadId }) => + execute: ({ threadId, includeArchived }) => sql` SELECT thread_id AS "threadId", @@ -1276,7 +1276,7 @@ const makeProjectionSnapshotQuery = Effect.gen(function* () { FROM projection_threads WHERE thread_id = ${threadId} AND deleted_at IS NULL - AND archived_at IS NULL + AND ${includeArchived ? sql`1 = 1` : sql`archived_at IS NULL`} LIMIT 1 `, }); @@ -1644,7 +1644,6 @@ const makeProjectionSnapshotQuery = Effect.gen(function* () { AND turns.turn_id = threads.latest_turn_id WHERE threads.thread_id = ${threadId} AND threads.deleted_at IS NULL - AND threads.archived_at IS NULL LIMIT 1 `, }); @@ -3190,10 +3189,13 @@ pending_approval_requests AS ( }); }); - const getThreadShellById: ProjectionSnapshotQueryShape["getThreadShellById"] = (threadId) => + const getThreadShellById: ProjectionSnapshotQueryShape["getThreadShellById"] = ( + threadId, + options, + ) => Effect.gen(function* () { const [threadRow, latestTurnRow, sessionRow, pullRequestRows] = yield* Effect.all([ - getActiveThreadRowById({ threadId }).pipe( + getThreadRowById({ threadId, includeArchived: options?.includeArchived === true }).pipe( Effect.mapError( toPersistenceSqlOrDecodeError( "ProjectionSnapshotQuery.getThreadShellById:getThread:query", @@ -3466,7 +3468,7 @@ pending_approval_requests AS ( latestTurnRow, sessionRow, ] = yield* Effect.all([ - getActiveThreadRowById({ threadId }).pipe( + getThreadRowById({ threadId, includeArchived: false }).pipe( Effect.mapError( toPersistenceSqlOrDecodeError( "ProjectionSnapshotQuery.getThreadDetailById:getThread:query", diff --git a/apps/server/src/orchestration/Layers/ProviderCommandReactor.test.ts b/apps/server/src/orchestration/Layers/ProviderCommandReactor.test.ts index 9bc701af0837..9c3fce6723a0 100644 --- a/apps/server/src/orchestration/Layers/ProviderCommandReactor.test.ts +++ b/apps/server/src/orchestration/Layers/ProviderCommandReactor.test.ts @@ -178,6 +178,7 @@ describe("ProviderCommandReactor", () => { readonly titleRegenerationBeforeStart?: "one" | "two"; readonly serverActivation?: Effect.Effect; readonly beforeReadySessionDispatch?: () => Effect.Effect; + readonly beforeStoppedSessionDispatch?: () => Effect.Effect; readonly beforeTurnStartDispatch?: () => Effect.Effect; readonly afterTurnStartDispatch?: () => Effect.Effect; readonly compactThreadEffect?: () => Effect.Effect; @@ -442,9 +443,11 @@ describe("ProviderCommandReactor", () => { const before = command.type === "thread.session.set" && command.session.status === "ready" ? input?.beforeReadySessionDispatch - : isReplay - ? input?.beforeTurnStartDispatch - : undefined; + : command.type === "thread.session.set" && command.session.status === "stopped" + ? input?.beforeStoppedSessionDispatch + : isReplay + ? input?.beforeTurnStartDispatch + : undefined; return (before?.() ?? Effect.void).pipe( Effect.andThen(engine.dispatch(command)), Effect.tap(() => @@ -4293,6 +4296,72 @@ describe("ProviderCommandReactor", () => { }), ); + effectIt.effect.each(["provider-stop", "session-update"] as const)( + "does not restore a waiting error cleared during %s", + (stage) => + Effect.gen(function* () { + const stopStarted = yield* Deferred.make(); + const releaseStop = yield* Deferred.make(); + const harness = yield* Effect.promise(() => + createHarness( + stage === "provider-stop" + ? { + stopSessionEffect: () => + Deferred.succeed(stopStarted, undefined).pipe( + Effect.andThen(Deferred.await(releaseStop)), + ), + } + : { + beforeStoppedSessionDispatch: () => + Deferred.succeed(stopStarted, undefined).pipe( + Effect.andThen(Deferred.await(releaseStop)), + ), + }, + ), + ); + const threadId = ThreadId.make("thread-1"); + const now = "2026-01-01T00:00:00.000Z"; + const session = { + threadId, + status: "ready" as const, + providerName: "codex", + providerInstanceId: ProviderInstanceId.make("codex_work"), + runtimeMode: "approval-required" as const, + activeTurnId: null, + lastError: "Usage limit reached. Automatic continuation scheduled.", + updatedAt: now, + }; + yield* harness.engine.dispatch({ + type: "thread.session.set", + commandId: CommandId.make("cmd-session-waiting-before-stop"), + threadId, + session, + createdAt: now, + }); + yield* harness.engine.dispatch({ + type: "thread.session.stop", + commandId: CommandId.make("cmd-stop-waiting-session"), + threadId, + createdAt: now, + }); + yield* Deferred.await(stopStarted); + yield* harness.engine.dispatch({ + type: "thread.session.set", + commandId: CommandId.make("cmd-clear-canceled-wait-during-stop"), + threadId, + session: { ...session, lastError: null }, + createdAt: now, + }); + yield* Deferred.succeed(releaseStop, undefined); + yield* Effect.promise(() => harness.drain()); + const thread = yield* harness.snapshotQuery + .getThreadShellById(threadId) + .pipe(Effect.map(Option.getOrThrow)); + expect(thread.session?.status).toBe("stopped"); + expect(thread.session?.lastError).toBeNull(); + }), + ); + effectIt.effect("stops a ready provider session after automatic settlement", () => Effect.gen(function* () { const sessionStopped = yield* Deferred.make(); diff --git a/apps/server/src/orchestration/Layers/ProviderCommandReactor.ts b/apps/server/src/orchestration/Layers/ProviderCommandReactor.ts index 6a42b7c67ae1..90fb27b101f2 100644 --- a/apps/server/src/orchestration/Layers/ProviderCommandReactor.ts +++ b/apps/server/src/orchestration/Layers/ProviderCommandReactor.ts @@ -390,6 +390,7 @@ const make = Effect.gen(function* () { const setThreadSession = (input: { readonly threadId: ThreadId; readonly session: OrchestrationSession; + readonly expectedSession?: OrchestrationSession; readonly createdAt: string; }) => serverCommandId("provider-session-set").pipe( @@ -399,6 +400,7 @@ const make = Effect.gen(function* () { commandId, threadId: input.threadId, session: input.session, + ...(input.expectedSession ? { expectedSession: input.expectedSession } : {}), createdAt: input.createdAt, }), ), @@ -1743,23 +1745,34 @@ const make = Effect.gen(function* () { ), ); }, - onSuccess: () => - setThreadSession({ - threadId: thread.id, - session: { + onSuccess: Effect.fnUntraced( + function* () { + const currentThread = yield* resolveThreadShell(thread.id); + if (!currentThread) return; + const session = currentThread.session; + yield* setThreadSession({ threadId: thread.id, - status: "stopped", - providerName: thread.session?.providerName ?? null, - ...(thread.session?.providerInstanceId !== undefined - ? { providerInstanceId: thread.session.providerInstanceId } - : {}), - runtimeMode: thread.session?.runtimeMode ?? DEFAULT_RUNTIME_MODE, - activeTurnId: null, - lastError: thread.session?.lastError ?? null, - updatedAt: now, - }, - createdAt: now, + ...(session ? { expectedSession: session } : {}), + session: { + threadId: thread.id, + status: "stopped", + providerName: session?.providerName ?? null, + ...(session?.providerInstanceId !== undefined + ? { providerInstanceId: session.providerInstanceId } + : {}), + runtimeMode: session?.runtimeMode ?? DEFAULT_RUNTIME_MODE, + activeTurnId: null, + lastError: session?.lastError ?? null, + updatedAt: now, + }, + createdAt: now, + }); + }, + Effect.retry({ + times: 1, + while: (error) => error._tag === "OrchestrationCommandInvariantError", }), + ), }), Effect.ensuring(clearStopping), ); diff --git a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts index d61739f72c21..c133e6620b0f 100644 --- a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts +++ b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts @@ -68,6 +68,7 @@ import { ProviderRuntimeIngestionService } from "../Services/ProviderRuntimeInge import { ProjectionSnapshotQuery } from "../Services/ProjectionSnapshotQuery.ts"; import { ServerConfig } from "../../config.ts"; import { ServerSettingsService } from "../../serverSettings.ts"; +import { UsageLimitContinuation } from "../UsageLimitContinuation.ts"; import * as NodeServices from "@effect/platform-node/NodeServices"; import { makeSqlStatementCounter } from "../../../integration/SqlStatementCounter.integration.ts"; @@ -270,6 +271,7 @@ describe("ProviderRuntimeIngestion", () => { serverSettings?: Partial; threadTitle?: string; workspaceSubdirectory?: string; + onUsageLimitFailure?: UsageLimitContinuation["Service"]["recordFailure"]; isGitRepository?: CheckpointStore.CheckpointStore["Service"]["isGitRepository"]; }) { const repositoryRoot = makeTempDir("t3-provider-project-"); @@ -321,6 +323,11 @@ describe("ProviderRuntimeIngestion", () => { sleep: (duration) => realClock.sleep(duration), }; const layer = ProviderRuntimeIngestionLive.pipe( + Layer.provide( + Layer.mock(UsageLimitContinuation)({ + recordFailure: options?.onUsageLimitFailure ?? (() => Effect.void), + }), + ), Layer.provide(Layer.succeed(Clock.Clock, shiftedClock)), Layer.provideMerge(orchestrationLayer), Layer.provideMerge(ingestionProjectionSnapshotLayer), @@ -440,6 +447,42 @@ describe("ProviderRuntimeIngestion", () => { }; } + it("preserves the failure boundary from before lifecycle projection", async () => { + const failureSequences: number[] = []; + const harness = await createHarness({ + onUsageLimitFailure: (_event, sequence) => + Effect.sync(() => { + failureSequences.push(sequence); + }), + }); + const base = { + provider: ProviderDriverKind.make("codex"), + providerInstanceId: ProviderInstanceId.make("codex"), + threadId: asThreadId("thread-1"), + turnId: asTurnId("turn-1"), + createdAt: "2026-01-01T00:00:00.000Z", + }; + await harness.emitAndDrain([ + { ...base, type: "turn.started", eventId: asEventId("quota-started") }, + ]); + const beforeFailure = (await harness.readModel()).snapshotSequence; + await harness.emitAndDrain([ + { + ...base, + type: "turn.completed", + eventId: asEventId("quota-failed"), + payload: { + state: "failed", + errorMessage: "Codex usage limit reached.", + usageLimit: {}, + }, + }, + ]); + + expect(failureSequences).toEqual([beforeFailure]); + expect((await harness.readModel()).snapshotSequence).toBeGreaterThan(beforeFailure); + }); + it("maps turn started/completed events into thread session updates", async () => { const harness = await createHarness(); const now = "2026-01-01T00:00:00.000Z"; diff --git a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts index 0db70e491235..036d2f5d5cf2 100644 --- a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts +++ b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts @@ -54,6 +54,7 @@ import { import { projectActivityPayload } from "../ActivityPayloadProjection.ts"; import { forkParked } from "../../serverActivation.ts"; import { ServerSettingsService } from "../../serverSettings.ts"; +import { UsageLimitContinuation } from "../UsageLimitContinuation.ts"; import { resolveProjectSettings } from "@t3tools/shared/projectSettings"; import { canReplaceThreadTitle } from "../threadTitles.ts"; @@ -133,6 +134,7 @@ type RuntimeIngestionInput = | { source: "runtime"; event: ProviderRuntimeEvent; + snapshotSequence: number; } | { source: "domain"; @@ -1028,6 +1030,7 @@ const make = Effect.gen(function* () { const projectionTurnRepository = yield* ProjectionTurnRepository; const projectionThreadActivityRepository = yield* ProjectionThreadActivityRepository; const serverSettingsService = yield* ServerSettingsService; + const usageLimitContinuation = yield* UsageLimitContinuation; const checkpointStore = yield* CheckpointStore.CheckpointStore; const providerCommandId = (event: ProviderRuntimeEvent, tag: string) => crypto.randomUUIDv4.pipe( @@ -1755,7 +1758,7 @@ const make = Effect.gen(function* () { }, ); - const processRuntimeEvent = (event: ProviderRuntimeEvent) => + const processRuntimeEvent = (event: ProviderRuntimeEvent, snapshotSequence: number) => Effect.gen(function* () { if ( event.type === "content.delta" && @@ -2592,6 +2595,13 @@ const make = Effect.gen(function* () { ), ), ).pipe(Effect.asVoid); + if ( + event.type === "turn.completed" && + event.payload.usageLimit && + shouldApplyThreadLifecycle + ) { + yield* usageLimitContinuation.recordFailure(event, snapshotSequence); + } }); const processDomainEvent = (_event: TurnStartRequestedDomainEvent) => Effect.void; @@ -2636,7 +2646,7 @@ const make = Effect.gen(function* () { const processInput = (input: RuntimeIngestionInput) => { switch (input.source) { case "runtime": - return processRuntimeEvent(input.event); + return processRuntimeEvent(input.event, input.snapshotSequence); case "domain": return processDomainEvent(input.event); case "diff": @@ -2689,7 +2699,12 @@ const make = Effect.gen(function* () { Stream.runForEach(providerService.streamEvents, (event) => event.type === "turn.diff.updated" ? diffWorker.enqueue(event) - : worker.enqueue({ source: "runtime", event }), + : // Include cancellations during queueing and failure projection in the recovery guard. + orchestrationEngine.latestSequence.pipe( + Effect.flatMap((snapshotSequence) => + worker.enqueue({ source: "runtime", event, snapshotSequence }), + ), + ), ), ); yield* forkParked( diff --git a/apps/server/src/orchestration/Services/ProjectionSnapshotQuery.ts b/apps/server/src/orchestration/Services/ProjectionSnapshotQuery.ts index eac3ede9c1ee..96d4a13af294 100644 --- a/apps/server/src/orchestration/Services/ProjectionSnapshotQuery.ts +++ b/apps/server/src/orchestration/Services/ProjectionSnapshotQuery.ts @@ -225,10 +225,12 @@ export interface ProjectionSnapshotQueryShape { ) => Effect.Effect, ProjectionRepositoryError>; /** - * Read a single active thread shell row by id. + * Read one thread shell by id, optionally including archived threads for cleanup. + * Deleted threads are always excluded. */ readonly getThreadShellById: ( threadId: ThreadId, + options?: { readonly includeArchived?: boolean }, ) => Effect.Effect, ProjectionRepositoryError>; /** Read the active thread and session facts used to ingest provider events. */ diff --git a/apps/server/src/orchestration/UsageLimitContinuation.test.ts b/apps/server/src/orchestration/UsageLimitContinuation.test.ts new file mode 100644 index 000000000000..1bf5fc3c5b66 --- /dev/null +++ b/apps/server/src/orchestration/UsageLimitContinuation.test.ts @@ -0,0 +1,646 @@ +import * as NodeServices from "@effect/platform-node/NodeServices"; +import { expect, it } from "@effect/vitest"; +import { + CommandId, + EventId, + MessageId, + ProjectId, + ProviderDriverKind, + ProviderInstanceId, + ThreadId, + TurnId, + type ProviderRuntimeTurnCompletedEvent, + type ServerProvider, + type ServerProviderUsageLimits, +} from "@t3tools/contracts"; +import * as DateTime from "effect/DateTime"; +import * as Effect from "effect/Effect"; +import * as Exit from "effect/Exit"; +import * as Fiber from "effect/Fiber"; +import * as Layer from "effect/Layer"; +import * as Option from "effect/Option"; +import * as Queue from "effect/Queue"; +import * as Scope from "effect/Scope"; +import { TestClock } from "effect/testing"; +import * as SqlClient from "effect/unstable/sql/SqlClient"; + +import * as ServerSecretStore from "../auth/ServerSecretStore.ts"; +import { ServerConfig } from "../config.ts"; +import { OrchestrationCommandReceiptRepositoryLive } from "../persistence/Layers/OrchestrationCommandReceipts.ts"; +import { OrchestrationEventStoreLive } from "../persistence/Layers/OrchestrationEventStore.ts"; +import { SqlitePersistenceMemory } from "../persistence/Layers/Sqlite.ts"; +import * as ProviderSessionRuntime from "../persistence/ProviderSessionRuntime.ts"; +import * as RepositoryIdentityResolver from "../project/RepositoryIdentityResolver.ts"; +import { ProviderSessionDirectoryLive } from "../provider/Layers/ProviderSessionDirectory.ts"; +import { ProviderRegistry } from "../provider/Services/ProviderRegistry.ts"; +import { ProviderSessionDirectory } from "../provider/Services/ProviderSessionDirectory.ts"; +import { readPendingUsageLimitContinuation } from "../provider/usageLimitContinuation.ts"; +import * as ServerSettings from "../serverSettings.ts"; +import { OrchestrationEngineLive } from "./Layers/OrchestrationEngine.ts"; +import { OrchestrationProjectionPipelineLive } from "./Layers/ProjectionPipeline.ts"; +import { OrchestrationProjectionSnapshotQueryLive } from "./Layers/ProjectionSnapshotQuery.ts"; +import { OrchestrationEngineService } from "./Services/OrchestrationEngine.ts"; +import { ProjectionSnapshotQuery } from "./Services/ProjectionSnapshotQuery.ts"; +import * as ThreadBackgroundLiveness from "./ThreadBackgroundLiveness.ts"; +import * as ThreadPlanProgress from "./ThreadPlanProgress.ts"; +import * as UsageLimitContinuation from "./UsageLimitContinuation.ts"; + +const NOW = "2026-09-18T00:00:00.000Z"; +const RESET = "2026-09-18T01:00:00.000Z"; +const instanceId = ProviderInstanceId.make("codex"); +const firstId = ThreadId.make("first"); +const modelSelection = { instanceId, model: "gpt-5.4" }; +const errorMessage = "Usage limit reached"; + +const testLayer = Layer.mergeAll( + OrchestrationEngineLive.pipe( + Layer.provide(OrchestrationProjectionSnapshotQueryLive), + Layer.provide(OrchestrationProjectionPipelineLive), + ), + OrchestrationProjectionSnapshotQueryLive, + ProviderSessionDirectoryLive.pipe(Layer.provide(ProviderSessionRuntime.layer)), +).pipe( + Layer.provideMerge(ServerSettings.layer.pipe(Layer.provide(ServerSecretStore.layer))), + Layer.provide(ThreadBackgroundLiveness.layer), + Layer.provide(ThreadPlanProgress.layer), + Layer.provide(OrchestrationEventStoreLive), + Layer.provide(OrchestrationCommandReceiptRepositoryLive), + Layer.provide(RepositoryIdentityResolver.layer), + Layer.provideMerge(SqlitePersistenceMemory), + Layer.provide(ServerConfig.layerTest(process.cwd(), { prefix: "t3code-usage-limit-test-" })), + Layer.provideMerge(NodeServices.layer), +); + +const available = (checkedAt: string): ServerProviderUsageLimits => ({ + checkedAt, + windows: [{ id: "session", kind: "session", label: "Session", usedPercent: 10 }], +}); +const provider = (usageLimits: ServerProviderUsageLimits): ServerProvider => ({ + instanceId, + driver: ProviderDriverKind.make("codex"), + enabled: true, + installed: true, + version: null, + status: "ready", + auth: { status: "authenticated" }, + checkedAt: usageLimits.checkedAt, + models: [], + slashCommands: [], + skills: [], + usageLimits, +}); + +const makeHarness = Effect.fn("makeUsageLimitHarness")(function* () { + yield* TestClock.setTime(Date.parse(NOW)); + const engine = yield* OrchestrationEngineService; + const snapshots = yield* ProjectionSnapshotQuery; + const directory = yield* ProviderSessionDirectory; + const settings = yield* ServerSettings.ServerSettingsService; + const probes = yield* Queue.unbounded(); + const probeFibers = yield* Queue.unbounded>(); + const responses = yield* Queue.unbounded>(); + let sequence = 0; + const commandId = () => CommandId.make(`command-${++sequence}`); + const now = DateTime.now.pipe(Effect.map(DateTime.formatIso)); + const read = (id = firstId) => + snapshots + .getSnapshot() + .pipe(Effect.map((snapshot) => snapshot.threads.find((thread) => thread.id === id)!)); + const pending = (id = firstId) => + directory + .getBinding(id) + .pipe( + Effect.map((binding) => + readPendingUsageLimitContinuation(Option.getOrThrow(binding).runtimePayload), + ), + ); + const projectId = ProjectId.make("project"); + yield* engine.dispatch({ + type: "project.create", + commandId: commandId(), + projectId, + title: "Project", + workspaceRoot: "/tmp/usage-limit-project", + createdAt: NOW, + }); + const create = Effect.fnUntraced(function* (id = firstId, account = instanceId) { + yield* engine.dispatch({ + type: "thread.create", + commandId: commandId(), + threadId: id, + projectId, + title: "Thread", + modelSelection: { ...modelSelection, instanceId: account }, + runtimeMode: "full-access", + interactionMode: "default", + branch: null, + worktreePath: null, + createdAt: NOW, + }); + yield* directory.upsert({ + threadId: id, + provider: ProviderDriverKind.make("codex"), + providerInstanceId: account, + resumeCursor: { threadId: `provider-${id}` }, + }); + }); + const fail = Effect.fnUntraced(function* (id = firstId, retryAt?: string) { + const thread = yield* read(id); + const turnId = TurnId.make(`turn-${++sequence}`); + const createdAt = yield* now; + const session = { + threadId: id, + providerName: "codex", + providerInstanceId: thread.modelSelection.instanceId, + runtimeMode: "full-access" as const, + updatedAt: createdAt, + }; + yield* engine.dispatch({ + type: "thread.session.set", + commandId: commandId(), + threadId: id, + createdAt, + session: { ...session, status: "running", activeTurnId: turnId, lastError: null }, + }); + const snapshotSequence = yield* engine.latestSequence; + yield* engine.dispatch({ + type: "thread.session.set", + commandId: commandId(), + threadId: id, + createdAt, + session: { ...session, status: "error", activeTurnId: null, lastError: errorMessage }, + }); + const event: ProviderRuntimeTurnCompletedEvent = { + type: "turn.completed", + eventId: EventId.make(`failure-${turnId}`), + provider: ProviderDriverKind.make("codex"), + providerInstanceId: session.providerInstanceId, + threadId: id, + turnId, + createdAt, + payload: { state: "failed", errorMessage, usageLimit: retryAt ? { retryAt } : {} }, + }; + return { event, snapshotSequence }; + }); + const start = Effect.gen(function* () { + const scope = yield* Scope.make(); + yield* Effect.addFinalizer(() => Scope.close(scope, Exit.void)); + const service = yield* UsageLimitContinuation.make.pipe( + Effect.provide( + Layer.mock(ProviderRegistry)({ + refreshInstance: (id) => + Effect.withFiber((fiber) => + Queue.offer(probeFibers, fiber).pipe( + Effect.andThen(Queue.offer(probes, id)), + Effect.andThen(Queue.take(responses)), + ), + ), + }), + ), + Scope.provide(scope), + ); + yield* service.start().pipe(Scope.provide(scope)); + yield* TestClock.adjust(0); + yield* service.drain; + return { ...service, stop: Scope.close(scope, Exit.void) }; + }); + yield* create(); + yield* settings.updateSettings({ + continueThreadsAfterUsageLimit: true, + usageLimitContinuationPrompt: "Carry on with the remaining work", + }); + return { + engine, + settings, + directory, + commandId, + read, + pending, + now, + create, + fail, + start, + probes, + probeFibers, + responses, + }; +}); + +type Harness = Effect.Success>; +type Service = UsageLimitContinuation.UsageLimitContinuation["Service"]; +const record = Effect.fnUntraced(function* ( + h: Harness, + service: Service, + id = firstId, + retryAt?: string, +) { + const failure = yield* h.fail(id, retryAt); + yield* service.recordFailure(failure.event, failure.snapshotSequence); + yield* service.drain; + return failure; +}); +const respond = Effect.fnUntraced(function* ( + h: Harness, + service: Service, + limits?: ReadonlyArray, +) { + const probeFiber = yield* Queue.take(h.probeFibers); + yield* Queue.offer(h.responses, limits ?? [provider(available(yield* h.now))]); + expect(Exit.isSuccess(yield* Fiber.await(probeFiber))).toBe(true); + yield* service.drain; +}); + +it.effect("waits for reset, coalesces account probes, and appends the current prompt once", () => + Effect.gen(function* () { + const h = yield* makeHarness(); + const secondId = ThreadId.make("second"); + yield* h.create(secondId); + const service = yield* h.start; + const failure = yield* record(h, service, firstId, RESET); + yield* service.recordFailure(failure.event, failure.snapshotSequence); + yield* record(h, service, secondId, RESET); + yield* h.settings.updateSettings({ usageLimitContinuationPrompt: "Finish the remaining work" }); + yield* TestClock.adjust("1 hour"); + expect(yield* Queue.size(h.probes)).toBe(0); + expect((yield* h.read()).messages).toEqual([]); + yield* TestClock.adjust("1 second"); + expect(yield* Queue.take(h.probes)).toBe(instanceId); + yield* h.engine.dispatch({ + type: "thread.meta.update", + commandId: h.commandId(), + threadId: secondId, + title: "Renamed while waiting", + }); + yield* h.engine.dispatch({ + type: "thread.activity.append", + commandId: h.commandId(), + threadId: firstId, + activity: { + id: EventId.make("checkpoint"), + kind: "checkpoint.completed", + tone: "info", + summary: "Checkpoint saved", + payload: null, + createdAt: RESET, + turnId: TurnId.make(failure.event.turnId!), + }, + createdAt: RESET, + }); + yield* respond(h, service); + for (const id of [firstId, secondId]) { + expect((yield* h.read(id)).messages).toMatchObject([ + { role: "user", text: "Finish the remaining work", attachments: [] }, + ]); + expect(yield* h.pending(id)).toBeUndefined(); + } + yield* service.recordFailure(failure.event, failure.snapshotSequence); + yield* service.drain; + yield* TestClock.adjust("5 minutes"); + expect((yield* h.read()).messages).toHaveLength(1); + expect(yield* Queue.size(h.probes)).toBe(0); + }).pipe(Effect.provide(testLayer)), +); + +it.effect("does not opt an old failure in when the setting is enabled later", () => + Effect.gen(function* () { + const h = yield* makeHarness(); + yield* h.settings.updateSettings({ continueThreadsAfterUsageLimit: false }); + const service = yield* h.start; + const failure = yield* record(h, service); + yield* h.settings.updateSettings({ continueThreadsAfterUsageLimit: true }); + yield* service.recordFailure(failure.event, failure.snapshotSequence); + yield* service.drain; + yield* TestClock.adjust("2 hours"); + expect(yield* h.pending()).toBeUndefined(); + expect((yield* h.read()).messages).toEqual([]); + expect(yield* Queue.size(h.probes)).toBe(0); + }).pipe(Effect.provide(testLayer)), +); + +const appendThreadActivity = Effect.fnUntraced(function* (h: Harness) { + for (let index = 0; index < 1_000; index++) { + yield* h.engine.dispatch({ + type: "thread.meta.update", + commandId: h.commandId(), + threadId: firstId, + title: `Thread ${index}`, + }); + } +}); + +it.effect("retries a failure after its pending continuation could not be persisted", () => + Effect.gen(function* () { + const h = yield* makeHarness(); + const sql = yield* SqlClient.SqlClient; + const service = yield* h.start; + yield* sql`CREATE TEMP TRIGGER reject_usage_wait + BEFORE UPDATE OF runtime_payload_json ON provider_session_runtime + BEGIN SELECT RAISE(FAIL, 'temporary write failure'); END`; + const failure = yield* record(h, service, firstId, RESET); + expect(yield* h.pending()).toBeUndefined(); + yield* sql`DROP TRIGGER reject_usage_wait`; + yield* service.recordFailure(failure.event, failure.snapshotSequence); + yield* service.drain; + expect((yield* h.pending())?.failedTurnId).toBe(failure.event.turnId); + yield* TestClock.adjust("3601 seconds"); + yield* Queue.take(h.probes); + yield* respond(h, service); + expect((yield* h.read()).messages).toHaveLength(1); + }).pipe(Effect.provide(testLayer)), +); + +it.effect("clears an older waiting banner when a rescheduled banner write fails", () => + Effect.gen(function* () { + const h = yield* makeHarness(); + const sql = yield* SqlClient.SqlClient; + const service = yield* h.start; + yield* record(h, service); + const waitingError = (yield* h.read()).session?.lastError; + yield* TestClock.adjust(0); + yield* Queue.take(h.probes); + yield* sql`CREATE TEMP TRIGGER reject_waiting_banner + BEFORE UPDATE OF last_error ON projection_thread_sessions + BEGIN SELECT RAISE(FAIL, 'temporary banner write failure'); END`; + yield* respond(h, service, [ + provider({ checkedAt: NOW, windows: [], unavailable: { reason: "probeFailed" } }), + ]); + expect((yield* h.pending())?.nextCheckAt).toBe("2026-09-18T00:05:00.000Z"); + expect((yield* h.read()).session?.lastError).toBe(waitingError); + yield* sql`DROP TRIGGER reject_waiting_banner`; + yield* h.settings.updateSettings({ continueThreadsAfterUsageLimit: false }); + yield* TestClock.adjust(0); + yield* service.drain; + expect(yield* h.pending()).toBeUndefined(); + expect((yield* h.read()).session?.lastError).toBe(errorMessage); + }).pipe(Effect.provide(testLayer)), +); + +it.effect("honors a stop after more than 1,000 events before the worker records the failure", () => + Effect.gen(function* () { + const h = yield* makeHarness(); + const service = yield* h.start; + const failure = yield* h.fail(firstId, RESET); + yield* appendThreadActivity(h); + yield* h.engine.dispatch({ + type: "thread.turn.interrupt", + commandId: h.commandId(), + threadId: firstId, + createdAt: NOW, + }); + yield* service.recordFailure(failure.event, failure.snapshotSequence); + yield* service.drain; + yield* TestClock.adjust("2 hours"); + expect(yield* h.pending()).toBeUndefined(); + expect((yield* h.read()).messages).toEqual([]); + expect(yield* Queue.size(h.probes)).toBe(0); + }).pipe(Effect.provide(testLayer)), +); + +it.effect.each([ + "thread.turn.interrupt", + "thread.session.stop", + "thread.archive", + "disable", + "provider-change", +] as const)( + "%s cancels an outstanding probe and cannot be undone by replaying the failure", + (action) => + Effect.gen(function* () { + const h = yield* makeHarness(); + const service = yield* h.start; + const failure = yield* record(h, service); + expect((yield* h.read()).session?.lastError).toContain("Automatic continuation"); + yield* TestClock.adjust(0); + yield* Queue.take(h.probes); + if (action === "disable") { + yield* h.settings.updateSettings({ continueThreadsAfterUsageLimit: false }); + } else if (action === "provider-change") { + yield* h.engine.dispatch({ + type: "thread.meta.update", + commandId: h.commandId(), + threadId: firstId, + modelSelection: { ...modelSelection, instanceId: ProviderInstanceId.make("other") }, + }); + } else { + yield* h.engine.dispatch({ + type: action, + commandId: h.commandId(), + threadId: firstId, + createdAt: NOW, + }); + } + yield* TestClock.adjust(0); + yield* service.drain; + expect(yield* h.pending()).toBeUndefined(); + expect((yield* h.read()).session?.lastError).toBe(errorMessage); + yield* h.settings.updateSettings({ continueThreadsAfterUsageLimit: true }); + yield* respond(h, service); + yield* service.recordFailure(failure.event, failure.snapshotSequence); + yield* service.drain; + yield* TestClock.adjust("2 hours"); + expect((yield* h.read()).messages).toEqual([]); + expect(yield* Queue.size(h.probes)).toBe(0); + expect((yield* h.settings.getSettings).usageLimitContinuationPrompt).toBe( + "Carry on with the remaining work", + ); + }).pipe(Effect.provide(testLayer)), +); + +it.effect( + "recovers the stored wait across worker restarts, including an idle session closing", + () => + Effect.gen(function* () { + const h = yield* makeHarness(); + const first = yield* h.start; + yield* record(h, first, firstId, RESET); + const stored = yield* h.pending(); + yield* first.stop; + const session = (yield* h.read()).session!; + yield* h.engine.dispatch({ + type: "thread.session.set", + commandId: h.commandId(), + threadId: firstId, + createdAt: NOW, + session: { ...session, status: "stopped", updatedAt: "2026-09-18T00:01:00.000Z" }, + }); + const second = yield* h.start; + yield* TestClock.adjust("59 minutes"); + expect(yield* Queue.size(h.probes)).toBe(0); + expect(yield* h.pending()).toEqual(stored); + yield* second.stop; + yield* TestClock.adjust("2 hours"); + const third = yield* h.start; + yield* Queue.take(h.probes); + yield* respond(h, third); + expect((yield* h.read()).messages).toHaveLength(1); + expect(yield* h.pending()).toBeUndefined(); + expect(Option.getOrThrow(yield* h.directory.getBinding(firstId)).resumeCursor).toEqual({ + threadId: "provider-first", + }); + }).pipe(Effect.provide(testLayer)), +); + +for (const [name, quota] of [ + ["stale", provider(available("2026-09-17T23:59:59.000Z"))], + ["failed", provider({ checkedAt: NOW, windows: [], unavailable: { reason: "probeFailed" } })], + [ + "another account's", + { ...provider(available(NOW)), instanceId: ProviderInstanceId.make("other") }, + ], +] as const) { + it.effect(`retries ${name} quota without continuing prematurely`, () => + Effect.gen(function* () { + const h = yield* makeHarness(); + const service = yield* h.start; + yield* record(h, service); + yield* TestClock.adjust(0); + yield* Queue.take(h.probes); + yield* respond(h, service, [quota]); + expect((yield* h.read()).messages).toEqual([]); + expect((yield* h.pending())?.nextCheckAt).toBe("2026-09-18T00:05:00.000Z"); + yield* TestClock.adjust("4 minutes"); + expect(yield* Queue.size(h.probes)).toBe(0); + yield* TestClock.adjust("1 minute"); + yield* Queue.take(h.probes); + yield* respond(h, service); + expect((yield* h.read()).messages).toHaveLength(1); + }).pipe(Effect.provide(testLayer)), + ); +} + +it.effect("waits for a new reset when the provider still reports exhausted usage", () => + Effect.gen(function* () { + const h = yield* makeHarness(); + const service = yield* h.start; + yield* record(h, service); + yield* TestClock.adjust(0); + yield* Queue.take(h.probes); + yield* respond(h, service, [ + provider({ + checkedAt: NOW, + windows: [ + { id: "session", kind: "session", label: "Session", usedPercent: 100, resetsAt: RESET }, + ], + }), + ]); + expect((yield* h.pending())?.nextCheckAt).toBe("2026-09-18T01:00:01.000Z"); + yield* TestClock.adjust("1 hour"); + expect((yield* h.read()).messages).toEqual([]); + expect(yield* Queue.size(h.probes)).toBe(0); + yield* TestClock.adjust("1 second"); + yield* Queue.take(h.probes); + yield* respond(h, service); + expect((yield* h.read()).messages).toHaveLength(1); + }).pipe(Effect.provide(testLayer)), +); + +it.effect("backs off when the automatic continuation itself exhausts quota", () => + Effect.gen(function* () { + const h = yield* makeHarness(); + const service = yield* h.start; + yield* record(h, service); + yield* TestClock.adjust(0); + yield* Queue.take(h.probes); + yield* respond(h, service); + yield* record(h, service); + yield* TestClock.adjust("4 minutes"); + expect((yield* h.read()).messages).toHaveLength(1); + expect(yield* Queue.size(h.probes)).toBe(0); + yield* TestClock.adjust("1 minute"); + yield* Queue.take(h.probes); + yield* respond(h, service); + expect((yield* h.read()).messages).toHaveLength(2); + }).pipe(Effect.provide(testLayer)), +); + +it.effect("restores the error and stops retrying an unsupported usage probe", () => + Effect.gen(function* () { + const h = yield* makeHarness(); + const service = yield* h.start; + yield* record(h, service); + yield* TestClock.adjust(0); + yield* Queue.take(h.probes); + yield* respond(h, service, [ + provider({ checkedAt: NOW, windows: [], unavailable: { reason: "unsupported" } }), + ]); + yield* TestClock.adjust("10 minutes"); + expect(yield* h.pending()).toBeUndefined(); + expect((yield* h.read()).session?.lastError).toBe(errorMessage); + expect((yield* h.read()).messages).toEqual([]); + expect(yield* Queue.size(h.probes)).toBe(0); + }).pipe(Effect.provide(testLayer)), +); + +it.effect.each(["stop", "disable", "new-turn"] as const)( + "rejects a continuation already queued at command admission after %s", + (action) => + Effect.gen(function* () { + const h = yield* makeHarness(); + const { event, snapshotSequence } = yield* h.fail(); + const session = (yield* h.read()).session!; + const expectedUsageLimit = { + turnId: TurnId.make(event.turnId!), + providerInstanceId: instanceId, + sessionUpdatedAt: session.updatedAt, + snapshotSequence, + }; + if (action === "disable") { + yield* h.settings.updateSettings({ continueThreadsAfterUsageLimit: false }); + } else if (action === "stop") { + yield* appendThreadActivity(h); + yield* h.engine.dispatch({ + type: "thread.turn.interrupt", + commandId: h.commandId(), + threadId: firstId, + createdAt: NOW, + }); + } else { + yield* h.fail(); + expectedUsageLimit.snapshotSequence = yield* h.engine.latestSequence; + } + const before = yield* h.engine.latestSequence; + const error = yield* h.engine + .dispatch({ + type: "thread.turn.start", + commandId: h.commandId(), + threadId: firstId, + message: { + messageId: MessageId.make("stale-continuation"), + role: "user", + text: "continue", + attachments: [], + }, + runtimeMode: "full-access", + interactionMode: "default", + expectedUsageLimit, + createdAt: NOW, + }) + .pipe(Effect.flip); + expect(error).toMatchObject({ + _tag: "OrchestrationCommandInvariantError", + detail: expect.stringContaining( + action === "disable" + ? "disabled" + : action === "stop" + ? "changed before" + : "no longer matches", + ), + }); + if (action !== "disable") { + const bannerError = yield* h.engine + .dispatch({ + type: "thread.session.set", + commandId: h.commandId(), + threadId: firstId, + session: { ...session, lastError: "Waiting for usage reset" }, + expectedUsageLimit, + createdAt: NOW, + }) + .pipe(Effect.flip); + expect(bannerError._tag).toBe("OrchestrationCommandInvariantError"); + } + expect(yield* h.engine.latestSequence).toBe(before); + expect((yield* h.read()).messages).toEqual([]); + expect((yield* h.read()).session?.lastError).toBe(errorMessage); + }).pipe(Effect.provide(testLayer)), +); diff --git a/apps/server/src/orchestration/UsageLimitContinuation.ts b/apps/server/src/orchestration/UsageLimitContinuation.ts new file mode 100644 index 000000000000..09b7e1bb80e9 --- /dev/null +++ b/apps/server/src/orchestration/UsageLimitContinuation.ts @@ -0,0 +1,419 @@ +import { + CommandId, + MessageId, + ModelSelection, + type OrchestrationThreadShell, + type ProviderInstanceId, + type ProviderRuntimeEvent, + type ServerProvider, + ThreadId, + TurnId, +} from "@t3tools/contracts"; +import { makeDrainableWorker } from "@t3tools/shared/DrainableWorker"; +import * as Clock from "effect/Clock"; +import * as Context from "effect/Context"; +import * as Crypto from "effect/Crypto"; +import * as DateTime from "effect/DateTime"; +import * as Effect from "effect/Effect"; +import * as Fiber from "effect/Fiber"; +import * as Layer from "effect/Layer"; +import * as Option from "effect/Option"; +import * as Schema from "effect/Schema"; +import * as Scope from "effect/Scope"; +import * as Stream from "effect/Stream"; + +import { ProviderRegistry } from "../provider/Services/ProviderRegistry.ts"; +import { ProviderSessionDirectory } from "../provider/Services/ProviderSessionDirectory.ts"; +import { usageLimitRetryAt } from "../provider/Layers/codexUsageLimits.ts"; +import { + type PendingUsageLimitContinuation, + readPendingUsageLimitContinuation, +} from "../provider/usageLimitContinuation.ts"; +import { ServerSettingsService } from "../serverSettings.ts"; +import { forkParked } from "../serverActivation.ts"; +import { OrchestrationEngineService } from "./Services/OrchestrationEngine.ts"; +import { ProjectionSnapshotQuery } from "./Services/ProjectionSnapshotQuery.ts"; +import { cancelsUsageLimitContinuation } from "./decider.ts"; + +type FailedTurn = Extract; +const RECHECK_MS = 5 * 60_000; +const WAITING_MESSAGE_PREFIX = "Usage limit reached. Automatic continuation will check usage at "; +const iso = (millis: number) => DateTime.formatIso(DateTime.makeUnsafe(millis)); +const sameModel = Schema.toEquivalence(ModelSelection); + +export class UsageLimitContinuation extends Context.Service< + UsageLimitContinuation, + { + readonly start: () => Effect.Effect; + readonly recordFailure: (event: FailedTurn, snapshotSequence: number) => Effect.Effect; + readonly drain: Effect.Effect; + } +>()("t3/orchestration/UsageLimitContinuation") {} + +function waitingMessage(pending: PendingUsageLimitContinuation): string { + return `${WAITING_MESSAGE_PREFIX}${pending.nextCheckAt}.`; +} + +/** @public Service construction is part of the canonical Effect module API. */ +export const make = Effect.gen(function* () { + const engine = yield* OrchestrationEngineService; + const query = yield* ProjectionSnapshotQuery; + const directory = yield* ProviderSessionDirectory; + const registry = yield* ProviderRegistry; + const settingsService = yield* ServerSettingsService; + const crypto = yield* Crypto.Crypto; + const scope = yield* Scope.Scope; + const pending = new Map(); + const probing = new Set(); + const handledFailures = new Map(); + let timer: Fiber.Fiber | undefined; + + const getThread = (threadId: ThreadId) => + query.getThreadShellById(threadId).pipe(Effect.map(Option.getOrUndefined)); + const eligible = ( + thread: OrchestrationThreadShell | undefined, + wait: PendingUsageLimitContinuation, + ) => + thread !== undefined && + thread.archivedAt === null && + thread.latestTurn?.turnId === wait.failedTurnId && + thread.latestTurn.state === "error" && + thread.session !== null && + thread.session.activeTurnId === null && + thread.session.status !== "starting" && + thread.session.status !== "running" && + thread.session.providerInstanceId === wait.providerInstanceId && + sameModel(thread.modelSelection, wait.modelSelection); + + const expected = ( + thread: OrchestrationThreadShell, + wait: PendingUsageLimitContinuation, + sequence: number, + ) => ({ + turnId: wait.failedTurnId, + providerInstanceId: wait.providerInstanceId, + sessionUpdatedAt: thread.session?.updatedAt ?? wait.failedAt, + snapshotSequence: sequence, + }); + + const setBanner = Effect.fn("UsageLimitContinuation.setBanner")(function* ( + threadId: ThreadId, + wait: PendingUsageLimitContinuation, + message: string, + ) { + const sequence = yield* engine.latestSequence; + const thread = yield* getThread(threadId); + if (!eligible(thread, wait) || !thread?.session) return; + const now = iso(yield* Clock.currentTimeMillis); + yield* engine + .dispatch({ + type: "thread.session.set", + commandId: CommandId.make(yield* crypto.randomUUIDv4), + threadId, + session: { ...thread.session, lastError: message, updatedAt: now }, + expectedUsageLimit: expected(thread, wait, sequence), + createdAt: now, + }) + .pipe(Effect.ignore); + }); + + const cancel = Effect.fn("UsageLimitContinuation.cancel")(function* (threadId: ThreadId) { + const wait = pending.get(threadId); + if (!wait) return; + yield* directory.setUsageLimitContinuation({ + threadId, + pending: null, + expectedFailedTurnId: wait.failedTurnId, + }); + pending.delete(threadId); + yield* Effect.gen(function* () { + const thread = yield* query + .getThreadShellById(threadId, { includeArchived: true }) + .pipe(Effect.map(Option.getOrUndefined)); + if ( + thread?.session?.lastError?.startsWith(WAITING_MESSAGE_PREFIX) !== true || + thread.latestTurn?.turnId !== wait.failedTurnId || + thread.session.activeTurnId !== null || + thread.session.status === "starting" || + thread.session.status === "running" + ) + return; + const now = iso(yield* Clock.currentTimeMillis); + yield* engine.dispatch({ + type: "thread.session.set", + commandId: CommandId.make(yield* crypto.randomUUIDv4), + threadId, + session: { ...thread.session, lastError: wait.errorMessage, updatedAt: now }, + expectedSession: thread.session, + createdAt: now, + }); + }).pipe( + // An explicit stop can update the session while its waiting banner is cleared. + Effect.retry({ + times: 1, + while: (error) => error._tag === "OrchestrationCommandInvariantError", + }), + Effect.ignore, + ); + }); + + const canceledSince = ( + threadId: ThreadId, + wait: PendingUsageLimitContinuation, + sequence: number, + ) => + engine + .readThreadEvents({ + threadId, + fromSequenceExclusive: wait.snapshotSequence, + toSequenceInclusive: sequence, + limit: sequence - wait.snapshotSequence, + }) + .pipe( + Stream.filter(cancelsUsageLimitContinuation), + Stream.runHead, + Effect.map(Option.isSome), + ); + + const save = Effect.fn("UsageLimitContinuation.save")(function* ( + threadId: ThreadId, + wait: PendingUsageLimitContinuation, + ) { + yield* directory.setUsageLimitContinuation({ threadId, pending: wait }); + pending.set(threadId, wait); + yield* setBanner(threadId, wait, waitingMessage(wait)); + }); + + type Attempt = { + threadId: ThreadId; + wait: PendingUsageLimitContinuation; + }; + const finishProbe = Effect.fn("UsageLimitContinuation.finishProbe")(function* ( + instanceId: ProviderInstanceId, + attempts: Attempt[], + startedAt: number, + providers: readonly ServerProvider[], + ) { + probing.delete(instanceId); + const now = yield* Clock.currentTimeMillis; + const limits = providers.find((provider) => provider.instanceId === instanceId)?.usageLimits; + for (const attempt of attempts) { + const { threadId, wait } = attempt; + if (pending.get(threadId) !== wait) continue; + const sequence = yield* engine.latestSequence; + const thread = yield* getThread(threadId); + if (!thread || !eligible(thread, wait)) { + yield* cancel(threadId); + continue; + } + const settings = yield* settingsService.getSettings; + if ( + !settings.continueThreadsAfterUsageLimit || + (yield* canceledSince(threadId, wait, yield* engine.latestSequence)) + ) { + yield* cancel(threadId); + continue; + } + if (limits?.unavailable?.reason === "unsupported") { + yield* cancel(threadId); + continue; + } + const fresh = + limits !== undefined && + !limits.unavailable && + Date.parse(limits.checkedAt) >= startedAt && + limits.windows.length > 0; + if (fresh && limits.windows.every((window) => window.usedPercent < 100)) { + const key = `usage-limit:${threadId}:${wait.failedTurnId}`; + const result = yield* engine + .dispatch({ + type: "thread.turn.start", + commandId: CommandId.make(key), + threadId, + message: { + messageId: MessageId.make(key), + role: "user", + text: settings.usageLimitContinuationPrompt, + attachments: [], + }, + modelSelection: thread.modelSelection, + runtimeMode: thread.runtimeMode, + interactionMode: thread.interactionMode, + expectedUsageLimit: expected(thread, wait, sequence), + createdAt: iso(now), + }) + .pipe(Effect.result); + // A rejected conditional command cannot be retried under the same receipt. + // Leave the ordinary error if the thread changed while the probe ran. + yield* cancel(threadId); + if (result._tag === "Failure") + yield* Effect.logDebug("usage-limit continuation was not admitted", result.failure); + } else { + const retryAt = fresh ? usageLimitRetryAt(limits.windows, iso(now)) : undefined; + yield* save(threadId, { + ...wait, + nextCheckAt: iso(retryAt ? Date.parse(retryAt) + 1000 : now + RECHECK_MS), + }); + } + } + }); + + const checkDue = Effect.fn("UsageLimitContinuation.checkDue")(function* () { + const now = yield* Clock.currentTimeMillis; + const groups = new Map(); + for (const [threadId, wait] of pending) { + if (probing.has(wait.providerInstanceId) || Date.parse(wait.nextCheckAt) > now) continue; + const sequence = yield* engine.latestSequence; + const thread = yield* getThread(threadId); + if (!eligible(thread, wait) || !thread || (yield* canceledSince(threadId, wait, sequence))) { + yield* cancel(threadId); + continue; + } + const group = groups.get(wait.providerInstanceId) ?? []; + group.push({ threadId, wait }); + groups.set(wait.providerInstanceId, group); + } + for (const [instanceId, attempts] of groups) { + probing.add(instanceId); + yield* registry.refreshInstance(instanceId).pipe( + Effect.orElseSucceed(() => []), + Effect.flatMap((providers) => enqueue(finishProbe(instanceId, attempts, now, providers))), + Effect.forkIn(scope), + ); + } + }); + + const armTimer: Effect.Effect = Effect.gen(function* () { + if (timer) yield* Fiber.interrupt(timer); + timer = undefined; + const times = [...pending.values()] + .filter((wait) => !probing.has(wait.providerInstanceId)) + .map((wait) => Date.parse(wait.nextCheckAt)); + if (times.length === 0) return; + const delay = Math.max(0, Math.min(...times) - (yield* Clock.currentTimeMillis)); + timer = yield* Effect.sleep(delay).pipe( + Effect.andThen(Effect.suspend(() => enqueue(checkDue()))), + Effect.forkIn(scope), + ); + }); + const worker = yield* makeDrainableWorker((work: Effect.Effect) => + work.pipe(Effect.andThen(armTimer)), + ); + const enqueue = (work: Effect.Effect): Effect.Effect => + worker.enqueue( + work.pipe( + Effect.catchCause((cause) => + Effect.gen(function* () { + const now = yield* Clock.currentTimeMillis; + for (const [threadId, wait] of pending) { + if (Date.parse(wait.nextCheckAt) <= now) + pending.set(threadId, { ...wait, nextCheckAt: iso(now + RECHECK_MS) }); + } + yield* Effect.logWarning("usage-limit continuation failed", cause); + }), + ), + ), + ); + + const recordFailure = Effect.fn("UsageLimitContinuation.recordFailure")(function* ( + event: FailedTurn, + snapshotSequence: number, + ) { + yield* enqueue( + Effect.gen(function* () { + if (!event.payload.usageLimit || !event.turnId || !event.providerInstanceId) return; + if (handledFailures.get(event.threadId) === event.turnId) return; + const thread = yield* getThread(event.threadId); + if (thread?.latestTurn?.turnId !== event.turnId || thread.latestTurn.state !== "error") + return; + handledFailures.set(event.threadId, TurnId.make(event.turnId)); + if (!(yield* settingsService.getSettings).continueThreadsAfterUsageLimit) return; + const binding = yield* directory.getBinding(event.threadId); + if (!thread || Option.isNone(binding) || binding.value.resumeCursor == null) return; + const now = yield* Clock.currentTimeMillis; + const retryAt = event.payload.usageLimit.retryAt; + let nextCheckAt = retryAt ? Math.max(now, Date.parse(retryAt) + 1000) : now; + if (!retryAt || Date.parse(retryAt) <= now) { + const detail = yield* query.getThreadDetailSnapshot(event.threadId, { turnLimit: 1 }); + const lastUser = Option.isSome(detail) + ? detail.value.thread.messages.findLast((message) => message.role === "user") + : undefined; + if (lastUser?.id.startsWith("usage-limit:")) + nextCheckAt = Math.max(now, Date.parse(event.createdAt) + RECHECK_MS); + } + const wait: PendingUsageLimitContinuation = { + failedTurnId: TurnId.make(event.turnId), + providerInstanceId: event.providerInstanceId, + modelSelection: thread.modelSelection, + errorMessage: event.payload.errorMessage ?? "Usage limit reached.", + failedAt: event.createdAt, + snapshotSequence, + nextCheckAt: iso(nextCheckAt), + }; + if ( + eligible(thread, wait) && + !(yield* canceledSince(event.threadId, wait, yield* engine.latestSequence)) + ) + yield* save(event.threadId, wait); + }).pipe( + Effect.tapError(() => + Effect.sync(() => { + if (handledFailures.get(event.threadId) === event.turnId) + handledFailures.delete(event.threadId); + }), + ), + ), + ); + }); + + const start = Effect.fn("UsageLimitContinuation.start")(function* () { + const events = yield* engine.subscribeDomainEvents; + const changes = yield* settingsService.subscribeChanges; + yield* forkParked( + enqueue( + Effect.gen(function* () { + const enabled = (yield* settingsService.getSettings).continueThreadsAfterUsageLimit; + for (const binding of yield* directory.listBindings()) { + const wait = readPendingUsageLimitContinuation(binding.runtimePayload); + if (!wait) continue; + pending.set(binding.threadId, wait); + handledFailures.set(binding.threadId, wait.failedTurnId); + if ( + !enabled || + binding.providerInstanceId !== wait.providerInstanceId || + (yield* canceledSince(binding.threadId, wait, yield* engine.latestSequence)) + ) + yield* cancel(binding.threadId); + } + }), + ), + ); + yield* forkParked( + Stream.runForEach(events, (event) => + cancelsUsageLimitContinuation(event) + ? enqueue( + Effect.suspend(() => { + const threadId = ThreadId.make(event.aggregateId); + const wait = pending.get(threadId); + return wait && event.sequence > wait.snapshotSequence + ? cancel(threadId) + : Effect.void; + }), + ) + : Effect.void, + ), + ); + yield* forkParked( + Stream.runForEach(changes, (settings) => + settings.continueThreadsAfterUsageLimit + ? Effect.void + : enqueue( + Effect.suspend(() => Effect.forEach([...pending.keys()], cancel, { discard: true })), + ), + ), + ); + }); + return { start, recordFailure, drain: worker.drain } satisfies UsageLimitContinuation["Service"]; +}); + +export const layer = Layer.effect(UsageLimitContinuation, make); diff --git a/apps/server/src/orchestration/decider.ts b/apps/server/src/orchestration/decider.ts index 0119c0e8599a..dba81d68c575 100644 --- a/apps/server/src/orchestration/decider.ts +++ b/apps/server/src/orchestration/decider.ts @@ -3,6 +3,7 @@ import { MAX_SCRIPT_ID_LENGTH, SCRIPT_RUN_COMMAND_PATTERN, MessageId, + OrchestrationSession, ThreadLinkedPullRequest, UserInputRequestedPayload, isImportedAgentSessionMessageId, @@ -53,6 +54,7 @@ const isScriptRunCommand = Schema.is(SCRIPT_RUN_COMMAND_PATTERN); const nowIso = Effect.map(DateTime.now, DateTime.formatIso); const decodeUserInputRequestedPayload = Schema.decodeUnknownOption(UserInputRequestedPayload); +const sessionsEqual = Schema.toEquivalence(OrchestrationSession); const threadPullRequestLinksEqual = Schema.toEquivalence(Schema.NullOr(ThreadLinkedPullRequest)); /** @@ -131,6 +133,60 @@ function hasQueuedTurnStartForThread( ); } +export function cancelsUsageLimitContinuation(event: OrchestrationEvent): boolean { + switch (event.type) { + case "thread.turn-start-requested": + case "thread.turn-interrupt-requested": + case "thread.session-stop-requested": + case "thread.archived": + case "thread.deleted": + case "thread.runtime-mode-set": + case "thread.interaction-mode-set": + return true; + case "thread.meta-updated": + return event.payload.modelSelection !== undefined || event.payload.worktreePath !== undefined; + default: + return false; + } +} + +function requireExpectedUsageLimit( + thread: OrchestrationThread, + command: Extract, +): Effect.Effect { + const expected = command.expectedUsageLimit; + if (expected === undefined) return Effect.void; + const turn = thread.latestTurn; + const session = thread.session; + if ( + thread.deletedAt !== null || + thread.archivedAt !== null || + turn?.turnId !== expected.turnId || + turn.state !== "error" || + session === null || + session.providerInstanceId !== expected.providerInstanceId || + thread.modelSelection.instanceId !== expected.providerInstanceId || + (command.type === "thread.session.set" && session.updatedAt !== expected.sessionUpdatedAt) || + session.activeTurnId !== null || + session.status === "starting" || + session.status === "running" || + thread.messages.some( + (message) => + message.role === "user" && + !isImportedAgentSessionMessageId(message.id) && + Date.parse(message.createdAt) > Date.parse(turn.completedAt ?? turn.requestedAt), + ) + ) { + return Effect.fail( + new OrchestrationCommandInvariantError({ + commandType: command.type, + detail: `Thread '${thread.id}' no longer matches the usage-limit failure.`, + }), + ); + } + return Effect.void; +} + function findPullRequestLink( thread: Pick, key: ThreadPullRequestKey, @@ -1377,6 +1433,7 @@ export const decideOrchestrationCommand = Effect.fn("decideOrchestrationCommand" command, threadId: command.threadId, }); + yield* requireExpectedUsageLimit(targetThread, command); const sourceProposedPlan = command.sourceProposedPlan; const sourceThread = sourceProposedPlan ? yield* requireThread({ @@ -1867,6 +1924,16 @@ export const decideOrchestrationCommand = Effect.fn("decideOrchestrationCommand" command, threadId: command.threadId, }); + yield* requireExpectedUsageLimit(thread, command); + if ( + command.expectedSession !== undefined && + (thread.session === null || !sessionsEqual(thread.session, command.expectedSession)) + ) { + return yield* new OrchestrationCommandInvariantError({ + commandType: command.type, + detail: `Thread '${thread.id}' session changed before the conditional update.`, + }); + } const sessionSetEvent: Omit = { ...(yield* withEventBase({ aggregateKind: "thread", diff --git a/apps/server/src/persistence/ProviderSessionRuntime.ts b/apps/server/src/persistence/ProviderSessionRuntime.ts index 2673512edf10..97df3ce83f36 100644 --- a/apps/server/src/persistence/ProviderSessionRuntime.ts +++ b/apps/server/src/persistence/ProviderSessionRuntime.ts @@ -18,6 +18,8 @@ import { ThreadId, } from "@t3tools/contracts"; +import { SetUsageLimitContinuationInput } from "../provider/usageLimitContinuation.ts"; + import { PersistenceDecodeError, type PersistenceErrorCorrelation, @@ -79,7 +81,7 @@ export class ProviderSessionRuntimeRepository extends Context.Service< * Insert or replace a provider runtime row. * * Upserts by canonical `threadId`, retaining imported transcript records - * from the current database row. + * and pending usage continuation from the current database row. */ readonly upsert: ( runtime: ProviderSessionRuntime, @@ -91,6 +93,10 @@ export class ProviderSessionRuntimeRepository extends Context.Service< input: RecordImportedTranscriptInput, ) => Effect.Effect; + readonly setUsageLimitContinuation: ( + input: SetUsageLimitContinuationInput, + ) => Effect.Effect; + /** * Read provider runtime state by canonical thread id. */ @@ -147,6 +153,10 @@ const GetRuntimeRequestSchema = Schema.Struct({ const DeleteRuntimeRequestSchema = GetRuntimeRequestSchema; +const SetUsageLimitContinuationRequestSchema = SetUsageLimitContinuationInput.mapFields( + Struct.assign({ pending: Schema.fromJsonString(SetUsageLimitContinuationInput.fields.pending) }), +); + const RecordImportedTranscriptRequestSchema = RecordImportedTranscriptInput.mapFields( Struct.assign({ source: Schema.fromJsonString(AgentSessionImportSource) }), ); @@ -171,7 +181,8 @@ export const make = Effect.gen(function* () { const sql = yield* SqlClient.SqlClient; // Runtime writes can carry stale payloads. Only recordImportedTranscript may - // change source records, so restore that field from the row being updated. + // change source records. The continuation worker likewise owns its marker; + // preserve both fields from the row being updated. const upsertRuntimeRow = SqlSchema.void({ Request: ProviderSessionRuntimeDbRowSchema, execute: (runtime) => @@ -198,7 +209,7 @@ export const make = Effect.gen(function* () { ${runtime.resumeCursor}, CASE WHEN json_type(${runtime.runtimePayload}) = 'object' - THEN json_remove(${runtime.runtimePayload}, '$.importedTranscripts') + THEN json_remove(${runtime.runtimePayload}, '$.importedTranscripts', '$.usageLimitContinuation') ELSE ${runtime.runtimePayload} END ) @@ -212,22 +223,18 @@ export const make = Effect.gen(function* () { last_seen_at = excluded.last_seen_at, resume_cursor_json = excluded.resume_cursor_json, runtime_payload_json = CASE - WHEN json_type( - CASE - WHEN json_valid(provider_session_runtime.runtime_payload_json) - THEN provider_session_runtime.runtime_payload_json - ELSE '{}' - END, - '$.importedTranscripts' - ) IS NOT NULL - THEN json_set( - CASE - WHEN json_type(excluded.runtime_payload_json) = 'object' - THEN excluded.runtime_payload_json - ELSE '{}' - END, - '$.importedTranscripts', - json_extract(provider_session_runtime.runtime_payload_json, '$.importedTranscripts') + WHEN json_valid(provider_session_runtime.runtime_payload_json) + AND ( + json_type(provider_session_runtime.runtime_payload_json, '$.importedTranscripts') IS NOT NULL + OR json_type(provider_session_runtime.runtime_payload_json, '$.usageLimitContinuation') IS NOT NULL + ) + THEN json_patch( + CASE WHEN json_type(excluded.runtime_payload_json) = 'object' + THEN excluded.runtime_payload_json ELSE '{}' END, + json_object( + 'importedTranscripts', json_extract(provider_session_runtime.runtime_payload_json, '$.importedTranscripts'), + 'usageLimitContinuation', json_extract(provider_session_runtime.runtime_payload_json, '$.usageLimitContinuation') + ) ) ELSE excluded.runtime_payload_json END @@ -260,7 +267,7 @@ export const make = Effect.gen(function* () { ${runtime.resumeCursor}, CASE WHEN json_type(${runtime.runtimePayload}) = 'object' - THEN json_remove(${runtime.runtimePayload}, '$.importedTranscripts') + THEN json_remove(${runtime.runtimePayload}, '$.importedTranscripts', '$.usageLimitContinuation') ELSE ${runtime.runtimePayload} END ) @@ -268,6 +275,26 @@ export const make = Effect.gen(function* () { `, }); + const setUsageLimitContinuationRow = SqlSchema.void({ + Request: SetUsageLimitContinuationRequestSchema, + execute: ({ threadId, pending, expectedFailedTurnId }) => + sql` + UPDATE provider_session_runtime + SET runtime_payload_json = CASE + WHEN ${pending} = 'null' THEN json_remove(runtime_payload_json, '$.usageLimitContinuation') + ELSE json_set( + CASE WHEN json_type(runtime_payload_json) = 'object' + THEN runtime_payload_json ELSE '{}' END, + '$.usageLimitContinuation', json(${pending}) + ) + END + WHERE thread_id = ${threadId} + AND (${pending} = 'null' OR provider_instance_id = json_extract(${pending}, '$.providerInstanceId')) + AND (${expectedFailedTurnId ?? null} IS NULL + OR json_extract(runtime_payload_json, '$.usageLimitContinuation.failedTurnId') = ${expectedFailedTurnId ?? null}) + `, + }); + const recordImportedTranscriptRow = SqlSchema.void({ Request: RecordImportedTranscriptRequestSchema, execute: ({ threadId, source }) => @@ -387,6 +414,18 @@ export const make = Effect.gen(function* () { ), ); + const setUsageLimitContinuation: ProviderSessionRuntimeRepository["Service"]["setUsageLimitContinuation"] = + (input) => + setUsageLimitContinuationRow(input).pipe( + Effect.mapError( + toPersistenceSqlOrDecodeError( + "ProviderSessionRuntimeRepository.setUsageLimitContinuation:query", + "ProviderSessionRuntimeRepository.setUsageLimitContinuation:encodeRequest", + { threadId: input.threadId }, + ), + ), + ); + const getByThreadId: ProviderSessionRuntimeRepository["Service"]["getByThreadId"] = (input) => getRuntimeRowByThreadId(input).pipe( Effect.mapError( @@ -466,6 +505,7 @@ export const make = Effect.gen(function* () { return { upsert, recordImportedTranscript, + setUsageLimitContinuation, getByThreadId, list, deleteByThreadId, diff --git a/apps/server/src/project/AgentSessionImporter.test.ts b/apps/server/src/project/AgentSessionImporter.test.ts index 4eb03a5cc036..413a84c27d9c 100644 --- a/apps/server/src/project/AgentSessionImporter.test.ts +++ b/apps/server/src/project/AgentSessionImporter.test.ts @@ -232,6 +232,7 @@ it.layer(NodeServices.layer)("AgentSessionImporter", (it) => { upsert: (binding) => Effect.sync(() => void bindings.push(binding)), getProvider: () => Effect.die("unused"), recordImportedTranscript: () => Effect.void, + setUsageLimitContinuation: () => Effect.die("unused"), getBinding: () => Effect.succeed(Option.none()), listThreadIds: () => Effect.die("unused"), listBindings: () => Effect.die("unused"), @@ -337,6 +338,7 @@ it.layer(NodeServices.layer)("AgentSessionImporter", (it) => { upsert: () => Effect.die("must not bind a scanner skip"), getProvider: () => Effect.die("unused"), recordImportedTranscript: () => Effect.die("unused"), + setUsageLimitContinuation: () => Effect.die("unused"), getBinding: () => Effect.die("must not read a scanner skip binding"), listThreadIds: () => Effect.die("unused"), listBindings: () => Effect.die("unused"), @@ -414,6 +416,7 @@ it.layer(NodeServices.layer)("AgentSessionImporter", (it) => { }, getProvider: () => Effect.die("unused"), recordImportedTranscript: () => Effect.void, + setUsageLimitContinuation: () => Effect.die("unused"), getBinding: () => Effect.succeed(bindings[0] === undefined ? Option.none() : Option.some(bindings[0])), listThreadIds: () => Effect.die("unused"), @@ -456,6 +459,7 @@ it.layer(NodeServices.layer)("AgentSessionImporter", (it) => { upsert: () => Effect.die("must not replace an active binding"), getProvider: () => Effect.die("unused"), recordImportedTranscript: () => Effect.void, + setUsageLimitContinuation: () => Effect.die("unused"), getBinding: () => Effect.succeed(Option.some(runningBinding)), listThreadIds: () => Effect.die("unused"), listBindings: () => Effect.die("unused"), @@ -511,6 +515,7 @@ it.layer(NodeServices.layer)("AgentSessionImporter", (it) => { upsert: () => Effect.die("must not bind malformed or wrong-project sessions"), getProvider: () => Effect.die("unused"), recordImportedTranscript: () => Effect.die("unused"), + setUsageLimitContinuation: () => Effect.die("unused"), getBinding: () => Effect.succeed(Option.none()), listThreadIds: () => Effect.die("unused"), listBindings: () => Effect.die("unused"), @@ -562,6 +567,7 @@ const integrationRuntimeRepository = ProviderSessionRuntime.layer.pipe( ); const integrationLayer = Layer.mergeAll( OrchestrationEngineLive.pipe( + Layer.provide(ServerSettingsService.layerTest()), Layer.provide(OrchestrationProjectionSnapshotQueryLive), Layer.provide(OrchestrationProjectionPipelineLive), ), diff --git a/apps/server/src/provider/Layers/CodexAdapter.test.ts b/apps/server/src/provider/Layers/CodexAdapter.test.ts index 9f464bdaa177..b09cafc798ba 100644 --- a/apps/server/src/provider/Layers/CodexAdapter.test.ts +++ b/apps/server/src/provider/Layers/CodexAdapter.test.ts @@ -223,6 +223,7 @@ function makeScopedRuntimeFactory(options?: { readonly failConstruction?: boolea const providerSessionDirectoryTestLayer = Layer.succeed(ProviderSessionDirectory, { upsert: () => Effect.void, recordImportedTranscript: () => Effect.die("unused"), + setUsageLimitContinuation: () => Effect.die("unused"), getProvider: () => Effect.die(new Error("ProviderSessionDirectory.getProvider is not used in test")), getBinding: () => Effect.succeed(Option.none()), @@ -2850,6 +2851,7 @@ function codexErrorNotification(input: { function codexRateLimitsNotification(input: { readonly id: string; readonly rateLimitReachedType?: string; + readonly limitId?: string; readonly primary?: { readonly usedPercent: number; readonly resetsInSeconds: number }; readonly secondary?: { readonly usedPercent: number; readonly resetsInSeconds: number }; }): ProviderEvent { @@ -2863,7 +2865,7 @@ function codexRateLimitsNotification(input: { method: "account/rateLimits/updated", payload: { rateLimits: { - limitId: "codex", + limitId: input.limitId ?? "codex", ...(input.rateLimitReachedType ? { rateLimitReachedType: input.rateLimitReachedType } : {}), ...(input.primary ? { @@ -2958,6 +2960,7 @@ usageLimitLayer("CodexAdapterLive usage limits", (it) => { } if (event.type === "turn.completed") { NodeAssert.equal(event.payload.errorMessage, expected); + NodeAssert.equal(event.payload.usageLimit, undefined); } } }), @@ -2990,6 +2993,9 @@ usageLimitLayer("CodexAdapterLive usage limits", (it) => { const events = Array.from(yield* Fiber.join(eventsFiber)); const completed = events.find((event) => event.type === "turn.completed"); + NodeAssert.deepStrictEqual(completed?.payload.usageLimit, { + retryAt: "2026-01-01T03:20:00.000Z", + }); NodeAssert.equal( completed?.payload.errorMessage, "Codex usage limit reached. The session limit resets in 3h 20m. Send the message again once the limit resets.", @@ -3059,6 +3065,36 @@ usageLimitLayer("CodexAdapterLive usage limits", (it) => { NodeAssert.equal(runtimeError?.payload.message, expected); const completed = events.find((event) => event.type === "turn.completed"); NodeAssert.equal(completed?.payload.errorMessage, expected); + NodeAssert.deepStrictEqual(completed?.payload.usageLimit, {}); + }), + ); + + it.effect("does not arm a model-specific stop using an earlier main allowance", () => + Effect.gen(function* () { + const { adapter, runtime } = yield* startUsageLimitRuntime(); + const eventsFiber = yield* adapter.streamEvents.pipe( + Stream.take(3), + Stream.runCollect, + Effect.forkChild, + ); + yield* runtime.emit( + codexRateLimitsNotification({ + id: "evt-main", + primary: { usedPercent: 100, resetsInSeconds: 600 }, + }), + ); + yield* runtime.emit( + codexRateLimitsNotification({ + id: "evt-spark", + limitId: "codex_spark", + primary: { usedPercent: 100, resetsInSeconds: 1200 }, + }), + ); + yield* runtime.emit(codexUsageLimitTurnFailed("evt-spark-stop")); + const events = Array.from(yield* Fiber.join(eventsFiber)); + const completed = events.find((event) => event.type === "turn.completed"); + NodeAssert.ok(completed); + NodeAssert.equal(completed.payload.usageLimit, undefined); }), ); diff --git a/apps/server/src/provider/Layers/CodexAdapter.ts b/apps/server/src/provider/Layers/CodexAdapter.ts index 0ecc9693ab04..b7b1bbce6dbe 100644 --- a/apps/server/src/provider/Layers/CodexAdapter.ts +++ b/apps/server/src/provider/Layers/CodexAdapter.ts @@ -75,6 +75,7 @@ import { type CodexRateLimitSnapshot, codexRateLimitsToUpdate, codexUsageLimitMessage, + codexUsageLimitRecovery, mergeCodexRateLimits, } from "./codexUsageLimits.ts"; const isCodexAppServerProcessExitedError = Schema.is(CodexErrors.CodexAppServerProcessExitedError); @@ -2314,6 +2315,7 @@ export const makeCodexAdapter = Effect.fn("makeCodexAdapter")(function* ( // after the stop and often sparse, so keep the session's merged view of // it and read it when a turn fails on the limit. let rateLimits: CodexRateLimitSnapshot | undefined; + let latestRateLimits: CodexRateLimitSnapshot | undefined; const sessionScope = yield* Scope.make("sequential"); let sessionScopeTransferred = false; yield* Effect.addFinalizer(() => @@ -2377,6 +2379,7 @@ export const makeCodexAdapter = Effect.fn("makeCodexAdapter")(function* ( event.payload, ); if (limitsPayload) { + latestRateLimits = limitsPayload.rateLimits; rateLimits = mergeCodexRateLimits(rateLimits, limitsPayload.rateLimits); } } else if (event.method === "error") { @@ -2391,6 +2394,7 @@ export const makeCodexAdapter = Effect.fn("makeCodexAdapter")(function* ( let usageLimitError: ProviderRuntimeEvent | undefined; let usageLimitMessage: string | undefined; + let usageLimit: { readonly retryAt?: string } | undefined; if (event.method === "turn/completed") { const completedPayload = readPayload( EffectCodexSchema.V2TurnCompletedNotification, @@ -2402,6 +2406,12 @@ export const makeCodexAdapter = Effect.fn("makeCodexAdapter")(function* ( : undefined; if (turnError?.codexErrorInfo === "usageLimitExceeded") { usageLimitMessage = codexUsageLimitMessage(rateLimits, event.createdAt); + usageLimit = codexUsageLimitRecovery( + latestRateLimits?.limitId && latestRateLimits.limitId !== "codex" + ? latestRateLimits + : rateLimits, + event.createdAt, + ); usageLimitError = { ...runtimeEventBase(event, event.threadId), type: "runtime.error", @@ -2421,6 +2431,7 @@ export const makeCodexAdapter = Effect.fn("makeCodexAdapter")(function* ( payload: { ...runtimeEvent.payload, ...(usageLimitMessage ? { errorMessage: usageLimitMessage } : {}), + ...(usageLimit ? { usageLimit } : {}), tokenUsage: completeCodexTurnTokenUsage( turnTokenUsage, String(runtimeEvent.turnId), diff --git a/apps/server/src/provider/Layers/OpenCodeAdapter.test.ts b/apps/server/src/provider/Layers/OpenCodeAdapter.test.ts index 96d6b10d3839..b4e48e423961 100644 --- a/apps/server/src/provider/Layers/OpenCodeAdapter.test.ts +++ b/apps/server/src/provider/Layers/OpenCodeAdapter.test.ts @@ -576,6 +576,7 @@ const OpenCodeRuntimeTestDouble: OpenCodeRuntimeShape = { const providerSessionDirectoryTestLayer = Layer.succeed(ProviderSessionDirectory, { upsert: () => Effect.void, recordImportedTranscript: () => Effect.die("unused"), + setUsageLimitContinuation: () => Effect.die("unused"), getProvider: () => Effect.die(new Error("ProviderSessionDirectory.getProvider is not used in test")), getBinding: () => Effect.succeed(Option.none()), diff --git a/apps/server/src/provider/Layers/ProviderService.test.ts b/apps/server/src/provider/Layers/ProviderService.test.ts index b9997e1df312..5611d594881e 100644 --- a/apps/server/src/provider/Layers/ProviderService.test.ts +++ b/apps/server/src/provider/Layers/ProviderService.test.ts @@ -4906,6 +4906,7 @@ const boundedListing = makeProviderServiceLayer({ directory: { upsert: () => Effect.void, recordImportedTranscript: () => Effect.die("unused"), + setUsageLimitContinuation: () => Effect.die("unused"), getProvider: () => Effect.die("ProviderService.listSessions does not use getProvider"), getBinding, listThreadIds, diff --git a/apps/server/src/provider/Layers/ProviderSessionDirectory.test.ts b/apps/server/src/provider/Layers/ProviderSessionDirectory.test.ts index 8b41bd3e518c..b34157b48f89 100644 --- a/apps/server/src/provider/Layers/ProviderSessionDirectory.test.ts +++ b/apps/server/src/provider/Layers/ProviderSessionDirectory.test.ts @@ -8,6 +8,7 @@ import { ProviderDriverKind, ProviderInstanceId, ThreadId, + TurnId, type AgentSessionImportSource, } from "@t3tools/contracts"; import { assert, expect, it } from "@effect/vitest"; @@ -22,6 +23,10 @@ import { } from "../../persistence/Layers/Sqlite.ts"; import * as ProviderSessionRuntime from "../../persistence/ProviderSessionRuntime.ts"; import { ProviderSessionDirectory } from "../Services/ProviderSessionDirectory.ts"; +import { + readPendingUsageLimitContinuation, + type PendingUsageLimitContinuation, +} from "../usageLimitContinuation.ts"; import { ProviderSessionDirectoryLive } from "./ProviderSessionDirectory.ts"; const importedSource = { @@ -145,6 +150,110 @@ it.layer(makeDirectoryLayer(SqlitePersistenceMemory))("ProviderSessionDirectoryL }), ); + it.effect( + "preserves pending usage waits across ordinary writes and never restores a cleared wait", + () => + Effect.gen(function* () { + const directory = yield* ProviderSessionDirectory; + const repository = yield* ProviderSessionRuntime.ProviderSessionRuntimeRepository; + const threadId = ThreadId.make("thread-usage-wait"); + const providerInstanceId = ProviderInstanceId.make("codex"); + const pending: PendingUsageLimitContinuation = { + failedTurnId: TurnId.make("turn-quota"), + providerInstanceId, + modelSelection: { instanceId: providerInstanceId, model: "gpt-5" }, + errorMessage: "Usage limit reached.", + failedAt: "2026-09-18T10:00:00.000Z", + nextCheckAt: "2026-09-18T11:00:00.000Z", + snapshotSequence: 12, + }; + yield* directory.upsert({ + threadId, + provider: ProviderDriverKind.make("codex"), + providerInstanceId, + resumeCursor: { threadId: "provider-thread" }, + runtimePayload: { cwd: "/tmp/project" }, + }); + yield* directory.recordImportedTranscript({ threadId, source: importedSource }); + const beforeWait = Option.getOrThrow(yield* repository.getByThreadId({ threadId })); + yield* directory.setUsageLimitContinuation({ threadId, pending }); + // A session write that read before the wait must preserve it. + yield* repository.upsert(beforeWait); + const withWait = Option.getOrThrow(yield* repository.getByThreadId({ threadId })); + expect(readPendingUsageLimitContinuation(withWait.runtimePayload)).toEqual(pending); + expect(withWait.resumeCursor).toEqual({ threadId: "provider-thread" }); + expect(withWait.runtimePayload).toMatchObject({ + cwd: "/tmp/project", + importedTranscripts: [importedSource], + }); + expect( + readPendingUsageLimitContinuation( + Option.getOrThrow(yield* directory.getBinding(threadId)).runtimePayload, + ), + ).toEqual(pending); + + const replacement = { ...pending, failedTurnId: TurnId.make("turn-quota-new") }; + yield* directory.setUsageLimitContinuation({ threadId, pending: replacement }); + yield* directory.setUsageLimitContinuation({ + threadId, + pending: null, + expectedFailedTurnId: pending.failedTurnId, + }); + expect( + readPendingUsageLimitContinuation( + Option.getOrThrow(yield* directory.getBinding(threadId)).runtimePayload, + ), + ).toEqual(replacement); + yield* directory.setUsageLimitContinuation({ + threadId, + pending: null, + expectedFailedTurnId: replacement.failedTurnId, + }); + // Neither a stale ordinary write nor a stale timer update may resurrect it. + yield* repository.upsert(withWait); + yield* directory.setUsageLimitContinuation({ + threadId, + pending, + expectedFailedTurnId: pending.failedTurnId, + }); + const cleared = Option.getOrThrow(yield* directory.getBinding(threadId)); + expect(readPendingUsageLimitContinuation(cleared.runtimePayload)).toBeUndefined(); + expect(cleared.runtimePayload).toEqual({ + cwd: "/tmp/project", + importedTranscripts: [importedSource], + }); + }), + ); + + it.effect("rejects a pending wait owned by another provider instance", () => + Effect.gen(function* () { + const directory = yield* ProviderSessionDirectory; + const threadId = ThreadId.make("thread-usage-instance"); + yield* directory.upsert({ + threadId, + provider: ProviderDriverKind.make("codex"), + providerInstanceId: ProviderInstanceId.make("codex-work"), + }); + yield* directory.setUsageLimitContinuation({ + threadId, + pending: { + failedTurnId: TurnId.make("turn-quota"), + providerInstanceId: ProviderInstanceId.make("codex-personal"), + modelSelection: { instanceId: ProviderInstanceId.make("codex-personal"), model: "gpt-5" }, + errorMessage: "Usage limit reached.", + failedAt: "2026-09-18T10:00:00.000Z", + nextCheckAt: "2026-09-18T11:00:00.000Z", + snapshotSequence: 1, + }, + }); + expect( + readPendingUsageLimitContinuation( + Option.getOrThrow(yield* directory.getBinding(threadId)).runtimePayload, + ), + ).toBeUndefined(); + }), + ); + it.effect("keeps the existing binding when an insert conflicts", () => Effect.gen(function* () { const directory = yield* ProviderSessionDirectory; diff --git a/apps/server/src/provider/Layers/ProviderSessionDirectory.ts b/apps/server/src/provider/Layers/ProviderSessionDirectory.ts index 29ec8d2ed168..12ea602b13ab 100644 --- a/apps/server/src/provider/Layers/ProviderSessionDirectory.ts +++ b/apps/server/src/provider/Layers/ProviderSessionDirectory.ts @@ -178,6 +178,15 @@ const makeProviderSessionDirectory = Effect.gen(function* () { Effect.mapError(toPersistenceError("ProviderSessionDirectory.recordImportedTranscript")), ); + const setUsageLimitContinuation: ProviderSessionDirectoryShape["setUsageLimitContinuation"] = ( + input, + ) => + repository + .setUsageLimitContinuation(input) + .pipe( + Effect.mapError(toPersistenceError("ProviderSessionDirectory.setUsageLimitContinuation")), + ); + const listThreadIds: ProviderSessionDirectoryShape["listThreadIds"] = () => repository.list().pipe( Effect.mapError(toPersistenceError("ProviderSessionDirectory.listThreadIds:list")), @@ -199,6 +208,7 @@ const makeProviderSessionDirectory = Effect.gen(function* () { return { upsert, recordImportedTranscript, + setUsageLimitContinuation, getProvider, getBinding, listThreadIds, diff --git a/apps/server/src/provider/Layers/codexUsageLimits.test.ts b/apps/server/src/provider/Layers/codexUsageLimits.test.ts index 37006789986c..d0ce06c63c4c 100644 --- a/apps/server/src/provider/Layers/codexUsageLimits.test.ts +++ b/apps/server/src/provider/Layers/codexUsageLimits.test.ts @@ -7,6 +7,7 @@ import { codexRateLimitsToUpdate, codexResetCreditsToContract, codexUsageLimitMessage, + codexUsageLimitRecovery, mergeCodexRateLimits, } from "./codexUsageLimits.ts"; @@ -297,3 +298,45 @@ describe("mergeCodexRateLimits", () => { ).toBe(main); }); }); + +describe("usage-limit continuation hints", () => { + const now = "2026-09-18T10:00:00.000Z"; + const seconds = Date.parse(now) / 1000; + + it("carries the exact latest exhausted reset and ignores windows with remaining allowance", () => { + expect( + codexUsageLimitRecovery( + { + limitId: "codex", + primary: { usedPercent: 100, resetsAt: seconds + 123 }, + secondary: { usedPercent: 90, resetsAt: seconds + 789 }, + }, + now, + ), + ).toEqual({ retryAt: "2026-09-18T10:02:03.000Z" }); + expect( + codexUsageLimitRecovery( + { + limitId: "codex", + primary: { usedPercent: 100, resetsAt: seconds + 123 }, + secondary: { usedPercent: 100, resetsAt: seconds + 789 }, + }, + now, + ), + ).toEqual({ retryAt: "2026-09-18T10:13:09.000Z" }); + }); + + it("requires a future reset for every exhausted window", () => { + for (const resetsAt of [undefined, seconds]) { + expect( + codexUsageLimitRecovery( + { + primary: { usedPercent: 100, resetsAt: seconds + 123 }, + secondary: { usedPercent: 100, ...(resetsAt === undefined ? {} : { resetsAt }) }, + }, + now, + ), + ).toEqual({}); + } + }); +}); diff --git a/apps/server/src/provider/Layers/codexUsageLimits.ts b/apps/server/src/provider/Layers/codexUsageLimits.ts index 74ce04a7895f..fc3dd16c084f 100644 --- a/apps/server/src/provider/Layers/codexUsageLimits.ts +++ b/apps/server/src/provider/Layers/codexUsageLimits.ts @@ -221,14 +221,55 @@ export function codexUsageLimitMessage( ): string { const atMs = Date.parse(atIso); const windows = snapshot && Number.isFinite(atMs) ? codexRateLimitsToWindows(snapshot) : []; - let reset = ""; - let latestResetMs = Number.NEGATIVE_INFINITY; + const latest = latestExhaustedWindow(windows, atIso); + const reset = latest?.resetsAt + ? ` The ${latest.kind} limit resets in ${formatCodexUsageLimitWait(Date.parse(latest.resetsAt) - atMs)}.` + : ""; + return `Codex usage limit reached.${reset}${codexUsageLimitNextStep(snapshot?.rateLimitReachedType)}`; +} + +function latestExhaustedWindow( + windows: ReadonlyArray, + atIso: string, +): ServerProviderUsageWindow | undefined { + const atMs = Date.parse(atIso); + let latest: ServerProviderUsageWindow | undefined; + let latestResetMs = atMs; for (const window of windows) { if (window.usedPercent < 100 || !window.resetsAt) continue; const resetMs = Date.parse(window.resetsAt); - if (!Number.isFinite(resetMs) || resetMs <= atMs || resetMs <= latestResetMs) continue; + if (!Number.isFinite(resetMs) || resetMs <= latestResetMs) continue; + latest = window; latestResetMs = resetMs; - reset = ` The ${window.kind} limit resets in ${formatCodexUsageLimitWait(resetMs - atMs)}.`; } - return `Codex usage limit reached.${reset}${codexUsageLimitNextStep(snapshot?.rateLimitReachedType)}`; + return latest; +} + +/** A timer is useful only when every exhausted window has a future reset. */ +export function usageLimitRetryAt( + windows: ReadonlyArray, + atIso: string, +): string | undefined { + const atMs = Date.parse(atIso); + if ( + !Number.isFinite(atMs) || + windows.some( + (window) => + window.usedPercent >= 100 && (!window.resetsAt || !(Date.parse(window.resetsAt) > atMs)), + ) + ) + return undefined; + return latestExhaustedWindow(windows, atIso)?.resetsAt; +} + +export function codexUsageLimitRecovery( + snapshot: CodexRateLimitSnapshot | undefined, + atIso: string, +): { readonly retryAt?: string } | undefined { + if (snapshot?.limitId && snapshot.limitId !== "codex") return undefined; + if (snapshot?.rateLimitReachedType && snapshot.rateLimitReachedType !== "rate_limit_reached") { + return undefined; + } + const retryAt = usageLimitRetryAt(snapshot ? codexRateLimitsToWindows(snapshot) : [], atIso); + return retryAt ? { retryAt } : {}; } diff --git a/apps/server/src/provider/Services/ProviderSessionDirectory.ts b/apps/server/src/provider/Services/ProviderSessionDirectory.ts index 9dbafd3e804e..d7216543416f 100644 --- a/apps/server/src/provider/Services/ProviderSessionDirectory.ts +++ b/apps/server/src/provider/Services/ProviderSessionDirectory.ts @@ -15,6 +15,8 @@ import type { ProviderValidationError, } from "../Errors.ts"; +import type { SetUsageLimitContinuationInput } from "../usageLimitContinuation.ts"; + export interface ProviderRuntimeBinding { readonly threadId: ThreadId; readonly provider: ProviderDriverKind; @@ -57,6 +59,10 @@ export interface ProviderSessionDirectoryShape { readonly source: AgentSessionImportSource; }) => Effect.Effect; + readonly setUsageLimitContinuation: ( + input: SetUsageLimitContinuationInput, + ) => Effect.Effect; + readonly getProvider: ( threadId: ThreadId, ) => Effect.Effect; diff --git a/apps/server/src/provider/usageLimitContinuation.ts b/apps/server/src/provider/usageLimitContinuation.ts new file mode 100644 index 000000000000..f3921ae0f91e --- /dev/null +++ b/apps/server/src/provider/usageLimitContinuation.ts @@ -0,0 +1,39 @@ +import { + IsoDateTime, + ModelSelection, + ProviderInstanceId, + ThreadId, + TurnId, +} from "@t3tools/contracts"; +import * as Option from "effect/Option"; +import * as Schema from "effect/Schema"; + +export const PendingUsageLimitContinuation = Schema.Struct({ + failedTurnId: TurnId, + providerInstanceId: ProviderInstanceId, + modelSelection: ModelSelection, + errorMessage: Schema.String, + failedAt: IsoDateTime, + nextCheckAt: IsoDateTime, + snapshotSequence: Schema.Int.check(Schema.isGreaterThanOrEqualTo(0)), +}); +export type PendingUsageLimitContinuation = typeof PendingUsageLimitContinuation.Type; + +export const SetUsageLimitContinuationInput = Schema.Struct({ + threadId: ThreadId, + pending: Schema.NullOr(PendingUsageLimitContinuation), + expectedFailedTurnId: Schema.optional(TurnId), +}); +export type SetUsageLimitContinuationInput = typeof SetUsageLimitContinuationInput.Type; + +const decodePayload = Schema.decodeUnknownOption( + Schema.Struct({ + usageLimitContinuation: PendingUsageLimitContinuation, + }), +); + +export function readPendingUsageLimitContinuation( + runtimePayload: unknown, +): PendingUsageLimitContinuation | undefined { + return Option.getOrUndefined(decodePayload(runtimePayload))?.usageLimitContinuation; +} diff --git a/apps/server/src/server.ts b/apps/server/src/server.ts index 3fd0bb7274a0..87f40082d6f7 100644 --- a/apps/server/src/server.ts +++ b/apps/server/src/server.ts @@ -55,6 +55,7 @@ import { ProviderInstanceRegistry } from "./provider/Services/ProviderInstanceRe import { ProviderRegistry } from "./provider/Services/ProviderRegistry.ts"; import { ProviderSessionReaperLive } from "./provider/Layers/ProviderSessionReaper.ts"; import { ProviderUsageLimitsIngestionLive } from "./provider/Layers/ProviderUsageLimitsIngestion.ts"; +import * as UsageLimitContinuation from "./orchestration/UsageLimitContinuation.ts"; import * as OpenCodeRuntime from "./provider/opencodeRuntime.ts"; import * as CheckpointDiffQuery from "./checkpointing/CheckpointDiffQuery.ts"; import * as CheckpointStore from "./checkpointing/CheckpointStore.ts"; @@ -246,6 +247,7 @@ const PlatformServicesLive = NodeServices.layer; const ReactorLayerLive = Layer.empty.pipe( Layer.provideMerge(OrchestrationReactorLive), Layer.provideMerge(ProviderRuntimeIngestionLive), + Layer.provideMerge(UsageLimitContinuation.layer), Layer.provideMerge(ProviderCommandReactorLive), Layer.provideMerge(CheckpointReactorLive), Layer.provideMerge(StorageCleanup.layer), @@ -454,7 +456,7 @@ const ProviderRuntimeLayerLive = ProviderSessionReaperLive.pipe( // telemetry instead of waiting for the next status probe. Layer.provideMerge(ProviderUsageLimitsIngestionLive), Layer.provideMerge(ProviderLayerLive), - Layer.provideMerge(OrchestrationLayerLive), + Layer.provideMerge(OrchestrationLayerLive.pipe(Layer.provide(ServerSettingsLayerLive))), ); const AntigravityInstallationRefreshLive = Layer.effectDiscard( diff --git a/apps/server/src/serverRuntimeStartup.reconcile.test.ts b/apps/server/src/serverRuntimeStartup.reconcile.test.ts index 37fd210ee6da..b733eb7b8871 100644 --- a/apps/server/src/serverRuntimeStartup.reconcile.test.ts +++ b/apps/server/src/serverRuntimeStartup.reconcile.test.ts @@ -152,6 +152,7 @@ it.effect("marks active running sessions that have persisted resume state", () = ), upsert: (binding) => Effect.sync(() => upserts.push(binding)), recordImportedTranscript: () => Effect.die("unused"), + setUsageLimitContinuation: () => Effect.die("unused"), getProvider: () => Effect.die("unused"), listThreadIds: () => Effect.die("unused"), listBindings: () => Effect.succeed([]), @@ -281,6 +282,7 @@ it.effect.each( ), ), recordImportedTranscript: () => Effect.die("unused"), + setUsageLimitContinuation: () => Effect.die("unused"), getProvider: () => Effect.die("unused"), listThreadIds: () => Effect.die("unused"), listBindings: () => Effect.succeed([]), @@ -414,6 +416,7 @@ it.effect("does not continue archived or deleted marked sessions", () => { }, upsert: () => Effect.void, recordImportedTranscript: () => Effect.die("unused"), + setUsageLimitContinuation: () => Effect.die("unused"), getProvider: () => Effect.die("unused"), listThreadIds: () => Effect.die("unused"), listBindings: () => Effect.succeed([]), @@ -470,6 +473,7 @@ it.effect("retries continuation preparation before settling a persistent failure ), upsert: () => Effect.void, recordImportedTranscript: () => Effect.die("unused"), + setUsageLimitContinuation: () => Effect.die("unused"), getProvider: () => Effect.die("unused"), listThreadIds: () => Effect.die("unused"), listBindings: () => Effect.succeed([]), @@ -542,6 +546,7 @@ it.effect("reconciles multiple active and archived orphans but skips live sessio ), upsert: (binding) => Effect.sync(() => upserts.push(binding)), recordImportedTranscript: () => Effect.die("unused"), + setUsageLimitContinuation: () => Effect.die("unused"), getProvider: () => Effect.die("unused"), listThreadIds: () => Effect.die("unused"), listBindings: () => Effect.succeed([]), @@ -622,6 +627,7 @@ it.effect( ), upsert: () => Effect.fail(writeFailure), recordImportedTranscript: () => Effect.die("unused"), + setUsageLimitContinuation: () => Effect.die("unused"), getProvider: () => Effect.die("unused"), listThreadIds: () => Effect.die("unused"), listBindings: () => Effect.succeed([]), @@ -660,6 +666,7 @@ it.effect("retries failed projections and continues after a persistent failure", getBinding: () => Effect.succeed(Option.none()), upsert: () => Effect.void, recordImportedTranscript: () => Effect.die("unused"), + setUsageLimitContinuation: () => Effect.die("unused"), getProvider: () => Effect.die("unused"), listThreadIds: () => Effect.die("unused"), listBindings: () => Effect.succeed([]), @@ -710,6 +717,7 @@ it.effect("does not fail startup when the live provider session inventory cannot getBinding: () => Effect.die("unused"), upsert: () => Effect.die("unused"), recordImportedTranscript: () => Effect.die("unused"), + setUsageLimitContinuation: () => Effect.die("unused"), getProvider: () => Effect.die("unused"), listThreadIds: () => Effect.die("unused"), listBindings: () => Effect.succeed([]), @@ -776,6 +784,7 @@ for (const scenario of [ upserts.push(binding); }), recordImportedTranscript: () => Effect.die("unused"), + setUsageLimitContinuation: () => Effect.die("unused"), getProvider: () => Effect.die("unused"), listThreadIds: () => Effect.die("unused"), listBindings: () => Effect.succeed([]), @@ -850,6 +859,7 @@ for (const preparedStatus of [ yield* Deferred.succeed(cleared, undefined); }), recordImportedTranscript: () => Effect.die("unused"), + setUsageLimitContinuation: () => Effect.die("unused"), getProvider: () => Effect.die("unused"), listThreadIds: () => Effect.die("unused"), listBindings: () => @@ -956,6 +966,7 @@ it.effect("settles failed opt-in recovery without retrying the provider turn", ( binding = next; }), recordImportedTranscript: () => Effect.die("unused"), + setUsageLimitContinuation: () => Effect.die("unused"), getProvider: () => Effect.die("unused"), listThreadIds: () => Effect.die("unused"), listBindings: () => Effect.succeed([]), diff --git a/apps/web/src/components/settings/SettingsPanels.tsx b/apps/web/src/components/settings/SettingsPanels.tsx index 180ab45361df..f60cb1cc10b0 100644 --- a/apps/web/src/components/settings/SettingsPanels.tsx +++ b/apps/web/src/components/settings/SettingsPanels.tsx @@ -9,6 +9,7 @@ import { type BackgroundActivityProfile, type DesktopUpdateChannel, ProviderDriverKind, + PROVIDER_SEND_TURN_MAX_INPUT_CHARS, type ProviderInstanceId, type ScopedThreadRef, type SidebarProjectGroupingMode, @@ -22,6 +23,7 @@ import { import { DEFAULT_ENVIRONMENT_IDENTIFICATION_MODE, DEFAULT_UNIFIED_SETTINGS, + UsageLimitContinuationPrompt, type DiffLayout, type EnvironmentIdentificationMode, MAX_APPEARANCE_CONTRAST, @@ -106,6 +108,7 @@ import { AlertDialogTitle, } from "../ui/alert-dialog"; import { Button } from "../ui/button"; +import { Textarea } from "../ui/textarea"; import { Collapsible, CollapsiblePanel, CollapsibleTrigger } from "../ui/collapsible"; import { Dialog, @@ -601,6 +604,12 @@ export function useSettingsRestore(onRestored?: () => void) { DEFAULT_UNIFIED_SETTINGS.continueThreadsAfterServerUpdate ? ["Continue threads after restarts"] : []), + ...(settings.continueThreadsAfterUsageLimit !== + DEFAULT_UNIFIED_SETTINGS.continueThreadsAfterUsageLimit || + settings.usageLimitContinuationPrompt !== + DEFAULT_UNIFIED_SETTINGS.usageLimitContinuationPrompt + ? ["Continue after usage limits"] + : []), ...(isBackgroundActivityDirty ? ["Background activity"] : []), ...(settings.defaultThreadEnvMode !== DEFAULT_UNIFIED_SETTINGS.defaultThreadEnvMode ? ["New thread mode"] @@ -672,6 +681,8 @@ export function useSettingsRestore(onRestored?: () => void) { settings.responseStreamingMode, settings.enableProviderUpdateChecks, settings.continueThreadsAfterServerUpdate, + settings.continueThreadsAfterUsageLimit, + settings.usageLimitContinuationPrompt, settings.sidebarAutoSettleAfterDays, settings.sidebarAutoSettleOnMerge, settings.sidebarProjectGroupingMode, @@ -776,6 +787,8 @@ export function useSettingsRestore(onRestored?: () => void) { responseStreamingMode: DEFAULT_UNIFIED_SETTINGS.responseStreamingMode, enableProviderUpdateChecks: DEFAULT_UNIFIED_SETTINGS.enableProviderUpdateChecks, continueThreadsAfterServerUpdate: DEFAULT_UNIFIED_SETTINGS.continueThreadsAfterServerUpdate, + continueThreadsAfterUsageLimit: DEFAULT_UNIFIED_SETTINGS.continueThreadsAfterUsageLimit, + usageLimitContinuationPrompt: DEFAULT_UNIFIED_SETTINGS.usageLimitContinuationPrompt, backgroundActivity: DEFAULT_UNIFIED_SETTINGS.backgroundActivity, backgroundActivityProfile: DEFAULT_UNIFIED_SETTINGS.backgroundActivityProfile, automaticGitFetchInterval: DEFAULT_UNIFIED_SETTINGS.automaticGitFetchInterval, @@ -2099,6 +2112,50 @@ function LegacyFeaturesSection() { ); } +const isUsageLimitContinuationPrompt = Schema.is(UsageLimitContinuationPrompt); + +function UsageLimitPromptEditor({ + value, + onSave, +}: { + value: string | null; + onSave: (value: string) => void; +}) { + const [draft, setDraft] = useState(null); + const [error, setError] = useState(null); + return ( +
+