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) => ( (null); + const resetAt = thread.runtime?.usageLimitResetAt ?? null; + const canSchedule = + resetAt !== null && + Date.parse(resetAt) > 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); + setError(null); + try { + 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); + } + } + 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} + {error ? ( + + {error} + + ) : null} + + ); +} diff --git a/apps/server/src/orchestration-v2/Orchestrator.ts b/apps/server/src/orchestration-v2/Orchestrator.ts index b6119d03ff85..f8b2fdc005d0 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,14 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio return { ...thread, ...(command.title === undefined ? {} : { title: command.title }), + ...(command.limitRecovery === undefined + ? {} + : { + 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 @@ -3754,6 +3786,50 @@ 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.requestId !== command.usageLimitRecoveryRequestId || + 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, + ...(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; + } + } + 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/UsageLimitRecoveryWorker.ts b/apps/server/src/orchestration-v2/UsageLimitRecoveryWorker.ts new file mode 100644 index 000000000000..dfadda069277 --- /dev/null +++ b/apps/server/src/orchestration-v2/UsageLimitRecoveryWorker.ts @@ -0,0 +1,109 @@ +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 * as ServerSettings from "../serverSettings.ts"; +import * as ProjectionStore from "./ProjectionStore.ts"; +import * as ThreadManagement from "./ThreadManagementService.ts"; + +/** The persisted run and reset form the identity of one recovery opportunity. */ +export function limitRecoveryCommand( + thread: OrchestrationV2ThreadShell, + 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; + const deliveryIdentity = `${identity}:${recovery.requestId ?? "legacy"}`; + return { + type: "message.dispatch", + 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" }, + createdBy: "user", + creationSource: "server", + }; +} + +const makeSweep = Effect.gen(function* () { + const projections = yield* ProjectionStore.ProjectionStoreV2; + const threads = yield* ThreadManagement.ThreadManagementService; + const settings = yield* ServerSettings.ServerSettingsService; + return Effect.fn("UsageLimitRecoveryWorker.sweep")(function* () { + const preferences = yield* settings.getSettings; + const snapshot = yield* projections.getShellSnapshot(); + const nowMs = DateTime.toEpochMillis(yield* DateTime.now); + for (const thread of snapshot.threads) { + const 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* makeSweep; + 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..50e4e6990c32 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 "./UsageLimitRecoveryWorker.ts"; import { SourceControlProviderRegistry } from "../sourceControl/SourceControlProviderRegistry.ts"; import * as NodeServices from "@effect/platform-node/NodeServices"; import { assert, it } from "@effect/vitest"; @@ -3014,3 +3015,255 @@ it.layer(SharedApplicationDataPlaneTestLayer)("shared application data plane", ( }), ); }); + +it.layer(TestLayer)("usage-limit recovery", (it) => { + it.effect.each([ + "resume", + "cancel", + "rearm", + "snooze-race", + "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 === "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", + 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.", + }, + }, + }, + ], + }); + } + if (scenario === "rearm") { + yield* orchestrator.dispatch(resume!); + yield* orchestrator.dispatch({ + type: "thread.metadata.update", + commandId: CommandId.make(`recovery:rearm:${scenario}`), + threadId, + limitRecovery: { runId: run.id, resetAt, autoResume: true }, + }); + yield* orchestrator.dispatch(resume!); + assert.lengthOf((yield* orchestrator.getThreadProjection(threadId)).runs, 1); + const rearmedShell = (yield* orchestrator.getShellSnapshot()).threads.find( + (thread) => thread.id === threadId, + )!; + const freshResume = limitRecoveryCommand( + rearmedShell, + 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); + } + 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..83b176495a75 100644 --- a/apps/server/src/orchestration-v2/runtimeLayer.ts +++ b/apps/server/src/orchestration-v2/runtimeLayer.ts @@ -1,3 +1,4 @@ +import * as UsageLimitRecoveryWorker from "./UsageLimitRecoveryWorker.ts"; import * as Layer from "effect/Layer"; import { OrchestrationEventInfrastructureLayerLive, @@ -295,6 +296,9 @@ export const OrchestrationV2ProductionLayerLive = Layer.mergeAll( threadLaunchProvided, threadLifecycleProvided, scheduledTaskProvided, + UsageLimitRecoveryWorker.workerLive.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..c6b1fb81f7fb 100644 --- a/apps/web/src/components/ChatView.tsx +++ b/apps/web/src/components/ChatView.tsx @@ -1,3 +1,4 @@ +import { usageLimitRecoveryBannerItem } from "./chat/UsageLimitRecoveryBanner"; import { resolveBackgroundDraftWorkspaceOptions, resolveDraftHeroState, @@ -2123,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; @@ -3975,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; @@ -6965,7 +6974,26 @@ 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, + 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]; @@ -6977,6 +7005,7 @@ export default function ChatView(props: ChatViewProps) { if (!localCheckoutBranchMismatch || !showBranchMismatchBanner || !activeBranchMismatchKey) { return [ ...feedbackBannerItems, + ...limitRecoveryItems, ...usageLimitsItems, ...projectCloneItems, ...systemComposerBannerItems, @@ -6988,6 +7017,7 @@ export default function ChatView(props: ChatViewProps) { } return [ ...feedbackBannerItems, + ...limitRecoveryItems, ...usageLimitsItems, ...projectCloneItems, ...systemComposerBannerItems, @@ -7037,6 +7067,7 @@ export default function ChatView(props: ChatViewProps) { }, [ activeBranchMismatchKey, feedbackBannerItems, + limitRecoveryBanner, handleRestoreThreadBranch, isRestoringThreadBranch, backgroundWorkBannerItem, @@ -10356,7 +10387,7 @@ export default function ChatView(props: ChatViewProps) { onOpenProviderSetup={openProviderSetup} /> Promise; +}; + +export function usageLimitRecoveryBannerItem(props: RecoveryProps): ComposerBannerStackItem { + const { runId, resetAt, stoppedAt } = props; + const canSchedule = resetAt !== null && Date.parse(resetAt) > Date.parse(stoppedAt); + 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", + 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/settings/SettingsPanels.tsx b/apps/web/src/components/settings/SettingsPanels.tsx index 36220f3bc143..52e6103fc96a 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..1c62d140c247 100644 --- a/packages/contracts/src/orchestrationV2.ts +++ b/packages/contracts/src/orchestrationV2.ts @@ -326,6 +326,14 @@ export const OrchestrationV2ProviderCapabilities = Schema.Struct({ }); export type OrchestrationV2ProviderCapabilities = typeof OrchestrationV2ProviderCapabilities.Type; +export const OrchestrationV2LimitRecovery = Schema.Struct({ + requestId: Schema.optional(CommandId), + runId: RunId, + resetAt: IsoDateTime, + autoResume: Schema.Boolean, +}); +export type OrchestrationV2LimitRecovery = typeof OrchestrationV2LimitRecovery.Type; + export const OrchestrationV2AppThread = Schema.Struct({ ...OrchestrationV2CreationFields, id: ThreadId, @@ -369,6 +377,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 +1525,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 +2339,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 +2429,8 @@ 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), + usageLimitRecoveryRequestId: Schema.optional(CommandId), /** 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)),