diff --git a/docs/non-compaction-retry-policy.md b/docs/non-compaction-retry-policy.md index 8c1142f60..693f9d5de 100644 --- a/docs/non-compaction-retry-policy.md +++ b/docs/non-compaction-retry-policy.md @@ -54,7 +54,7 @@ Current retryable categories include: The normalized classifier recognizes the transient categories above from structured flags/status and provider-aware text patterns. Classifier refusals remain a separate typed `stopDetails` decision. -Beyond `isRetryableError(...)`, empty generic aborts may enter the same retry engine when no user, dispose, or streaming-edit-guard abort is in progress. An interrupted turn whose tool calls already have matching results can also be continued safely: the failed assistant/tool-result sequence is preserved so completed side effects are not replayed. Resolved stream stalls use the same preserve-and-continue path. +Beyond `isRetryableError(...)`, empty generic aborts may enter the same retry engine when no user, dispose, or streaming-edit-guard abort is in progress. An interrupted turn whose tool calls already have matching results can also be continued safely: the failed assistant/tool-result sequence is preserved so completed side effects are not replayed. Resolved stream stalls and HTTP/2 stream resets (`NGHTTP2_INTERNAL_ERROR`, `NGHTTP2_REFUSED_STREAM`, `HTTP2StreamReset`) use the same preserve-and-continue path. Cursor idle-stall recovery still requires the exec-resolved marker (the Connect stream may still be open); an HTTP/2 RST does not, because the stream is already dead. Retry state is owned by `TurnRecovery`: @@ -233,4 +233,4 @@ A new retry chain can still start later on a future retryable error after counte - Retry strips the failing assistant error from **runtime context** before re-continue, but session history still keeps that error entry. - `RpcSessionState` currently exposes `autoCompactionEnabled` but not an `autoRetryEnabled` field; RPC callers must track their own toggle state or query settings through other APIs. - Model fallback changes append temporary `model_change` entries and may later restore the primary model when its cooldown expires, depending on `retry.fallbackRevertPolicy`. -- Usage-aware fallback runs before a provider request when both `retry.modelFallback` and `retry.usageAwareFallback` are enabled. Unknown/unmapped usage fails open. At the reserve threshold, `"confirm"` asks interactive sessions and keeps the current model when declined; sessions without a confirmation UI automatically apply an eligible configured fallback. `"auto"` applies an eligible fallback without asking. `"fail-closed"` rejects reserve or depleted usage instead of spending it or selecting a fallback. Depleted usage under the other policies applies an eligible fallback without a reserve confirmation. +- Usage-aware fallback runs before a provider request when both `retry.modelFallback` and `retry.usageAwareFallback` are enabled. Unknown/unmapped usage fails open. At the reserve threshold, `"confirm"` asks interactive sessions and keeps the current model when declined; sessions without a confirmation UI automatically apply an eligible configured fallback. `"auto"` applies an eligible fallback without asking. `"fail-closed"` rejects reserve or depleted usage instead of spending it or selecting a fallback. Depleted usage under the other policies applies an eligible fallback without a reserve confirmation. \ No newline at end of file diff --git a/packages/ai/CHANGELOG.md b/packages/ai/CHANGELOG.md index 1befacee9..bdd03977a 100644 --- a/packages/ai/CHANGELOG.md +++ b/packages/ai/CHANGELOG.md @@ -15,6 +15,8 @@ - Fixed `omp usage invalidate` to discard stale OAuth and API-key usage snapshots, then force a cache-bypassing, per-provider serialized refresh with a broker request budget sized for the full unfiltered account batch, so upgraded subscriptions do not silently retain pre-change quota data. - Fixed quota reporting and Cookie capture guidance for China (Beijing) Alibaba Token Plan credentials ([#8509](https://github.com/can1357/oh-my-pi/issues/8509)). +- Fixed `omp usage invalidate` to discard stale OAuth and API-key usage snapshots, then force a cache-bypassing, per-provider serialized refresh so upgraded subscriptions do not silently retain pre-change quota data. +- Classified Cursor HTTP/2 `NGHTTP2_INTERNAL_ERROR` / `NGHTTP2_REFUSED_STREAM` stream closes as transient so session recovery can continue instead of treating the RST as a hard stop. ## [17.3.3] - 2026-08-14 diff --git a/packages/ai/src/error/flags.ts b/packages/ai/src/error/flags.ts index 6167e6f5a..f7c121a9b 100644 --- a/packages/ai/src/error/flags.ts +++ b/packages/ai/src/error/flags.ts @@ -106,7 +106,7 @@ const TRANSIENT_ENVELOPE_PATTERN = /anthropic stream envelope error:/i; const TRANSIENT_ENVELOPE_BEFORE_START_PATTERN = /before message_start/i; export const STREAM_READ_ERROR_PATTERN = /stream[_ -]?read[_ -]?error/i; export const TRANSIENT_TRANSPORT_PATTERN = - /\b(?:no[_ -]?capacity|(?:high|peak)[ _-]?demand|(?:at|over|insufficient)[ _-]?capacity|capacity[ _-]?(?:exceeded|exhausted)|peak[ _-]?load)\b|overloaded|provider.?returned.?error|rate.?limit|too many requests|429|500|502|503|504|service.?unavailable|server.?error|internal.?error|retry your request|network.?error|connection.?error|connection.?refused|unable.?to.?connect\.\s*is the computer able to access the url\?|other side closed|fetch failed|upstream.?connect|upstream.?request.?failed|reset before headers|socket hang up|timed? out|timeout|terminated|retry delay|stream stall|no error details in response|HTTP2(?:StreamReset|RefusedStream|EnhanceYourCalm)|malformed.?function.?call/i; + /\b(?:no[_ -]?capacity|(?:high|peak)[ _-]?demand|(?:at|over|insufficient)[ _-]?capacity|capacity[ _-]?(?:exceeded|exhausted)|peak[ _-]?load)\b|overloaded|provider.?returned.?error|rate.?limit|too many requests|429|500|502|503|504|service.?unavailable|server.?error|internal.?error|retry your request|network.?error|connection.?error|connection.?refused|unable.?to.?connect\.\s*is the computer able to access the url\?|other side closed|fetch failed|upstream.?connect|upstream.?request.?failed|reset before headers|socket hang up|timed? out|timeout|terminated|retry delay|stream stall|no error details in response|HTTP2(?:StreamReset|RefusedStream|EnhanceYourCalm)|nghttp2_(?:internal_error|refused_stream)|stream closed with error code nghttp2_(?:internal_error|refused_stream)|malformed.?function.?call/i; const AUTH_FAILURE_PATTERN = /\b(?:401|403|unauthorized|forbidden|authentication|auth[_ ]?unavailable|no auth available|(?:invalid|no)[_ ]?api[_ ]?key)\b/i; const MALFORMED_FUNCTION_CALL_PATTERN = /\bmalformed.?function.?call\b/i; diff --git a/packages/ai/test/error-id.test.ts b/packages/ai/test/error-id.test.ts index 668218794..8823684e7 100644 --- a/packages/ai/test/error-id.test.ts +++ b/packages/ai/test/error-id.test.ts @@ -130,6 +130,24 @@ describe("error-id classification", () => { expect(assistant.errorId).toBe(id); }); + it("classifies Cursor NGHTTP2 stream resets as transient", () => { + for (const errorMessage of [ + "Stream closed with error code NGHTTP2_INTERNAL_ERROR", + "Stream closed with error code NGHTTP2_REFUSED_STREAM", + "Connect error failed_precondition: Error: Stream closed with error code NGHTTP2_REFUSED_STREAM", + ]) { + const assistant = message({ + api: "cursor-agent", + provider: "cursor", + model: "composer-2.5", + errorMessage, + }); + const id = AIError.classifyMessage(assistant); + expect(AIError.is(id, AIError.Flag.Transient)).toBe(true); + expect(AIError.retriable(id)).toBe(true); + } + }); + it("merges existing cause-chain kinds with finalized error text kinds", () => { const assistant = message({ errorId: AIError.create(AIError.Flag.ThinkingLoop), diff --git a/packages/coding-agent/CHANGELOG.md b/packages/coding-agent/CHANGELOG.md index 13686d399..f0b811a92 100644 --- a/packages/coding-agent/CHANGELOG.md +++ b/packages/coding-agent/CHANGELOG.md @@ -41,6 +41,9 @@ - Fixed interrupted `vibe_wait` calls being reported as elapsed timeout windows while preserving the per-call wait timeout. - Optimized checkpoint/rewind prompt rendering: the rewind instruction moved from the permanent checkpoint tool result (stale after rewind) into a transient `` that is branch-cut away on rewind, leaving only the goal; the rewind-report prompt is now forward-looking with no negation guard; the rewind tool description collapsed to one line. ([#8499](https://github.com/can1357/oh-my-pi/issues/8499)) +### Fixed + +- Continued Cursor turns that died with `NGHTTP2_INTERNAL_ERROR` / `NGHTTP2_REFUSED_STREAM` after tool calls already had results, instead of leaving the agent idle until the user typed "continue". HTTP/2 stream resets now use the same preserve-and-continue path as idle stream stalls, without requiring the Cursor exec-resolved marker that MCP/todo blocks never carry. ## [17.3.3] - 2026-08-14 diff --git a/packages/coding-agent/src/session/turn-recovery.ts b/packages/coding-agent/src/session/turn-recovery.ts index f1220df61..617931b93 100644 --- a/packages/coding-agent/src/session/turn-recovery.ts +++ b/packages/coding-agent/src/session/turn-recovery.ts @@ -74,6 +74,9 @@ const EMPTY_STOP_MAX_RETRIES = 3; const SIBLING_UNBLOCK_BUFFER_MS = 1_000; const NON_WHITESPACE_RE = /\S/; const USAGE_PREFLIGHT_BLOCKED_PREFIX = "Usage preflight blocked:"; +const STREAM_STALL_ERROR_RE = /stream stall/i; +const HTTP2_STREAM_RESET_ERROR_RE = + /stream closed with error code\s+nghttp2_(?:internal_error|refused_stream)|nghttp2_(?:internal_error|refused_stream)|HTTP2(?:StreamReset|RefusedStream)/i; function hasNonWhitespace(value: string): boolean { return NON_WHITESPACE_RE.test(value); @@ -1103,10 +1106,10 @@ export class TurnRecovery { } /** - * Classify a reasonless abort or stream stall whose emitted tool calls all - * have results. The failed assistant/tool-result pair stays in context so - * continuation cannot replay completed side effects; synthetic results tell - * the next turn that an unexecuted call must be reissued. + * Classify a reasonless abort, idle stream stall, or HTTP/2 stream reset whose + * emitted tool calls all have results. The failed assistant/tool-result pair + * stays in context so continuation cannot replay completed side effects; + * synthetic results tell the next turn that an unexecuted call must be reissued. */ classifyResolvedInterruptedToolTurn(message: AssistantMessage): "reasonless-abort" | "stream-stall" | undefined { const id = this.#classifyRetryMessage(message); @@ -1118,20 +1121,28 @@ export class TurnRecovery { !this.#host.isDisposed() && !this.#host.streamingEditAbortTriggered() && ((message.stopReason === "aborted" && AIError.is(id, AIError.Flag.Abort)) || genericAbort); + const errorMessage = message.errorMessage ?? ""; const streamStall = + message.stopReason === "error" && STREAM_STALL_ERROR_RE.test(errorMessage) && AIError.retriable(id); + const transportReset = message.stopReason === "error" && - message.errorMessage?.toLowerCase().includes("stream stall") === true && - AIError.retriable(id); - if (!reasonlessAbort && !streamStall) return undefined; + HTTP2_STREAM_RESET_ERROR_RE.test(errorMessage) && + AIError.retriable(id) && + !this.#host.abortInProgress() && + !this.#host.isDisposed() && + !this.#host.streamingEditAbortTriggered(); + if (!reasonlessAbort && !streamStall && !transportReset) return undefined; if (reasonlessAbort && genericAbort) message.errorId = AIError.create(AIError.Flag.Abort); - // The Cursor server-execution marker gate applies only to the stream-stall + // The Cursor server-execution marker gate applies only to the idle stream-stall // path: an unmarked/unresolved Cursor block there means the server has not // finished executing, so resuming would race it. A reasonless abort instead // ends the turn and the agent loop pairs every un-run call (Cursor's unmarked // `todo`/MCP blocks included) with a synthetic `executed: false` result, so // the tool-result reconciliation below is the safety gate and the marker is - // irrelevant. + // irrelevant. An HTTP/2 RST_STREAM / NGHTTP2_* close also ends the Connect + // stream, so there is no in-flight server exec to race — unmarked MCP/todo + // blocks are safe to continue once every emitted call has a result. const resolvedToolCallIds: string[] = []; for (const block of message.content) { if (block.type !== "toolCall") continue; diff --git a/packages/coding-agent/test/agent-session-retry-cap.test.ts b/packages/coding-agent/test/agent-session-retry-cap.test.ts index 3b2b27b54..e7902693e 100644 --- a/packages/coding-agent/test/agent-session-retry-cap.test.ts +++ b/packages/coding-agent/test/agent-session-retry-cap.test.ts @@ -1263,6 +1263,129 @@ describe("AgentSession retry delay cap", () => { expect(retryEndEvents).toContainEqual(expect.objectContaining({ success: true, attempt: 1 })); expect(lastAssistant(session).content).toContainEqual({ type: "text", text: "Recovered after Cursor stall" }); }); + + it("resumes a Cursor HTTP/2 stream reset after an unmarked MCP tool result", async () => { + const resetMessage = "Stream closed with error code NGHTTP2_INTERNAL_ERROR"; + const model = createMockModel({ + id: "composer-2.5", + provider: "cursor", + }); + authStorage.setRuntimeApiKey("cursor", "cursor-test-key"); + const toolCall: ToolCall = { + type: "toolCall", + id: "cursor-mcp-1", + name: "mcp__databricks_production_execute_sql", + arguments: { query: "SELECT 1" }, + }; + const toolResult: ToolResultMessage = { + role: "toolResult", + toolCallId: toolCall.id, + toolName: toolCall.name, + content: [{ type: "text", text: "1" }], + isError: false, + timestamp: Date.now(), + }; + let streamCalls = 0; + let resumedWithToolResult = false; + const agent = new Agent({ + getApiKey: requestedModel => `${requestedModel.provider}-test-key`, + initialState: { + model, + systemPrompt: ["Test"], + tools: [], + messages: [], + }, + cursorOnToolResult: message => message, + streamFn: (_requestedModel, context, options) => { + streamCalls += 1; + if (streamCalls > 1) { + resumedWithToolResult = context.messages.some( + message => message.role === "toolResult" && message.toolCallId === toolCall.id, + ); + model.push({ content: ["Recovered after Cursor HTTP/2 reset"] }); + return model.stream(model, context, options); + } + + const stream = new AssistantMessageEventStream(); + queueMicrotask(async () => { + await options?.cursorOnToolResult?.(toolResult); + const partial: AssistantMessage = { + role: "assistant", + content: [toolCall], + api: model.api, + provider: model.provider, + model: model.id, + usage: { + input: 0, + output: 0, + cacheRead: 0, + cacheWrite: 0, + totalTokens: 0, + cost: { input: 0, output: 0, cacheRead: 0, cacheWrite: 0, total: 0 }, + }, + stopReason: "stop", + timestamp: Date.now(), + }; + stream.push({ type: "start", partial }); + stream.push({ type: "toolcall_start", contentIndex: 0, partial }); + stream.push({ + type: "toolcall_delta", + contentIndex: 0, + delta: JSON.stringify(toolCall.arguments), + partial, + }); + stream.push({ type: "toolcall_end", contentIndex: 0, toolCall, partial }); + stream.push({ + type: "error", + reason: "error", + error: { + ...partial, + stopReason: "error", + errorMessage: resetMessage, + }, + }); + }); + return stream; + }, + }); + + const settings = Settings.isolated({ + "compaction.enabled": false, + "retry.baseDelayMs": 5, + "retry.maxRetries": 1, + }); + settings.setModelRole("default", `${model.provider}/${model.id}`); + session = new AgentSession({ + agent, + sessionManager: SessionManager.inMemory(), + settings, + modelRegistry, + }); + const retryStartEvents: AutoRetryStartEvent[] = []; + const retryEndEvents: AutoRetryEndEvent[] = []; + session.subscribe(event => { + if (event.type === "auto_retry_start") retryStartEvents.push(event); + if (event.type === "auto_retry_end") retryEndEvents.push(event); + }); + + await session.prompt("Run the query"); + await session.waitForIdle(); + + expect(streamCalls).toBe(2); + expect(resumedWithToolResult).toBe(true); + expect( + session.agent.state.messages.some( + message => message.role === "toolResult" && message.toolCallId === toolCall.id, + ), + ).toBe(true); + expect(retryStartEvents).toHaveLength(1); + expect(retryEndEvents).toContainEqual(expect.objectContaining({ success: true, attempt: 1 })); + expect(lastAssistant(session).content).toContainEqual({ + type: "text", + text: "Recovered after Cursor HTTP/2 reset", + }); + }); + it("resumes a Cursor reasonless abort after an unmarked client-side tool call", async () => { const model = createMockModel({ id: "composer-2.5", diff --git a/packages/coding-agent/test/turn-recovery-replay-unsafe.test.ts b/packages/coding-agent/test/turn-recovery-replay-unsafe.test.ts index 73f95ffe7..a880bc4e1 100644 --- a/packages/coding-agent/test/turn-recovery-replay-unsafe.test.ts +++ b/packages/coding-agent/test/turn-recovery-replay-unsafe.test.ts @@ -2,6 +2,7 @@ import { afterAll, beforeAll, describe, expect, it } from "bun:test"; import type { AgentMessage, SyntheticToolResultDetails } from "@oh-my-pi/pi-agent-core"; import type { AssistantMessage, ToolResultMessage } from "@oh-my-pi/pi-ai"; import * as AIError from "@oh-my-pi/pi-ai/error"; +import { kCursorExecResolved } from "@oh-my-pi/pi-ai/utils/block-symbols"; import { getBundledModel } from "@oh-my-pi/pi-catalog/models"; import type { Model, Usage } from "@oh-my-pi/pi-catalog/types"; import { ModelRegistry } from "@oh-my-pi/pi-coding-agent/config/model-registry"; @@ -534,4 +535,102 @@ describe("TurnRecovery replay-unsafe output classification", () => { expect(recoveryFor(message, [syntheticResult("call-1")]).isRetryableError(message)).toBe(false); }); }); + + describe("HTTP/2 stream reset after resolved tool calls", () => { + const nghttp2Internal = "Stream closed with error code NGHTTP2_INTERNAL_ERROR"; + const nghttp2Refused = "Stream closed with error code NGHTTP2_REFUSED_STREAM"; + const stallMessage = "Provider stream stalled while waiting for the next event"; + + function cursorMessage(content: AssistantMessage["content"], errorMessage: string): AssistantMessage { + const message = makeMessage(content, model); + message.provider = "cursor"; + message.errorMessage = errorMessage; + return message; + } + + function execToolCall(id: string, marked = false): AssistantMessage["content"][number] { + const block: AssistantMessage["content"][number] = { + type: "toolCall", + id, + name: "bash", + arguments: { command: "pwd" }, + }; + if (marked) (block as { [kCursorExecResolved]?: true })[kCursorExecResolved] = true; + return block; + } + + function mcpToolCall(id: string): AssistantMessage["content"][number] { + return { + type: "toolCall", + id, + name: "mcp__databricks_production_execute_sql", + arguments: { query: "SELECT 1" }, + }; + } + + function realResult(toolCallId: string, toolName = "bash"): ToolResultMessage { + return { + role: "toolResult", + toolCallId, + toolName, + content: [{ type: "text", text: "/workspace" }], + isError: false, + timestamp: Date.now(), + }; + } + + function recoveryForReset(message: AssistantMessage, tail: readonly AgentMessage[]): TurnRecovery { + return new TurnRecovery(createHost(model, modelRegistry, { messages: [message as AgentMessage, ...tail] })); + } + + it("continues a Cursor NGHTTP2_INTERNAL_ERROR after a marked exec result", () => { + const message = cursorMessage([execToolCall("call-1", true)], nghttp2Internal); + const recovery = recoveryForReset(message, [realResult("call-1")]); + expect(recovery.isRetryableError(message)).toBe(false); + expect(recovery.classifyResolvedInterruptedToolTurn(message)).toBe("stream-stall"); + }); + + it("continues a Cursor NGHTTP2_REFUSED_STREAM after a marked exec result", () => { + const message = cursorMessage([execToolCall("call-1", true)], nghttp2Refused); + expect(recoveryForReset(message, [realResult("call-1")]).classifyResolvedInterruptedToolTurn(message)).toBe( + "stream-stall", + ); + }); + + it("continues a Cursor HTTP/2 reset after an unmarked MCP result", () => { + const message = cursorMessage([mcpToolCall("mcp-1")], nghttp2Internal); + const recovery = recoveryForReset(message, [realResult("mcp-1", "mcp__databricks_production_execute_sql")]); + expect(recovery.classifyResolvedInterruptedToolTurn(message)).toBe("stream-stall"); + }); + + it("does not continue a Cursor idle stall after an unmarked MCP call", () => { + const message = cursorMessage([mcpToolCall("mcp-1")], stallMessage); + const recovery = recoveryForReset(message, [realResult("mcp-1", "mcp__databricks_production_execute_sql")]); + expect(recovery.classifyResolvedInterruptedToolTurn(message)).toBeUndefined(); + }); + + it("does not continue an HTTP/2 reset whose tool call has no result", () => { + const message = cursorMessage([execToolCall("call-1", true)], nghttp2Internal); + expect(recoveryForReset(message, []).classifyResolvedInterruptedToolTurn(message)).toBeUndefined(); + }); + + it("does not continue an HTTP/2 CANCEL reset", () => { + const message = cursorMessage([execToolCall("call-1", true)], "Stream closed with error code NGHTTP2_CANCEL"); + expect( + recoveryForReset(message, [realResult("call-1")]).classifyResolvedInterruptedToolTurn(message), + ).toBeUndefined(); + }); + + it("matches a Connect-wrapped NGHTTP2 close", () => { + const message = cursorMessage( + [mcpToolCall("mcp-1")], + "Connect error failed_precondition: Error: Stream closed with error code NGHTTP2_INTERNAL_ERROR", + ); + expect( + recoveryForReset(message, [ + realResult("mcp-1", "mcp__databricks_production_execute_sql"), + ]).classifyResolvedInterruptedToolTurn(message), + ).toBe("stream-stall"); + }); + }); });