From d5452d930fe05f804e768d0afa16e609f18a771d Mon Sep 17 00:00:00 2001 From: Julius Marminge Date: Sat, 19 Sep 2026 22:24:16 -0700 Subject: [PATCH 1/9] feat(v2): show provider limit stops as Limited Port the usage-limit classification and presentation from [#10550](https://github.com/pingdotgg/t3code/pull/10550) to V2's structured provider failures. Related original detection work is in [#7165](https://github.com/pingdotgg/t3code/pull/7165), [#10321](https://github.com/pingdotgg/t3code/pull/10321), and [#10473](https://github.com/pingdotgg/t3code/pull/10473). Codex, Claude, and Grok classify explicit limit signals at the adapter boundary. Shell and detail summaries derive the failure from the current failed root turn, preserving thread isolation and clearing stale state when a new run, new attempt, recovery, or replacement error supersedes it. Web, desktop, and mobile show Limited with warning styling and preserve attention notifications. No auto-resume or legacy session migration. Validation: 637 focused tests passed; scoped typechecks, lint, and format checks. Browser verification omitted at the maintainer's request. Adapted from Vitalii Yehorov's original limit detection and presentation work. Co-authored-by: Vitalii Yehorov --- .../features/threads/thread-list-v2-items.tsx | 7 +- .../src/features/threads/thread-work-log.tsx | 22 ++- .../src/features/threads/threadListV2.test.ts | 29 ++++ .../src/features/threads/threadListV2.ts | 11 +- .../src/features/threads/work-log-layout.tsx | 3 +- apps/mobile/src/lib/threadActivity.test.ts | 25 +++ apps/mobile/src/lib/threadActivity.ts | 14 +- .../orchestration-v2/Adapters/AcpAdapterV2.ts | 11 +- .../Adapters/ClaudeAdapterV2.test.ts | 50 ++++++ .../Adapters/ClaudeAdapterV2.ts | 36 ++++- .../Adapters/CodexAdapterV2.test.ts | 100 ++++++++++++ .../Adapters/CodexAdapterV2.ts | 69 +++++++-- .../Adapters/GrokAdapterV2.test.ts | 30 ++++ .../Adapters/GrokAdapterV2.ts | 15 +- .../orchestration-v2/ProjectionStore.test.ts | 146 ++++++++++++++++++ .../src/orchestration-v2/ProjectionStore.ts | 36 ++++- .../src/orchestration-v2/ProviderFailure.ts | 4 +- .../src/provider/acp/XAiAcpExtension.ts | 2 +- apps/web/src/components/ChatView.tsx | 5 + apps/web/src/components/Sidebar.logic.test.ts | 28 ++++ apps/web/src/components/Sidebar.logic.ts | 16 +- apps/web/src/components/Sidebar.tsx | 45 ++++-- .../ThreadNotificationCoordinator.test.tsx | 6 +- .../ThreadNotificationCoordinator.tsx | 10 +- .../src/components/chat/MessagesTimeline.tsx | 6 +- .../src/components/chat/ThreadErrorBanner.tsx | 10 +- apps/web/src/session-logic.test.ts | 7 + apps/web/src/session-logic.ts | 13 +- docs/user/thread-sidebar.md | 4 + packages/client-runtime/src/state/models.ts | 3 + .../src/state/threadExecution.test.ts | 55 ++++++- .../src/state/threadExecution.ts | 9 +- .../contracts/src/orchestrationV2.test.ts | 8 + packages/contracts/src/orchestrationV2.ts | 2 + packages/shared/package.json | 4 + .../shared/src/orchestrationV2ThreadError.ts | 45 ++++++ 36 files changed, 807 insertions(+), 79 deletions(-) create mode 100644 packages/shared/src/orchestrationV2ThreadError.ts diff --git a/apps/mobile/src/features/threads/thread-list-v2-items.tsx b/apps/mobile/src/features/threads/thread-list-v2-items.tsx index c647382bb2c2..5ebb192e51de 100644 --- a/apps/mobile/src/features/threads/thread-list-v2-items.tsx +++ b/apps/mobile/src/features/threads/thread-list-v2-items.tsx @@ -70,6 +70,7 @@ const STATUS_LABEL_BY_STATUS: Partial< input: { label: "Input", className: "text-adaptive-indigo-600-300" }, working: { label: "Working", className: "text-adaptive-sky-600-400" }, failed: { label: "Failed", className: "text-danger-foreground" }, + limited: { label: "Limited", className: "text-warning-foreground" }, }; function threadTimeLabel(thread: EnvironmentThreadShell): string { @@ -914,13 +915,15 @@ export const ThreadListV2Row = memo(function ThreadListV2Row(props: { ) : null} - {status === "failed" && thread.runtime?.lastError ? ( + {(status === "failed" || status === "limited") && thread.runtime?.lastError ? ( diff --git a/apps/mobile/src/features/threads/thread-work-log.tsx b/apps/mobile/src/features/threads/thread-work-log.tsx index 4e00addd164b..9355c10b108f 100644 --- a/apps/mobile/src/features/threads/thread-work-log.tsx +++ b/apps/mobile/src/features/threads/thread-work-log.tsx @@ -888,7 +888,11 @@ const ThreadWorkLogRow = memo(function ThreadWorkLogRow( const accessiblePreview = [previewText, answerPreview].filter(Boolean).join(": "); const displayText = workEntryRowLabel(row.workEntry, expanded); const isSystemNotice = row.projectedItem.item.type === "system_notice"; - const iconIsDestructive = !isSystemNotice && (row.icon === "alert" || row.icon === "warning"); + const isUsageLimit = + row.projectedItem.item.type === "error" && + row.projectedItem.item.failure.class === "usage_limit"; + const iconIsDestructive = + !isSystemNotice && !isUsageLimit && (row.icon === "alert" || row.icon === "warning"); const failed = row.status === "failure"; const toolIcon = row.workEntry.toolIcon ?? row.workEntry.toolSource?.icon; const icon = reasoning ? "brain" : (toolPresentation?.icon ?? workRowSymbolName(row.icon)); @@ -942,16 +946,20 @@ const ThreadWorkLogRow = memo(function ThreadWorkLogRow( icon={icon} color={props.iconSubtleColor} colorClassName={ - iconIsDestructive - ? "accent-adaptive-rose-600-400" - : failed - ? "accent-danger-foreground/40" - : undefined + isUsageLimit + ? "accent-warning-foreground" + : iconIsDestructive + ? "accent-adaptive-rose-600-400" + : failed + ? "accent-danger-foreground/40" + : undefined } /> )} - + {isSystemNotice ? row.summary : displayText} {answerPreview ? ( { }); describe("resolveThreadListV2Status", () => { + it("distinguishes usage limits from ordinary failures and clears the label after recovery", () => { + const thread = makeThread({ + id: ThreadId.make("limited"), + title: "Limited", + runtime: { + status: "failed", + activeRunId: null, + providerInstanceId: ProviderInstanceId.make("codex"), + providerName: "Codex", + lastError: "Plan limit reached", + lastErrorClass: "usage_limit", + updatedAt: NOW, + }, + }); + expect(resolveThreadListV2Status(thread)).toBe("limited"); + expect( + resolveThreadListV2Status({ + ...thread, + runtime: { ...thread.runtime!, lastErrorClass: null }, + }), + ).toBe("failed"); + expect( + resolveThreadListV2Status({ + ...thread, + runtime: { ...thread.runtime!, status: "completed" }, + }), + ).toBe("ready"); + }); + it("prioritizes approval over a running runtime", () => { const thread = makeThread({ id: ThreadId.make("t"), diff --git a/apps/mobile/src/features/threads/threadListV2.ts b/apps/mobile/src/features/threads/threadListV2.ts index e056deeac283..b01ee350b13d 100644 --- a/apps/mobile/src/features/threads/threadListV2.ts +++ b/apps/mobile/src/features/threads/threadListV2.ts @@ -59,7 +59,14 @@ export function resolveThreadListV2ProviderDrivers( * The orchestrator v2 presentation bridge parks runtime at idle when the * post-settlement background roster is nonempty. */ -export type ThreadListV2Status = "approval" | "input" | "working" | "waiting" | "failed" | "ready"; +export type ThreadListV2Status = + | "approval" + | "input" + | "working" + | "waiting" + | "failed" + | "limited" + | "ready"; export type ThreadListV2SwipeAction = "archive" | "settle" | "unsettle" | "snooze" | "unsnooze"; export function resolveThreadListV2SnoozeMenuSelection(input: { @@ -197,7 +204,7 @@ export function resolveThreadListV2Status( return "waiting"; } if (thread.runtime?.status === "failed") { - return "failed"; + return thread.runtime.lastErrorClass === "usage_limit" ? "limited" : "failed"; } return "ready"; } diff --git a/apps/mobile/src/features/threads/work-log-layout.tsx b/apps/mobile/src/features/threads/work-log-layout.tsx index d8d6dcd014d7..bd02feb00a39 100644 --- a/apps/mobile/src/features/threads/work-log-layout.tsx +++ b/apps/mobile/src/features/threads/work-log-layout.tsx @@ -55,7 +55,7 @@ export function WorkLogLabel({ tone = "default", }: { children: ReactNode; - tone?: "default" | "danger"; + tone?: "default" | "danger" | "warning"; }) { return ( {children} diff --git a/apps/mobile/src/lib/threadActivity.test.ts b/apps/mobile/src/lib/threadActivity.test.ts index 8468a38d8638..5520e86c27fa 100644 --- a/apps/mobile/src/lib/threadActivity.test.ts +++ b/apps/mobile/src/lib/threadActivity.test.ts @@ -338,6 +338,31 @@ describe("buildThreadFeed", () => { expect(presented.some((entry) => entry.type === "run-fold")).toBe(false); }); + it("presents a usage-limit stop as a warning while preserving its explanation", () => { + const message = "Plan usage limit reached. Try again after reset."; + const entries = buildThreadFeed([ + projected( + { + ...base("item-limit", "2026-06-20T00:00:02.000Z", 1), + type: "error", + status: "failed", + title: "Usage limit reached", + failure: { class: "usage_limit", message, code: "usageLimitExceeded", retryable: null }, + }, + 0, + ), + ]); + const activity = entries.flatMap((entry) => + entry.type === "activity-group" ? entry.activities : [], + )[0]; + expect(activity).toMatchObject({ + summary: "Usage limit reached", + status: "neutral", + icon: "warning", + }); + expect(activity?.getFullDetail()).toContain(message); + }); + it("presents provider retries as visible work-log activity", () => { const retryBase = { ...base("item-provider-retry", "2026-06-20T00:00:02.000Z", 1), diff --git a/apps/mobile/src/lib/threadActivity.ts b/apps/mobile/src/lib/threadActivity.ts index 37eb7b0f8285..7651fdf99cc8 100644 --- a/apps/mobile/src/lib/threadActivity.ts +++ b/apps/mobile/src/lib/threadActivity.ts @@ -362,7 +362,8 @@ function itemIsProminent(item: OrchestrationV2TurnItem): boolean { function itemStatus(item: OrchestrationV2TurnItem): ThreadFeedActivity["status"] { if (item.type === "notification") return item.outcome === "failed" ? "failure" : null; if (item.type === "error") { - if (item.status === "failed") return "failure"; + if (item.status === "failed") + return item.failure.class === "usage_limit" ? "neutral" : "failure"; return item.status === "completed" ? "success" : "neutral"; } if (!itemIsToolLike(item)) return null; @@ -432,7 +433,7 @@ function itemIcon(item: OrchestrationV2TurnItem): ThreadFeedActivity["icon"] { case "system_notice": return "warning"; case "error": - return "alert"; + return item.failure.class === "usage_limit" ? "warning" : "alert"; case "checkpoint": case "proposed_plan": case "todo_list": @@ -486,7 +487,7 @@ function itemSummary( case "run_interrupt_result": return "Run interrupted"; case "error": - return "Provider error"; + return item.failure.class === "usage_limit" ? "Usage limit reached" : "Provider error"; case "handoff": return "Context handed off"; case "fork": @@ -671,7 +672,12 @@ function toFeedActivity( logo: toolPresentation?.logo ?? null, toolLike: itemIsToolLike(item), prominent: itemIsProminent(item), - status: workEntryDisplayIndicatesToolFailure(workEntry) ? "failure" : itemStatus(item), + status: + item.type === "error" && item.failure.class === "usage_limit" + ? itemStatus(item) + : workEntryDisplayIndicatesToolFailure(workEntry) + ? "failure" + : itemStatus(item), lifecycleStatus: itemLifecycleStatus(item), workEntry, projectedItem: row, diff --git a/apps/server/src/orchestration-v2/Adapters/AcpAdapterV2.ts b/apps/server/src/orchestration-v2/Adapters/AcpAdapterV2.ts index b72f631a6aa7..36f177eba3b3 100644 --- a/apps/server/src/orchestration-v2/Adapters/AcpAdapterV2.ts +++ b/apps/server/src/orchestration-v2/Adapters/AcpAdapterV2.ts @@ -192,6 +192,8 @@ export interface AcpAdapterV2ExtensionContext { } export interface AcpAdapterV2Flavor { + /** Interprets provider-specific prompt errors before they cross into orchestration. */ + readonly promptFailure?: (cause: unknown) => OrchestrationV2ProviderFailure; readonly driver: ProviderDriverKind; readonly capabilities: OrchestrationV2ProviderCapabilities; readonly clientCapabilitiesMeta?: Record; @@ -6679,10 +6681,11 @@ export function makeAcpAdapterV2(options: AcpAdapterV2Options): ProviderAdapterV yield* finalizeTurn( context, context.interrupted ? "interrupted" : "failed", - makeProviderFailure({ - cause: Cause.squash(cause), - class: "provider_error", - }), + flavor.promptFailure?.(Cause.squash(cause)) ?? + makeProviderFailure({ + cause: Cause.squash(cause), + class: "provider_error", + }), ).pipe( Effect.andThen( Effect.logWarning("orchestration-v2.acp-prompt-failed", { diff --git a/apps/server/src/orchestration-v2/Adapters/ClaudeAdapterV2.test.ts b/apps/server/src/orchestration-v2/Adapters/ClaudeAdapterV2.test.ts index 4be6afe87b6c..ee2fb5d4a416 100644 --- a/apps/server/src/orchestration-v2/Adapters/ClaudeAdapterV2.test.ts +++ b/apps/server/src/orchestration-v2/Adapters/ClaudeAdapterV2.test.ts @@ -2327,6 +2327,7 @@ describe("ClaudeAdapterV2 background wake turns", () => { assert.equal(terminal.status, "failed"); if (terminal.status !== "failed") return; assert.include(terminal.failure.message, expected); + assert.equal(terminal.failure.class, recovered ? "provider_error" : "usage_limit"); }).pipe(Effect.provide(Layer.merge(idAllocatorLayer, NodeServices.layer))), ); @@ -2380,6 +2381,7 @@ describe("ClaudeAdapterV2 background wake turns", () => { const terminal = yield* Queue.take(harness.terminalReceipts); assert.equal(terminal.status, "failed"); if (terminal.status !== "failed") return; + assert.equal(terminal.failure.class, expectedLimit ? "usage_limit" : "provider_error"); assert.equal( terminal.failure.message, expectedLimit @@ -2389,6 +2391,50 @@ describe("ClaudeAdapterV2 background wake turns", () => { }).pipe(Effect.provide(Layer.merge(idAllocatorLayer, NodeServices.layer))), ); + it.effect.each([429, 401, 529])( + "classifies the current Claude API status %s after rate-limit evidence", + (apiErrorStatus) => + Effect.gen(function* () { + const harness = yield* makeWakeHarness; + yield* harness.runtime.startTurn( + makeClaudeTestTurnInput({ + threadId: harness.threadId, + providerThread: harness.providerThread, + now: yield* DateTime.now, + attemptId: RunAttemptId.make(`attempt-status-${apiErrorStatus}`), + text: "Continue.", + attachments: [], + }), + ); + yield* Queue.offer( + harness.sdkMessages, + makeAssistantErrorFrame({ + uuid: "00000000-0000-4000-8000-000000000650", + error: "rate_limit", + }), + ); + yield* Queue.offer( + harness.sdkMessages, + makeResultFrame({ + uuid: "00000000-0000-4000-8000-000000000651", + result: "API Error", + terminalReason: "api_error", + isError: true, + apiErrorStatus, + }), + ); + const terminal = yield* Queue.take(harness.terminalReceipts); + assert.equal(terminal.status, "failed"); + if (terminal.status !== "failed") return; + assert.equal( + terminal.failure.class, + apiErrorStatus === 429 ? "usage_limit" : "provider_error", + ); + if (apiErrorStatus !== 429) + assert.notInclude(terminal.failure.message.toLowerCase(), "usage limit"); + }).pipe(Effect.provide(Layer.merge(idAllocatorLayer, NodeServices.layer))), + ); + it.effect("surfaces a Claude safety model fallback without failing the turn", () => Effect.gen(function* () { const harness = yield* makeWakeHarness; @@ -2970,6 +3016,10 @@ describe("ClaudeAdapterV2 background wake turns", () => { assert.equal(terminal.status, "failed"); if (terminal.status !== "failed") return; assert.isNotEmpty(terminal.failure.message); + assert.equal( + terminal.failure.class, + terminalReason === "blocking_limit" ? "usage_limit" : "provider_error", + ); assert.isFalse( harness.events.some( (event) => diff --git a/apps/server/src/orchestration-v2/Adapters/ClaudeAdapterV2.ts b/apps/server/src/orchestration-v2/Adapters/ClaudeAdapterV2.ts index 2becedbc87ee..fca48f3881d6 100644 --- a/apps/server/src/orchestration-v2/Adapters/ClaudeAdapterV2.ts +++ b/apps/server/src/orchestration-v2/Adapters/ClaudeAdapterV2.ts @@ -2101,6 +2101,7 @@ function terminalStatusFromResult( // The SDK reports API-level failures (401 auth, 529 overloaded, …) as // subtype "success" with is_error set; the turn produced no real work. return isOverloadedResult(message) || + message.api_error_status === 429 || terminalResultError(message.terminal_reason, failureHint) !== undefined || (message.is_error && failureHint !== undefined) ? "failed" @@ -2138,16 +2139,25 @@ function isClaudeTaskNotificationOriginResult(message: SDKMessage): message is S function providerFailureFromResult( message: SDKResultMessage, failureHint?: string, + usageLimited = false, ): OrchestrationV2ProviderFailure | null { + const failureClass = + message.terminal_reason === "blocking_limit" || + (message.subtype === "success" && message.api_error_status === 429) || + usageLimited + ? "usage_limit" + : "provider_error"; const listedError = resultUserFacingError(message); const structuredError = isOverloadedResult(message) ? "Claude API is overloaded (529). Try again shortly." - : terminalResultError(message.terminal_reason, failureHint); + : message.subtype === "success" && message.api_error_status === 429 + ? "Claude API rate limit reached. Try again later." + : terminalResultError(message.terminal_reason, failureHint); if (message.subtype !== "success") { return makeProviderFailure({ message: listedError ?? structuredError ?? message.errors.join("\n"), code: message.subtype, - class: "provider_error", + class: failureClass, }); } if (!message.is_error && structuredError === undefined) { @@ -2160,7 +2170,7 @@ function providerFailureFromResult( apiErrorStatus === null ? (message.terminal_reason ?? "sdk_result_error") : `api_error_${apiErrorStatus}`, - class: "provider_error", + class: failureClass, retryable: apiErrorStatus === 429 || apiErrorStatus === 529 ? true : null, }); } @@ -2173,7 +2183,12 @@ function providerFailureFromApiRetry(message: SDKAPIRetryMessage): Orchestration message.error_status === null ? message.error : `api_error_${Math.trunc(message.error_status)}`, - class: message.error_status === null ? "transport_error" : "provider_error", + class: + message.error_status === 429 + ? "usage_limit" + : message.error_status === null + ? "transport_error" + : "provider_error", retryable: true, }); } @@ -5215,14 +5230,23 @@ export function makeClaudeAdapterV2( next.delete(context.providerTurnId); return next; }); + const usageLimited = + context.authenticationFailureMessage === undefined && + (context.rejectedRateLimitTypes.size > 0 || context.latestAssistantRateLimited) && + (message.subtype !== "success" || + message.api_error_status == null || + message.api_error_status === 429) && + (message.terminal_reason == null || + message.terminal_reason === "api_error" || + message.terminal_reason === "blocking_limit"); const failureHint = context.authenticationFailureMessage ?? - (context.rejectedRateLimitTypes.size > 0 || context.latestAssistantRateLimited + (usageLimited ? "Claude usage limit reached. Send the message again once the limit resets." : undefined); const resultFailure = interrupted ? null - : providerFailureFromResult(message, failureHint); + : providerFailureFromResult(message, failureHint, usageLimited); yield* finalizeActiveTurn({ context, status: interrupted ? "interrupted" : terminalStatusFromResult(message, failureHint), diff --git a/apps/server/src/orchestration-v2/Adapters/CodexAdapterV2.test.ts b/apps/server/src/orchestration-v2/Adapters/CodexAdapterV2.test.ts index d8fa9118cceb..10da6ac82043 100644 --- a/apps/server/src/orchestration-v2/Adapters/CodexAdapterV2.test.ts +++ b/apps/server/src/orchestration-v2/Adapters/CodexAdapterV2.test.ts @@ -4862,6 +4862,106 @@ describe("CodexAdapterV2 post-settle continuation", () => { ), ); + for (const scenario of [ + { + name: "usage", + code: "usageLimitExceeded", + notification: false, + expectedClass: "usage_limit", + }, + { name: "rate", code: "rateLimitExceeded", notification: false, expectedClass: "usage_limit" }, + { + name: "ordinary", + code: "contextWindowExceeded", + notification: false, + expectedClass: "provider_error", + }, + { + name: "notification", + code: "usageLimitExceeded", + notification: true, + expectedClass: "usage_limit", + }, + { + name: "replacement", + code: "usageLimitExceeded", + notification: true, + expectedClass: "provider_error", + }, + { name: "retry", code: "usageLimitExceeded", notification: true, expectedClass: "usage_limit" }, + ] as const) { + it.effect(`classifies Codex terminal failures from ${scenario.name} evidence`, () => + Effect.scoped( + Effect.gen(function* () { + const nativeThreadId = `native-limit-${scenario.name}`; + const nativeTurnId = `turn-limit-${scenario.name}`; + const message = "Provider stopped this request."; + const transcript = makeCodexReplayTranscript({ + scenario: `codex-limit-${scenario.name}`, + entries: [ + ...codexReplayPreamble({ nativeThreadId, nativeTurnId, prompt: "Continue." }), + ...(scenario.notification + ? [ + { + type: "emit_inbound" as const, + label: "error", + frame: { + method: "error", + params: { + threadId: nativeThreadId, + turnId: nativeTurnId, + willRetry: scenario.name === "retry", + error: { + message, + codexErrorInfo: scenario.code, + additionalDetails: null, + }, + }, + }, + }, + ] + : []), + { + type: "emit_inbound", + label: "turn/completed", + frame: { + method: "turn/completed", + params: { + threadId: nativeThreadId, + turn: { + ...makeCodexReplayTurn({ id: nativeTurnId, status: "failed" }), + error: { + message: scenario.name === "replacement" ? "A different failure." : message, + ...(scenario.notification ? {} : { codexErrorInfo: scenario.code }), + }, + }, + }, + }, + }, + ], + }); + const harness = yield* makeCodexReplayHarness(transcript); + yield* harness.runtime.startTurn( + makeCodexTestTurnInput({ + threadId: harness.threadId, + providerThread: harness.providerThread, + now: yield* DateTime.now, + text: "Continue.", + attemptId: RunAttemptId.make(`attempt-limit-${scenario.name}`), + }), + ); + yield* harness.firstTerminal; + const terminal = harness.terminalEvents()[0]; + assert.equal(terminal?.status, "failed"); + if (terminal?.status !== "failed") return; + assert.equal(terminal.failure.class, scenario.expectedClass); + assert.equal(terminal.threadDisposition, "reusable"); + if (scenario.name === "retry") assert.equal(terminal.retry?.attempt, 1); + }).pipe(Effect.provide(Layer.merge(idAllocatorLayer, NodeServices.layer))), + ), + ); + } + const FAILED_SCENARIO = "codex-failed-mid-command"; const FAILED_NATIVE_THREAD = "native-codex-failed-thread"; const FAILED_NATIVE_TURN = "native-codex-failed-turn"; diff --git a/apps/server/src/orchestration-v2/Adapters/CodexAdapterV2.ts b/apps/server/src/orchestration-v2/Adapters/CodexAdapterV2.ts index eb979952e469..559f0cdb73d1 100644 --- a/apps/server/src/orchestration-v2/Adapters/CodexAdapterV2.ts +++ b/apps/server/src/orchestration-v2/Adapters/CodexAdapterV2.ts @@ -959,6 +959,10 @@ function codexErrorInfoCode(value: unknown): string | null { } interface ActiveCodexTurnContext { + latestProviderFailure?: { + readonly nativeMessage: string; + readonly failure: OrchestrationV2ProviderFailure; + }; readonly nativeStartReady?: Deferred.Deferred; readonly input: ProviderAdapterV2TurnInput; readonly projectionAppThread: OrchestrationV2AppThread; @@ -978,6 +982,7 @@ interface ActiveCodexTurnContext { } interface ActiveCodexProviderRetry { + readonly nativeMessage: string; readonly retry: OrchestrationV2ProviderRetry; readonly failure: OrchestrationV2ProviderFailure; readonly startedAt: DateTime.Utc; @@ -3630,13 +3635,26 @@ export function makeCodexAdapterV2(adapterOptions: CodexAdapterV2Options): Provi yield* client.handleServerNotification("error", (payload) => Effect.gen(function* () { - if (!payload.willRetry) { - return; - } const context = yield* awaitActiveTurn(payload.turnId); if (context === undefined) { return; } + const notificationCode = codexErrorInfoCode(payload.error.codexErrorInfo); + if (!payload.willRetry) { + context.latestProviderFailure = { + nativeMessage: payload.error.message, + failure: makeProviderFailure({ + message: payload.error.additionalDetails?.trim() || payload.error.message, + code: notificationCode, + class: + notificationCode === "usageLimitExceeded" || + notificationCode === "rateLimitExceeded" + ? "usage_limit" + : "provider_error", + }), + }; + return; + } const updatedAt = yield* DateTime.now; const previous = (yield* Ref.get(providerRetries)).get(context.providerTurnId); const progress = parseCodexRetryProgress(payload.error.message); @@ -3654,15 +3672,18 @@ export function makeCodexAdapterV2(adapterOptions: CodexAdapterV2Options): Provi : additionalDetails, code, class: - code?.startsWith("http") === true || code?.startsWith("responseStream") === true - ? "transport_error" - : "provider_error", + code === "usageLimitExceeded" || code === "rateLimitExceeded" + ? "usage_limit" + : code?.startsWith("http") === true || code?.startsWith("responseStream") === true + ? "transport_error" + : "provider_error", retryable: true, }); const itemOrdinal = previous?.itemOrdinal ?? (yield* resolveItemOrdinal(context, `terminal-failure:${context.providerTurnId}`)); const state: ActiveCodexProviderRetry = { + nativeMessage: payload.error.message, retry, failure, startedAt: previous?.startedAt ?? updatedAt, @@ -4579,10 +4600,27 @@ export function makeCodexAdapterV2(adapterOptions: CodexAdapterV2Options): Provi readonly context: ActiveCodexTurnContext; readonly status: OrchestrationV2ProviderTurn["status"]; readonly failureMessage?: string; + readonly failureCode?: string | null; readonly providerRetry?: ActiveCodexProviderRetry; }): Effect.fn.Return { const terminalStatus = providerTurnStatusToTerminal(input.status); if (terminalStatus === "failed") { + const previousFailure = input.context.latestProviderFailure ?? input.providerRetry; + const failure = + input.failureCode === undefined && + previousFailure !== undefined && + (input.failureMessage === undefined || + input.failureMessage === previousFailure.nativeMessage) + ? previousFailure.failure + : makeProviderFailure({ + message: input.failureMessage, + code: input.failureCode, + class: + input.failureCode === "usageLimitExceeded" || + input.failureCode === "rateLimitExceeded" + ? "usage_limit" + : "provider_error", + }); return { type: "turn.terminal", driver: CODEX_PROVIDER, @@ -4594,13 +4632,7 @@ export function makeCodexAdapterV2(adapterOptions: CodexAdapterV2Options): Provi `terminal-failure:${input.context.providerTurnId}`, ), status: terminalStatus, - failure: - input.failureMessage === undefined && input.providerRetry !== undefined - ? input.providerRetry.failure - : makeProviderFailure({ - message: input.failureMessage, - class: "provider_error", - }), + failure, ...(input.providerRetry === undefined ? {} : { @@ -4629,6 +4661,7 @@ export function makeCodexAdapterV2(adapterOptions: CodexAdapterV2Options): Provi readonly nativeTurnId: string; readonly status: OrchestrationV2ProviderTurn["status"]; readonly failureMessage?: string; + readonly failureCode?: string | null; readonly providerRetry?: ActiveCodexProviderRetry; }) { const event = yield* makeRootTerminalEvent(input); @@ -4677,6 +4710,7 @@ export function makeCodexAdapterV2(adapterOptions: CodexAdapterV2Options): Provi readonly status: OrchestrationV2ProviderTurn["status"]; readonly completedAt: DateTime.Utc; readonly failureMessage?: string; + readonly failureCode?: string | null; }) => turnTerminalizationPermit.withPermits(1)( Effect.gen(function* () { @@ -4929,7 +4963,14 @@ export function makeCodexAdapterV2(adapterOptions: CodexAdapterV2Options): Provi completedAt: codexTimestamp(payload.turn.completedAt), ...(payload.turn.error?.message === undefined ? {} - : { failureMessage: payload.turn.error.message }), + : { + failureMessage: payload.turn.error.message, + ...(payload.turn.error.codexErrorInfo == null + ? {} + : { + failureCode: codexErrorInfoCode(payload.turn.error.codexErrorInfo), + }), + }), }); }), ); diff --git a/apps/server/src/orchestration-v2/Adapters/GrokAdapterV2.test.ts b/apps/server/src/orchestration-v2/Adapters/GrokAdapterV2.test.ts index 8524ca5d3cb6..99d780c804a5 100644 --- a/apps/server/src/orchestration-v2/Adapters/GrokAdapterV2.test.ts +++ b/apps/server/src/orchestration-v2/Adapters/GrokAdapterV2.test.ts @@ -1,3 +1,5 @@ +import * as EffectAcpErrors from "effect-acp/errors"; +import { xAiRateLimitedErrorCode } from "../../provider/acp/XAiAcpExtension.ts"; import { assert, describe, it } from "@effect/vitest"; import * as Effect from "effect/Effect"; import type * as EffectAcpSchema from "effect-acp/compat"; @@ -63,6 +65,34 @@ describe("acpSubagentStatusBlocksTurnSettlement", () => { }); describe("GrokAdapterV2 capabilities", () => { + it("preserves Grok's rate-limit stop and distinguishes other prompt failures", () => { + const flavor = makeGrokAcpAdapterFlavor({ + makeRuntime: () => Effect.never, + } as unknown as GrokAdapterV2Options); + const limit = flavor.promptFailure?.( + new EffectAcpErrors.AcpRequestError({ + code: xAiRateLimitedErrorCode, + errorMessage: "Grok usage limit reached. Try again later.", + }), + ); + assert.equal(limit?.class, "usage_limit"); + assert.equal(limit?.code, String(xAiRateLimitedErrorCode)); + assert.equal(limit?.message, "Grok usage limit reached. Try again later."); + assert.equal( + flavor.promptFailure?.( + new EffectAcpErrors.AcpRequestError({ + code: -32603, + errorMessage: "Internal error", + }), + ).class, + "provider_error", + ); + assert.equal( + flavor.promptFailure?.(new Error("Rate limit mentioned in an ordinary error")).class, + "provider_error", + ); + }); + it("wires hard Stop teardown but soft non-Stop interrupts in the constructor flavor", () => { const flavor = makeGrokAcpAdapterFlavor({ makeRuntime: () => Effect.never, diff --git a/apps/server/src/orchestration-v2/Adapters/GrokAdapterV2.ts b/apps/server/src/orchestration-v2/Adapters/GrokAdapterV2.ts index b3c73d9bf791..125d78982178 100644 --- a/apps/server/src/orchestration-v2/Adapters/GrokAdapterV2.ts +++ b/apps/server/src/orchestration-v2/Adapters/GrokAdapterV2.ts @@ -1,3 +1,5 @@ +import { makeProviderFailure } from "../ProviderFailure.ts"; +import { xAiRateLimitedErrorCode } from "../../provider/acp/XAiAcpExtension.ts"; import { HostProcessEnvironment, HostProcessPlatform } from "@t3tools/shared/hostProcess"; import { defaultInstanceIdForDriver, @@ -12,7 +14,7 @@ import * as Layer from "effect/Layer"; import * as Schema from "effect/Schema"; import type * as Scope from "effect/Scope"; import { ChildProcessSpawner } from "effect/unstable/process"; -import type * as EffectAcpErrors from "effect-acp/errors"; +import * as EffectAcpErrors from "effect-acp/errors"; import { ServerConfig } from "../../config.ts"; import { makeAcpNativeLoggerFactory } from "../../provider/acp/AcpNativeLogging.ts"; @@ -253,6 +255,17 @@ export function makeGrokAcpAdapterFlavor(options: GrokAdapterV2Options): AcpAdap environment: options.environment, childProcessSpawner: options.childProcessSpawner, })), + promptFailure: (cause) => + makeProviderFailure({ + cause, + ...(Schema.is(EffectAcpErrors.AcpRequestError)(cause) + ? { + message: cause.errorMessage, + code: String(cause.code), + class: cause.code === xAiRateLimitedErrorCode ? "usage_limit" : "provider_error", + } + : { class: "provider_error" }), + }), registerExtensions: registerGrokAcpExtensions, extractSubagentUpdate: extractXAiAcpSubagentUpdate, extractSubagentEndNotice: extractXAiAcpSubagentEndNotice, diff --git a/apps/server/src/orchestration-v2/ProjectionStore.test.ts b/apps/server/src/orchestration-v2/ProjectionStore.test.ts index e38e88b3e176..40edf27afc54 100644 --- a/apps/server/src/orchestration-v2/ProjectionStore.test.ts +++ b/apps/server/src/orchestration-v2/ProjectionStore.test.ts @@ -1918,6 +1918,152 @@ it.layer(TestLayer)("ProjectionStoreV2", (it) => { }), ); + it.effect("projects only the latest failed root turn's limit into SQL and memory shells", () => + Effect.gen(function* () { + const store = yield* ProjectionStoreV2; + const threadId = yield* addRolledBackRecoveryCandidate("limit-shell"); + const otherThreadId = yield* addRolledBackRecoveryCandidate("other-limit-shell"); + const original = (yield* store.getThreadProjection(threadId)).runs[0]!; + const now = yield* DateTime.now; + const limitItem = { + id: TurnItemId.make("limit-shell:error"), + threadId, + runId: original.id, + nodeId: original.rootNodeId, + providerThreadId: null, + providerTurnId: null, + nativeItemRef: null, + parentItemId: null, + ordinal: 2, + status: "failed" as const, + title: "Usage limit reached", + startedAt: now, + completedAt: now, + updatedAt: now, + type: "error" as const, + failure: { + class: "usage_limit" as const, + message: "Plan limit reached.", + code: "usageLimitExceeded", + retryable: null, + }, + }; + const applyRun = (status: typeof original.status, rootNodeId = original.rootNodeId) => + store.apply({ + id: EventId.make(`event:limit-shell:run:${status}:${rootNodeId}`), + type: "run.updated", + threadId, + occurredAt: now, + payload: { ...original, rootNodeId, status }, + }); + const assertSummary = Effect.fnUntraced(function* ( + lastError: string | null, + lastErrorClass: string | null, + ) { + const projection = yield* store.getThreadProjection(threadId); + const memoryShell = threadShellFromProjection(projection); + const shells = yield* store.getShellSnapshot(); + const sqlShell = shells.threads.find((row) => row.id === threadId)!; + for (const shell of [memoryShell, sqlShell]) { + assert.equal(shell.lastError, lastError); + assert.equal(shell.lastErrorClass, lastErrorClass); + } + assert.isNull(shells.threads.find((row) => row.id === otherThreadId)!.lastErrorClass); + }); + yield* store.apply({ + id: EventId.make("event:limit-shell:error"), + type: "turn-item.updated", + threadId, + occurredAt: now, + payload: limitItem, + }); + yield* applyRun("failed"); + yield* assertSummary("Plan limit reached.", "usage_limit"); + const session = { + id: ProviderSessionId.make("session:limit-shell:shared"), + driver, + providerInstanceId, + status: "ready" as const, + cwd: "/workspace", + model: modelSelection.model, + capabilities: CodexProviderCapabilitiesV2, + createdAt: now, + updatedAt: now, + lastError: null, + }; + for (const boundThreadId of [threadId, otherThreadId]) { + yield* store.apply({ + id: EventId.make(`event:limit-shell:bind:${boundThreadId}`), + type: "provider-session.attached", + threadId: boundThreadId, + driver, + providerInstanceId, + occurredAt: now, + payload: session, + }); + } + yield* assertSummary("Plan limit reached.", "usage_limit"); + yield* store.apply({ + id: EventId.make("event:limit-shell:session-failed"), + type: "provider-session.updated", + threadId, + occurredAt: now, + payload: { ...session, status: "error", lastError: "Provider process exited." }, + }); + yield* assertSummary("Provider process exited.", null); + yield* store.apply({ + id: EventId.make("event:limit-shell:session-recovered"), + type: "provider-session.updated", + threadId, + occurredAt: now, + payload: session, + }); + yield* assertSummary("Plan limit reached.", "usage_limit"); + // A failed child is visible in history but does not replace the root's reason. + yield* store.apply({ + id: EventId.make("event:limit-shell:child-error"), + type: "turn-item.updated", + threadId, + occurredAt: now, + payload: { + ...limitItem, + id: TurnItemId.make("limit-shell:child-error"), + nodeId: NodeId.make("child-node"), + ordinal: 3, + failure: { ...limitItem.failure, class: "provider_error", message: "Child failed." }, + }, + }); + yield* assertSummary("Plan limit reached.", "usage_limit"); + // A later ordinary root error replaces the limit classification. + yield* store.apply({ + id: EventId.make("event:limit-shell:replacement"), + type: "turn-item.updated", + threadId, + occurredAt: now, + payload: { + ...limitItem, + id: TurnItemId.make("limit-shell:replacement"), + ordinal: 4, + failure: { ...limitItem.failure, class: "provider_error", message: "Provider failed." }, + }, + }); + yield* assertSummary("Provider failed.", "provider_error"); + for (const status of [ + "running", + "completed", + "interrupted", + "cancelled", + "rolled_back", + ] as const) { + yield* applyRun(status); + yield* assertSummary(null, null); + } + // A new attempt's root cannot inherit an earlier attempt's limit. + yield* applyRun("failed", NodeId.make("new-attempt-root")); + yield* assertSummary(null, null); + }), + ); + it.effect("projects one shared provider session into multiple thread bindings", () => Effect.gen(function* () { const projectionStore = yield* ProjectionStoreV2; diff --git a/apps/server/src/orchestration-v2/ProjectionStore.ts b/apps/server/src/orchestration-v2/ProjectionStore.ts index cbef23b92134..35cfcd06b8b2 100644 --- a/apps/server/src/orchestration-v2/ProjectionStore.ts +++ b/apps/server/src/orchestration-v2/ProjectionStore.ts @@ -1,3 +1,7 @@ +import { + latestRootProviderFailure, + threadErrorSummary, +} from "@t3tools/shared/orchestrationV2ThreadError"; import { threadPullRequestsOf } from "@t3tools/shared/threadPullRequests"; import type { OrchestrationV2AppThread, @@ -726,6 +730,7 @@ type ShellThreadRow = { readonly activity_run_status: string | null; readonly activity_run_started_at: string | null; readonly last_error: string | null; + readonly terminal_failure_payload_json: string | null; readonly pending_request_payload_json: string | null; readonly latest_user_message_at: string | null; readonly has_actionable_proposed_plan: number; @@ -1182,7 +1187,10 @@ export function threadShellFromProjection( activityRunStatus: activityRun?.status ?? null, activityRunStartedAt: activityRun?.startedAt ?? activityRun?.requestedAt ?? null, status: latestRun?.status ?? "idle", - lastError: providerSession?.lastError ?? null, + ...threadErrorSummary( + latestRootProviderFailure(latestRun, projection.turnItems), + providerSession?.lastError ?? null, + ), pendingRuntimeRequest: pendingRuntimeRequest === null ? null @@ -1272,6 +1280,7 @@ type ShellThreadState = { readonly activityRunStatus: ShellActivityRunStatus | null; readonly activityRunStartedAt: DateTime.Utc | null; readonly lastError: string | null; + readonly lastErrorClass: OrchestrationV2ThreadShell["lastErrorClass"]; readonly pendingRuntimeRequest: OrchestrationV2ThreadProjection["runtimeRequests"][number] | null; readonly latestUserMessageAt: DateTime.Utc | null; readonly hasActionableProposedPlan: boolean; @@ -1407,6 +1416,7 @@ function shellFromState(input: { activityRunStartedAt: input.state.activityRunStartedAt, status: input.state.latestRunStatus, lastError: input.state.lastError, + lastErrorClass: input.state.lastErrorClass, pendingRuntimeRequest: input.state.pendingRuntimeRequest === null ? null @@ -3993,6 +4003,21 @@ export const layer: Layer.Layer = ORDER BY session.updated_at DESC, session.provider_session_id DESC LIMIT 1 ) AS last_error, + ( + SELECT item.payload_json + FROM orchestration_v2_projection_turn_items item + INNER JOIN orchestration_v2_projection_runs r ON r.run_id = item.run_id + WHERE r.run_id = ( + SELECT latest.run_id FROM orchestration_v2_projection_runs latest + WHERE latest.thread_id = t.thread_id + ORDER BY latest.ordinal DESC, latest.run_id DESC LIMIT 1 + ) + AND r.status = 'failed' + AND item.type = 'error' AND item.status = 'failed' + AND item.node_id IS json_extract(r.payload_json, '$.rootNodeId') + ORDER BY item.updated_at DESC, item.ordinal DESC, item.turn_item_id DESC + LIMIT 1 + ) AS terminal_failure_payload_json, ( SELECT request.payload_json FROM orchestration_v2_projection_runtime_requests request @@ -4288,6 +4313,10 @@ export const layer: Layer.Layer = row.pending_request_payload_json === null ? null : yield* decodeRuntimeRequestPayload(row.pending_request_payload_json); + const terminalFailureItem = + row.terminal_failure_payload_json === null + ? null + : yield* decodeTurnItemPayload(row.terminal_failure_payload_json); const latestRunId = row.latest_run_id === null ? null : RunId.make(row.latest_run_id); const latestRunStatus = shellStatusFromStoredRunStatus(row.latest_run_status); const pendingBackgroundTasks = [ @@ -4334,7 +4363,10 @@ export const layer: Layer.Layer = row.activity_run_status === "waiting" ? row.activity_run_status : null, - lastError: row.last_error, + ...threadErrorSummary( + terminalFailureItem?.type === "error" ? terminalFailureItem.failure : null, + row.last_error, + ), pendingRuntimeRequest, latestUserMessageAt: row.latest_user_message_at === null diff --git a/apps/server/src/orchestration-v2/ProviderFailure.ts b/apps/server/src/orchestration-v2/ProviderFailure.ts index a0e649d25ae9..0400eaf05ff2 100644 --- a/apps/server/src/orchestration-v2/ProviderFailure.ts +++ b/apps/server/src/orchestration-v2/ProviderFailure.ts @@ -137,7 +137,7 @@ export function makeProviderFailureTurnItem(input: { parentItemId: null, ordinal: input.itemOrdinal, status: "failed", - title: "Provider error", + title: input.failure.class === "usage_limit" ? "Usage limit reached" : "Provider error", startedAt: input.retryStartedAt ?? input.occurredAt, completedAt: input.occurredAt, updatedAt: input.occurredAt, @@ -170,7 +170,7 @@ export function makeProviderRetryTurnItem(input: { if (input.status === "completed") { title = "Provider recovered"; } else if (input.status === "failed") { - title = "Provider error"; + title = input.failure.class === "usage_limit" ? "Usage limit reached" : "Provider error"; } else if (input.status === "interrupted" || input.status === "cancelled") { title = "Provider retry stopped"; } diff --git a/apps/server/src/provider/acp/XAiAcpExtension.ts b/apps/server/src/provider/acp/XAiAcpExtension.ts index 3f5851f1c269..8afc286a74ed 100644 --- a/apps/server/src/provider/acp/XAiAcpExtension.ts +++ b/apps/server/src/provider/acp/XAiAcpExtension.ts @@ -15,7 +15,7 @@ import type * as AcpSessionRuntime from "./AcpSessionRuntime.ts"; import type { AcpToolCallState } from "./AcpRuntimeModel.ts"; const xAiStopReasonMissingMetaKey = "xAiStopReasonMissing"; -const xAiRateLimitedErrorCode = -32003; +export const xAiRateLimitedErrorCode = -32003; const completedXAiPromptIdLimit = 128; const XAiPromptCompleteNotification = Schema.Struct({ diff --git a/apps/web/src/components/ChatView.tsx b/apps/web/src/components/ChatView.tsx index 56fc29c89a2d..f2a43c9dfb21 100644 --- a/apps/web/src/components/ChatView.tsx +++ b/apps/web/src/components/ChatView.tsx @@ -10357,6 +10357,11 @@ export default function ChatView(props: ChatViewProps) { /> { setThreadError(activeThread.id, null); dismissThreadErrorBannerForSession(threadErrorBannerKey); diff --git a/apps/web/src/components/Sidebar.logic.test.ts b/apps/web/src/components/Sidebar.logic.test.ts index 674c9aa3b8df..6a3e38805eac 100644 --- a/apps/web/src/components/Sidebar.logic.test.ts +++ b/apps/web/src/components/Sidebar.logic.test.ts @@ -894,6 +894,34 @@ describe("resolveSidebarThreadStatus", () => { ).toBe("working"); }); + it("keeps usage-limit stops Limited and visible until the thread recovers", () => { + const limited = { + ...runtime, + status: "failed" as const, + lastError: "Plan limit reached", + lastErrorClass: "usage_limit" as const, + }; + expect(resolveSidebarThreadStatus({ ...idle, runtime: limited })).toBe("limited"); + expect( + resolveSidebarThreadStatus({ ...idle, runtime: { ...limited, status: "running" } }), + ).toBe("working"); + expect( + resolveSidebarThreadStatus({ ...idle, runtime: { ...limited, status: "completed" } }), + ).toBe("ready"); + expect(resolveSidebarV2TopStatus({ status: "limited", isUnread: false, isWoke: false })).toBe( + "limited", + ); + expect( + shouldRecedeSidebarThread({ + status: "limited", + isUnread: false, + isWoke: false, + isActive: false, + isSelected: false, + }), + ).toBe(false); + }); + it("reports failed only while the latest run failed", () => { expect( resolveSidebarThreadStatus({ diff --git a/apps/web/src/components/Sidebar.logic.ts b/apps/web/src/components/Sidebar.logic.ts index 418e03b5d564..536fa2c8f616 100644 --- a/apps/web/src/components/Sidebar.logic.ts +++ b/apps/web/src/components/Sidebar.logic.ts @@ -877,7 +877,14 @@ export function resolveThreadRowClassName(input: { // false Done. // Unread completion is tracked separately: it describes whether a ready // thread needs attention, not what the thread is currently doing. -export type SidebarThreadStatus = "approval" | "input" | "working" | "waiting" | "failed" | "ready"; +export type SidebarThreadStatus = + | "approval" + | "input" + | "working" + | "waiting" + | "failed" + | "limited" + | "ready"; export function shouldRecedeSidebarThread(input: { status: SidebarThreadStatus; @@ -916,7 +923,7 @@ export function resolveSidebarThreadStatus(thread: SidebarThreadStatusInput): Si return "waiting"; } if (thread.runtime?.status === "failed") { - return "failed"; + return thread.runtime.lastErrorClass === "usage_limit" ? "limited" : "failed"; } return "ready"; } @@ -925,6 +932,7 @@ export type SidebarV2TopStatusKind = | "approval" | "done" | "failed" + | "limited" | "input" | "waiting" | "woke" @@ -947,8 +955,8 @@ export function resolveSidebarV2TopStatus(input: { if (input.status === "input") { return "input"; } - if (input.status === "failed") { - return "failed"; + if (input.status === "failed" || input.status === "limited") { + return input.status; } if (input.isWoke) { return "woke"; diff --git a/apps/web/src/components/Sidebar.tsx b/apps/web/src/components/Sidebar.tsx index 47ddda96d63e..484664aff397 100644 --- a/apps/web/src/components/Sidebar.tsx +++ b/apps/web/src/components/Sidebar.tsx @@ -506,9 +506,20 @@ function SidebarThreadTooltip({ ) : null} {thread.runtime?.lastError ? ( -
+
-
Error occurred
+
+ {thread.runtime.lastErrorClass === "usage_limit" + ? "Usage limit reached" + : "Error occurred"} +
) : null}
@@ -1233,25 +1244,31 @@ const SidebarThreadRow = memo(function SidebarThreadRow(props: { icon: "input" as const, className: "text-indigo-600 dark:text-indigo-300", } - : status === "failed" + : status === "limited" ? { - label: "Failed", + label: "Limited", icon: "failed" as const, - className: "text-red-700 dark:text-red-300", + className: "text-amber-700 dark:text-amber-300", } - : isWoke + : status === "failed" ? { - label: "Woke", - icon: "woke" as const, - className: "text-amber-700 dark:text-amber-300", + label: "Failed", + icon: "failed" as const, + className: "text-red-700 dark:text-red-300", } - : isUnread + : isWoke ? { - label: "Done", - icon: "done" as const, - className: "text-emerald-700 dark:text-emerald-300", + label: "Woke", + icon: "woke" as const, + className: "text-amber-700 dark:text-amber-300", } - : null; + : isUnread + ? { + label: "Done", + icon: "done" as const, + className: "text-emerald-700 dark:text-emerald-300", + } + : null; const isWokeStatus = topStatus?.icon === "woke"; const branchMismatch = resolveLocalCheckoutBranchMismatch({ diff --git a/apps/web/src/components/ThreadNotificationCoordinator.test.tsx b/apps/web/src/components/ThreadNotificationCoordinator.test.tsx index e25ab858a61e..6db13a3e1e03 100644 --- a/apps/web/src/components/ThreadNotificationCoordinator.test.tsx +++ b/apps/web/src/components/ThreadNotificationCoordinator.test.tsx @@ -18,6 +18,7 @@ const state = vi.hoisted(() => ({ approval: false, sessionError: false, turnError: false, + limited: false, subagent: false, add: vi.fn( (_toast: { title: string; description: string; actionProps: { onClick: () => void } }) => @@ -57,9 +58,10 @@ function mockThreadShell() { activeRunId: null, status: state.completedAt ? "completed" - : state.sessionError || state.turnError + : state.sessionError || state.turnError || state.limited ? "failed" : "running", + lastErrorClass: state.limited ? "usage_limit" : null, pendingRuntimeRequest: state.input ? { id: "request-1", kind: "user_input", createdAt: SHELL_NOW } : state.approval @@ -147,6 +149,7 @@ beforeEach(() => { approval: false, sessionError: false, turnError: false, + limited: false, subagent: false, }); vi.stubGlobal("IS_REACT_ACT_ENVIRONMENT", true); @@ -218,6 +221,7 @@ describe("thread notifications", () => { ["approval", "Approval needed"], ["sessionError", "Thread failed"], ["turnError", "Thread failed"], + ["limited", "Usage limit reached"], ] as const)("uses the same %s event for in-app and desktop alerts", async (event, title) => { state.mode = "notifications-and-sound"; await render(); diff --git a/apps/web/src/components/ThreadNotificationCoordinator.tsx b/apps/web/src/components/ThreadNotificationCoordinator.tsx index b5c7980c0306..0bde089a8898 100644 --- a/apps/web/src/components/ThreadNotificationCoordinator.tsx +++ b/apps/web/src/components/ThreadNotificationCoordinator.tsx @@ -115,7 +115,7 @@ function EnvironmentNotifications({ if (status === "ready" && thread.latestRun?.status === "failed") status = "failed"; const prior = previous.current.get(thread.id); const attention = - status === "input" || status === "approval" || status === "failed" + status === "input" || status === "approval" || status === "failed" || status === "limited" ? `${thread.latestRun?.runId ?? ""}:${status}` : null; const completedAt = Date.parse(thread.latestRun?.completedAt ?? ""); @@ -139,9 +139,11 @@ function EnvironmentNotifications({ ? "Thread completed" : status === "approval" ? "Approval needed" - : status === "failed" - ? "Thread failed" - : "Input needed"; + : status === "limited" + ? "Usage limit reached" + : status === "failed" + ? "Thread failed" + : "Input needed"; if (hasNotificationSound(mode)) { void playNotificationSound(kind, () => hasNotificationSound(getClientSettings().notificationMode), diff --git a/apps/web/src/components/chat/MessagesTimeline.tsx b/apps/web/src/components/chat/MessagesTimeline.tsx index 82880be8f6a5..1c0577557479 100644 --- a/apps/web/src/components/chat/MessagesTimeline.tsx +++ b/apps/web/src/components/chat/MessagesTimeline.tsx @@ -2632,7 +2632,7 @@ function v2EventPresentation(item: OrchestrationV2TurnItem): { tone: item.status === "completed" ? "success" - : item.status === "running" + : item.status === "running" || item.failure.class === "usage_limit" ? "warning" : "danger", icon: CircleAlertIcon, @@ -4873,8 +4873,10 @@ const SimpleWorkEntryRow = memo(function SimpleWorkEntryRow(props: { }; const iconConfig = workToneIcon(workEntry.tone); const showWarningIndicator = workEntry.sourceActivityKind === "runtime.warning"; - const showFailedIndicator = workEntryDisplayIndicatesToolFailure(workEntry); + const showFailedIndicator = + !showWarningIndicator && workEntryDisplayIndicatesToolFailure(workEntry); const showDestructiveRowStyle = + !showWarningIndicator && showFailedIndicator && (workEntrySignalsSevereFailure(workEntry) || !workLogEntryIsToolLike(workEntry)); const entryIconName = diff --git a/apps/web/src/components/chat/ThreadErrorBanner.tsx b/apps/web/src/components/chat/ThreadErrorBanner.tsx index 29a6a6db1846..1a01b37fcf04 100644 --- a/apps/web/src/components/chat/ThreadErrorBanner.tsx +++ b/apps/web/src/components/chat/ThreadErrorBanner.tsx @@ -1,3 +1,4 @@ +import type { OrchestrationV2ProviderFailureClass } from "@t3tools/contracts"; import { memo } from "react"; import { Alert, AlertAction, AlertDescription } from "../ui/alert"; import { Button } from "../ui/button"; @@ -36,18 +37,21 @@ export function isThreadErrorBannerDismissedForSession(bannerKey: string | null) export const ThreadErrorBanner = memo(function ThreadErrorBanner({ error, onDismiss, + errorClass, }: { error: string | null; + errorClass?: OrchestrationV2ProviderFailureClass | null; onDismiss?: () => void; }) { if (!error) return null; + const variant = errorClass === "usage_limit" ? "warning" : "error"; return (
@@ -61,7 +65,7 @@ export const ThreadErrorBanner = memo(function ThreadErrorBanner({ {onDismiss && ( )} diff --git a/apps/web/src/session-logic.test.ts b/apps/web/src/session-logic.test.ts index cede1974d0c3..d4438390277a 100644 --- a/apps/web/src/session-logic.test.ts +++ b/apps/web/src/session-logic.test.ts @@ -99,6 +99,13 @@ describe("V2 session presentation", () => { }, } satisfies Extract; + expect( + providerErrorPresentation({ + ...retryItem, + status: "failed", + failure: { ...retryItem.failure, class: "usage_limit" }, + }), + ).toMatchObject({ label: "Usage limit reached after 2/10 retries" }); expect(providerErrorPresentation(retryItem)).toEqual({ label: "Retrying provider (2/10)", detail: "Claude API overloaded. Retrying in 1.5s.", diff --git a/apps/web/src/session-logic.ts b/apps/web/src/session-logic.ts index 4c8f68652687..1a0ceba7c574 100644 --- a/apps/web/src/session-logic.ts +++ b/apps/web/src/session-logic.ts @@ -360,7 +360,10 @@ export function providerErrorPresentation( ): { readonly label: string; readonly detail: string } { if (item.retry === undefined) { return { - label: item.title?.trim() || "Provider error", + label: + item.failure.class === "usage_limit" + ? "Usage limit reached" + : item.title?.trim() || "Provider error", detail: item.failure.message, }; } @@ -374,7 +377,7 @@ export function providerErrorPresentation( : item.status === "completed" ? `Provider recovered (${progress} retries)` : item.status === "failed" - ? `Provider error after ${progress} retries` + ? `${item.failure.class === "usage_limit" ? "Usage limit reached" : "Provider error"} after ${progress} retries` : `Provider retry stopped (${progress})`; const retryDelay = item.status === "running" && item.retry.retryDelayMs !== null && item.retry.retryDelayMs > 0 @@ -478,7 +481,11 @@ function projectedWorkEntry(row: OrchestrationV2ProjectedTurnItem): WorkLogEntry return { ...common, ...presentation, - ...(item.retry === undefined ? { sourceActivityKind: "runtime.error" } : {}), + ...(item.failure.class === "usage_limit" + ? { sourceActivityKind: "runtime.warning" } + : item.retry === undefined + ? { sourceActivityKind: "runtime.error" } + : {}), toolData: item, }; } diff --git a/docs/user/thread-sidebar.md b/docs/user/thread-sidebar.md index 042442807e36..4bd5d0edab71 100644 --- a/docs/user/thread-sidebar.md +++ b/docs/user/thread-sidebar.md @@ -134,6 +134,10 @@ for custom configuration. ## Inspect agent work +**Limited** means the provider stopped on a usage or rate limit. The conversation +keeps the provider's explanation. Retry after the limit resets, or switch to +another provider instance. + On web and desktop, use **Agents** to follow work delegated to subagents. Expand a tool call in the conversation to see its full command and output. diff --git a/packages/client-runtime/src/state/models.ts b/packages/client-runtime/src/state/models.ts index a34c938b3a60..3b992caec60c 100644 --- a/packages/client-runtime/src/state/models.ts +++ b/packages/client-runtime/src/state/models.ts @@ -5,6 +5,7 @@ import type { MessageId, OrchestrationProjectShell, OrchestrationV2RunStatus, + OrchestrationV2ProviderFailureClass, OrchestrationV2ThreadProjection, OrchestrationV2ThreadShell, PlanId, @@ -53,6 +54,7 @@ export interface ThreadRuntimeSummary { readonly providerInstanceId: ProviderInstanceId; readonly providerName: string | null; readonly lastError: string | null; + readonly lastErrorClass?: OrchestrationV2ProviderFailureClass | null; readonly updatedAt: string; } @@ -173,6 +175,7 @@ function shellRuntime(thread: OrchestrationV2ThreadShell): ThreadRuntimeSummary providerInstanceId: thread.providerInstanceId, providerName: null, lastError: thread.lastError ?? null, + lastErrorClass: thread.lastErrorClass ?? null, updatedAt: iso(thread.updatedAt), }; } diff --git a/packages/client-runtime/src/state/threadExecution.test.ts b/packages/client-runtime/src/state/threadExecution.test.ts index 6df118f80699..c9d4bd9c0e26 100644 --- a/packages/client-runtime/src/state/threadExecution.test.ts +++ b/packages/client-runtime/src/state/threadExecution.test.ts @@ -1,4 +1,10 @@ -import { MessageId, RunId, type OrchestrationV2RunStatus } from "@t3tools/contracts"; +import { + TurnItemId, + NodeId, + MessageId, + RunId, + type OrchestrationV2RunStatus, +} from "@t3tools/contracts"; import * as DateTime from "effect/DateTime"; import { describe, expect, it } from "vite-plus/test"; @@ -34,6 +40,53 @@ function run(id: string, ordinal: number, status: OrchestrationV2RunStatus) { } describe("thread execution presentation", () => { + it("derives the current root failure without inheriting errors from children or previous runs", () => { + const failed = { ...run("limited", 1, "failed"), rootNodeId: NodeId.make("root") }; + const item = { + id: TurnItemId.make("limit-error"), + threadId: v2Projection.thread.id, + runId: failed.id, + nodeId: failed.rootNodeId, + providerThreadId: null, + providerTurnId: null, + nativeItemRef: null, + parentItemId: null, + ordinal: 1, + type: "error" as const, + status: "failed" as const, + title: "Usage limit reached", + startedAt: now, + completedAt: now, + updatedAt: now, + failure: { + class: "usage_limit" as const, + message: "Plan limit reached", + code: "usageLimitExceeded", + retryable: null, + }, + }; + const projection = { ...v2Projection, runs: [failed], turnItems: [item] }; + expect(deriveThreadRuntime(projection)).toMatchObject({ + lastError: "Plan limit reached", + lastErrorClass: "usage_limit", + }); + expect( + deriveThreadRuntime({ + ...projection, + turnItems: [{ ...item, nodeId: NodeId.make("child") }], + }), + ).toMatchObject({ lastError: null, lastErrorClass: null }); + expect( + deriveThreadRuntime({ + ...projection, + runs: [{ ...failed, rootNodeId: NodeId.make("new-root") }], + }), + ).toMatchObject({ lastError: null, lastErrorClass: null }); + expect( + deriveThreadRuntime({ ...projection, runs: [failed, run("new", 2, "running")] }), + ).toMatchObject({ status: "running", lastError: null, lastErrorClass: null }); + }); + it("keeps live activity attached to an executing run when a newer run is queued", () => { const runningRun = run("run-running", 1, "running"); const queuedRun = run("run-queued", 2, "queued"); diff --git a/packages/client-runtime/src/state/threadExecution.ts b/packages/client-runtime/src/state/threadExecution.ts index 83995c073dd4..9a7012cc57d5 100644 --- a/packages/client-runtime/src/state/threadExecution.ts +++ b/packages/client-runtime/src/state/threadExecution.ts @@ -1,3 +1,7 @@ +import { + latestRootProviderFailure, + threadErrorSummary, +} from "@t3tools/shared/orchestrationV2ThreadError"; import type { OrchestrationV2ThreadProjection } from "@t3tools/contracts"; import { derivePendingBackgroundWork } from "@t3tools/shared/orchestrationV2PendingBackgroundWork"; import * as DateTime from "effect/DateTime"; @@ -92,7 +96,10 @@ export function deriveThreadRuntime( : null, providerInstanceId: projection.thread.providerInstanceId, providerName: providerSession?.driver ?? null, - lastError: providerSession?.lastError ?? null, + ...threadErrorSummary( + latestRootProviderFailure(latestRunProjection, projection.turnItems), + providerSession?.lastError ?? null, + ), updatedAt: DateTime.formatIso(projection.updatedAt), }; } diff --git a/packages/contracts/src/orchestrationV2.test.ts b/packages/contracts/src/orchestrationV2.test.ts index 313a0b24df54..ea34526d294a 100644 --- a/packages/contracts/src/orchestrationV2.test.ts +++ b/packages/contracts/src/orchestrationV2.test.ts @@ -585,6 +585,14 @@ describe("orchestration V2 contracts", () => { expect(errorItem.type).toBe("error"); if (errorItem.type !== "error") throw new Error("expected error item"); expect(errorItem.failure.message).toBe("Invalid reasoning effort."); + const usageLimited = decodeOrchestrationV2TurnItem({ + ...errorItem, + startedAt: now, + completedAt: now, + updatedAt: now, + failure: { ...errorItem.failure, class: "usage_limit" }, + }); + expect(usageLimited.type === "error" && usageLimited.failure.class).toBe("usage_limit"); expect(() => decodeOrchestrationV2TurnItem({ ...errorItem, diff --git a/packages/contracts/src/orchestrationV2.ts b/packages/contracts/src/orchestrationV2.ts index 582b341597d7..d6a6ea061204 100644 --- a/packages/contracts/src/orchestrationV2.ts +++ b/packages/contracts/src/orchestrationV2.ts @@ -955,6 +955,7 @@ export const OrchestrationV2FileChangeDetail = Schema.Struct({ export type OrchestrationV2FileChangeDetail = typeof OrchestrationV2FileChangeDetail.Type; export const OrchestrationV2ProviderFailureClass = Schema.Literals([ + "usage_limit", "provider_error", "transport_error", "permission_error", @@ -1486,6 +1487,7 @@ export const OrchestrationV2ThreadShell = Schema.Struct({ ), status: OrchestrationV2ShellThreadStatus, lastError: Schema.optional(Schema.NullOr(Schema.String)), + lastErrorClass: Schema.optional(Schema.NullOr(OrchestrationV2ProviderFailureClass)), pendingRuntimeRequest: Schema.NullOr(OrchestrationV2PendingRuntimeRequestSummary), latestVisibleMessage: Schema.NullOr(OrchestrationV2LatestVisibleMessageSummary), latestUserMessageAt: Schema.NullOr(Schema.DateTimeUtc), diff --git a/packages/shared/package.json b/packages/shared/package.json index 9b76a0795a1b..6e476bb2029f 100644 --- a/packages/shared/package.json +++ b/packages/shared/package.json @@ -3,6 +3,10 @@ "private": true, "type": "module", "exports": { + "./orchestrationV2ThreadError": { + "types": "./src/orchestrationV2ThreadError.ts", + "import": "./src/orchestrationV2ThreadError.ts" + }, "./legacyCliLauncher": { "types": "./src/legacyCliLauncher.ts", "import": "./src/legacyCliLauncher.ts" diff --git a/packages/shared/src/orchestrationV2ThreadError.ts b/packages/shared/src/orchestrationV2ThreadError.ts new file mode 100644 index 000000000000..d22f2385a06e --- /dev/null +++ b/packages/shared/src/orchestrationV2ThreadError.ts @@ -0,0 +1,45 @@ +import type { + OrchestrationV2ProviderFailure, + OrchestrationV2Run, + OrchestrationV2TurnItem, +} from "@t3tools/contracts"; +import * as DateTime from "effect/DateTime"; + +/** Only a failed root turn of the current run owns the thread's failure state. */ +export function latestRootProviderFailure( + run: OrchestrationV2Run | null, + turnItems: ReadonlyArray, +): OrchestrationV2ProviderFailure | null { + if (run?.status !== "failed") return null; + let latest: Extract | null = null; + for (const item of turnItems) { + if ( + item.type !== "error" || + item.status !== "failed" || + item.runId !== run.id || + item.nodeId !== run.rootNodeId + ) + continue; + if ( + latest === null || + DateTime.toEpochMillis(item.updatedAt) > DateTime.toEpochMillis(latest.updatedAt) || + (DateTime.toEpochMillis(item.updatedAt) === DateTime.toEpochMillis(latest.updatedAt) && + (item.ordinal > latest.ordinal || (item.ordinal === latest.ordinal && item.id > latest.id))) + ) { + latest = item; + } + } + return latest?.failure ?? null; +} + +/** A distinct session failure supersedes the turn's classification. */ +export function threadErrorSummary( + failure: OrchestrationV2ProviderFailure | null, + sessionError: string | null, +) { + return { + lastError: sessionError ?? failure?.message ?? null, + lastErrorClass: + sessionError !== null && sessionError !== failure?.message ? null : (failure?.class ?? null), + }; +} From 4daa4466ceb1085f59b8cf9064bba40c86c19c28 Mon Sep 17 00:00:00 2001 From: Julius Marminge Date: Sat, 19 Sep 2026 23:02:27 -0700 Subject: [PATCH 2/9] fix(v2): clear limit warnings after provider recovery --- .../src/features/threads/thread-work-log.tsx | 3 +- apps/mobile/src/lib/threadActivity.test.ts | 167 +++++++++--------- apps/mobile/src/lib/threadActivity.ts | 6 +- apps/web/src/session-logic.test.ts | 22 +++ apps/web/src/session-logic.ts | 2 +- 5 files changed, 118 insertions(+), 82 deletions(-) diff --git a/apps/mobile/src/features/threads/thread-work-log.tsx b/apps/mobile/src/features/threads/thread-work-log.tsx index 9355c10b108f..b70a13784da1 100644 --- a/apps/mobile/src/features/threads/thread-work-log.tsx +++ b/apps/mobile/src/features/threads/thread-work-log.tsx @@ -890,7 +890,8 @@ const ThreadWorkLogRow = memo(function ThreadWorkLogRow( const isSystemNotice = row.projectedItem.item.type === "system_notice"; const isUsageLimit = row.projectedItem.item.type === "error" && - row.projectedItem.item.failure.class === "usage_limit"; + row.projectedItem.item.failure.class === "usage_limit" && + row.projectedItem.item.status !== "completed"; const iconIsDestructive = !isSystemNotice && !isUsageLimit && (row.icon === "alert" || row.icon === "warning"); const failed = row.status === "failure"; diff --git a/apps/mobile/src/lib/threadActivity.test.ts b/apps/mobile/src/lib/threadActivity.test.ts index 5520e86c27fa..7d9d267496d1 100644 --- a/apps/mobile/src/lib/threadActivity.test.ts +++ b/apps/mobile/src/lib/threadActivity.test.ts @@ -363,88 +363,97 @@ describe("buildThreadFeed", () => { expect(activity?.getFullDetail()).toContain(message); }); - it("presents provider retries as visible work-log activity", () => { - const retryBase = { - ...base("item-provider-retry", "2026-06-20T00:00:02.000Z", 1), - type: "error" as const, - failure: { - class: "transport_error" as const, - message: "The response stream disconnected.", - code: "responseStreamDisconnected", - retryable: true, - }, - retry: { - attempt: 2, - maxAttempts: 5, - retryDelayMs: null, - }, - }; - const runningFeed = buildThreadFeed([ - projected( - { - ...retryBase, - status: "running", - title: "Provider retry", - completedAt: null, - }, - 0, - ), - ]); - const recoveredFeed = buildThreadFeed([ - projected( - { - ...retryBase, - status: "completed", - title: "Provider recovered", + it.each(["transport_error", "usage_limit"] as const)( + "presents %s retries and clears warning markers on recovery", + (failureClass) => { + const retryBase = { + ...base("item-provider-retry", "2026-06-20T00:00:02.000Z", 1), + type: "error" as const, + failure: { + class: failureClass, + message: "The response stream disconnected.", + code: "responseStreamDisconnected", + retryable: true, }, - 0, - ), - ]); - const failedFeed = buildThreadFeed([ - projected( - { - ...retryBase, - status: "failed", - title: "Provider retry failed", + retry: { + attempt: 2, + maxAttempts: 5, + retryDelayMs: null, }, - 0, - ), - projected(command("2026-06-20T00:00:03.000Z"), 1), - ]); - const runningActivity = runningFeed.find((entry) => entry.type === "activity-group") - ?.activities[0]; - const recoveredActivity = recoveredFeed.find((entry) => entry.type === "activity-group") - ?.activities[0]; - if (runningActivity === undefined || recoveredActivity === undefined) { - throw new Error("Expected provider retry work-log activities."); - } + }; + const runningFeed = buildThreadFeed([ + projected( + { + ...retryBase, + status: "running", + title: "Provider retry", + completedAt: null, + }, + 0, + ), + ]); + const recoveredFeed = buildThreadFeed([ + projected( + { + ...retryBase, + status: "completed", + title: "Provider recovered", + }, + 0, + ), + ]); + if (failureClass === "usage_limit") { + const recoveredActivity = recoveredFeed.flatMap((entry) => + entry.type === "activity-group" ? entry.activities : [], + )[0]; + expect(recoveredActivity).toMatchObject({ status: "success", icon: "check" }); + } + const failedFeed = buildThreadFeed([ + projected( + { + ...retryBase, + status: "failed", + title: "Provider retry failed", + }, + 0, + ), + projected(command("2026-06-20T00:00:03.000Z"), 1), + ]); + const runningActivity = runningFeed.find((entry) => entry.type === "activity-group") + ?.activities[0]; + const recoveredActivity = recoveredFeed.find((entry) => entry.type === "activity-group") + ?.activities[0]; + if (runningActivity === undefined || recoveredActivity === undefined) { + throw new Error("Expected provider retry work-log activities."); + } - expect(runningActivity).toMatchObject({ - summary: "Provider retry", - status: "neutral", - toolLike: false, - }); - expect(threadFeedActivityIsVisible(runningActivity)).toBe(true); - expect(recoveredActivity).toMatchObject({ - summary: "Provider recovered", - status: "success", - toolLike: false, - }); - const failedPresentation = deriveThreadFeedPresentation( - failedFeed, - { runId, status: "running", startedAt: null, completedAt: null }, - new Set(), - ); - expect(failedPresentation.map((entry) => entry.type)).toEqual([ - "activity-group", - "work-toggle", - ]); - expect( - failedPresentation[0]?.type === "activity-group" - ? failedPresentation[0].activities[0]?.summary - : null, - ).toBe("Provider retry failed"); - }); + expect(runningActivity).toMatchObject({ + summary: "Provider retry", + status: "neutral", + toolLike: false, + }); + expect(threadFeedActivityIsVisible(runningActivity)).toBe(true); + expect(recoveredActivity).toMatchObject({ + summary: "Provider recovered", + status: "success", + toolLike: false, + }); + const failedPresentation = deriveThreadFeedPresentation( + failedFeed, + { runId, status: "running", startedAt: null, completedAt: null }, + new Set(), + ); + expect(failedPresentation.map((entry) => entry.type)).toEqual([ + "activity-group", + "work-toggle", + ]); + expect( + failedPresentation[0]?.type === "activity-group" + ? failedPresentation[0].activities[0]?.summary + : null, + ).toBe("Provider retry failed"); + }, + ); it.each(["pending", "running", "completed"] as const)( "omits %s task progress without hiding adjacent conversation items", diff --git a/apps/mobile/src/lib/threadActivity.ts b/apps/mobile/src/lib/threadActivity.ts index 7651fdf99cc8..0604fde4f098 100644 --- a/apps/mobile/src/lib/threadActivity.ts +++ b/apps/mobile/src/lib/threadActivity.ts @@ -433,7 +433,11 @@ function itemIcon(item: OrchestrationV2TurnItem): ThreadFeedActivity["icon"] { case "system_notice": return "warning"; case "error": - return item.failure.class === "usage_limit" ? "warning" : "alert"; + return item.failure.class === "usage_limit" + ? item.status === "completed" + ? "check" + : "warning" + : "alert"; case "checkpoint": case "proposed_plan": case "todo_list": diff --git a/apps/web/src/session-logic.test.ts b/apps/web/src/session-logic.test.ts index d4438390277a..cfa0787cd275 100644 --- a/apps/web/src/session-logic.test.ts +++ b/apps/web/src/session-logic.test.ts @@ -106,6 +106,28 @@ describe("V2 session presentation", () => { failure: { ...retryItem.failure, class: "usage_limit" }, }), ).toMatchObject({ label: "Usage limit reached after 2/10 retries" }); + const recoveredLimit = { + ...retryItem, + status: "completed" as const, + completedAt: now, + failure: { ...retryItem.failure, class: "usage_limit" as const }, + }; + const [recoveredEntry] = deriveTimelineEntriesFromVisibleTurnItems({ + visibleTurnItems: [ + { + item: recoveredLimit, + position: 0, + visibility: "local", + sourceThreadId: recoveredLimit.threadId, + sourceItemId: recoveredLimit.id, + }, + ], + optimisticMessages: [], + }); + if (recoveredEntry?.kind !== "work") throw new Error("Expected recovered provider work"); + expect(recoveredEntry.entry.label).toBe("Provider recovered (2/10 retries)"); + expect(recoveredEntry.entry.sourceActivityKind).not.toBe("runtime.warning"); + expect(workEntryDisplayIndicatesToolFailure(recoveredEntry.entry)).toBe(false); expect(providerErrorPresentation(retryItem)).toEqual({ label: "Retrying provider (2/10)", detail: "Claude API overloaded. Retrying in 1.5s.", diff --git a/apps/web/src/session-logic.ts b/apps/web/src/session-logic.ts index 1a0ceba7c574..059710305382 100644 --- a/apps/web/src/session-logic.ts +++ b/apps/web/src/session-logic.ts @@ -481,7 +481,7 @@ function projectedWorkEntry(row: OrchestrationV2ProjectedTurnItem): WorkLogEntry return { ...common, ...presentation, - ...(item.failure.class === "usage_limit" + ...(item.failure.class === "usage_limit" && item.status !== "completed" ? { sourceActivityKind: "runtime.warning" } : item.retry === undefined ? { sourceActivityKind: "runtime.error" } From ae0a1ed2728756a6dda9da0a119e24f016fe5903 Mon Sep 17 00:00:00 2001 From: Julius Marminge Date: Sat, 19 Sep 2026 23:09:56 -0700 Subject: [PATCH 3/9] feat(v2): carry provider limit reset times --- .../Adapters/ClaudeAdapterV2.ts | 21 ++++++++++++++++++- .../Adapters/CodexAdapterV2.ts | 19 +++++++++++++++-- .../orchestration-v2/ProjectionStore.test.ts | 5 +++++ .../src/orchestration-v2/ProjectionStore.ts | 2 ++ .../src/orchestration-v2/ProviderFailure.ts | 8 ++++++- .../provider/Layers/codexUsageLimits.test.ts | 17 +++++++++++++++ .../src/provider/Layers/codexUsageLimits.ts | 16 ++++++++++++++ packages/client-runtime/src/state/models.ts | 2 ++ packages/contracts/src/orchestrationV2.ts | 3 +++ .../shared/src/orchestrationV2ThreadError.ts | 4 ++++ 10 files changed, 93 insertions(+), 4 deletions(-) diff --git a/apps/server/src/orchestration-v2/Adapters/ClaudeAdapterV2.ts b/apps/server/src/orchestration-v2/Adapters/ClaudeAdapterV2.ts index fca48f3881d6..d62db9ac3308 100644 --- a/apps/server/src/orchestration-v2/Adapters/ClaudeAdapterV2.ts +++ b/apps/server/src/orchestration-v2/Adapters/ClaudeAdapterV2.ts @@ -2357,6 +2357,7 @@ interface ActiveClaudeTurnContext { readonly announcedUsageLimits: Set; authenticationFailureMessage: string | undefined; readonly rejectedRateLimitTypes: Set; + readonly rateLimitResetTimes: Map; latestAssistantRateLimited: boolean; readonly subagentsByTaskId: Map; readonly subagentsByToolUseId: Map; @@ -4598,12 +4599,20 @@ export function makeClaudeAdapterV2( if (context !== null) { if (blocked) { context.rejectedRateLimitTypes.add(limitType); + const resetMs = (rateLimitInfo.resetsAt ?? NaN) * 1000; + context.rateLimitResetTimes.set( + limitType, + Number.isFinite(resetMs) && resetMs > 0 && resetMs < 8.64e15 + ? DateTime.formatIso(DateTime.makeUnsafe(resetMs)) + : null, + ); } else if ( rateLimitInfo.status === "allowed" || rateLimitInfo.status === "allowed_warning" || overageAllowed ) { context.rejectedRateLimitTypes.delete(limitType); + context.rateLimitResetTimes.delete(limitType); } } // Rejected windows pause the SDK without ending its turn. Overage @@ -5244,15 +5253,24 @@ export function makeClaudeAdapterV2( (usageLimited ? "Claude usage limit reached. Send the message again once the limit resets." : undefined); + const resetTimes = Array.from(context.rateLimitResetTimes.values()); + const resetAt = + resetTimes.length > 0 && resetTimes.every((time) => time !== null) + ? resetTimes.reduce((latest, time) => (time! > latest ? time! : latest), "") + : null; const resultFailure = interrupted ? null : providerFailureFromResult(message, failureHint, usageLimited); + const terminalFailure = + resultFailure?.class === "usage_limit" + ? { ...resultFailure, resetAt } + : resultFailure; yield* finalizeActiveTurn({ context, status: interrupted ? "interrupted" : terminalStatusFromResult(message, failureHint), completedAt, result: message, - ...(resultFailure === null ? {} : { failure: resultFailure }), + ...(terminalFailure === null ? {} : { failure: terminalFailure }), }); } }); @@ -5723,6 +5741,7 @@ export function makeClaudeAdapterV2( announcedUsageLimits: new Set(), authenticationFailureMessage: undefined, rejectedRateLimitTypes: new Set(), + rateLimitResetTimes: new Map(), latestAssistantRateLimited: false, subagentsByTaskId: new Map(), subagentsByToolUseId: new Map(), diff --git a/apps/server/src/orchestration-v2/Adapters/CodexAdapterV2.ts b/apps/server/src/orchestration-v2/Adapters/CodexAdapterV2.ts index 559f0cdb73d1..73c016386844 100644 --- a/apps/server/src/orchestration-v2/Adapters/CodexAdapterV2.ts +++ b/apps/server/src/orchestration-v2/Adapters/CodexAdapterV2.ts @@ -12,7 +12,12 @@ import { type CodexTurnTokenUsageState, } from "../../provider/CodexTurnTokenUsage.ts"; import type { ServerProviderShape } from "../../provider/Services/ServerProvider.ts"; -import { codexRateLimitsToUpdate } from "../../provider/Layers/codexUsageLimits.ts"; +import { + codexRateLimitsToUpdate, + mergeCodexRateLimits, + codexUsageLimitResetAt, + type CodexRateLimitSnapshot, +} from "../../provider/Layers/codexUsageLimits.ts"; import { CodexSettings, defaultInstanceIdForDriver, ProviderDriverKind } from "@t3tools/contracts"; import { HostProcessEnvironment } from "@t3tools/shared/hostProcess"; import { getModelSelectionStringOptionValue, modelSelectionsEqual } from "@t3tools/shared/model"; @@ -1480,6 +1485,7 @@ export function makeCodexAdapterV2(adapterOptions: CodexAdapterV2Options): Provi now, }); const events = yield* Queue.unbounded(); + const rateLimitSnapshot = yield* Ref.make(undefined); const activeTurns = yield* Ref.make(new Map()); const turnTokenUsageByThread = new Map(); const usageStateForThread = (nativeThreadId: string) => { @@ -3553,6 +3559,9 @@ export function makeCodexAdapterV2(adapterOptions: CodexAdapterV2Options): Provi yield* client.handleServerNotification("account/rateLimits/updated", (payload) => Effect.gen(function* () { + yield* Ref.update(rateLimitSnapshot, (previous) => + mergeCodexRateLimits(previous, payload.rateLimits), + ); const update = codexRateLimitsToUpdate(payload.rateLimits); if (update && adapterOptions.onUsageLimits) { const checkedAt = DateTime.formatIso(yield* DateTime.now); @@ -4632,7 +4641,13 @@ export function makeCodexAdapterV2(adapterOptions: CodexAdapterV2Options): Provi `terminal-failure:${input.context.providerTurnId}`, ), status: terminalStatus, - failure, + failure: + failure.class === "usage_limit" + ? { + ...failure, + resetAt: codexUsageLimitResetAt(yield* Ref.get(rateLimitSnapshot)), + } + : failure, ...(input.providerRetry === undefined ? {} : { diff --git a/apps/server/src/orchestration-v2/ProjectionStore.test.ts b/apps/server/src/orchestration-v2/ProjectionStore.test.ts index 40edf27afc54..d14665da437f 100644 --- a/apps/server/src/orchestration-v2/ProjectionStore.test.ts +++ b/apps/server/src/orchestration-v2/ProjectionStore.test.ts @@ -1944,6 +1944,7 @@ it.layer(TestLayer)("ProjectionStoreV2", (it) => { failure: { class: "usage_limit" as const, message: "Plan limit reached.", + resetAt: "2099-01-01T00:00:00.000Z", code: "usageLimitExceeded", retryable: null, }, @@ -1967,6 +1968,10 @@ it.layer(TestLayer)("ProjectionStoreV2", (it) => { for (const shell of [memoryShell, sqlShell]) { assert.equal(shell.lastError, lastError); assert.equal(shell.lastErrorClass, lastErrorClass); + assert.equal( + shell.usageLimitResetAt, + lastErrorClass === "usage_limit" ? "2099-01-01T00:00:00.000Z" : null, + ); } assert.isNull(shells.threads.find((row) => row.id === otherThreadId)!.lastErrorClass); }); diff --git a/apps/server/src/orchestration-v2/ProjectionStore.ts b/apps/server/src/orchestration-v2/ProjectionStore.ts index 35cfcd06b8b2..4b6bd99d6f16 100644 --- a/apps/server/src/orchestration-v2/ProjectionStore.ts +++ b/apps/server/src/orchestration-v2/ProjectionStore.ts @@ -1281,6 +1281,7 @@ type ShellThreadState = { readonly activityRunStartedAt: DateTime.Utc | null; readonly lastError: string | null; readonly lastErrorClass: OrchestrationV2ThreadShell["lastErrorClass"]; + readonly usageLimitResetAt: OrchestrationV2ThreadShell["usageLimitResetAt"]; readonly pendingRuntimeRequest: OrchestrationV2ThreadProjection["runtimeRequests"][number] | null; readonly latestUserMessageAt: DateTime.Utc | null; readonly hasActionableProposedPlan: boolean; @@ -1417,6 +1418,7 @@ function shellFromState(input: { status: input.state.latestRunStatus, lastError: input.state.lastError, lastErrorClass: input.state.lastErrorClass, + usageLimitResetAt: input.state.usageLimitResetAt, pendingRuntimeRequest: input.state.pendingRuntimeRequest === null ? null diff --git a/apps/server/src/orchestration-v2/ProviderFailure.ts b/apps/server/src/orchestration-v2/ProviderFailure.ts index 0400eaf05ff2..89a33ea2f78f 100644 --- a/apps/server/src/orchestration-v2/ProviderFailure.ts +++ b/apps/server/src/orchestration-v2/ProviderFailure.ts @@ -10,7 +10,7 @@ import type { RunId, ThreadId, } from "@t3tools/contracts"; -import type * as DateTime from "effect/DateTime"; +import * as DateTime from "effect/DateTime"; import type { IdAllocatorV2Shape } from "./IdAllocator.ts"; @@ -94,6 +94,7 @@ export function makeProviderFailure(input: { readonly code?: string | null | undefined; readonly class?: OrchestrationV2ProviderFailureClass; readonly retryable?: boolean | null; + readonly resetAt?: string | null; }): OrchestrationV2ProviderFailure { const rawMessage = input.message ?? DEFAULT_PROVIDER_FAILURE_MESSAGE; const message = boundedText(rawMessage, MAX_PROVIDER_FAILURE_MESSAGE_LENGTH); @@ -106,6 +107,11 @@ export function makeProviderFailure(input: { message: message || DEFAULT_PROVIDER_FAILURE_MESSAGE, code, retryable: input.retryable ?? null, + ...(input.class === "usage_limit" && + input.resetAt != null && + Number.isFinite(Date.parse(input.resetAt)) + ? { resetAt: DateTime.formatIso(DateTime.makeUnsafe(input.resetAt)) } + : {}), }; } diff --git a/apps/server/src/provider/Layers/codexUsageLimits.test.ts b/apps/server/src/provider/Layers/codexUsageLimits.test.ts index 37006789986c..fe34190ebb6d 100644 --- a/apps/server/src/provider/Layers/codexUsageLimits.test.ts +++ b/apps/server/src/provider/Layers/codexUsageLimits.test.ts @@ -297,3 +297,20 @@ describe("mergeCodexRateLimits", () => { ).toBe(main); }); }); + +import { codexUsageLimitResetAt } from "./codexUsageLimits.ts"; + +describe("codexUsageLimitResetAt", () => { + it("waits for every exhausted window and never invents an unknown reset", () => { + expect( + codexUsageLimitResetAt({ + primary: { usedPercent: 100, resetsAt: 2000000000 }, + secondary: { usedPercent: 100, resetsAt: 2000100000 }, + }), + ).toBe("2033-05-19T07:46:40.000Z"); + expect(codexUsageLimitResetAt({ primary: { usedPercent: 100 } })).toBeNull(); + expect( + codexUsageLimitResetAt({ primary: { usedPercent: 50, resetsAt: 2000000000 } }), + ).toBeNull(); + }); +}); diff --git a/apps/server/src/provider/Layers/codexUsageLimits.ts b/apps/server/src/provider/Layers/codexUsageLimits.ts index 74ce04a7895f..f72ca02be835 100644 --- a/apps/server/src/provider/Layers/codexUsageLimits.ts +++ b/apps/server/src/provider/Layers/codexUsageLimits.ts @@ -232,3 +232,19 @@ export function codexUsageLimitMessage( } return `Codex usage limit reached.${reset}${codexUsageLimitNextStep(snapshot?.rateLimitReachedType)}`; } + +/** All exhausted windows must reset before a continuation can run. */ +export function codexUsageLimitResetAt( + snapshot: CodexRateLimitSnapshot | undefined, +): string | null { + if (!snapshot) return null; + const windows = codexRateLimitsToWindows(snapshot).filter((window) => window.usedPercent >= 100); + if (windows.length === 0 || windows.some((window) => !window.resetsAt)) return null; + return windows.reduce( + (latest, window) => + latest === null || Date.parse(window.resetsAt!) > Date.parse(latest) + ? window.resetsAt! + : latest, + null, + ); +} diff --git a/packages/client-runtime/src/state/models.ts b/packages/client-runtime/src/state/models.ts index 3b992caec60c..8b3756e95b8c 100644 --- a/packages/client-runtime/src/state/models.ts +++ b/packages/client-runtime/src/state/models.ts @@ -55,6 +55,7 @@ export interface ThreadRuntimeSummary { readonly providerName: string | null; readonly lastError: string | null; readonly lastErrorClass?: OrchestrationV2ProviderFailureClass | null; + readonly usageLimitResetAt?: string | null; readonly updatedAt: string; } @@ -176,6 +177,7 @@ function shellRuntime(thread: OrchestrationV2ThreadShell): ThreadRuntimeSummary providerName: null, lastError: thread.lastError ?? null, lastErrorClass: thread.lastErrorClass ?? null, + usageLimitResetAt: thread.usageLimitResetAt ?? null, updatedAt: iso(thread.updatedAt), }; } diff --git a/packages/contracts/src/orchestrationV2.ts b/packages/contracts/src/orchestrationV2.ts index d6a6ea061204..8960f6f85732 100644 --- a/packages/contracts/src/orchestrationV2.ts +++ b/packages/contracts/src/orchestrationV2.ts @@ -979,6 +979,8 @@ export const OrchestrationV2ProviderFailure = Schema.Struct({ message: OrchestrationV2ProviderFailureMessage, code: Schema.NullOr(OrchestrationV2ProviderFailureCode), retryable: Schema.NullOr(Schema.Boolean), + /** Reported reset time; absent when the provider cannot name one. */ + resetAt: Schema.optional(Schema.NullOr(IsoDateTime)), }); export type OrchestrationV2ProviderFailure = typeof OrchestrationV2ProviderFailure.Type; @@ -1488,6 +1490,7 @@ export const OrchestrationV2ThreadShell = Schema.Struct({ status: OrchestrationV2ShellThreadStatus, lastError: Schema.optional(Schema.NullOr(Schema.String)), lastErrorClass: Schema.optional(Schema.NullOr(OrchestrationV2ProviderFailureClass)), + usageLimitResetAt: Schema.optional(Schema.NullOr(IsoDateTime)), pendingRuntimeRequest: Schema.NullOr(OrchestrationV2PendingRuntimeRequestSummary), latestVisibleMessage: Schema.NullOr(OrchestrationV2LatestVisibleMessageSummary), latestUserMessageAt: Schema.NullOr(Schema.DateTimeUtc), diff --git a/packages/shared/src/orchestrationV2ThreadError.ts b/packages/shared/src/orchestrationV2ThreadError.ts index d22f2385a06e..b1b299f07e58 100644 --- a/packages/shared/src/orchestrationV2ThreadError.ts +++ b/packages/shared/src/orchestrationV2ThreadError.ts @@ -37,7 +37,11 @@ export function threadErrorSummary( failure: OrchestrationV2ProviderFailure | null, sessionError: string | null, ) { + const currentFailure = + sessionError !== null && sessionError !== failure?.message ? null : failure; return { + usageLimitResetAt: + currentFailure?.class === "usage_limit" ? (currentFailure.resetAt ?? null) : null, lastError: sessionError ?? failure?.message ?? null, lastErrorClass: sessionError !== null && sessionError !== failure?.message ? null : (failure?.class ?? null), From 7fb3e3e0b74ba8a6e6609123f53e70ecaf62ce75 Mon Sep 17 00:00:00 2001 From: Julius Marminge Date: Sat, 19 Sep 2026 23:25:49 -0700 Subject: [PATCH 4/9] fix(v2): retain late provider reset updates --- .../Adapters/CodexAdapterV2.test.ts | 48 ++++++++++- .../Adapters/CodexAdapterV2.ts | 79 +++++++++++++++++-- .../provider/Layers/codexUsageLimits.test.ts | 2 +- 3 files changed, 121 insertions(+), 8 deletions(-) diff --git a/apps/server/src/orchestration-v2/Adapters/CodexAdapterV2.test.ts b/apps/server/src/orchestration-v2/Adapters/CodexAdapterV2.test.ts index 10da6ac82043..982e09513dae 100644 --- a/apps/server/src/orchestration-v2/Adapters/CodexAdapterV2.test.ts +++ b/apps/server/src/orchestration-v2/Adapters/CodexAdapterV2.test.ts @@ -4888,6 +4888,18 @@ describe("CodexAdapterV2 post-settle continuation", () => { notification: true, expectedClass: "provider_error", }, + { + name: "known-reset", + code: "usageLimitExceeded", + notification: false, + expectedClass: "usage_limit", + }, + { + name: "late-reset", + code: "usageLimitExceeded", + notification: false, + expectedClass: "usage_limit", + }, { name: "retry", code: "usageLimitExceeded", notification: true, expectedClass: "usage_limit" }, ] as const) { it.effect(`classifies Codex terminal failures from ${scenario.name} evidence`, () => @@ -4896,10 +4908,25 @@ describe("CodexAdapterV2 post-settle continuation", () => { const nativeThreadId = `native-limit-${scenario.name}`; const nativeTurnId = `turn-limit-${scenario.name}`; const message = "Provider stopped this request."; + const resetAt = "2033-05-19T07:20:00.000Z"; + const snapshot = { + type: "emit_inbound" as const, + label: "account/rateLimits/updated", + frame: { + method: "account/rateLimits/updated", + params: { + rateLimits: { + limitId: "codex", + primary: { usedPercent: 100, resetsAt: 2000100000, windowDurationMins: 300 }, + }, + }, + }, + }; const transcript = makeCodexReplayTranscript({ scenario: `codex-limit-${scenario.name}`, entries: [ ...codexReplayPreamble({ nativeThreadId, nativeTurnId, prompt: "Continue." }), + ...(scenario.name === "known-reset" ? [snapshot] : []), ...(scenario.notification ? [ { @@ -4938,9 +4965,17 @@ describe("CodexAdapterV2 post-settle continuation", () => { }, }, }, + ...(scenario.name === "late-reset" ? [snapshot] : []), ], }); - const harness = yield* makeCodexReplayHarness(transcript); + const resetReceipt = yield* Deferred.make(); + const harness = yield* makeCodexReplayHarness(transcript, (event) => + event.type === "turn_item.updated" && + event.turnItem.type === "error" && + event.turnItem.failure.resetAt === resetAt + ? Deferred.succeed(resetReceipt, undefined) + : Effect.void, + ); yield* harness.runtime.startTurn( makeCodexTestTurnInput({ threadId: harness.threadId, @@ -4956,6 +4991,17 @@ describe("CodexAdapterV2 post-settle continuation", () => { if (terminal?.status !== "failed") return; assert.equal(terminal.failure.class, scenario.expectedClass); assert.equal(terminal.threadDisposition, "reusable"); + if (scenario.name === "known-reset") assert.equal(terminal.failure.resetAt, resetAt); + if (scenario.name === "late-reset") { + yield* Deferred.await(resetReceipt); + const item = harness.events.find( + (event) => + event.type === "turn_item.updated" && + event.turnItem.type === "error" && + event.turnItem.failure.resetAt === resetAt, + ); + assert.isDefined(item); + } if (scenario.name === "retry") assert.equal(terminal.retry?.attempt, 1); }).pipe(Effect.provide(Layer.merge(idAllocatorLayer, NodeServices.layer))), ), diff --git a/apps/server/src/orchestration-v2/Adapters/CodexAdapterV2.ts b/apps/server/src/orchestration-v2/Adapters/CodexAdapterV2.ts index 73c016386844..c7ba5610eb49 100644 --- a/apps/server/src/orchestration-v2/Adapters/CodexAdapterV2.ts +++ b/apps/server/src/orchestration-v2/Adapters/CodexAdapterV2.ts @@ -43,6 +43,7 @@ import type { ProviderApprovalDecision, ProviderApprovalOption, ProviderRequestKind, + ProviderThreadId, ProviderTurnId, ProviderInstanceId, RuntimeMode, @@ -98,7 +99,11 @@ import { type ProviderContinuationRequest, ProviderContinuationRequests, } from "../ProviderContinuationRequests.ts"; -import { makeProviderFailure, makeProviderRetryTurnItem } from "../ProviderFailure.ts"; +import { + makeProviderFailure, + makeProviderFailureTurnItem, + makeProviderRetryTurnItem, +} from "../ProviderFailure.ts"; import { turnScopedSelectionTransition } from "../ProviderSelectionTransition.ts"; import { isProviderNativeImageAttachment, @@ -1486,6 +1491,9 @@ export function makeCodexAdapterV2(adapterOptions: CodexAdapterV2Options): Provi }); const events = yield* Queue.unbounded(); const rateLimitSnapshot = yield* Ref.make(undefined); + const limitedTurnItems = yield* Ref.make( + new Map>(), + ); const activeTurns = yield* Ref.make(new Map()); const turnTokenUsageByThread = new Map(); const usageStateForThread = (nativeThreadId: string) => { @@ -1632,6 +1640,11 @@ export function makeCodexAdapterV2(adapterOptions: CodexAdapterV2Options): Provi subagent: null, startedAt: input.startedAt, }; + yield* Ref.update(limitedTurnItems, (current) => { + const next = new Map(current); + next.delete(context.providerThread.id); + return next; + }); beginTurnTokenUsage(context); yield* Ref.update(activeTurns, (current) => { const updated = new Map(current); @@ -3562,6 +3575,23 @@ export function makeCodexAdapterV2(adapterOptions: CodexAdapterV2Options): Provi yield* Ref.update(rateLimitSnapshot, (previous) => mergeCodexRateLimits(previous, payload.rateLimits), ); + const resetAt = codexUsageLimitResetAt(yield* Ref.get(rateLimitSnapshot)); + for (const item of (yield* Ref.get(limitedTurnItems)).values()) { + if (item.failure.resetAt === resetAt) continue; + const updated = { + ...item, + updatedAt: yield* DateTime.now, + failure: { ...item.failure, resetAt }, + }; + yield* Ref.update(limitedTurnItems, (current) => + new Map(current).set(item.providerThreadId!, updated), + ); + yield* emitProviderEvent({ + type: "turn_item.updated", + driver: CODEX_PROVIDER, + turnItem: updated, + }); + } const update = codexRateLimitsToUpdate(payload.rateLimits); if (update && adapterOptions.onUsageLimits) { const checkedAt = DateTime.formatIso(yield* DateTime.now); @@ -4670,6 +4700,40 @@ export function makeCodexAdapterV2(adapterOptions: CodexAdapterV2Options): Provi }, ); + const emitRootTerminal = Effect.fnUntraced(function* ( + context: ActiveCodexTurnContext, + event: CodexRootTerminalEvent, + ) { + const current = + event.status === "failed" && event.failure.class === "usage_limit" + ? { + ...event, + failure: { + ...event.failure, + resetAt: codexUsageLimitResetAt(yield* Ref.get(rateLimitSnapshot)), + }, + } + : event; + yield* emitProviderEvent(current); + if (current.status === "failed" && current.failure.class === "usage_limit") { + const item = makeProviderFailureTurnItem({ + idAllocator, + driver: CODEX_PROVIDER, + threadId: context.input.threadId, + runId: context.input.runId, + nodeId: context.input.rootNodeId, + providerThreadId: current.providerThreadId, + providerTurnId: current.providerTurnId, + itemOrdinal: current.failureItemOrdinal, + failure: current.failure, + occurredAt: yield* DateTime.now, + }); + yield* Ref.update(limitedTurnItems, (items) => + new Map(items).set(context.providerThread.id, item), + ); + } + }); + const emitOrDeferRootTerminal = Effect.fn("CodexAdapterV2.emitOrDeferRootTerminal")( function* (input: { readonly context: ActiveCodexTurnContext; @@ -4691,7 +4755,7 @@ export function makeCodexAdapterV2(adapterOptions: CodexAdapterV2Options): Provi }); return; } - yield* emitProviderEvent(event); + yield* emitRootTerminal(input.context, event); }, ); @@ -4700,7 +4764,10 @@ export function makeCodexAdapterV2(adapterOptions: CodexAdapterV2Options): Provi const activeTurnContexts = Array.from((yield* Ref.get(activeTurns)).values()); const readyEvents = yield* Ref.modify(deferredRootTerminals, (current) => { const updated = new Map(current); - const ready: Array = []; + const ready: Array<{ + context: ActiveCodexTurnContext; + event: CodexRootTerminalEvent; + }> = []; for (const [nativeTurnId, deferred] of current) { if ( !activeTurnContexts.some((candidate) => @@ -4708,13 +4775,13 @@ export function makeCodexAdapterV2(adapterOptions: CodexAdapterV2Options): Provi ) ) { updated.delete(nativeTurnId); - ready.push(deferred.event); + ready.push(deferred); } } return [ready, updated] as const; }); - for (const event of readyEvents) { - yield* emitProviderEvent(event); + for (const ready of readyEvents) { + yield* emitRootTerminal(ready.context, ready.event); } }, ); diff --git a/apps/server/src/provider/Layers/codexUsageLimits.test.ts b/apps/server/src/provider/Layers/codexUsageLimits.test.ts index fe34190ebb6d..8d310a9fff73 100644 --- a/apps/server/src/provider/Layers/codexUsageLimits.test.ts +++ b/apps/server/src/provider/Layers/codexUsageLimits.test.ts @@ -307,7 +307,7 @@ describe("codexUsageLimitResetAt", () => { primary: { usedPercent: 100, resetsAt: 2000000000 }, secondary: { usedPercent: 100, resetsAt: 2000100000 }, }), - ).toBe("2033-05-19T07:46:40.000Z"); + ).toBe("2033-05-19T07:20:00.000Z"); expect(codexUsageLimitResetAt({ primary: { usedPercent: 100 } })).toBeNull(); expect( codexUsageLimitResetAt({ primary: { usedPercent: 50, resetsAt: 2000000000 } }), From 26bb18fad2104439142f1b0f02abf03edaadd5a3 Mon Sep 17 00:00:00 2001 From: Julius Marminge Date: Sat, 19 Sep 2026 23:31:22 -0700 Subject: [PATCH 5/9] fix(v2): keep reset data after quota becomes available --- apps/server/src/orchestration-v2/Adapters/CodexAdapterV2.ts | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/apps/server/src/orchestration-v2/Adapters/CodexAdapterV2.ts b/apps/server/src/orchestration-v2/Adapters/CodexAdapterV2.ts index c7ba5610eb49..789a48313a86 100644 --- a/apps/server/src/orchestration-v2/Adapters/CodexAdapterV2.ts +++ b/apps/server/src/orchestration-v2/Adapters/CodexAdapterV2.ts @@ -3577,7 +3577,8 @@ export function makeCodexAdapterV2(adapterOptions: CodexAdapterV2Options): Provi ); const resetAt = codexUsageLimitResetAt(yield* Ref.get(rateLimitSnapshot)); for (const item of (yield* Ref.get(limitedTurnItems)).values()) { - if (item.failure.resetAt === resetAt) continue; + // Keep the stopped turn's reset when a later snapshot reports available quota. + if (resetAt === null || item.failure.resetAt === resetAt) continue; const updated = { ...item, updatedAt: yield* DateTime.now, From 444be0e9b44757f4c4a91186aaa7bd672f2ef51a Mon Sep 17 00:00:00 2001 From: Julius Marminge Date: Sat, 19 Sep 2026 23:36:17 -0700 Subject: [PATCH 6/9] fix(v2): bind reset data to the stopped turn --- apps/server/src/orchestration-v2/Adapters/CodexAdapterV2.ts | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/apps/server/src/orchestration-v2/Adapters/CodexAdapterV2.ts b/apps/server/src/orchestration-v2/Adapters/CodexAdapterV2.ts index 789a48313a86..1caac8985da7 100644 --- a/apps/server/src/orchestration-v2/Adapters/CodexAdapterV2.ts +++ b/apps/server/src/orchestration-v2/Adapters/CodexAdapterV2.ts @@ -3577,8 +3577,8 @@ export function makeCodexAdapterV2(adapterOptions: CodexAdapterV2Options): Provi ); const resetAt = codexUsageLimitResetAt(yield* Ref.get(rateLimitSnapshot)); for (const item of (yield* Ref.get(limitedTurnItems)).values()) { - // Keep the stopped turn's reset when a later snapshot reports available quota. - if (resetAt === null || item.failure.resetAt === resetAt) continue; + // Fill late reset data once; later account windows do not change this stopped turn. + if (resetAt === null || item.failure.resetAt != null) continue; const updated = { ...item, updatedAt: yield* DateTime.now, From 6922e6ef4c04d4eee8499a1dfad8cd8111875038 Mon Sep 17 00:00:00 2001 From: Julius Marminge Date: Sat, 19 Sep 2026 23:51:17 -0700 Subject: [PATCH 7/9] fix(v2): preserve captured Codex limit failures --- .../Adapters/CodexAdapterV2.test.ts | 100 +++++++++++++++++- .../Adapters/CodexAdapterV2.ts | 9 +- 2 files changed, 102 insertions(+), 7 deletions(-) diff --git a/apps/server/src/orchestration-v2/Adapters/CodexAdapterV2.test.ts b/apps/server/src/orchestration-v2/Adapters/CodexAdapterV2.test.ts index 982e09513dae..c789e2110455 100644 --- a/apps/server/src/orchestration-v2/Adapters/CodexAdapterV2.test.ts +++ b/apps/server/src/orchestration-v2/Adapters/CodexAdapterV2.test.ts @@ -4900,6 +4900,18 @@ describe("CodexAdapterV2 post-settle continuation", () => { notification: false, expectedClass: "usage_limit", }, + { + name: "matching-details", + code: "usageLimitExceeded", + notification: true, + expectedClass: "usage_limit", + }, + { + name: "deferred-reset", + code: "usageLimitExceeded", + notification: false, + expectedClass: "usage_limit", + }, { name: "retry", code: "usageLimitExceeded", notification: true, expectedClass: "usage_limit" }, ] as const) { it.effect(`classifies Codex terminal failures from ${scenario.name} evidence`, () => @@ -4926,7 +4938,45 @@ describe("CodexAdapterV2 post-settle continuation", () => { scenario: `codex-limit-${scenario.name}`, entries: [ ...codexReplayPreamble({ nativeThreadId, nativeTurnId, prompt: "Continue." }), - ...(scenario.name === "known-reset" ? [snapshot] : []), + ...(scenario.name === "known-reset" || scenario.name === "deferred-reset" + ? [snapshot] + : []), + ...(scenario.name === "deferred-reset" + ? [ + { + type: "emit_inbound" as const, + label: "item/completed/subAgentActivity-started", + frame: { + method: "item/completed", + params: { + threadId: nativeThreadId, + turnId: nativeTurnId, + item: { + type: "subAgentActivity", + id: "limit-child-spawn", + kind: "started", + agentThreadId: "native-limit-child", + agentPath: "/root/limit_child", + }, + }, + }, + }, + { + type: "emit_inbound" as const, + label: "turn/started/child", + frame: { + method: "turn/started", + params: { + threadId: "native-limit-child", + turn: makeCodexReplayTurn({ + id: "limit-child-turn", + status: "inProgress", + }), + }, + }, + }, + ] + : []), ...(scenario.notification ? [ { @@ -4941,7 +4991,10 @@ describe("CodexAdapterV2 post-settle continuation", () => { error: { message, codexErrorInfo: scenario.code, - additionalDetails: null, + additionalDetails: + scenario.name === "matching-details" + ? "Detailed provider allowance explanation." + : null, }, }, }, @@ -4959,13 +5012,49 @@ describe("CodexAdapterV2 post-settle continuation", () => { ...makeCodexReplayTurn({ id: nativeTurnId, status: "failed" }), error: { message: scenario.name === "replacement" ? "A different failure." : message, - ...(scenario.notification ? {} : { codexErrorInfo: scenario.code }), + ...(scenario.notification && scenario.name !== "matching-details" + ? {} + : { codexErrorInfo: scenario.code }), }, }, }, }, }, ...(scenario.name === "late-reset" ? [snapshot] : []), + ...(scenario.name === "deferred-reset" + ? [ + { + ...snapshot, + frame: { + method: "account/rateLimits/updated", + params: { + rateLimits: { + limitId: "codex", + primary: { + usedPercent: 100, + resetsAt: 2000200000, + windowDurationMins: 300, + }, + }, + }, + }, + }, + { + type: "emit_inbound" as const, + label: "turn/completed/child", + frame: { + method: "turn/completed", + params: { + threadId: "native-limit-child", + turn: makeCodexReplayTurn({ + id: "limit-child-turn", + status: "completed", + }), + }, + }, + }, + ] + : []), ], }); const resetReceipt = yield* Deferred.make(); @@ -4991,7 +5080,10 @@ describe("CodexAdapterV2 post-settle continuation", () => { if (terminal?.status !== "failed") return; assert.equal(terminal.failure.class, scenario.expectedClass); assert.equal(terminal.threadDisposition, "reusable"); - if (scenario.name === "known-reset") assert.equal(terminal.failure.resetAt, resetAt); + if (scenario.name === "known-reset" || scenario.name === "deferred-reset") + assert.equal(terminal.failure.resetAt, resetAt); + if (scenario.name === "matching-details") + assert.equal(terminal.failure.message, "Detailed provider allowance explanation."); if (scenario.name === "late-reset") { yield* Deferred.await(resetReceipt); const item = harness.events.find( diff --git a/apps/server/src/orchestration-v2/Adapters/CodexAdapterV2.ts b/apps/server/src/orchestration-v2/Adapters/CodexAdapterV2.ts index 1caac8985da7..7146de4f1d77 100644 --- a/apps/server/src/orchestration-v2/Adapters/CodexAdapterV2.ts +++ b/apps/server/src/orchestration-v2/Adapters/CodexAdapterV2.ts @@ -4647,10 +4647,11 @@ export function makeCodexAdapterV2(adapterOptions: CodexAdapterV2Options): Provi if (terminalStatus === "failed") { const previousFailure = input.context.latestProviderFailure ?? input.providerRetry; const failure = - input.failureCode === undefined && previousFailure !== undefined && (input.failureMessage === undefined || - input.failureMessage === previousFailure.nativeMessage) + input.failureMessage === previousFailure.nativeMessage) && + (input.failureCode === undefined || + input.failureCode === previousFailure.failure.code) ? previousFailure.failure : makeProviderFailure({ message: input.failureMessage, @@ -4711,7 +4712,9 @@ export function makeCodexAdapterV2(adapterOptions: CodexAdapterV2Options): Provi ...event, failure: { ...event.failure, - resetAt: codexUsageLimitResetAt(yield* Ref.get(rateLimitSnapshot)), + resetAt: + event.failure.resetAt ?? + codexUsageLimitResetAt(yield* Ref.get(rateLimitSnapshot)), }, } : event; From cfa8774a61aea0165681c428ab2f0358a04c7e8d Mon Sep 17 00:00:00 2001 From: Julius Marminge Date: Sun, 20 Sep 2026 17:58:41 -0700 Subject: [PATCH 8/9] fix(server): parameterize run IDs in turn-start history --- .../src/orchestration-v2/ProjectionStore.test.ts | 12 ++++++++++++ apps/server/src/orchestration-v2/ProjectionStore.ts | 2 +- 2 files changed, 13 insertions(+), 1 deletion(-) diff --git a/apps/server/src/orchestration-v2/ProjectionStore.test.ts b/apps/server/src/orchestration-v2/ProjectionStore.test.ts index 2ceb0b9a9723..c44cbc55d708 100644 --- a/apps/server/src/orchestration-v2/ProjectionStore.test.ts +++ b/apps/server/src/orchestration-v2/ProjectionStore.test.ts @@ -299,6 +299,18 @@ it.effect("memory recovery selection includes unfinished items from missing runs ); it.layer(TestLayer)("ProjectionStoreV2", (it) => { + it.effect("limits turn-start history to the requested runs, including an empty selection", () => + Effect.gen(function* () { + const store = yield* ProjectionStoreV2; + const threadId = yield* addRolledBackRecoveryCandidate("selected-turn-start-history"); + const runId = (yield* store.getThreadProjection(threadId)).runs[0]!.id; + const history = yield* store.getTurnStartHistory(threadId); + assert.isNotEmpty(history); + assert.deepEqual(yield* store.getTurnStartHistory(threadId, [runId]), history); + assert.deepEqual(yield* store.getTurnStartHistory(threadId, []), []); + assert.deepEqual(yield* store.getTurnStartHistory(threadId, [RunId.make("run:other")]), []); + }), + ); it.effect("preserves stored provider usage when a terminal update omits it", () => Effect.gen(function* () { const projectionStore = yield* ProjectionStoreV2; diff --git a/apps/server/src/orchestration-v2/ProjectionStore.ts b/apps/server/src/orchestration-v2/ProjectionStore.ts index dce217a9f7e5..2faaec6dff33 100644 --- a/apps/server/src/orchestration-v2/ProjectionStore.ts +++ b/apps/server/src/orchestration-v2/ProjectionStore.ts @@ -3422,7 +3422,7 @@ export const layer: Layer.Layer = WHERE thread_id = ${threadId} AND type IN ('user_message','assistant_message','command_execution','error', 'run_interrupt_result','file_change','proposed_plan') - AND (${runIds === undefined ? 1 : 0} OR run_id IN (SELECT value FROM json_each(${JSON.stringify(runIds ?? [])}))) + AND ${runIds === undefined ? sql`1` : sql`run_id IN ${sql.in(runIds)}`} ORDER BY ordinal ASC, turn_item_id ASC `; return yield* decodeRows(decodeTurnItemPayload, threadId)(rows); From 86d8cd45be4774129bd7e101a105e199442f63f0 Mon Sep 17 00:00:00 2001 From: Julius Marminge Date: Sun, 20 Sep 2026 18:33:43 -0700 Subject: [PATCH 9/9] fix(server): load reused checkpoint scopes at turn start --- .../src/orchestration-v2/ProjectionStore.ts | 2 +- .../orchestration-v2/TurnStartReads.test.ts | 51 ++++++++++++++++++- 2 files changed, 51 insertions(+), 2 deletions(-) diff --git a/apps/server/src/orchestration-v2/ProjectionStore.ts b/apps/server/src/orchestration-v2/ProjectionStore.ts index 2faaec6dff33..64ca61890156 100644 --- a/apps/server/src/orchestration-v2/ProjectionStore.ts +++ b/apps/server/src/orchestration-v2/ProjectionStore.ts @@ -3351,7 +3351,7 @@ export const layer: Layer.Layer = `.pipe(Effect.flatMap(decodeRows(decodeProviderTurnPayload, threadId))); const checkpointScopes = yield* sql` SELECT payload_json FROM orchestration_v2_projection_checkpoint_scopes - WHERE thread_id = ${threadId} AND node_id IN (SELECT json_extract(payload_json, '$.rootNodeId') FROM orchestration_v2_projection_runs WHERE thread_id = ${threadId} AND run_id = ${runId}) ORDER BY ordinal_within_parent ASC, scope_id ASC + WHERE thread_id = ${threadId} AND scope_id IN ${sql.in(nodes.flatMap((node) => (node.checkpointScopeId === null ? [] : [node.checkpointScopeId])))} ORDER BY ordinal_within_parent ASC, scope_id ASC `.pipe(Effect.flatMap(decodeRows(decodeCheckpointScopePayload, threadId))); const contextHandoffs = yield* sql` SELECT payload_json FROM orchestration_v2_projection_context_handoffs diff --git a/apps/server/src/orchestration-v2/TurnStartReads.test.ts b/apps/server/src/orchestration-v2/TurnStartReads.test.ts index da8136caf814..cbb5096a3ac4 100644 --- a/apps/server/src/orchestration-v2/TurnStartReads.test.ts +++ b/apps/server/src/orchestration-v2/TurnStartReads.test.ts @@ -1,5 +1,7 @@ import { assert, it } from "@effect/vitest"; import { + CheckpointScopeId, + NodeId, EventId, MessageId, ProjectId, @@ -68,7 +70,7 @@ it.effect.each(["sqlite", "memory"] as const)( modelSelection: thread.modelSelection, providerThreadId: null, userMessageId: messageId, - rootNodeId: null, + rootNodeId: NodeId.make("startup:root"), activeAttemptId: null, status: "queued", requestedAt: now, @@ -96,6 +98,49 @@ it.effect.each(["sqlite", "memory"] as const)( }); yield* putThread(thread); yield* putRun(run); + const scopeId = CheckpointScopeId.make("startup:shared-scope"); + yield* store.apply({ + id: EventId.make("startup:root-event"), + type: "node.updated", + threadId, + occurredAt: now, + payload: { + id: run.rootNodeId!, + threadId, + runId, + parentNodeId: null, + rootNodeId: run.rootNodeId!, + kind: "root_turn", + status: "pending", + countsForRun: true, + providerThreadId: null, + providerTurnId: null, + nativeItemRef: null, + runtimeRequestId: null, + checkpointScopeId: scopeId, + startedAt: null, + completedAt: null, + }, + }); + yield* store.apply({ + id: EventId.make("startup:scope-event"), + type: "checkpoint-scope.created", + threadId, + occurredAt: now, + payload: { + id: scopeId, + threadId, + runId: RunId.make("run:old"), + nodeId: NodeId.make("startup:earlier-root"), + parentScopeId: null, + providerThreadId: null, + kind: "root_run", + ordinalWithinParent: 0, + advancesAppRunCount: true, + cwd: "/repo/worktree", + createdAt: now, + }, + }); const message = { id: messageId, threadId, @@ -163,6 +208,10 @@ it.effect.each(["sqlite", "memory"] as const)( } const context = yield* store.getTurnStartContext(threadId, runId); assert.equal(context.thread.id, threadId); + assert.deepEqual( + context.checkpointScopes.map((scope) => scope.id), + [scopeId], + ); assert.equal(context.messages.find((m) => m.id === messageId)?.text, "Continue"); assert.isTrue(context.hasConversation); assert.deepEqual(context.turnItems, []);