From a3a9fb32139a53133e297e27bb37a48c168a3250 Mon Sep 17 00:00:00 2001 From: Julius Marminge Date: Sat, 19 Sep 2026 23:17:09 -0700 Subject: [PATCH 01/10] feat(v2): resume limited threads at reset --- .../settings/SettingsThreadsRouteScreen.tsx | 15 +- .../features/threads/ThreadDetailScreen.tsx | 6 + .../threads/UsageLimitRecoveryCard.tsx | 66 ++++++ .../src/orchestration-v2/Orchestrator.ts | 61 ++++++ .../src/orchestration-v2/ProjectionStore.ts | 2 + .../UsageLimitRecoveryService.ts | 105 ++++++++++ .../src/orchestration-v2/runtimeLayer.test.ts | 192 ++++++++++++++++++ .../src/orchestration-v2/runtimeLayer.ts | 4 + apps/web/src/components/ChatView.tsx | 19 ++ .../chat/UsageLimitRecoveryCard.tsx | 61 ++++++ .../components/settings/SettingsPanels.tsx | 16 ++ .../src/components/settings/settingsSearch.ts | 6 + docs/user/thread-sidebar.md | 7 + .../client-runtime/src/operations/commands.ts | 5 +- packages/client-runtime/src/state/models.ts | 2 + .../src/state/sharedSettings.test.ts | 1 + .../src/state/sharedSettings.ts | 1 + packages/contracts/src/orchestrationV2.ts | 11 + packages/contracts/src/settings.ts | 2 + 19 files changed, 580 insertions(+), 2 deletions(-) create mode 100644 apps/mobile/src/features/threads/UsageLimitRecoveryCard.tsx create mode 100644 apps/server/src/orchestration-v2/UsageLimitRecoveryService.ts create mode 100644 apps/web/src/components/chat/UsageLimitRecoveryCard.tsx diff --git a/apps/mobile/src/features/settings/SettingsThreadsRouteScreen.tsx b/apps/mobile/src/features/settings/SettingsThreadsRouteScreen.tsx index bd8ec3252155..eef44fa9707e 100644 --- a/apps/mobile/src/features/settings/SettingsThreadsRouteScreen.tsx +++ b/apps/mobile/src/features/settings/SettingsThreadsRouteScreen.tsx @@ -86,7 +86,9 @@ function AutoSettleSettingsRows() { return null; } - const writeToAll = (patch: Partial) => { + const writeToAll = ( + patch: Partial & { autoResumeLimitedThreads?: boolean }, + ) => { if (writeInFlight.current) return; const writes = planMobileScopedSettingsPatch(syncTargets, projectSelected, patch); if (writes.length === 0) return; @@ -164,6 +166,17 @@ function AutoSettleSettingsRows() { onClear={clearProjectOverrides} /> ) : null} + {!projectSelected ? ( + + writeToAll({ autoResumeLimitedThreads: value })} + /> + + ) : null} ) : null} + {props.feedbackSubmissions.map((submission) => ( Date.parse(thread.latestRun?.completedAt ?? thread.updatedAt); + const runId = thread.latestRun?.runId; + const recovery = thread.limitRecovery; + const scheduled = + recovery?.runId === runId && recovery?.resetAt === resetAt && recovery?.autoResume; + if ( + thread.runtime?.status !== "failed" || + thread.runtime.lastErrorClass !== "usage_limit" || + !runId + ) + return null; + async function toggle() { + if (!resetAt || !runId || !canSchedule) return; + setPending(true); + try { + await updateMetadata({ + environmentId, + input: { threadId: thread.id, limitRecovery: { runId, resetAt, autoResume: !scheduled } }, + }); + } finally { + setPending(false); + } + } + return ( + + + {resetAt + ? `Usage limit resets ${DateTime.toDateUtc(DateTime.makeUnsafe(resetAt)).toLocaleString()}.` + : "The provider did not report a reset time. Retry manually when your limit is available."} + + {canSchedule ? ( + void toggle()} + className="self-start rounded-lg bg-subtle px-3 py-2 active:opacity-70" + > + + {scheduled ? "Cancel auto-resume" : "Resume at reset"} + + + ) : null} + + ); +} diff --git a/apps/server/src/orchestration-v2/Orchestrator.ts b/apps/server/src/orchestration-v2/Orchestrator.ts index b6119d03ff85..09f113941175 100644 --- a/apps/server/src/orchestration-v2/Orchestrator.ts +++ b/apps/server/src/orchestration-v2/Orchestrator.ts @@ -1,3 +1,4 @@ +import { latestRootProviderFailure } from "@t3tools/shared/orchestrationV2ThreadError"; import { threadPullRequestsOf } from "@t3tools/shared/threadPullRequests"; import { normalizeThreadPullRequestKey, @@ -70,6 +71,7 @@ import { import { applyToProjection, emptyProjection, + threadShellFromProjection, isTurnItemAtOrBeforeRun, ProjectionStoreV2, type ProjectionCheckpointContext, @@ -2225,6 +2227,28 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio : null; const now = yield* DateTime.now; + if (command.type === "thread.metadata.update" && command.limitRecovery != null) { + const projection = yield* loadProjectionForCommand(command); + const run = projection.runs.at(-1) ?? null; + const failure = latestRootProviderFailure(run, projection.turnItems); + if ( + thread.archivedAt !== null || + thread.settledOverride === "settled" || + run?.id !== command.limitRecovery.runId || + failure?.class !== "usage_limit" || + failure.resetAt !== command.limitRecovery.resetAt || + Date.parse(command.limitRecovery.resetAt) <= + DateTime.toEpochMillis(run.completedAt ?? run.requestedAt) || + projection.runtimeRequests.some((request) => request.status === "pending") || + projection.runs.some((candidate) => candidate.status === "queued") + ) { + return yield* new OrchestratorDispatchError({ + commandId: command.commandId, + commandType: command.type, + cause: "The provider limit changed before recovery could be configured.", + }); + } + } let snoozedUntil: DateTime.Utc | null = null; if (command.type === "thread.snooze") { const projection = yield* loadProjectionForCommand(command); @@ -2382,6 +2406,9 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio return { ...thread, ...(command.title === undefined ? {} : { title: command.title }), + ...(command.limitRecovery === undefined + ? {} + : { limitRecovery: command.limitRecovery }), ...(command.branch === undefined ? {} : { branch: command.branch }), ...(command.worktreePath === undefined ? {} : { worktreePath: command.worktreePath }), ...(command.linkedPullRequest === undefined @@ -3754,6 +3781,40 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio ) => Effect.gen(function* () { let projection = yield* getProjectionWithPendingEvents(command.threadId, events); + if (command.usageLimitContinuationOfRunId !== undefined) { + const run = projection.runs.at(-1) ?? null; + const failure = latestRootProviderFailure(run, projection.turnItems); + const recovery = projection.thread.limitRecovery; + const now = yield* DateTime.now; + if ( + run?.id !== command.usageLimitContinuationOfRunId || + failure?.class !== "usage_limit" || + threadShellFromProjection(projection).lastErrorClass !== "usage_limit" || + !recovery?.autoResume || + recovery.runId !== run.id || + recovery.resetAt !== failure.resetAt || + Date.parse(recovery.resetAt) > DateTime.toEpochMillis(now) || + projection.thread.archivedAt !== null || + projection.thread.deletedAt !== null || + projection.thread.settledOverride === "settled" || + projection.thread.providerInstanceId !== run.providerInstanceId || + projection.runtimeRequests.some((request) => request.status === "pending") || + (projection.thread.snoozedUntil != null && + DateTime.toEpochMillis(projection.thread.snoozedUntil) > DateTime.toEpochMillis(now)) + ) { + yield* emit( + events, + command, + )({ + type: "thread.metadata-updated", + threadId: command.threadId, + occurredAt: now, + payload: projection.thread, + }); + return; + } + } + if (command.restartContinuationOfRunId !== undefined) { const source = projection.runs.find((run) => run.id === command.restartContinuationOfRunId); if ( diff --git a/apps/server/src/orchestration-v2/ProjectionStore.ts b/apps/server/src/orchestration-v2/ProjectionStore.ts index 64ca61890156..a8e90be26d38 100644 --- a/apps/server/src/orchestration-v2/ProjectionStore.ts +++ b/apps/server/src/orchestration-v2/ProjectionStore.ts @@ -1274,6 +1274,7 @@ export function threadShellFromProjection( pinOrderKey: projection.thread.pinOrderKey ?? null, lastVisitedAt: projection.thread.lastVisitedAt, titleRegeneration: projection.thread.titleRegeneration ?? null, + limitRecovery: projection.thread.limitRecovery ?? null, deletedAt: projection.thread.deletedAt, }; } @@ -1495,6 +1496,7 @@ function shellFromState(input: { pinOrderKey: input.state.thread.pinOrderKey ?? null, lastVisitedAt: input.state.thread.lastVisitedAt, titleRegeneration: input.state.thread.titleRegeneration ?? null, + limitRecovery: input.state.thread.limitRecovery ?? null, deletedAt: input.state.thread.deletedAt, }; } diff --git a/apps/server/src/orchestration-v2/UsageLimitRecoveryService.ts b/apps/server/src/orchestration-v2/UsageLimitRecoveryService.ts new file mode 100644 index 000000000000..c6a353866e76 --- /dev/null +++ b/apps/server/src/orchestration-v2/UsageLimitRecoveryService.ts @@ -0,0 +1,105 @@ +import { + CommandId, + MessageId, + type OrchestrationV2ThreadShell, + type OrchestrationV2Command, +} from "@t3tools/contracts"; +import * as DateTime from "effect/DateTime"; +import * as Effect from "effect/Effect"; +import * as Layer from "effect/Layer"; +import * as Schedule from "effect/Schedule"; +import { ServerSettingsService } from "../serverSettings.ts"; +import { ProjectionStoreV2 } from "./ProjectionStore.ts"; +import { ThreadManagementService } from "./ThreadManagementService.ts"; + +/** The persisted run and reset form the identity of one recovery opportunity. */ +export function limitRecoveryCommand( + thread: OrchestrationV2ThreadShell, + autoResume: boolean, + nowMs: number, +): OrchestrationV2Command | null { + if ( + thread.status !== "failed" || + thread.lastErrorClass !== "usage_limit" || + !thread.latestRunId || + !thread.usageLimitResetAt || + thread.archivedAt !== null || + thread.settledOverride === "settled" || + thread.pendingRuntimeRequest !== null + ) + return null; + const resetMs = Date.parse(thread.usageLimitResetAt); + // An already-expired window reported with a fresh failure cannot start a retry loop. + if ( + !Number.isFinite(resetMs) || + resetMs <= DateTime.toEpochMillis(thread.latestRunCompletedAt ?? thread.updatedAt) + ) + return null; + const identity = `${thread.id}:${thread.latestRunId}:${resetMs}`; + const recovery = thread.limitRecovery; + if (recovery?.runId !== thread.latestRunId || recovery.resetAt !== thread.usageLimitResetAt) { + if (!autoResume) return null; + return { + type: "thread.metadata.update", + commandId: CommandId.make(`limit-arm:${identity}`), + threadId: thread.id, + limitRecovery: { runId: thread.latestRunId, resetAt: thread.usageLimitResetAt, autoResume }, + }; + } + if ( + !recovery.autoResume || + resetMs > nowMs || + (thread.snoozedUntil != null && DateTime.toEpochMillis(thread.snoozedUntil) > nowMs) + ) + return null; + return { + type: "message.dispatch", + commandId: CommandId.make(`limit-resume:${identity}`), + messageId: MessageId.make(`limit-resume:${identity}`), + threadId: thread.id, + usageLimitContinuationOfRunId: thread.latestRunId, + text: "Continue where you left off.", + attachments: [], + dispatchMode: { type: "start_immediately" }, + createdBy: "user", + creationSource: "server", + }; +} + +export const make = Effect.gen(function* () { + const projections = yield* ProjectionStoreV2; + const threads = yield* ThreadManagementService; + const settings = yield* ServerSettingsService; + return Effect.fn("UsageLimitRecoveryService.sweep")(function* () { + const preferences = yield* settings.getSettings; + const snapshot = yield* projections.getShellSnapshot(); + const nowMs = DateTime.toEpochMillis(yield* DateTime.now); + for (const thread of snapshot.threads) { + const command = limitRecoveryCommand(thread, preferences.autoResumeLimitedThreads, nowMs); + if (command === null) continue; + yield* threads.dispatch(command).pipe( + Effect.catchCause((cause) => + Effect.logWarning("orchestration-v2.limit-recovery.dispatch-failed", { + threadId: thread.id, + cause, + }), + ), + ); + } + }); +}); + +// The schedule is derived from persisted failures and thread recovery choices, +// so restarts need no timer restoration and disconnected clients need not run it. +export const workerLive = Layer.effectDiscard( + Effect.gen(function* () { + const sweep = yield* make; + yield* sweep().pipe( + Effect.catchCause((cause) => + Effect.logWarning("orchestration-v2.limit-recovery.sweep-failed", { cause }), + ), + Effect.repeat(Schedule.spaced("30 seconds")), + Effect.forkScoped, + ); + }), +); diff --git a/apps/server/src/orchestration-v2/runtimeLayer.test.ts b/apps/server/src/orchestration-v2/runtimeLayer.test.ts index f3adb7418846..5378ea488280 100644 --- a/apps/server/src/orchestration-v2/runtimeLayer.test.ts +++ b/apps/server/src/orchestration-v2/runtimeLayer.test.ts @@ -1,3 +1,4 @@ +import { limitRecoveryCommand } from "./UsageLimitRecoveryService.ts"; import { SourceControlProviderRegistry } from "../sourceControl/SourceControlProviderRegistry.ts"; import * as NodeServices from "@effect/platform-node/NodeServices"; import { assert, it } from "@effect/vitest"; @@ -3014,3 +3015,194 @@ it.layer(SharedApplicationDataPlaneTestLayer)("shared application data plane", ( }), ); }); + +it.layer(TestLayer)("usage-limit recovery", (it) => { + it.effect.each(["resume", "cancel", "new-message", "archive", "settle", "replacement"] as const)( + "guards a scheduled usage-limit continuation against %s", + (scenario) => + Effect.gen(function* () { + const orchestrator = yield* OrchestratorV2; + const events = yield* EventSinkV2; + const threadId = ThreadId.make(`recovery:${scenario}`); + const projects = yield* ProjectionProjectRepository; + const projectId = ProjectId.make(`recovery:project:${scenario}`); + const projectAt = DateTime.formatIso(yield* DateTime.now); + yield* projects.upsert({ + projectId, + title: "Recovery project", + workspaceRoot: process.cwd(), + defaultModelSelection: modelSelection, + defaultThreadEnvMode: null, + autoPull: false, + scripts: [], + createdAt: projectAt, + updatedAt: projectAt, + deletedAt: null, + }); + yield* orchestrator.dispatch({ + type: "thread.create", + commandId: CommandId.make(`recovery:create:${scenario}`), + threadId, + projectId, + title: "Limited thread", + modelSelection, + runtimeMode: "full-access", + interactionMode: "default", + branch: null, + worktreePath: null, + createdBy: "user", + creationSource: "web", + }); + yield* orchestrator.dispatch({ + type: "message.dispatch", + commandId: CommandId.make(`recovery:message:${scenario}`), + threadId, + messageId: MessageId.make(`recovery:message:${scenario}`), + text: "Work on this.", + attachments: [], + dispatchMode: { type: "defer_start" }, + createdBy: "user", + creationSource: "web", + }); + const projection = yield* orchestrator.getThreadProjection(threadId); + const run = projection.runs[0]!; + const now = yield* DateTime.now; + const resetAt = DateTime.formatIso(DateTime.add(now, { minutes: 1 })); + yield* events.write({ + commandId: CommandId.make(`recovery:failure:${scenario}`), + events: [ + { + id: EventId.make(`recovery:run:${scenario}`), + type: "run.updated", + threadId, + occurredAt: now, + payload: { ...run, status: "failed", completedAt: now }, + }, + { + id: EventId.make(`recovery:error:${scenario}`), + type: "turn-item.updated", + threadId, + occurredAt: now, + payload: { + id: TurnItemId.make(`recovery:error:${scenario}`), + type: "error", + threadId, + runId: run.id, + nodeId: run.rootNodeId, + providerThreadId: null, + providerTurnId: null, + nativeItemRef: null, + parentItemId: null, + ordinal: 2, + status: "failed", + title: "Usage limit reached", + startedAt: now, + completedAt: now, + updatedAt: now, + failure: { + class: "usage_limit", + message: "Plan limit reached.", + code: "usageLimitExceeded", + retryable: null, + resetAt, + }, + }, + }, + ], + }); + const shell = (yield* orchestrator.getShellSnapshot()).threads.find( + (thread) => thread.id === threadId, + )!; + assert.isNull(limitRecoveryCommand(shell, false, DateTime.toEpochMillis(now))); + const arm = limitRecoveryCommand(shell, true, DateTime.toEpochMillis(now)); + assert.isNotNull(arm); + yield* orchestrator.dispatch(arm!); + const armedShell = (yield* orchestrator.getShellSnapshot()).threads.find( + (thread) => thread.id === threadId, + )!; + assert.deepEqual(armedShell.limitRecovery, { runId: run.id, resetAt, autoResume: true }); + assert.isNull(limitRecoveryCommand(armedShell, true, DateTime.toEpochMillis(now))); + yield* orchestrator.dispatch({ + type: "message.dispatch", + commandId: CommandId.make(`recovery:early:${scenario}`), + messageId: MessageId.make(`recovery:early:${scenario}`), + threadId, + usageLimitContinuationOfRunId: run.id, + text: "Continue where you left off.", + attachments: [], + dispatchMode: { type: "start_immediately" }, + createdBy: "user", + creationSource: "server", + }); + assert.lengthOf((yield* orchestrator.getThreadProjection(threadId)).runs, 1); + yield* TestClock.adjust("1 minute"); + const resume = limitRecoveryCommand( + armedShell, + true, + DateTime.toEpochMillis(yield* DateTime.now), + ); + assert.isNotNull(resume); + if (scenario === "cancel") + yield* orchestrator.dispatch({ + type: "thread.metadata.update", + commandId: CommandId.make(`recovery:cancel:${scenario}`), + threadId, + limitRecovery: { runId: run.id, resetAt, autoResume: false }, + }); + if (scenario === "archive") + yield* orchestrator.dispatch({ + type: "thread.archive", + commandId: CommandId.make(`recovery:archive:${scenario}`), + threadId, + }); + if (scenario === "new-message") + yield* orchestrator.dispatch({ + type: "message.dispatch", + commandId: CommandId.make(`recovery:new-message:${scenario}`), + threadId, + messageId: MessageId.make(`recovery:new-message:${scenario}`), + text: "I will continue manually.", + attachments: [], + dispatchMode: { type: "defer_start" }, + createdBy: "user", + creationSource: "web", + }); + if (scenario === "settle") + yield* orchestrator.dispatch({ + type: "thread.settle", + commandId: CommandId.make(`recovery:settle:${scenario}`), + threadId, + }); + if (scenario === "replacement") { + const current = yield* orchestrator.getThreadProjection(threadId); + const error = current.turnItems.find((item) => item.type === "error")!; + if (error.type !== "error") throw new Error("Expected provider error"); + yield* events.write({ + commandId: CommandId.make(`recovery:replacement:${scenario}`), + events: [ + { + id: EventId.make(`recovery:replacement:${scenario}`), + type: "turn-item.updated", + threadId, + occurredAt: yield* DateTime.now, + payload: { + ...error, + failure: { + ...error.failure, + class: "provider_error", + message: "A replacement failure.", + }, + }, + }, + ], + }); + } + const before = yield* orchestrator.getThreadProjection(threadId); + yield* orchestrator.dispatch(resume!); + yield* orchestrator.dispatch(resume!); + const after = yield* orchestrator.getThreadProjection(threadId); + assert.lengthOf(after.runs, before.runs.length + (scenario === "resume" ? 1 : 0)); + assert.lengthOf(after.messages, before.messages.length + (scenario === "resume" ? 1 : 0)); + }), + ); +}); diff --git a/apps/server/src/orchestration-v2/runtimeLayer.ts b/apps/server/src/orchestration-v2/runtimeLayer.ts index 5fd7579fb22f..4400151b963b 100644 --- a/apps/server/src/orchestration-v2/runtimeLayer.ts +++ b/apps/server/src/orchestration-v2/runtimeLayer.ts @@ -1,3 +1,4 @@ +import { workerLive as usageLimitRecoveryWorkerLive } from "./UsageLimitRecoveryService.ts"; import * as Layer from "effect/Layer"; import { OrchestrationEventInfrastructureLayerLive, @@ -295,6 +296,9 @@ export const OrchestrationV2ProductionLayerLive = Layer.mergeAll( threadLaunchProvided, threadLifecycleProvided, scheduledTaskProvided, + usageLimitRecoveryWorkerLive.pipe( + Layer.provide(Layer.mergeAll(projectionStoreLayer, threadManagementProvided)), + ), providerContinuationWorkerProvided, agentSessionImporterProvided, ).pipe(Layer.provideMerge(OrchestrationLayerLive)); diff --git a/apps/web/src/components/ChatView.tsx b/apps/web/src/components/ChatView.tsx index f2a43c9dfb21..36b4d720d8f6 100644 --- a/apps/web/src/components/ChatView.tsx +++ b/apps/web/src/components/ChatView.tsx @@ -1,3 +1,4 @@ +import { UsageLimitRecoveryCard } from "./chat/UsageLimitRecoveryCard"; import { resolveBackgroundDraftWorkspaceOptions, resolveDraftHeroState, @@ -10355,6 +10356,24 @@ export default function ChatView(props: ChatViewProps) { onDismiss={() => setDismissedProviderStatusBannerKey(providerStatusBannerKey)} onOpenProviderSetup={openProviderSetup} /> + {serverRuntime?.status === "failed" && + serverRuntime.lastErrorClass === "usage_limit" && + activeThreadShell?.latestRun ? ( + { + const result = await updateThreadMetadata({ + environmentId, + input: { threadId: activeThread.id, limitRecovery }, + }); + if (result._tag === "Failure") throw squashAtomCommandFailure(result); + }} + /> + ) : null} Promise; +}) { + const [pending, setPending] = useState(false); + const [error, setError] = useState(null); + const canSchedule = resetAt !== null && Date.parse(resetAt) > Date.parse(stoppedAt); + const scheduled = + recovery?.runId === runId && recovery.resetAt === resetAt && recovery.autoResume; + async function toggle() { + if (resetAt === null || !canSchedule) return; + setPending(true); + setError(null); + try { + await onChange({ runId, resetAt, autoResume: !scheduled }); + } catch (cause) { + setError(cause instanceof Error ? cause.message : "Could not change limit recovery."); + } finally { + setPending(false); + } + } + return ( +
+

+ {resetAt + ? `Usage limit resets ${DateTime.toDateUtc(DateTime.makeUnsafe(resetAt)).toLocaleString()}.` + : "The provider did not report a reset time. Retry manually when your limit is available."} +

+ {canSchedule ? ( + + ) : null} + {error ? ( +

+ {error} +

+ ) : null} +
+ ); +} diff --git a/apps/web/src/components/settings/SettingsPanels.tsx b/apps/web/src/components/settings/SettingsPanels.tsx index 36220f3bc143..2adfdd67817a 100644 --- a/apps/web/src/components/settings/SettingsPanels.tsx +++ b/apps/web/src/components/settings/SettingsPanels.tsx @@ -2228,6 +2228,22 @@ export function GeneralSettingsPanel() { } /> + + updateSettings({ autoResumeLimitedThreads: Boolean(checked) }) + } + aria-label="Auto-resume limited threads" + /> + } + /> {supportsAutoSettlement ? ( <> { expect( Object.keys(pickSharedServerSettings(DEFAULT_SERVER_SETTINGS, restartCapabilities)).sort(), ).toEqual([ + "autoResumeLimitedThreads", "continueThreadsAfterServerUpdate", "newWorktreesStartFromOrigin", "sidebarAutoSettleAfterDays", diff --git a/packages/client-runtime/src/state/sharedSettings.ts b/packages/client-runtime/src/state/sharedSettings.ts index 9845764d5e07..ed8d36276e6a 100644 --- a/packages/client-runtime/src/state/sharedSettings.ts +++ b/packages/client-runtime/src/state/sharedSettings.ts @@ -25,6 +25,7 @@ const SHARED_SERVER_SETTING_KEYS = [ "continueThreadsAfterServerUpdate", "sidebarAutoSettleAfterDays", "sidebarAutoSettleOnMerge", + "autoResumeLimitedThreads", "newWorktreesStartFromOrigin", "sourceControlWritingStyle", "textGenerationModelSelection", diff --git a/packages/contracts/src/orchestrationV2.ts b/packages/contracts/src/orchestrationV2.ts index 8960f6f85732..c9c6681dd999 100644 --- a/packages/contracts/src/orchestrationV2.ts +++ b/packages/contracts/src/orchestrationV2.ts @@ -326,6 +326,13 @@ export const OrchestrationV2ProviderCapabilities = Schema.Struct({ }); export type OrchestrationV2ProviderCapabilities = typeof OrchestrationV2ProviderCapabilities.Type; +export const OrchestrationV2LimitRecovery = Schema.Struct({ + runId: RunId, + resetAt: IsoDateTime, + autoResume: Schema.Boolean, +}); +export type OrchestrationV2LimitRecovery = typeof OrchestrationV2LimitRecovery.Type; + export const OrchestrationV2AppThread = Schema.Struct({ ...OrchestrationV2CreationFields, id: ThreadId, @@ -369,6 +376,7 @@ export const OrchestrationV2AppThread = Schema.Struct({ unsettledAt: Schema.optional(Schema.NullOr(Schema.DateTimeUtc)), snoozedUntil: Schema.optional(Schema.NullOr(Schema.DateTimeUtc)), snoozedAt: Schema.optional(Schema.NullOr(Schema.DateTimeUtc)), + limitRecovery: Schema.optional(Schema.NullOr(OrchestrationV2LimitRecovery)), pinnedAt: Schema.optional(Schema.NullOr(Schema.DateTimeUtc)), // Fractional-index slot in the user-arranged pinned order. Optional so // payloads from pre-reorder servers still decode. @@ -1516,6 +1524,7 @@ export const OrchestrationV2ThreadShell = Schema.Struct({ unsettledAt: Schema.optional(Schema.NullOr(Schema.DateTimeUtc)), snoozedUntil: Schema.optional(Schema.NullOr(Schema.DateTimeUtc)), snoozedAt: Schema.optional(Schema.NullOr(Schema.DateTimeUtc)), + limitRecovery: Schema.optional(Schema.NullOr(OrchestrationV2LimitRecovery)), /** Omitted by servers that predate thread pinning. */ pinnedAt: Schema.optional(Schema.NullOr(Schema.DateTimeUtc)), /** Slot in the user-arranged pinned order; omitted by pre-reorder servers. */ @@ -2329,6 +2338,7 @@ export const OrchestrationV2Command = Schema.Union([ expectedWorktreePath: Schema.optional(Schema.NullOr(TrimmedNonEmptyString)), /** Reject unless no message or run has landed on this thread. */ expectedEmpty: Schema.optional(Schema.Boolean), + limitRecovery: Schema.optional(Schema.NullOr(OrchestrationV2LimitRecovery)), /** Link (object) or unlink (null) a pull request (#8160); absent leaves it unchanged. */ linkedPullRequest: Schema.optional(Schema.NullOr(ThreadLinkedPullRequest)), }), @@ -2418,6 +2428,7 @@ export const OrchestrationV2Command = Schema.Union([ modelSelection: Schema.optional(ModelSelection), sourcePlanRef: Schema.optional(Schema.Struct({ threadId: ThreadId, planId: PlanId })), restartContinuationOfRunId: Schema.optional(RunId), + usageLimitContinuationOfRunId: Schema.optional(RunId), /** Resolve untargeted delivery against the server's serialized thread state. */ deliveryIntent: Schema.optional(Schema.Literals(["auto", "steer", "restart"])), delegatedCompletion: Schema.optional( diff --git a/packages/contracts/src/settings.ts b/packages/contracts/src/settings.ts index fd341e7ff0b4..a3821686b925 100644 --- a/packages/contracts/src/settings.ts +++ b/packages/contracts/src/settings.ts @@ -1200,6 +1200,7 @@ export const ServerSettings = Schema.Struct({ sidebarAutoSettleAfterDays: Schema.NullOr(SidebarAutoSettleAfterDays).pipe( Schema.withDecodingDefault(Effect.succeed(DEFAULT_SIDEBAR_AUTO_SETTLE_AFTER_DAYS)), ), + autoResumeLimitedThreads: Schema.Boolean.pipe(Schema.withDecodingDefault(Effect.succeed(false))), sidebarAutoSettleOnMerge: Schema.Boolean.pipe(Schema.withDecodingDefault(Effect.succeed(true))), backgroundActivity: BackgroundActivitySettings, // Legacy flat fields retained for old settings files and old clients. New @@ -1530,6 +1531,7 @@ export const ServerSettingsPatch = Schema.Struct({ deviceHosts: Schema.optionalKey(SshDeviceHostConfigs), sidebarAutoSettleAfterDays: Schema.optionalKey(Schema.NullOr(SidebarAutoSettleAfterDays)), sidebarAutoSettleOnMerge: Schema.optionalKey(Schema.Boolean), + autoResumeLimitedThreads: Schema.optionalKey(Schema.Boolean), backgroundActivity: Schema.optionalKey( Schema.Struct({ schemaVersion: Schema.optionalKey(Schema.Literal(1)), From 67f4c9002cb70e278e3b636ad22f05fb7fbd91ae Mon Sep 17 00:00:00 2001 From: Julius Marminge Date: Sat, 19 Sep 2026 23:31:44 -0700 Subject: [PATCH 02/10] fix(web): keep limit recovery above the composer --- apps/web/src/components/ChatView.tsx | 38 ++++++++++--------- .../chat/UsageLimitRecoveryCard.tsx | 2 +- 2 files changed, 21 insertions(+), 19 deletions(-) diff --git a/apps/web/src/components/ChatView.tsx b/apps/web/src/components/ChatView.tsx index 36b4d720d8f6..deeb547e4ff1 100644 --- a/apps/web/src/components/ChatView.tsx +++ b/apps/web/src/components/ChatView.tsx @@ -10356,24 +10356,6 @@ export default function ChatView(props: ChatViewProps) { onDismiss={() => setDismissedProviderStatusBannerKey(providerStatusBannerKey)} onOpenProviderSetup={openProviderSetup} /> - {serverRuntime?.status === "failed" && - serverRuntime.lastErrorClass === "usage_limit" && - activeThreadShell?.latestRun ? ( - { - const result = await updateThreadMetadata({ - environmentId, - input: { threadId: activeThread.id, limitRecovery }, - }); - if (result._tag === "Failure") throw squashAtomCommandFailure(result); - }} - /> - ) : null} + {serverRuntime?.status === "failed" && + serverRuntime.lastErrorClass === "usage_limit" && + activeThreadShell?.latestRun ? ( + { + const result = await updateThreadMetadata({ + environmentId, + input: { threadId: activeThread.id, limitRecovery }, + }); + if (result._tag === "Failure") throw squashAtomCommandFailure(result); + }} + /> + ) : null} {isDraftHeroState ? (
+

{resetAt ? `Usage limit resets ${DateTime.toDateUtc(DateTime.makeUnsafe(resetAt)).toLocaleString()}.` From 4be5d817acfb1a7192d5caf5c5e72b9ddd6dddae Mon Sep 17 00:00:00 2001 From: Julius Marminge Date: Sat, 19 Sep 2026 23:37:09 -0700 Subject: [PATCH 03/10] fix(v2): invalidate cancelled recovery deliveries --- .../src/orchestration-v2/Orchestrator.ts | 8 +- .../UsageLimitRecoveryService.ts | 10 +- .../src/orchestration-v2/runtimeLayer.test.ts | 347 ++++++++++-------- .../components/settings/SettingsPanels.tsx | 4 +- docs/user/thread-sidebar.md | 2 +- packages/contracts/src/orchestrationV2.ts | 2 + 6 files changed, 210 insertions(+), 163 deletions(-) diff --git a/apps/server/src/orchestration-v2/Orchestrator.ts b/apps/server/src/orchestration-v2/Orchestrator.ts index 09f113941175..e3a1fd32b1f6 100644 --- a/apps/server/src/orchestration-v2/Orchestrator.ts +++ b/apps/server/src/orchestration-v2/Orchestrator.ts @@ -2408,7 +2408,12 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio ...(command.title === undefined ? {} : { title: command.title }), ...(command.limitRecovery === undefined ? {} - : { limitRecovery: command.limitRecovery }), + : { + limitRecovery: + command.limitRecovery === null + ? null + : { ...command.limitRecovery, requestId: command.commandId }, + }), ...(command.branch === undefined ? {} : { branch: command.branch }), ...(command.worktreePath === undefined ? {} : { worktreePath: command.worktreePath }), ...(command.linkedPullRequest === undefined @@ -3791,6 +3796,7 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio failure?.class !== "usage_limit" || threadShellFromProjection(projection).lastErrorClass !== "usage_limit" || !recovery?.autoResume || + recovery.requestId !== command.usageLimitRecoveryRequestId || recovery.runId !== run.id || recovery.resetAt !== failure.resetAt || Date.parse(recovery.resetAt) > DateTime.toEpochMillis(now) || diff --git a/apps/server/src/orchestration-v2/UsageLimitRecoveryService.ts b/apps/server/src/orchestration-v2/UsageLimitRecoveryService.ts index c6a353866e76..4d6cb5e30c85 100644 --- a/apps/server/src/orchestration-v2/UsageLimitRecoveryService.ts +++ b/apps/server/src/orchestration-v2/UsageLimitRecoveryService.ts @@ -52,12 +52,16 @@ export function limitRecoveryCommand( (thread.snoozedUntil != null && DateTime.toEpochMillis(thread.snoozedUntil) > nowMs) ) return null; + const deliveryIdentity = `${identity}:${recovery.requestId ?? "legacy"}`; return { type: "message.dispatch", - commandId: CommandId.make(`limit-resume:${identity}`), - messageId: MessageId.make(`limit-resume:${identity}`), + commandId: CommandId.make(`limit-resume:${deliveryIdentity}`), + messageId: MessageId.make(`limit-resume:${deliveryIdentity}`), threadId: thread.id, usageLimitContinuationOfRunId: thread.latestRunId, + ...(recovery.requestId === undefined + ? {} + : { usageLimitRecoveryRequestId: recovery.requestId }), text: "Continue where you left off.", attachments: [], dispatchMode: { type: "start_immediately" }, @@ -66,7 +70,7 @@ export function limitRecoveryCommand( }; } -export const make = Effect.gen(function* () { +const make = Effect.gen(function* () { const projections = yield* ProjectionStoreV2; const threads = yield* ThreadManagementService; const settings = yield* ServerSettingsService; diff --git a/apps/server/src/orchestration-v2/runtimeLayer.test.ts b/apps/server/src/orchestration-v2/runtimeLayer.test.ts index 5378ea488280..bafc48a7f28d 100644 --- a/apps/server/src/orchestration-v2/runtimeLayer.test.ts +++ b/apps/server/src/orchestration-v2/runtimeLayer.test.ts @@ -3017,192 +3017,227 @@ it.layer(SharedApplicationDataPlaneTestLayer)("shared application data plane", ( }); it.layer(TestLayer)("usage-limit recovery", (it) => { - it.effect.each(["resume", "cancel", "new-message", "archive", "settle", "replacement"] as const)( - "guards a scheduled usage-limit continuation against %s", - (scenario) => - Effect.gen(function* () { - const orchestrator = yield* OrchestratorV2; - const events = yield* EventSinkV2; - const threadId = ThreadId.make(`recovery:${scenario}`); - const projects = yield* ProjectionProjectRepository; - const projectId = ProjectId.make(`recovery:project:${scenario}`); - const projectAt = DateTime.formatIso(yield* DateTime.now); - yield* projects.upsert({ - projectId, - title: "Recovery project", - workspaceRoot: process.cwd(), - defaultModelSelection: modelSelection, - defaultThreadEnvMode: null, - autoPull: false, - scripts: [], - createdAt: projectAt, - updatedAt: projectAt, - deletedAt: null, + it.effect.each([ + "resume", + "cancel", + "rearm", + "new-message", + "archive", + "settle", + "replacement", + ] as const)("guards a scheduled usage-limit continuation against %s", (scenario) => + Effect.gen(function* () { + const orchestrator = yield* OrchestratorV2; + const events = yield* EventSinkV2; + const threadId = ThreadId.make(`recovery:${scenario}`); + const projects = yield* ProjectionProjectRepository; + const projectId = ProjectId.make(`recovery:project:${scenario}`); + const projectAt = DateTime.formatIso(yield* DateTime.now); + yield* projects.upsert({ + projectId, + title: "Recovery project", + workspaceRoot: process.cwd(), + defaultModelSelection: modelSelection, + defaultThreadEnvMode: null, + autoPull: false, + scripts: [], + createdAt: projectAt, + updatedAt: projectAt, + deletedAt: null, + }); + yield* orchestrator.dispatch({ + type: "thread.create", + commandId: CommandId.make(`recovery:create:${scenario}`), + threadId, + projectId, + title: "Limited thread", + modelSelection, + runtimeMode: "full-access", + interactionMode: "default", + branch: null, + worktreePath: null, + createdBy: "user", + creationSource: "web", + }); + yield* orchestrator.dispatch({ + type: "message.dispatch", + commandId: CommandId.make(`recovery:message:${scenario}`), + threadId, + messageId: MessageId.make(`recovery:message:${scenario}`), + text: "Work on this.", + attachments: [], + dispatchMode: { type: "defer_start" }, + createdBy: "user", + creationSource: "web", + }); + const projection = yield* orchestrator.getThreadProjection(threadId); + const run = projection.runs[0]!; + const now = yield* DateTime.now; + const resetAt = DateTime.formatIso(DateTime.add(now, { minutes: 1 })); + yield* events.write({ + commandId: CommandId.make(`recovery:failure:${scenario}`), + events: [ + { + id: EventId.make(`recovery:run:${scenario}`), + type: "run.updated", + threadId, + occurredAt: now, + payload: { ...run, status: "failed", completedAt: now }, + }, + { + id: EventId.make(`recovery:error:${scenario}`), + type: "turn-item.updated", + threadId, + occurredAt: now, + payload: { + id: TurnItemId.make(`recovery:error:${scenario}`), + type: "error", + threadId, + runId: run.id, + nodeId: run.rootNodeId, + providerThreadId: null, + providerTurnId: null, + nativeItemRef: null, + parentItemId: null, + ordinal: 2, + status: "failed", + title: "Usage limit reached", + startedAt: now, + completedAt: now, + updatedAt: now, + failure: { + class: "usage_limit", + message: "Plan limit reached.", + code: "usageLimitExceeded", + retryable: null, + resetAt, + }, + }, + }, + ], + }); + const shell = (yield* orchestrator.getShellSnapshot()).threads.find( + (thread) => thread.id === threadId, + )!; + assert.isNull(limitRecoveryCommand(shell, false, DateTime.toEpochMillis(now))); + const arm = limitRecoveryCommand(shell, true, DateTime.toEpochMillis(now)); + assert.isNotNull(arm); + yield* orchestrator.dispatch(arm!); + const armedShell = (yield* orchestrator.getShellSnapshot()).threads.find( + (thread) => thread.id === threadId, + )!; + assert.deepEqual(armedShell.limitRecovery, { + runId: run.id, + resetAt, + autoResume: true, + requestId: arm!.commandId, + }); + assert.isNull(limitRecoveryCommand(armedShell, true, DateTime.toEpochMillis(now))); + yield* orchestrator.dispatch({ + type: "message.dispatch", + commandId: CommandId.make(`recovery:early:${scenario}`), + messageId: MessageId.make(`recovery:early:${scenario}`), + threadId, + usageLimitContinuationOfRunId: run.id, + text: "Continue where you left off.", + attachments: [], + dispatchMode: { type: "start_immediately" }, + createdBy: "user", + creationSource: "server", + }); + assert.lengthOf((yield* orchestrator.getThreadProjection(threadId)).runs, 1); + yield* TestClock.adjust("1 minute"); + const resume = limitRecoveryCommand( + armedShell, + true, + DateTime.toEpochMillis(yield* DateTime.now), + ); + assert.isNotNull(resume); + if (scenario === "cancel" || scenario === "rearm") + yield* orchestrator.dispatch({ + type: "thread.metadata.update", + commandId: CommandId.make(`recovery:cancel:${scenario}`), + threadId, + limitRecovery: { runId: run.id, resetAt, autoResume: false }, }); + if (scenario === "archive") yield* orchestrator.dispatch({ - type: "thread.create", - commandId: CommandId.make(`recovery:create:${scenario}`), + type: "thread.archive", + commandId: CommandId.make(`recovery:archive:${scenario}`), threadId, - projectId, - title: "Limited thread", - modelSelection, - runtimeMode: "full-access", - interactionMode: "default", - branch: null, - worktreePath: null, - createdBy: "user", - creationSource: "web", }); + if (scenario === "new-message") yield* orchestrator.dispatch({ type: "message.dispatch", - commandId: CommandId.make(`recovery:message:${scenario}`), + commandId: CommandId.make(`recovery:new-message:${scenario}`), threadId, - messageId: MessageId.make(`recovery:message:${scenario}`), - text: "Work on this.", + messageId: MessageId.make(`recovery:new-message:${scenario}`), + text: "I will continue manually.", attachments: [], dispatchMode: { type: "defer_start" }, createdBy: "user", creationSource: "web", }); - const projection = yield* orchestrator.getThreadProjection(threadId); - const run = projection.runs[0]!; - const now = yield* DateTime.now; - const resetAt = DateTime.formatIso(DateTime.add(now, { minutes: 1 })); + if (scenario === "settle") + yield* orchestrator.dispatch({ + type: "thread.settle", + commandId: CommandId.make(`recovery:settle:${scenario}`), + threadId, + }); + if (scenario === "replacement") { + const current = yield* orchestrator.getThreadProjection(threadId); + const error = current.turnItems.find((item) => item.type === "error")!; + if (error.type !== "error") throw new Error("Expected provider error"); yield* events.write({ - commandId: CommandId.make(`recovery:failure:${scenario}`), + commandId: CommandId.make(`recovery:replacement:${scenario}`), events: [ { - id: EventId.make(`recovery:run:${scenario}`), - type: "run.updated", - threadId, - occurredAt: now, - payload: { ...run, status: "failed", completedAt: now }, - }, - { - id: EventId.make(`recovery:error:${scenario}`), + id: EventId.make(`recovery:replacement:${scenario}`), type: "turn-item.updated", threadId, - occurredAt: now, + occurredAt: yield* DateTime.now, payload: { - id: TurnItemId.make(`recovery:error:${scenario}`), - type: "error", - threadId, - runId: run.id, - nodeId: run.rootNodeId, - providerThreadId: null, - providerTurnId: null, - nativeItemRef: null, - parentItemId: null, - ordinal: 2, - status: "failed", - title: "Usage limit reached", - startedAt: now, - completedAt: now, - updatedAt: now, + ...error, failure: { - class: "usage_limit", - message: "Plan limit reached.", - code: "usageLimitExceeded", - retryable: null, - resetAt, + ...error.failure, + class: "provider_error", + message: "A replacement failure.", }, }, }, ], }); - const shell = (yield* orchestrator.getShellSnapshot()).threads.find( - (thread) => thread.id === threadId, - )!; - assert.isNull(limitRecoveryCommand(shell, false, DateTime.toEpochMillis(now))); - const arm = limitRecoveryCommand(shell, true, DateTime.toEpochMillis(now)); - assert.isNotNull(arm); - yield* orchestrator.dispatch(arm!); - const armedShell = (yield* orchestrator.getShellSnapshot()).threads.find( - (thread) => thread.id === threadId, - )!; - assert.deepEqual(armedShell.limitRecovery, { runId: run.id, resetAt, autoResume: true }); - assert.isNull(limitRecoveryCommand(armedShell, true, DateTime.toEpochMillis(now))); + } + if (scenario === "rearm") { + yield* orchestrator.dispatch(resume!); yield* orchestrator.dispatch({ - type: "message.dispatch", - commandId: CommandId.make(`recovery:early:${scenario}`), - messageId: MessageId.make(`recovery:early:${scenario}`), + type: "thread.metadata.update", + commandId: CommandId.make(`recovery:rearm:${scenario}`), threadId, - usageLimitContinuationOfRunId: run.id, - text: "Continue where you left off.", - attachments: [], - dispatchMode: { type: "start_immediately" }, - createdBy: "user", - creationSource: "server", + limitRecovery: { runId: run.id, resetAt, autoResume: true }, }); + yield* orchestrator.dispatch(resume!); assert.lengthOf((yield* orchestrator.getThreadProjection(threadId)).runs, 1); - yield* TestClock.adjust("1 minute"); - const resume = limitRecoveryCommand( - armedShell, + const rearmedShell = (yield* orchestrator.getShellSnapshot()).threads.find( + (thread) => thread.id === threadId, + )!; + const freshResume = limitRecoveryCommand( + rearmedShell, true, DateTime.toEpochMillis(yield* DateTime.now), ); - assert.isNotNull(resume); - if (scenario === "cancel") - yield* orchestrator.dispatch({ - type: "thread.metadata.update", - commandId: CommandId.make(`recovery:cancel:${scenario}`), - threadId, - limitRecovery: { runId: run.id, resetAt, autoResume: false }, - }); - if (scenario === "archive") - yield* orchestrator.dispatch({ - type: "thread.archive", - commandId: CommandId.make(`recovery:archive:${scenario}`), - threadId, - }); - if (scenario === "new-message") - yield* orchestrator.dispatch({ - type: "message.dispatch", - commandId: CommandId.make(`recovery:new-message:${scenario}`), - threadId, - messageId: MessageId.make(`recovery:new-message:${scenario}`), - text: "I will continue manually.", - attachments: [], - dispatchMode: { type: "defer_start" }, - createdBy: "user", - creationSource: "web", - }); - if (scenario === "settle") - yield* orchestrator.dispatch({ - type: "thread.settle", - commandId: CommandId.make(`recovery:settle:${scenario}`), - threadId, - }); - if (scenario === "replacement") { - const current = yield* orchestrator.getThreadProjection(threadId); - const error = current.turnItems.find((item) => item.type === "error")!; - if (error.type !== "error") throw new Error("Expected provider error"); - yield* events.write({ - commandId: CommandId.make(`recovery:replacement:${scenario}`), - events: [ - { - id: EventId.make(`recovery:replacement:${scenario}`), - type: "turn-item.updated", - threadId, - occurredAt: yield* DateTime.now, - payload: { - ...error, - failure: { - ...error.failure, - class: "provider_error", - message: "A replacement failure.", - }, - }, - }, - ], - }); - } - const before = yield* orchestrator.getThreadProjection(threadId); - yield* orchestrator.dispatch(resume!); - yield* orchestrator.dispatch(resume!); - const after = yield* orchestrator.getThreadProjection(threadId); - assert.lengthOf(after.runs, before.runs.length + (scenario === "resume" ? 1 : 0)); - assert.lengthOf(after.messages, before.messages.length + (scenario === "resume" ? 1 : 0)); - }), + assert.isNotNull(freshResume); + assert.notEqual(freshResume!.commandId, resume!.commandId); + yield* orchestrator.dispatch(freshResume!); + yield* orchestrator.dispatch(freshResume!); + assert.lengthOf((yield* orchestrator.getThreadProjection(threadId)).runs, 2); + } + const before = yield* orchestrator.getThreadProjection(threadId); + yield* orchestrator.dispatch(resume!); + yield* orchestrator.dispatch(resume!); + const after = yield* orchestrator.getThreadProjection(threadId); + assert.lengthOf(after.runs, before.runs.length + (scenario === "resume" ? 1 : 0)); + assert.lengthOf(after.messages, before.messages.length + (scenario === "resume" ? 1 : 0)); + }), ); }); diff --git a/apps/web/src/components/settings/SettingsPanels.tsx b/apps/web/src/components/settings/SettingsPanels.tsx index 2adfdd67817a..52e6103fc96a 100644 --- a/apps/web/src/components/settings/SettingsPanels.tsx +++ b/apps/web/src/components/settings/SettingsPanels.tsx @@ -2230,8 +2230,8 @@ export function GeneralSettingsPanel() { Date: Sat, 19 Sep 2026 23:45:32 -0700 Subject: [PATCH 04/10] fix(mobile): surface failed limit recovery changes --- .../src/features/threads/UsageLimitRecoveryCard.tsx | 13 ++++++++++++- 1 file changed, 12 insertions(+), 1 deletion(-) diff --git a/apps/mobile/src/features/threads/UsageLimitRecoveryCard.tsx b/apps/mobile/src/features/threads/UsageLimitRecoveryCard.tsx index 8a583d3e0ba5..145cec4d958d 100644 --- a/apps/mobile/src/features/threads/UsageLimitRecoveryCard.tsx +++ b/apps/mobile/src/features/threads/UsageLimitRecoveryCard.tsx @@ -1,3 +1,4 @@ +import { squashAtomCommandFailure } from "@t3tools/client-runtime/state/runtime"; import type { EnvironmentThreadShell } from "@t3tools/client-runtime/state/shell"; import type { EnvironmentId } from "@t3tools/contracts"; import * as DateTime from "effect/DateTime"; @@ -16,6 +17,7 @@ export function UsageLimitRecoveryCard({ }) { const updateMetadata = useAtomCommand(threadEnvironment.updateMetadata); const [pending, setPending] = useState(false); + const [error, setError] = useState(null); const resetAt = thread.runtime?.usageLimitResetAt ?? null; const canSchedule = resetAt !== null && @@ -33,11 +35,15 @@ export function UsageLimitRecoveryCard({ async function toggle() { if (!resetAt || !runId || !canSchedule) return; setPending(true); + setError(null); try { - await updateMetadata({ + const result = await updateMetadata({ environmentId, input: { threadId: thread.id, limitRecovery: { runId, resetAt, autoResume: !scheduled } }, }); + if (result._tag === "Failure") throw squashAtomCommandFailure(result); + } catch (cause) { + setError(cause instanceof Error ? cause.message : "Could not change limit recovery."); } finally { setPending(false); } @@ -61,6 +67,11 @@ export function UsageLimitRecoveryCard({ ) : null} + {error ? ( + + {error} + + ) : null} ); } From 49a2b7407853ee37dc047a46bf07246a1d04ba6b Mon Sep 17 00:00:00 2001 From: Julius Marminge Date: Sun, 20 Sep 2026 00:33:23 -0700 Subject: [PATCH 05/10] fix(web): show limit recovery in the banner stack --- apps/web/src/components/ChatView.tsx | 51 ++++++++------ .../chat/UsageLimitRecoveryBanner.tsx | 69 +++++++++++++++++++ .../chat/UsageLimitRecoveryCard.tsx | 61 ---------------- 3 files changed, 98 insertions(+), 83 deletions(-) create mode 100644 apps/web/src/components/chat/UsageLimitRecoveryBanner.tsx delete mode 100644 apps/web/src/components/chat/UsageLimitRecoveryCard.tsx diff --git a/apps/web/src/components/ChatView.tsx b/apps/web/src/components/ChatView.tsx index deeb547e4ff1..cdcd918232cc 100644 --- a/apps/web/src/components/ChatView.tsx +++ b/apps/web/src/components/ChatView.tsx @@ -1,4 +1,4 @@ -import { UsageLimitRecoveryCard } from "./chat/UsageLimitRecoveryCard"; +import { usageLimitRecoveryBannerItem } from "./chat/UsageLimitRecoveryBanner"; import { resolveBackgroundDraftWorkspaceOptions, resolveDraftHeroState, @@ -6966,7 +6966,27 @@ export default function ChatView(props: ChatViewProps) { }), [feedbackSubmissions, routeThreadKey], ); + const limitRecoveryBanner = + serverRuntime?.status === "failed" && + serverRuntime.lastErrorClass === "usage_limit" && + activeThreadShell?.latestRun + ? usageLimitRecoveryBannerItem({ + runId: activeThreadShell.latestRun.runId, + resetAt: serverRuntime.usageLimitResetAt ?? null, + stoppedAt: activeThreadShell.latestRun.completedAt ?? activeThreadShell.updatedAt, + recovery: activeThreadShell.limitRecovery ?? null, + explanation: serverRuntime.lastError, + onChange: async (limitRecovery) => { + const result = await updateThreadMetadata({ + environmentId, + input: { threadId: activeThreadShell.id, limitRecovery }, + }); + if (result._tag === "Failure") throw squashAtomCommandFailure(result); + }, + }) + : null; const composerBannerItems = useMemo(() => { + const limitRecoveryItems = limitRecoveryBanner === null ? [] : [limitRecoveryBanner]; const backgroundWorkItems = backgroundWorkBannerItem === null ? [] : [backgroundWorkBannerItem]; const resumeCompactionItems = resumeCompactionBannerItem === null ? [] : [resumeCompactionBannerItem]; @@ -6978,6 +6998,7 @@ export default function ChatView(props: ChatViewProps) { if (!localCheckoutBranchMismatch || !showBranchMismatchBanner || !activeBranchMismatchKey) { return [ ...feedbackBannerItems, + ...limitRecoveryItems, ...usageLimitsItems, ...projectCloneItems, ...systemComposerBannerItems, @@ -6989,6 +7010,7 @@ export default function ChatView(props: ChatViewProps) { } return [ ...feedbackBannerItems, + ...limitRecoveryItems, ...usageLimitsItems, ...projectCloneItems, ...systemComposerBannerItems, @@ -7038,6 +7060,7 @@ export default function ChatView(props: ChatViewProps) { }, [ activeBranchMismatchKey, feedbackBannerItems, + limitRecoveryBanner, handleRestoreThreadBranch, isRestoringThreadBranch, backgroundWorkBannerItem, @@ -10357,7 +10380,11 @@ export default function ChatView(props: ChatViewProps) { onOpenProviderSetup={openProviderSetup} /> - {serverRuntime?.status === "failed" && - serverRuntime.lastErrorClass === "usage_limit" && - activeThreadShell?.latestRun ? ( - { - const result = await updateThreadMetadata({ - environmentId, - input: { threadId: activeThread.id, limitRecovery }, - }); - if (result._tag === "Failure") throw squashAtomCommandFailure(result); - }} - /> - ) : null} {isDraftHeroState ? (

Promise; +}; + +export function usageLimitRecoveryBannerItem(props: RecoveryProps): ComposerBannerStackItem { + const { runId, resetAt, stoppedAt, recovery, explanation } = props; + const canSchedule = resetAt !== null && Date.parse(resetAt) > Date.parse(stoppedAt); + const scheduled = + recovery?.runId === runId && recovery.resetAt === resetAt && recovery.autoResume; + return { + id: `usage-limit-recovery:${runId}`, + variant: "warning", + priority: "urgent", + icon: , + title: "Usage limit reached", + description: resetAt + ? `Resets ${new Date(resetAt).toLocaleString()}` + : "Reset time unavailable; retry manually", + children: ( +
+ {explanation ?

{explanation}

: null} + {scheduled ?

Auto-resume is scheduled for the reset.

: null} +
+ ), + actions: canSchedule ? : null, + }; +} + +function RecoveryActions({ runId, resetAt, recovery, onChange }: RecoveryProps) { + const [pending, setPending] = useState(false); + const [error, setError] = useState(null); + const scheduled = + recovery?.runId === runId && recovery.resetAt === resetAt && recovery.autoResume; + async function toggle() { + if (resetAt === null) return; + setPending(true); + setError(null); + try { + await onChange({ runId, resetAt, autoResume: !scheduled }); + } catch (cause) { + setError(cause instanceof Error ? cause.message : "Could not change limit recovery."); + } finally { + setPending(false); + } + } + return ( +
+ + {error ? ( +

+ {error} +

+ ) : null} +
+ ); +} diff --git a/apps/web/src/components/chat/UsageLimitRecoveryCard.tsx b/apps/web/src/components/chat/UsageLimitRecoveryCard.tsx deleted file mode 100644 index 213748296db3..000000000000 --- a/apps/web/src/components/chat/UsageLimitRecoveryCard.tsx +++ /dev/null @@ -1,61 +0,0 @@ -import * as DateTime from "effect/DateTime"; -import { type OrchestrationV2LimitRecovery, type RunId } from "@t3tools/contracts"; -import { useState } from "react"; -import { Button } from "../ui/button"; - -export function UsageLimitRecoveryCard({ - runId, - resetAt, - stoppedAt, - recovery, - onChange, -}: { - runId: RunId; - resetAt: string | null; - stoppedAt: string; - recovery: OrchestrationV2LimitRecovery | null; - onChange: (recovery: OrchestrationV2LimitRecovery) => Promise; -}) { - const [pending, setPending] = useState(false); - const [error, setError] = useState(null); - const canSchedule = resetAt !== null && Date.parse(resetAt) > Date.parse(stoppedAt); - const scheduled = - recovery?.runId === runId && recovery.resetAt === resetAt && recovery.autoResume; - async function toggle() { - if (resetAt === null || !canSchedule) return; - setPending(true); - setError(null); - try { - await onChange({ runId, resetAt, autoResume: !scheduled }); - } catch (cause) { - setError(cause instanceof Error ? cause.message : "Could not change limit recovery."); - } finally { - setPending(false); - } - } - return ( -
-

- {resetAt - ? `Usage limit resets ${DateTime.toDateUtc(DateTime.makeUnsafe(resetAt)).toLocaleString()}.` - : "The provider did not report a reset time. Retry manually when your limit is available."} -

- {canSchedule ? ( - - ) : null} - {error ? ( -

- {error} -

- ) : null} -
- ); -} From a51df13555634453de9a3e78cb03881a1c8a2a6c Mon Sep 17 00:00:00 2001 From: Julius Marminge Date: Sun, 20 Sep 2026 00:43:53 -0700 Subject: [PATCH 06/10] fix(web): preserve timeline spacing with recovery banners --- apps/web/src/components/ChatView.tsx | 16 ++++++++++------ 1 file changed, 10 insertions(+), 6 deletions(-) diff --git a/apps/web/src/components/ChatView.tsx b/apps/web/src/components/ChatView.tsx index cdcd918232cc..ed57bbe54b35 100644 --- a/apps/web/src/components/ChatView.tsx +++ b/apps/web/src/components/ChatView.tsx @@ -2124,6 +2124,14 @@ export default function ChatView(props: ChatViewProps) { widthStorageKey: `t3code:preview-panel-width:${activeThreadKey}`, }); const activeThreadShell = useThreadShell(isServerThread ? activeThreadRef : null); + const timelineThreadError = + serverRuntime?.status === "failed" && + serverRuntime.lastErrorClass === "usage_limit" && + activeThreadShell?.latestRun && + visibleThreadError === serverRuntime.lastError + ? null + : visibleThreadError; + const [timelineAnchor, setTimelineAnchor] = useState<{ readonly threadKey: string | null; readonly messageId: MessageId | null; @@ -3976,7 +3984,7 @@ export default function ChatView(props: ChatViewProps) { ) ? activeProviderStatus : null; - const hasTimelineTopBanner = Boolean(visibleThreadError) || visibleProviderStatus !== null; + const hasTimelineTopBanner = Boolean(timelineThreadError) || visibleProviderStatus !== null; const activeProjectCwd = activeProject?.workspaceRoot ?? null; const activeThreadWorktreePath = activeThread?.worktreePath ?? null; const activeWorkspaceRoot = activeThreadWorktreePath ?? activeProjectCwd ?? undefined; @@ -10380,11 +10388,7 @@ export default function ChatView(props: ChatViewProps) { onOpenProviderSetup={openProviderSetup} /> Date: Sun, 20 Sep 2026 12:04:29 -0700 Subject: [PATCH 07/10] fix(web): keep usage limit banner concise --- apps/web/src/components/ChatView.tsx | 1 - .../src/components/chat/UsageLimitRecoveryBanner.tsx | 11 +---------- 2 files changed, 1 insertion(+), 11 deletions(-) diff --git a/apps/web/src/components/ChatView.tsx b/apps/web/src/components/ChatView.tsx index ed57bbe54b35..c6b1fb81f7fb 100644 --- a/apps/web/src/components/ChatView.tsx +++ b/apps/web/src/components/ChatView.tsx @@ -6983,7 +6983,6 @@ export default function ChatView(props: ChatViewProps) { resetAt: serverRuntime.usageLimitResetAt ?? null, stoppedAt: activeThreadShell.latestRun.completedAt ?? activeThreadShell.updatedAt, recovery: activeThreadShell.limitRecovery ?? null, - explanation: serverRuntime.lastError, onChange: async (limitRecovery) => { const result = await updateThreadMetadata({ environmentId, diff --git a/apps/web/src/components/chat/UsageLimitRecoveryBanner.tsx b/apps/web/src/components/chat/UsageLimitRecoveryBanner.tsx index c29802299531..2d2992cc092a 100644 --- a/apps/web/src/components/chat/UsageLimitRecoveryBanner.tsx +++ b/apps/web/src/components/chat/UsageLimitRecoveryBanner.tsx @@ -9,15 +9,12 @@ type RecoveryProps = { resetAt: string | null; stoppedAt: string; recovery: OrchestrationV2LimitRecovery | null; - explanation: string | null; onChange: (recovery: OrchestrationV2LimitRecovery) => Promise; }; export function usageLimitRecoveryBannerItem(props: RecoveryProps): ComposerBannerStackItem { - const { runId, resetAt, stoppedAt, recovery, explanation } = props; + const { runId, resetAt, stoppedAt } = props; const canSchedule = resetAt !== null && Date.parse(resetAt) > Date.parse(stoppedAt); - const scheduled = - recovery?.runId === runId && recovery.resetAt === resetAt && recovery.autoResume; return { id: `usage-limit-recovery:${runId}`, variant: "warning", @@ -27,12 +24,6 @@ export function usageLimitRecoveryBannerItem(props: RecoveryProps): ComposerBann description: resetAt ? `Resets ${new Date(resetAt).toLocaleString()}` : "Reset time unavailable; retry manually", - children: ( -
- {explanation ?

{explanation}

: null} - {scheduled ?

Auto-resume is scheduled for the reset.

: null} -
- ), actions: canSchedule ? : null, }; } From ae6fa23b94be580595ce1c456b5dd5350483f036 Mon Sep 17 00:00:00 2001 From: Julius Marminge Date: Sun, 20 Sep 2026 18:10:35 -0700 Subject: [PATCH 08/10] fix(server): retry continuations after a snooze race --- .../src/orchestration-v2/Orchestrator.ts | 11 +++++++- .../src/orchestration-v2/runtimeLayer.test.ts | 26 +++++++++++++++++++ docs/user/thread-sidebar.md | 3 ++- 3 files changed, 38 insertions(+), 2 deletions(-) diff --git a/apps/server/src/orchestration-v2/Orchestrator.ts b/apps/server/src/orchestration-v2/Orchestrator.ts index e3a1fd32b1f6..f8b2fdc005d0 100644 --- a/apps/server/src/orchestration-v2/Orchestrator.ts +++ b/apps/server/src/orchestration-v2/Orchestrator.ts @@ -3815,7 +3815,16 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio type: "thread.metadata-updated", threadId: command.threadId, occurredAt: now, - payload: projection.thread, + payload: { + ...projection.thread, + ...(recovery?.autoResume && + recovery.requestId === command.usageLimitRecoveryRequestId && + recovery.runId === command.usageLimitContinuationOfRunId && + projection.thread.snoozedUntil != null && + DateTime.toEpochMillis(projection.thread.snoozedUntil) > DateTime.toEpochMillis(now) + ? { limitRecovery: { ...recovery, requestId: command.commandId } } + : {}), + }, }); return; } diff --git a/apps/server/src/orchestration-v2/runtimeLayer.test.ts b/apps/server/src/orchestration-v2/runtimeLayer.test.ts index bafc48a7f28d..8a36a6e7dd5c 100644 --- a/apps/server/src/orchestration-v2/runtimeLayer.test.ts +++ b/apps/server/src/orchestration-v2/runtimeLayer.test.ts @@ -3021,6 +3021,7 @@ it.layer(TestLayer)("usage-limit recovery", (it) => { "resume", "cancel", "rearm", + "snooze-race", "new-message", "archive", "settle", @@ -3153,6 +3154,31 @@ it.layer(TestLayer)("usage-limit recovery", (it) => { DateTime.toEpochMillis(yield* DateTime.now), ); assert.isNotNull(resume); + if (scenario === "snooze-race") { + const wakeAt = DateTime.formatIso(DateTime.add(yield* DateTime.now, { minutes: 1 })); + yield* orchestrator.dispatch({ + type: "thread.snooze", + commandId: CommandId.make("recovery:raced-snooze"), + threadId, + snoozedUntil: wakeAt, + }); + yield* orchestrator.dispatch(resume!); + assert.lengthOf((yield* orchestrator.getThreadProjection(threadId)).runs, 1); + yield* TestClock.adjust("1 minute"); + const current = (yield* orchestrator.getShellSnapshot()).threads.find( + (thread) => thread.id === threadId, + )!; + const freshResume = limitRecoveryCommand( + current, + true, + DateTime.toEpochMillis(yield* DateTime.now), + ); + assert.isNotNull(freshResume); + assert.notEqual(freshResume!.commandId, resume!.commandId); + yield* orchestrator.dispatch(freshResume!); + yield* orchestrator.dispatch(freshResume!); + assert.lengthOf((yield* orchestrator.getThreadProjection(threadId)).runs, 2); + } if (scenario === "cancel" || scenario === "rearm") yield* orchestrator.dispatch({ type: "thread.metadata.update", diff --git a/docs/user/thread-sidebar.md b/docs/user/thread-sidebar.md index ed0f635ae9c6..022e2912db88 100644 --- a/docs/user/thread-sidebar.md +++ b/docs/user/thread-sidebar.md @@ -140,7 +140,8 @@ another provider instance. When the provider reports a reset time, choose **Resume at reset** to schedule a continuation. You can cancel it from the thread. Enable **Auto-resume limited -threads** in thread behavior settings to schedule limit stops by default. +threads** in **Settings → General** on web and desktop, or **Settings → Thread +behavior** on mobile, to schedule limit stops by default. The environment must be running when the reset arrives; it resumes overdue continuations after a restart. Sending a new message, archiving, or settling the thread prevents a pending continuation from starting. From 96c0d276a5afa7ab23a833fe87f2189580c9c05a Mon Sep 17 00:00:00 2001 From: Julius Marminge Date: Sun, 20 Sep 2026 18:15:46 -0700 Subject: [PATCH 09/10] refactor(server): qualify recovery service dependencies --- .../orchestration-v2/UsageLimitRecoveryService.ts | 12 ++++++------ apps/server/src/orchestration-v2/runtimeLayer.ts | 4 ++-- 2 files changed, 8 insertions(+), 8 deletions(-) diff --git a/apps/server/src/orchestration-v2/UsageLimitRecoveryService.ts b/apps/server/src/orchestration-v2/UsageLimitRecoveryService.ts index 4d6cb5e30c85..4f1d13017854 100644 --- a/apps/server/src/orchestration-v2/UsageLimitRecoveryService.ts +++ b/apps/server/src/orchestration-v2/UsageLimitRecoveryService.ts @@ -8,9 +8,9 @@ import * as DateTime from "effect/DateTime"; import * as Effect from "effect/Effect"; import * as Layer from "effect/Layer"; import * as Schedule from "effect/Schedule"; -import { ServerSettingsService } from "../serverSettings.ts"; -import { ProjectionStoreV2 } from "./ProjectionStore.ts"; -import { ThreadManagementService } from "./ThreadManagementService.ts"; +import * as ServerSettings from "../serverSettings.ts"; +import * as ProjectionStore from "./ProjectionStore.ts"; +import * as ThreadManagement from "./ThreadManagementService.ts"; /** The persisted run and reset form the identity of one recovery opportunity. */ export function limitRecoveryCommand( @@ -71,9 +71,9 @@ export function limitRecoveryCommand( } const make = Effect.gen(function* () { - const projections = yield* ProjectionStoreV2; - const threads = yield* ThreadManagementService; - const settings = yield* ServerSettingsService; + const projections = yield* ProjectionStore.ProjectionStoreV2; + const threads = yield* ThreadManagement.ThreadManagementService; + const settings = yield* ServerSettings.ServerSettingsService; return Effect.fn("UsageLimitRecoveryService.sweep")(function* () { const preferences = yield* settings.getSettings; const snapshot = yield* projections.getShellSnapshot(); diff --git a/apps/server/src/orchestration-v2/runtimeLayer.ts b/apps/server/src/orchestration-v2/runtimeLayer.ts index 4400151b963b..7f42d6565104 100644 --- a/apps/server/src/orchestration-v2/runtimeLayer.ts +++ b/apps/server/src/orchestration-v2/runtimeLayer.ts @@ -1,4 +1,4 @@ -import { workerLive as usageLimitRecoveryWorkerLive } from "./UsageLimitRecoveryService.ts"; +import * as UsageLimitRecoveryService from "./UsageLimitRecoveryService.ts"; import * as Layer from "effect/Layer"; import { OrchestrationEventInfrastructureLayerLive, @@ -296,7 +296,7 @@ export const OrchestrationV2ProductionLayerLive = Layer.mergeAll( threadLaunchProvided, threadLifecycleProvided, scheduledTaskProvided, - usageLimitRecoveryWorkerLive.pipe( + UsageLimitRecoveryService.workerLive.pipe( Layer.provide(Layer.mergeAll(projectionStoreLayer, threadManagementProvided)), ), providerContinuationWorkerProvided, From 9eea14d5902c92d514dcb47e88051aaae5aa3b04 Mon Sep 17 00:00:00 2001 From: Julius Marminge Date: Sun, 20 Sep 2026 18:22:28 -0700 Subject: [PATCH 10/10] refactor(server): name limit recovery as a worker --- ...eLimitRecoveryService.ts => UsageLimitRecoveryWorker.ts} | 6 +++--- apps/server/src/orchestration-v2/runtimeLayer.test.ts | 2 +- apps/server/src/orchestration-v2/runtimeLayer.ts | 4 ++-- 3 files changed, 6 insertions(+), 6 deletions(-) rename apps/server/src/orchestration-v2/{UsageLimitRecoveryService.ts => UsageLimitRecoveryWorker.ts} (96%) diff --git a/apps/server/src/orchestration-v2/UsageLimitRecoveryService.ts b/apps/server/src/orchestration-v2/UsageLimitRecoveryWorker.ts similarity index 96% rename from apps/server/src/orchestration-v2/UsageLimitRecoveryService.ts rename to apps/server/src/orchestration-v2/UsageLimitRecoveryWorker.ts index 4f1d13017854..dfadda069277 100644 --- a/apps/server/src/orchestration-v2/UsageLimitRecoveryService.ts +++ b/apps/server/src/orchestration-v2/UsageLimitRecoveryWorker.ts @@ -70,11 +70,11 @@ export function limitRecoveryCommand( }; } -const make = Effect.gen(function* () { +const makeSweep = Effect.gen(function* () { const projections = yield* ProjectionStore.ProjectionStoreV2; const threads = yield* ThreadManagement.ThreadManagementService; const settings = yield* ServerSettings.ServerSettingsService; - return Effect.fn("UsageLimitRecoveryService.sweep")(function* () { + return Effect.fn("UsageLimitRecoveryWorker.sweep")(function* () { const preferences = yield* settings.getSettings; const snapshot = yield* projections.getShellSnapshot(); const nowMs = DateTime.toEpochMillis(yield* DateTime.now); @@ -97,7 +97,7 @@ const make = Effect.gen(function* () { // so restarts need no timer restoration and disconnected clients need not run it. export const workerLive = Layer.effectDiscard( Effect.gen(function* () { - const sweep = yield* make; + const sweep = yield* makeSweep; yield* sweep().pipe( Effect.catchCause((cause) => Effect.logWarning("orchestration-v2.limit-recovery.sweep-failed", { cause }), diff --git a/apps/server/src/orchestration-v2/runtimeLayer.test.ts b/apps/server/src/orchestration-v2/runtimeLayer.test.ts index 8a36a6e7dd5c..50e4e6990c32 100644 --- a/apps/server/src/orchestration-v2/runtimeLayer.test.ts +++ b/apps/server/src/orchestration-v2/runtimeLayer.test.ts @@ -1,4 +1,4 @@ -import { limitRecoveryCommand } from "./UsageLimitRecoveryService.ts"; +import { limitRecoveryCommand } from "./UsageLimitRecoveryWorker.ts"; import { SourceControlProviderRegistry } from "../sourceControl/SourceControlProviderRegistry.ts"; import * as NodeServices from "@effect/platform-node/NodeServices"; import { assert, it } from "@effect/vitest"; diff --git a/apps/server/src/orchestration-v2/runtimeLayer.ts b/apps/server/src/orchestration-v2/runtimeLayer.ts index 7f42d6565104..83b176495a75 100644 --- a/apps/server/src/orchestration-v2/runtimeLayer.ts +++ b/apps/server/src/orchestration-v2/runtimeLayer.ts @@ -1,4 +1,4 @@ -import * as UsageLimitRecoveryService from "./UsageLimitRecoveryService.ts"; +import * as UsageLimitRecoveryWorker from "./UsageLimitRecoveryWorker.ts"; import * as Layer from "effect/Layer"; import { OrchestrationEventInfrastructureLayerLive, @@ -296,7 +296,7 @@ export const OrchestrationV2ProductionLayerLive = Layer.mergeAll( threadLaunchProvided, threadLifecycleProvided, scheduledTaskProvided, - UsageLimitRecoveryService.workerLive.pipe( + UsageLimitRecoveryWorker.workerLive.pipe( Layer.provide(Layer.mergeAll(projectionStoreLayer, threadManagementProvided)), ), providerContinuationWorkerProvided,