diff --git a/src/adapters/cursor.ts b/src/adapters/cursor.ts index e495f2a6713..5cfef1b8ada 100644 --- a/src/adapters/cursor.ts +++ b/src/adapters/cursor.ts @@ -48,6 +48,7 @@ import { import { runCursorTurnWithRetry } from "./cursor/transport-retry"; import { cursorRequestHasShellAlias, cursorRequestUsesCodeMode } from "./cursor/tool-definitions"; import { + CURSOR_OUTPUT_GUARD_MAX_HOLD_BYTES, CURSOR_ECHO_RETRY_CONTINUATION_TEXT, CURSOR_ROUTING_COMMENTARY_RETRY_TEXT, CursorEnvelopeEchoSniffer, @@ -315,6 +316,8 @@ export function createCursorAdapter(provider: OcxProviderConfig, deps: CursorAda ? new CursorRoutingCommentarySniffer() : undefined; let guardHeld: AdapterEvent[] = []; + let guardHeldBytes = 0; + const guardEncoder = new TextEncoder(); // Exactly-once observation: every client-bound text delta passes through here // exactly once — held deltas only on release, ordinary deltas at emit time. const emitTextObserved = (event: AdapterEvent): void => { @@ -327,7 +330,49 @@ export function createCursorAdapter(provider: OcxProviderConfig, deps: CursorAda emitTextObserved(held); } guardHeld = []; + guardHeldBytes = 0; }; + // A single frame can carry a multi-megabyte payload (the transport accepts up to the + // 16 MiB Cursor message bound), so the serialized size is projected — object overhead + // plus raw payload length — BEFORE any encoded copy exists. Escapes only inflate the + // exact figure, making the raw length a safe lower bound for the overflow decision. + const GUARD_EVENT_OVERHEAD_BYTES = 64; + const projectedGuardEventBytes = (event: AdapterEvent): number => + GUARD_EVENT_OVERHEAD_BYTES + + (event.type === "text_delta" + ? Buffer.byteLength(event.text, "utf8") + : event.type === "thinking_delta" + ? Buffer.byteLength(event.thinking, "utf8") + : 0); + const holdGuardEvent = (event: AdapterEvent) => { + if (guardHeldBytes + projectedGuardEventBytes(event) > CURSOR_OUTPUT_GUARD_MAX_HOLD_BYTES) { + // Too large to retain even unescaped: settle the sniffers, release what was held, + // and pass this event through without ever encoding it. + echoSniffer?.finish(); + routingCommentarySniffer?.finish(); + releaseGuardHeld(); + if (event.type !== "heartbeat") emittedOutput = true; + emitTextObserved(event); + return false; + } + guardHeld.push(event); + // Count the complete retained representation, including per-event overhead, so an + // upstream cannot evade the cap with empty or non-text reasoning frames. + guardHeldBytes += guardEncoder.encode(JSON.stringify(event)).byteLength; + if (guardHeldBytes <= CURSOR_OUTPUT_GUARD_MAX_HOLD_BYTES) return true; + echoSniffer?.finish(); + routingCommentarySniffer?.finish(); + releaseGuardHeld(); + return false; + }; + // Each sniffer settles from a bounded leading window (40 B / 512 B respectively), so + // feeding an oversized delta whole would retain megabytes it never inspects. The + // bounded prefix still covers every decision path — including marker prefixes and + // routing claims — while the tail falls through to the aggregate cap. + const ECHO_SNIFF_FEED_MAX_CHARS = 512; + const ROUTING_SNIFF_FEED_MAX_CHARS = 2048; + const boundedSniffText = (text: string, maxChars: number): string => + text.length > maxChars ? text.slice(0, maxChars) : text; const guardsSettled = () => (!echoSniffer || echoSniffer.settled) && (!routingCommentarySniffer || routingCommentarySniffer.settled); @@ -368,27 +413,30 @@ export function createCursorAdapter(provider: OcxProviderConfig, deps: CursorAda } if (!guardsSettled()) { if (event.type === "text_delta") { - guardHeld.push(event); + // Classify the delta before the aggregate-cap check: an oversized first + // delta must still pass the armed sniffers (echo/hallucination detection is + // prefix-based), so the cap cannot disarm them before they see the text. if (echoSniffer && !echoSniffer.settled) { - const decision = echoSniffer.feed(event.text); + const decision = echoSniffer.feed(boundedSniffText(event.text, ECHO_SNIFF_FEED_MAX_CHARS)); if (decision.kind === "echo") { guardHeld = []; throw new CursorToolResultEchoError(decision.marker); } } if (routingCommentarySniffer && !routingCommentarySniffer.settled) { - const decision = routingCommentarySniffer.feed(event.text); + const decision = routingCommentarySniffer.feed(boundedSniffText(event.text, ROUTING_SNIFF_FEED_MAX_CHARS)); if (decision.kind === "hallucination") { guardHeld = []; throw new CursorRoutingCommentaryError(); } } + if (!holdGuardEvent(event)) continue; if (guardsSettled()) releaseGuardHeld(); continue; } else if (event.type === "thinking_delta" || event.type === "heartbeat") { // Reasoning before first text stays ordered; liveness still passes through. if (event.type === "thinking_delta") { - guardHeld.push(event); + holdGuardEvent(event); continue; } } else { diff --git a/src/adapters/cursor/envelope-echo.ts b/src/adapters/cursor/envelope-echo.ts index 4691731fa01..51cd69f9ccf 100644 --- a/src/adapters/cursor/envelope-echo.ts +++ b/src/adapters/cursor/envelope-echo.ts @@ -76,7 +76,7 @@ export const MAX_MIDSTREAM_SCAN_LENGTH = 512 * 1024; const MAX_MIDSTREAM_FINDINGS = 8; const MAX_ROUTING_COMMENTARY_BYTES = 512; /** Aggregate quarantine cap: past this, flush and disarm. */ -const MAX_HOLD_BYTES = 8 * 1024; +export const CURSOR_OUTPUT_GUARD_MAX_HOLD_BYTES = 8 * 1024; const encoder = new TextEncoder(); export class CursorToolResultEchoError extends Error { @@ -130,16 +130,22 @@ export class CursorMidstreamEchoObserver { private totalLength = 0; private disarmed = false; private lineDisarmed = false; - private corruptionWatch: { finding: MidstreamEchoFinding; remaining: number; window: string } | undefined; + private readonly corruptionWatches: Array<{ + finding: MidstreamEchoFinding; + remaining: number; + window: string; + }> = []; private readonly recorded: MidstreamEchoFinding[] = []; feed(textDelta: string): void { - if (this.disarmed && !this.corruptionWatch) return; + if (this.disarmed && this.corruptionWatches.length === 0) return; let index = 0; while (index < textDelta.length) { const newline = textDelta.indexOf("\n", index); const segment = newline === -1 ? textDelta.slice(index) : textDelta.slice(index, newline); - if (this.corruptionWatch) this.watchCorruption(segment + (newline === -1 ? "" : "\n")); + if (this.corruptionWatches.length > 0) { + this.watchCorruption(segment + (newline === -1 ? "" : "\n")); + } if (!this.disarmed && !this.lineDisarmed && segment.length > 0) { this.lineBuffer += segment; if (this.lineBuffer.length > MAX_MIDSTREAM_LINE_INDENT + 32) { @@ -160,9 +166,7 @@ export class CursorMidstreamEchoObserver { } findings(): readonly MidstreamEchoFinding[] { - if (this.corruptionWatch) { - this.settleCorruption(); - } + while (this.corruptionWatches.length > 0) this.settleCorruption(0); return this.recorded; } @@ -187,12 +191,24 @@ export class CursorMidstreamEchoObserver { this.lineDisarmed = true; return; } - const finding: MidstreamEchoFinding = { - marker, - offset: this.lineStartOffset, - callIdCorrupt: false, - }; - this.corruptionWatch = { finding, remaining: MIDSTREAM_CORRUPTION_WINDOW, window: "" }; + // A new marker ends the previous marker's corruption window: the text + // between two markers belongs to the earlier finding only. Without this, + // every open watch consumed the same following text, so one corrupt + // call-id after a second marker also marked the first, clean finding + // corrupt (clean-then-corrupt cross-contamination). + while (this.corruptionWatches.length > 0) this.settleCorruption(0); + if (this.recorded.length + this.corruptionWatches.length < MAX_MIDSTREAM_FINDINGS) { + const finding: MidstreamEchoFinding = { + marker, + offset: this.lineStartOffset, + callIdCorrupt: false, + }; + this.corruptionWatches.push({ + finding, + remaining: MIDSTREAM_CORRUPTION_WINDOW, + window: "", + }); + } this.lineDisarmed = true; return; } @@ -203,24 +219,28 @@ export class CursorMidstreamEchoObserver { } private watchCorruption(text: string): void { - const watch = this.corruptionWatch; - if (!watch) return; - const take = Math.min(watch.remaining, text.length); - watch.window += text.slice(0, take); - watch.remaining -= take; - if (watch.remaining <= 0) this.settleCorruption(); + for (const watch of this.corruptionWatches) { + const take = Math.min(watch.remaining, text.length); + watch.window += text.slice(0, take); + watch.remaining -= take; + } + let index = 0; + while (index < this.corruptionWatches.length) { + if (this.corruptionWatches[index]!.remaining <= 0) this.settleCorruption(index); + else index += 1; + } } - private settleCorruption(): void { - const watch = this.corruptionWatch; + private settleCorruption(index: number): void { + const watch = this.corruptionWatches[index]; if (!watch) return; const window = watch.window; watch.finding.callIdCorrupt = /fc_[0-9a-f]+[ \t]+mar-/.test(window) || /call_id: \S+[ \t]+\S+_0\b/.test(window); - if (this.recorded.length < MAX_MIDSTREAM_FINDINGS) this.recorded.push(watch.finding); + this.recorded.push(watch.finding); // Window text is discarded here; only booleans/offsets survive. - this.corruptionWatch = undefined; + this.corruptionWatches.splice(index, 1); } } @@ -251,7 +271,11 @@ export class CursorEnvelopeEchoSniffer { const stillPrefix = ECHO_MARKERS.some(marker => probe.length < marker.length && marker.startsWith(probe), ); - if (stillPrefix && this.byteCount <= MAX_SNIFF_BYTES && this.buffered.length < MAX_HOLD_BYTES) { + if ( + stillPrefix + && this.byteCount <= MAX_SNIFF_BYTES + && this.buffered.length < CURSOR_OUTPUT_GUARD_MAX_HOLD_BYTES + ) { return { kind: "hold" }; } this.done = true; @@ -316,7 +340,7 @@ export class CursorRoutingCommentarySniffer { && lineBreakCount < 2; if ( this.byteCount < MAX_ROUTING_COMMENTARY_BYTES - && this.buffered.length < MAX_HOLD_BYTES + && this.buffered.length < CURSOR_OUTPUT_GUARD_MAX_HOLD_BYTES && (lineBreakCount === 0 || pendingFailureClaim) && (hasRoutingHint || this.byteCount < 64) ) { diff --git a/structure/providers/cursor.md b/structure/providers/cursor.md index 7252130843a..ce21712e910 100644 --- a/structure/providers/cursor.md +++ b/structure/providers/cursor.md @@ -192,6 +192,15 @@ Translated Chat request construction uses the [inline-image budget](../transport ## Mid-stream envelope echo +Held quarantine output is bounded by the aggregate `CURSOR_OUTPUT_GUARD_MAX_HOLD_BYTES` (8 KiB) +budget in `src/adapters/cursor.ts`. Text deltas are fed to the armed echo and +routing-commentary sniffers BEFORE the cap check, so a single oversized first delta cannot +disarm the guards without being classified; each sniffer reads only the bounded leading window +its decision needs. Retained bytes are projected from payload length before any serialized +copy exists, so a multi-megabyte frame cannot force a same-size encoded allocation. An event +that cannot fit the remaining budget settles both sniffers, releases the held events, and is +emitted directly. + The prefix sniffer only watches the opening bytes of a turn. An external model that writes real prose first and then pastes a replayed `[Tool Result]` envelope defeats it, so that text reaches the client and is stored as assistant output. `CursorMidstreamEchoObserver` records those diff --git a/tests/providers/cursor/cursor-envelope-echo-retry.test.ts b/tests/providers/cursor/cursor-envelope-echo-retry.test.ts index daeeb77176a..092d301b049 100644 --- a/tests/providers/cursor/cursor-envelope-echo-retry.test.ts +++ b/tests/providers/cursor/cursor-envelope-echo-retry.test.ts @@ -96,8 +96,12 @@ describe("cursor external output quarantine + corrective retry (devlog 260826 ga observer.feed("Leading text about progress.\n"); observer.feed(RUN03_SPECIMEN); const findings = observer.findings(); - expect(findings).toHaveLength(1); - expect(findings[0]!.callIdCorrupt).toBe(true); + expect(findings).toHaveLength(2); + expect(findings[0]!.marker).toBe("[Tool Result]"); + // The corrupt call-id sits after the second (duplicated) marker, so it is + // attributed to that marker's window — the first marker's span is clean. + expect(findings[0]!.callIdCorrupt).toBe(false); + expect(findings[1]!.callIdCorrupt).toBe(true); }); test("clean call-id lines do not flag corruption", () => { @@ -105,8 +109,8 @@ describe("cursor external output quarantine + corrective retry (devlog 260826 ga observer.feed("Leading text.\n"); observer.feed("[Tool Result]\n[tool_result]\ncall_id: call-1\nfc_63367283-2aec-9a25_1\noutput:\nok\n"); const findings = observer.findings(); - expect(findings).toHaveLength(1); - expect(findings[0]!.callIdCorrupt).toBe(false); + expect(findings).toHaveLength(2); + expect(findings.every(finding => !finding.callIdCorrupt)).toBe(true); }); test("a marker fragmented across delta boundaries still fires", () => { @@ -119,6 +123,26 @@ describe("cursor external output quarantine + corrective retry (devlog 260826 ga expect(findings[0]!.callIdCorrupt).toBe(true); }); + test("closely spaced markers retain independent corruption findings", () => { + const observer = new CursorMidstreamEchoObserver(); + observer.feed("lead\n[Tool Result]\nfc_63367283 mar-broken_0\n[Tool Error]\nclean\n"); + expect(observer.findings()).toEqual([ + { marker: "[Tool Result]", offset: 5, callIdCorrupt: true }, + { marker: "[Tool Error]", offset: 44, callIdCorrupt: false }, + ]); + }); + + test("a clean marker is not contaminated by corruption after the next marker", () => { + const observer = new CursorMidstreamEchoObserver(); + observer.feed("lead\n[Tool Result]\nclean\n[Tool Error]\nfc_123 mar-broken_0\n"); + const findings = observer.findings(); + expect(findings).toHaveLength(2); + expect(findings[0]!.marker).toBe("[Tool Result]"); + expect(findings[0]!.callIdCorrupt).toBe(false); + expect(findings[1]!.marker).toBe("[Tool Error]"); + expect(findings[1]!.callIdCorrupt).toBe(true); + }); + test("a mid-line marker mention does not fire", () => { const observer = new CursorMidstreamEchoObserver(); observer.feed("first line\nThe string [Tool Result] appeared in the transcript I reviewed.\n"); @@ -313,6 +337,133 @@ describe("cursor external output quarantine + corrective retry (devlog 260826 ga expect(text).toBe("[note] leading bracket but not an envelope"); }); + test("reasoning-only quarantine is capped and disarms before unbounded retention", async () => { + let attempt = 0; + const factory = () => ({ + async *run() { + attempt += 1; + for (let i = 0; i < 100; i += 1) { + yield { type: "thinking", thinking: "x".repeat(128) } satisfies CursorServerMessage; + } + // Once the aggregate hold cap flushes, later marker-like text is ordinary output rather + // than evidence for a retry whose preceding reasoning has already reached the client. + yield { type: "text", text: ECHO_TEXT } satisfies CursorServerMessage; + yield { type: "done", usage: { inputTokens: 1, outputTokens: 1 } } satisfies CursorServerMessage; + }, + writeClient() {}, + }); + const adapter = createCursorAdapter( + { ...provider, apiKey: "cursor-token" }, + { createTransport: factory as never }, + ); + const events: AdapterEvent[] = []; + await adapter.runTurn?.( + toolResultBody("cursor/kimi-k3"), + { headers: new Headers() }, + event => events.push(event), + ); + + expect(attempt).toBe(1); + expect(events.filter(event => event.type === "thinking_delta")).toHaveLength(100); + expect(events.filter(event => event.type === "text_delta")).not.toHaveLength(0); + }); + + test("an oversized first text delta is still classified by the echo sniffer", async () => { + // One text delta larger than the aggregate hold cap whose leading bytes are the + // echoed envelope marker. + let attempt = 0; + const runRequests: CursorRunRequest[] = []; + const oversizedFactory = () => ({ + async *run(request: CursorRunRequest) { + runRequests.push(request); + attempt += 1; + if (attempt === 1) { + yield { type: "text", text: ECHO_TEXT + "x".repeat(32 * 1024) } satisfies CursorServerMessage; + yield { type: "done", usage: { inputTokens: 1, outputTokens: 1 } } satisfies CursorServerMessage; + return; + } + yield { type: "text", text: "STATE A17" } satisfies CursorServerMessage; + yield { type: "done", usage: { inputTokens: 1, outputTokens: 1 } } satisfies CursorServerMessage; + }, + writeClient() {}, + }); + const adapter = createCursorAdapter({ ...provider, apiKey: "cursor-token" }, { createTransport: oversizedFactory as never }); + const events: AdapterEvent[] = []; + await adapter.runTurn?.(toolResultBody("cursor/kimi-k3"), { headers: new Headers() }, event => events.push(event)); + expect(attempt).toBe(2); + const text = events.filter(e => e.type === "text_delta").map(e => (e as { text: string }).text).join(""); + expect(text).toBe("STATE A17"); + }); + + test("an oversized first text delta is still classified by the routing sniffer", async () => { + let attempt = 0; + const factory = () => ({ + async *run() { + attempt += 1; + if (attempt === 1) { + // Routing claim padded past the 8 KiB aggregate cap in a single delta. + yield { + type: "text", + text: "네이티브 셸은 차단됐으니 exec_command 경로로 읽겠습니다. " + "x".repeat(32 * 1024), + } satisfies CursorServerMessage; + yield { type: "done", usage: { inputTokens: 1, outputTokens: 1 } } satisfies CursorServerMessage; + return; + } + yield { type: "text", text: "READ_OK" } satisfies CursorServerMessage; + yield { type: "done", usage: { inputTokens: 1, outputTokens: 1 } } satisfies CursorServerMessage; + }, + writeClient() {}, + }); + const body = { + modelId: "cursor/kimi-k3-1m", + context: { + messages: [{ role: "user", content: "Read the file and report its first line.", timestamp: 1 }], + tools: [{ + name: "exec", + description: "Run JavaScript code to orchestrate nested tool calls.", + parameters: {}, + freeform: true, + }], + }, + stream: false, + options: {}, + _cursorConversationId: "cursor_routing_oversized", + _cursorIdentityScope: "acct-routing-commentary", + } as OcxParsedRequest; + const adapter = createCursorAdapter({ ...provider, apiKey: "cursor-token" }, { createTransport: factory as never }); + const events: AdapterEvent[] = []; + await adapter.runTurn?.(body, { headers: new Headers() }, event => events.push(event)); + expect(attempt).toBe(2); + const text = events.filter(e => e.type === "text_delta").map(e => (e as { text: string }).text).join(""); + expect(text).toBe("READ_OK"); + }); + + test("a single reasoning frame larger than the cap flushes without unbounded retention", async () => { + let attempt = 0; + const bigThinking = "y".repeat(64 * 1024); + const factory = () => ({ + async *run() { + attempt += 1; + yield { type: "thinking", thinking: bigThinking } satisfies CursorServerMessage; + yield { type: "text", text: "post-thought answer" } satisfies CursorServerMessage; + yield { type: "done", usage: { inputTokens: 1, outputTokens: 1 } } satisfies CursorServerMessage; + }, + writeClient() {}, + }); + const adapter = createCursorAdapter({ ...provider, apiKey: "cursor-token" }, { createTransport: factory as never }); + const events: AdapterEvent[] = []; + await adapter.runTurn?.( + toolResultBody("cursor/kimi-k3"), + { headers: new Headers() }, + event => events.push(event), + ); + expect(attempt).toBe(1); + const thinking = events.filter(e => e.type === "thinking_delta").map(e => (e as { thinking: string }).thinking).join(""); + expect(thinking).toBe(bigThinking); + const text = events.filter(e => e.type === "text_delta").map(e => (e as { text: string }).text).join(""); + expect(text).toBe("post-thought answer"); + }); + test("plain user turns (no trailing toolResult) never arm the sniffer", async () => { let attempt = 0; const factory = () => ({ @@ -511,4 +662,3 @@ describe("Cursor midstream envelope-echo remint", () => { clearCursorIncompleteToolRemintForTests(); }); }); -