Merge PR #8516: fix(session): resume Cursor turns after HTTP/2 stream reset (@joseotaviorf)

This commit is contained in:
can1357
2026-08-16 02:03:17 +02:00
8 changed files with 268 additions and 12 deletions
+2 -2
View File
@@ -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.
+2
View File
@@ -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
+1 -1
View File
@@ -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;
+18
View File
@@ -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),
+3
View File
@@ -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 `<system-notice>` 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
@@ -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;
@@ -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",
@@ -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");
});
});
});