diff --git a/packages/coding-agent/src/debug/raw-sse-buffer.ts b/packages/coding-agent/src/debug/raw-sse-buffer.ts index 1bb6d0b3e..b62614666 100644 --- a/packages/coding-agent/src/debug/raw-sse-buffer.ts +++ b/packages/coding-agent/src/debug/raw-sse-buffer.ts @@ -119,9 +119,16 @@ export class RawSseDebugBuffer { #records: RawSseDebugRecord[] = []; // Parallel to `#records`: `#recordChars[i]` is the precomputed char count // for `#records[i]`. Kept in lockstep by `#append` (push both) and - // `#enforceLimits` (shift both). See the comment above the class for why - // this is a sidecar array instead of a per-record property. + // `#enforceLimits` (advance `#head` to evict, then `slice` both together + // when compacting). See the comment above the class for why this is a + // sidecar array instead of a per-record property. #recordChars: number[] = []; + // Head-index ring over `#records`/`#recordChars`: index of the oldest live + // record. Eviction advances `#head` (amortized O(1)) rather than an O(n) + // front `shift()`; the dead `[0, #head)` prefix is reclaimed lazily by + // `#enforceLimits`. Live count is `#records.length - #head`; the live + // records are `#records[#head ..]`. + #head = 0; #totalChars = 0; #droppedRecords = 0; #droppedChars = 0; @@ -181,7 +188,7 @@ export class RawSseDebugBuffer { snapshot(): RawSseDebugSnapshot { return { - records: [...this.#records], + records: this.#records.slice(this.#head), droppedRecords: this.#droppedRecords, droppedChars: this.#droppedChars, totalEvents: this.#totalEvents, @@ -190,9 +197,12 @@ export class RawSseDebugBuffer { } toRawText(): string { - // Reads the live array directly: `rawRecordText` only computes a string - // from each record, so no caller-visible mutation is possible. - const body = this.#records.map(rawRecordText).join("\n"); + // Reads the live window directly: `rawRecordText` only computes a string + // from each record, so no caller-visible mutation is possible. With a + // non-empty dead prefix we map a slice past `#head`; `#head === 0` (the + // common case) maps `#records` in place with no extra copy. + const live = this.#head === 0 ? this.#records : this.#records.slice(this.#head); + const body = live.map(rawRecordText).join("\n"); if (this.#droppedRecords === 0) return body; const dropped = `: omp-debug-dropped records=${this.#droppedRecords} chars=${this.#droppedChars}\n\n`; return body.length > 0 ? `${dropped}${body}` : dropped; @@ -208,14 +218,25 @@ export class RawSseDebugBuffer { } #enforceLimits(): void { - while (this.#records.length > MAX_RAW_SSE_EVENTS || this.#totalChars > MAX_RAW_SSE_CHARS) { - if (this.#records.length === 0) return; - this.#records.shift(); - const chars = this.#recordChars.shift() ?? 0; + while (this.#records.length - this.#head > MAX_RAW_SSE_EVENTS || this.#totalChars > MAX_RAW_SSE_CHARS) { + if (this.#records.length - this.#head === 0) break; + const chars = this.#recordChars[this.#head] ?? 0; + this.#head += 1; this.#totalChars = Math.max(0, this.#totalChars - chars); this.#droppedRecords += 1; this.#droppedChars += chars; } + // Reclaim the consumed `[0, #head)` prefix once it grows large: one O(n) + // memmove amortized over many O(1) evictions, bounding the backing arrays + // to ~2x the live window. `#head >= MAX_RAW_SSE_EVENTS` covers the + // full-record-count steady state; `#head > liveCount` covers a small live + // window held by a few large records under the char budget. + const liveCount = this.#records.length - this.#head; + if (this.#head >= MAX_RAW_SSE_EVENTS || this.#head > liveCount) { + this.#records = this.#records.slice(this.#head); + this.#recordChars = this.#recordChars.slice(this.#head); + this.#head = 0; + } } #emit(): void { diff --git a/packages/coding-agent/test/debug/raw-sse-buffer.test.ts b/packages/coding-agent/test/debug/raw-sse-buffer.test.ts index ad5e66b04..f393e2feb 100644 --- a/packages/coding-agent/test/debug/raw-sse-buffer.test.ts +++ b/packages/coding-agent/test/debug/raw-sse-buffer.test.ts @@ -96,4 +96,78 @@ describe("RawSseDebugBuffer", () => { expect(buffer.toRawText()).toContain("requestId=req_pre_viewer"); expect(buffer.toRawText()).toContain("event: message_stop"); }); + + it("keeps oldest-first order and exact droppedRecords well past MAX_RAW_SSE_EVENTS", () => { + const buffer = new RawSseDebugBuffer(); + // > 2x MAX_RAW_SSE_EVENTS (1000) so the head-index ring compacts at least + // once; a corrupt slice would scramble order or counts. + const APPENDS = 2_300; + for (let i = 1; i <= APPENDS; i++) { + buffer.recordEvent({ event: null, data: "{}", raw: [`data: ${i}`] }, model); + } + + const snapshot = buffer.snapshot(); + // 1000 newest survive (sequences 1301..2300); the oldest 1300 are evicted. + expect(snapshot.records).toHaveLength(1_000); + expect(snapshot.droppedRecords).toBe(1_300); + expect(snapshot.totalEvents).toBe(APPENDS); + const last = snapshot.records[snapshot.records.length - 1]; + expect(snapshot.records[0].sequence).toBe(1_301); + expect(last.sequence).toBe(2_300); + // Contiguous + oldest-first across the compaction boundary. + expect(snapshot.records.every((record, idx) => record.sequence === 1_301 + idx)).toBe(true); + expect(rawSseRecordLines(snapshot.records[0])[0]).toBe("data: 1301"); + expect(rawSseRecordLines(last)[0]).toBe("data: 2300"); + }); + + it("evicts by char budget with exact droppedChars and live char total", () => { + const buffer = new RawSseDebugBuffer(); + // One 63_998-char line → originalChars 63_999 (≤ the 64_000 per-event cap, + // so no trim) → 64_000 chars per record. 8 records exactly fill + // MAX_RAW_SSE_CHARS (512_000); each further append evicts the oldest. + const PER_RECORD_CHARS = 64_000; + const LINE_LEN = 63_998; + const APPENDS = 12; + for (let i = 1; i <= APPENDS; i++) { + const line = `data: ${i}`.padEnd(LINE_LEN, "x"); + buffer.recordEvent({ event: null, data: "{}", raw: [line] }, model); + } + + const snapshot = buffer.snapshot(); + // 12 appended, 8 fit the budget → the oldest 4 (sequences 1..4) drop. + expect(snapshot.records).toHaveLength(8); + expect(snapshot.droppedRecords).toBe(4); + expect(snapshot.droppedChars).toBe(4 * PER_RECORD_CHARS); + expect(snapshot.records[0].sequence).toBe(5); + expect(snapshot.records[snapshot.records.length - 1].sequence).toBe(12); + // #totalChars is not exposed; derive live chars from the conservation + // invariant (sum appended − droppedChars) and confirm eviction stopped + // exactly at budget rather than over- or under-evicting. + const liveChars = APPENDS * PER_RECORD_CHARS - snapshot.droppedChars; + expect(liveChars).toBe(8 * PER_RECORD_CHARS); + expect(liveChars).toBeLessThanOrEqual(512_000); + }); + + it("emits the dropped header and oldest-first body in toRawText after eviction", () => { + const buffer = new RawSseDebugBuffer(); + const APPENDS = 1_005; // 5 past MAX_RAW_SSE_EVENTS → oldest 5 drop + for (let i = 1; i <= APPENDS; i++) { + buffer.recordEvent({ event: null, data: "{}", raw: [`data: ${i}`] }, model); + } + + // Dropped records 1..5: each "data: N" is 7 chars → 9 chars/record → 45. + const text = buffer.toRawText(); + expect(text.startsWith(": omp-debug-dropped records=5 chars=45\n\n")).toBe(true); + + // Body data lines, oldest-first, are exactly records 6..1005 — no dropped + // record leaks and order survives the head-index window. + const dataLines = text + .split("\n") + .filter(line => line.startsWith("data: ")) + .map(line => Number(line.slice("data: ".length))); + expect(dataLines).toHaveLength(1_000); + expect(dataLines[0]).toBe(6); + expect(dataLines.at(-1)).toBe(1_005); + expect(dataLines.every((n, idx) => n === 6 + idx)).toBe(true); + }); });