From 3317c1eb7b05eddf79ecdc4317f601ce43843ce6 Mon Sep 17 00:00:00 2001 From: Brandon Date: Thu, 1 Oct 2026 10:10:57 -0400 Subject: [PATCH] fix(shell): retain final output with bounded UTF-8 head and tail --- docs/issue-247-validation.md | 39 ++++ src/commands/console_input.ts | 5 +- src/core/bounded_output.ts | 73 ++++++++ src/core/shell_session.ts | 19 +- src/core/tool_executor.ts | 27 +-- test/bounded_output.test.ts | 69 ++++++++ test/console_input.test.ts | 57 ++++-- test/console_shell.test.ts | 15 ++ test/shell_output_capture.test.ts | 283 ++++++++++++++++++++++++++++++ 9 files changed, 548 insertions(+), 39 deletions(-) create mode 100644 docs/issue-247-validation.md create mode 100644 src/core/bounded_output.ts create mode 100644 test/bounded_output.test.ts create mode 100644 test/shell_output_capture.test.ts diff --git a/docs/issue-247-validation.md b/docs/issue-247-validation.md new file mode 100644 index 0000000..3fc6abd --- /dev/null +++ b/docs/issue-247-validation.md @@ -0,0 +1,39 @@ +# Issue #247: final shell output retention + +Reviewed against the original issue on 2026-10-01, based on main +`892b27b405964b779b7a7d7215f3458278a3f038`. + +The original one-shot implementation retained a prefix up to 64 MiB and then +discarded all subsequent chunks. The persistent shell retained only its first +8,000 characters. Both could drop the real final summary. + +## Acceptance evidence + +| Original requirement | Implementation and regression evidence | +| --- | --- | +| Bounded head and rolling tail; continue draining both pipes | `BoundedOutput` owns two fixed buffers (7,920 payload bytes within an 8,000-byte body budget). Both execution paths continue reading and streaming stdout/stderr. Tests count every emitted byte across 66 MiB real child output. | +| Actual final bytes survive beyond 64 MiB | `run_shell`, `run_tests`, direct user execution, and persistent user/model execution retain initial output plus distinct final stdout/stderr markers and authoritative exit 7. | +| Single oversized chunk | A direct 65 MiB append retains `HEAD_247` and `FINAL SUMMARY: 1 failed`, with fixed retained capacity. | +| Incremental UTF-8, interleaved pipes | Independent pipe decoders; unit tests cover two-, three-, and four-byte cuts, wrapping, and exact omission accounting. Child fixtures stagger individual UTF-8 bytes across both pipes; persistent callbacks never split surrogate pairs. | +| Explicit omissions without invented summary | Rendered notices count omitted decoded UTF-8 bytes; final text is copied from actual output. Malformed input is decoded using Node's normal replacement semantics, so counts describe normalized UTF-8 rather than malformed raw bytes. | +| Cancellation and timeout distinctions | Real post-cap termination preserves summaries and codes 130/124 for both execution paths. Persistent sessions visibly lose state, refuse execution until reset, and produce clean output afterward. | +| Completion and isolation | Capture renders after pipe draining/decoder flush. Persistent protocol markers never appear in normal capture or callbacks; next commands have clean output. | + +Explicit `/shell-result` sharing also uses bounded head/tail capture so a long +Unicode command cannot cause a second head-only truncation of the final summary. +File/web/search `capHeadTail` behavior remains unchanged. + +## Validation + +Seven new integration regressions failed on the unmodified implementation and +passed after the fix. Independent review included 5,000 randomized Unicode +capture cases without a prefix/suffix, byte-count, boundary, or budget failure. + +Local environment: Linux, Node 24.19.0. The checked-in tests also run in the +repository's Windows CI job; persistent Bash tests explicitly skip unsupported +platforms. + +Final local suite and hosted-check results are recorded in the PR and issue +closure comment. The initial full-suite run exposed two fixed-delay console +routing fixtures; those scenarios passed independently, and their synchronization +was hardened to observe completion rather than assume wall-clock timing. diff --git a/src/commands/console_input.ts b/src/commands/console_input.ts index 21566b6..54f68fe 100644 --- a/src/commands/console_input.ts +++ b/src/commands/console_input.ts @@ -1,6 +1,7 @@ import { ShellSession, type ShellCommandEvent } from "../core/shell_session.js"; import { randomUUID } from "node:crypto"; import { ToolExecutor } from "../core/tool_executor.js"; +import { BoundedOutput } from "../core/bounded_output.js"; import { sanitizeServerText } from "../core/transport.js"; /** Classify before history, prompt rewriting, or the busy queue. */ @@ -58,7 +59,9 @@ export class ConsoleShell { } }); if (fallback) this.event({ ...fallback, state: result.exitCode === 130 ? "cancelled" : "completed", exitCode: result.exitCode }); const full = `!${input.command}\ncwd: ${this.session.cwd}\nexit: ${result.exitCode}\n${result.output}`; - this.result = Buffer.from(full).subarray(0, 8192).toString("utf8").replace(/\ufffd$/, ""); + const shared = new BoundedOutput(8192); + shared.append(full); + this.result = shared.render(); // Stream once; retain the bounded capture for explicit sharing. Refusal and // state-loss explanations still render even when some output was streamed. const visible = streamed ? result.output.split("\n", 1)[0]! : result.output; diff --git a/src/core/bounded_output.ts b/src/core/bounded_output.ts new file mode 100644 index 0000000..3e1eefd --- /dev/null +++ b/src/core/bounded_output.ts @@ -0,0 +1,73 @@ +/** Bounded UTF-8 capture of decoded output, in observed pipe-event order. + * Callers decode each pipe independently before appending. The fixed buffers + * own their bytes: even a huge input cannot leave a retained backing buffer. + */ +export class BoundedOutput { + private readonly head: Buffer; + private readonly tail: Buffer; + private headLength = 0; + private tailLength = 0; + private tailNext = 0; + private totalBytes = 0; + + constructor(maxBytes = 8000) { + if (!Number.isSafeInteger(maxBytes) || maxBytes < 128) throw new RangeError("output budget must be at least 128 bytes"); + // Reserve space for the omission notice, including a safe-integer count. + const capacity = maxBytes - 80; + this.head = Buffer.alloc(Math.floor(capacity / 3)); + this.tail = Buffer.alloc(capacity - this.head.length); + } + + /** Fixed allocation, independent of the number/size of incoming chunks. */ + get capacityBytes(): number { return this.head.length + this.tail.length; } + get retainedBytes(): number { return this.headLength + this.tailLength; } + + append(text: string): void { + const bytes = Buffer.from(text, "utf8"); + this.totalBytes += bytes.length; + let offset = 0; + if (this.headLength < this.head.length) { + const count = Math.min(bytes.length, this.head.length - this.headLength); + bytes.copy(this.head, this.headLength, 0, count); + this.headLength += count; + offset = count; + } + const remaining = bytes.length - offset; + if (remaining >= this.tail.length) { + bytes.copy(this.tail, 0, bytes.length - this.tail.length); + this.tailLength = this.tail.length; + this.tailNext = 0; + } else if (remaining > 0) { + const first = Math.min(remaining, this.tail.length - this.tailNext); + bytes.copy(this.tail, this.tailNext, offset, offset + first); + bytes.copy(this.tail, 0, offset + first); + this.tailNext = (this.tailNext + remaining) % this.tail.length; + this.tailLength = Math.min(this.tail.length, this.tailLength + remaining); + } + } + + render(): string { + const tail = this.tailLength < this.tail.length + ? this.tail.subarray(0, this.tailLength) + : Buffer.concat([this.tail.subarray(this.tailNext), this.tail.subarray(0, this.tailNext)]); + const head = this.head.subarray(0, this.headLength); + if (this.totalBytes <= this.capacityBytes) return Buffer.concat([head, tail]).toString("utf8"); + + // Never manufacture replacement characters at either elision boundary. + // Appended text contains complete code points; only our cuts can split one. + let headEnd = head.length; + if (headEnd > 0) { + let start = headEnd - 1; + while (start > 0 && (head[start]! & 0xc0) === 0x80) start--; + const lead = head[start]!; + const width = lead < 0x80 ? 1 : lead < 0xe0 ? 2 : lead < 0xf0 ? 3 : 4; + if (headEnd - start < width) headEnd = start; + } + let tailStart = 0; + while (tailStart < tail.length && (tail[tailStart]! & 0xc0) === 0x80) tailStart++; + const omitted = this.totalBytes - headEnd - (tail.length - tailStart); + return head.subarray(0, headEnd).toString("utf8") + + `\n…[${omitted} UTF-8 bytes elided]…\n` + + tail.subarray(tailStart).toString("utf8"); + } +} diff --git a/src/core/shell_session.ts b/src/core/shell_session.ts index 50b1c49..d1b6b0f 100644 --- a/src/core/shell_session.ts +++ b/src/core/shell_session.ts @@ -4,6 +4,7 @@ import { realpathSync, statSync } from "node:fs"; import { resolve, sep } from "node:path"; import type { Readable } from "node:stream"; import { StringDecoder } from "node:string_decoder"; +import { BoundedOutput } from "./bounded_output.js"; import { childEnv } from "./child_env.js"; import type { RunOptions, ToolResult } from "./tool_executor.js"; @@ -157,7 +158,7 @@ export class ShellSession { emit("running"); return new Promise((settle) => { const marker = `\x1e${commandId}\x1f`; - let output = ""; + const output = new BoundedOutput(); let completed = false; let control = ""; let code: number | null = null; @@ -165,11 +166,9 @@ export class ShellSession { let outDone = false; let errDone = false; const retain = (text: string): void => { - if (output.length < 8000) { - const kept = text.slice(0, 8000 - output.length); - output += kept; - options.onOutput?.(kept); - } + if (!text) return; + output.append(text); + options.onOutput?.(text); }; const streamReader = (stream: Readable, end: () => void): (() => void) => { let pending = ""; @@ -183,7 +182,9 @@ export class ShellSession { end(); } else { // A marker may span chunks. Drain everything except its suffix. - const safe = Math.max(0, pending.length - marker.length + 1); + let safe = Math.max(0, pending.length - marker.length + 1); + // The marker lookbehind must not divide an astral code point. + if (safe > 0 && /[\uD800-\uDBFF]/.test(pending[safe - 1]!)) safe--; retain(pending.slice(0, safe)); pending = pending.slice(safe); } @@ -197,7 +198,7 @@ export class ShellSession { clearTimeout(timer); options.signal?.removeEventListener("abort", abort); cleanupOut(); cleanupErr(); - if (state !== "completed") result.output += output; + result.output += output.render(); fd.off("data", onControl); this.failActive = null; emit(state, result.exitCode); @@ -213,7 +214,7 @@ export class ShellSession { return; } this.cwd = physical; - finish({ output: `[exit ${code}]\n${output}`, exitCode: code }, "completed"); + finish({ output: `[exit ${code}]\n`, exitCode: code }, "completed"); }; const cleanupOut = streamReader(child.stdout!, () => { outDone = true; check(); }); const cleanupErr = streamReader(child.stderr!, () => { errDone = true; check(); }); diff --git a/src/core/tool_executor.ts b/src/core/tool_executor.ts index 499ebe2..0994dbc 100644 --- a/src/core/tool_executor.ts +++ b/src/core/tool_executor.ts @@ -9,6 +9,8 @@ import { spawn, spawnSync } from "node:child_process"; import { closeSync, constants as fsConstants, existsSync, mkdirSync, openSync, readFileSync, readdirSync, realpathSync, statSync, writeFileSync } from "node:fs"; import { dirname, relative, resolve, sep } from "node:path"; +import { StringDecoder } from "node:string_decoder"; +import { BoundedOutput } from "./bounded_output.js"; import type { ToolName } from "./brain_protocol.js"; import { validateToolCall } from "./tool_registry.js"; import { GitCommitGuard, SpawnGitRunner } from "./git_commit_guard.js"; @@ -257,18 +259,19 @@ export class ToolExecutor { stdio: ["ignore", "pipe", "pipe"], }); - let out = ""; - let bytes = 0; - const CAP = 64 * 1024 * 1024; - const absorb = (chunk: Buffer): void => { - options.onOutput?.(chunk.toString("utf8")); - bytes += chunk.length; - // Keep draining past the cap so the pipe never blocks the child, but - // stop retaining; capHeadTail trims the ends at the boundary anyway. - if (bytes <= CAP) out += chunk.toString("utf8"); + const output = new BoundedOutput(MAX_OUTPUT); + const absorb = (text: string): void => { + if (!text) return; + output.append(text); + options.onOutput?.(text); }; - child.stdout?.on("data", absorb); - child.stderr?.on("data", absorb); + // stdout and stderr may interleave mid-codepoint. Each owns a decoder, + // while capture and live output receive every complete decoded chunk. + for (const pipe of [child.stdout, child.stderr]) { + const decoder = new StringDecoder("utf8"); + pipe?.on("data", (chunk: Buffer) => absorb(decoder.write(chunk))); + pipe?.on("end", () => absorb(decoder.end())); + } let settled = false; let verdict: "timeout" | "aborted" | null = null; @@ -324,7 +327,7 @@ export class ToolExecutor { // 'close' rather than 'exit': it fires once the pipes are drained, so a // test summary arriving with the exit is not lost. child.on("close", (code, sig) => { - const body = capHeadTail(out, MAX_OUTPUT); + const body = output.render(); if (verdict === "timeout") { finish({ output: `[timeout after ${Math.round(timeoutMs / 1000)}s]\n${body}`, exitCode: 124 }); return; diff --git a/test/bounded_output.test.ts b/test/bounded_output.test.ts new file mode 100644 index 0000000..e7fb1cd --- /dev/null +++ b/test/bounded_output.test.ts @@ -0,0 +1,69 @@ +import { test } from "node:test"; +import assert from "node:assert/strict"; +import { BoundedOutput } from "../src/core/bounded_output.js"; + +function verify(capture: BoundedOutput, source: string, budget: number): void { + const rendered = capture.render(); + assert.ok(Buffer.byteLength(rendered) <= budget); + assert.ok(capture.retainedBytes <= capture.capacityBytes); + const match = /\n…\[(\d+) UTF-8 bytes elided\]…\n/.exec(rendered); + if (!match) { assert.equal(rendered, source); return; } + const head = rendered.slice(0, match.index); + const tail = rendered.slice(match.index + match[0].length); + assert.ok(source.startsWith(head)); + assert.ok(source.endsWith(tail)); + assert.equal(Number(match[1]), Buffer.byteLength(source) - Buffer.byteLength(head + tail)); + assert.doesNotMatch(head + tail, /\ufffd|[\uD800-\uDBFF](?![\uDC00-\uDFFF])|(? { + const capture = new BoundedOutput(128); + const source = "a".repeat(15) + "šŸ˜€ā‚¬ē•Œ"; + capture.append(source); + verify(capture, source, 128); + assert.equal(capture.render(), source); +}); + +test("one chunk beyond 64 MiB retains the genuine head and tail in fixed owned buffers", () => { + const capture = new BoundedOutput(); + const source = "HEAD_247" + "x".repeat(65 * 1024 * 1024) + "FINAL SUMMARY: 1 failed"; + capture.append(source); + verify(capture, source, 8000); + assert.match(capture.render(), /^HEAD_247/); + assert.ok(capture.render().endsWith("FINAL SUMMARY: 1 failed")); + assert.equal(capture.capacityBytes, 7920); + assert.equal(capture.retainedBytes, 7920); +}); + +test("rolling tail wraps correctly across many chunks and repeated renders", () => { + const capture = new BoundedOutput(256); + let source = ""; + for (let index = 0; index < 1500; index++) { + const chunk = `${index}:` + "ā‚¬šŸ˜€ē•Œz".repeat(index % 11); + capture.append(chunk); + source += chunk; + verify(capture, source, 256); + } +}); + +test("UTF-8 head and tail cuts omit complete partial boundary bytes accurately", () => { + for (const point of ["Ć©", "€", "šŸ˜€"]) { + for (let offset = 0; offset < 8; offset++) { + const capture = new BoundedOutput(128); + const source = "x".repeat(offset) + point.repeat(100) + "FINAL"; + for (const character of source) capture.append(character); + verify(capture, source, 128); + assert.ok(capture.render().endsWith("FINAL")); + } + } +}); + +test("byte count describes normalized decoded UTF-8, including actual replacement characters", () => { + const capture = new BoundedOutput(128); + const source = "\ufffd".repeat(100); + capture.append(source); + const rendered = capture.render(); + const match = /\[(\d+) UTF-8 bytes elided\]/.exec(rendered)!; + const retained = rendered.replace(/\n…\[\d+ UTF-8 bytes elided\]…\n/, ""); + assert.equal(Number(match[1]), Buffer.byteLength(source) - Buffer.byteLength(retained)); +}); diff --git a/test/console_input.test.ts b/test/console_input.test.ts index e1b138a..ae0b09c 100644 --- a/test/console_input.test.ts +++ b/test/console_input.test.ts @@ -66,51 +66,74 @@ for (const tty of [false, true]) { Object.defineProperty(process, 'stdin', { value: input }); Object.defineProperty(process.stdout, 'isTTY', { value: ${tty} }); Object.defineProperty(process.stdout, 'columns', { value: 100 }); + let observed = ''; + const write = process.stdout.write.bind(process.stdout); + process.stdout.write = (chunk, ...args) => { observed += String(chunk); return write(chunk, ...args); }; let calls = 0; let sharing = false; + let releaseModel = null; globalThis.fetch = async (_url, options) => { const body = String(options?.body ?? ''); if (body.includes('SHELL_ONLY') || body.includes('QUEUED_SHELL')) throw new Error('shell output leaked'); if (body.includes('SHARE_ALLOWED') && !sharing) throw new Error('implicit sharing'); if (sharing && !body.includes('SHARE_ALLOWED')) throw new Error('explicit sharing missing output'); calls++; process.stdout.write('MODEL_CALL_' + calls + '\\n'); - await new Promise(r => setTimeout(r, 100)); + await new Promise(resolve => { releaseModel = resolve; }); return new Response('data: {"type":"delta","text":"model response"}\\n\\ndata: {"type":"done","uvt":0,"cents":0}\\n\\n', { headers: {'content-type':'text/event-stream'} }); }; const tokens = { get: async () => 'test-token' }; const ctx = { cfg: { backend:'cloud', baseUrl:'https://stub.test', defaultModel:'', defaultEffort:'', permissionMode:'ask', autoApply:false, telemetry:false }, flags: { cwd:${JSON.stringify(cwd)}, json:false, yes:false }, tokens, api:new ApiClient('https://stub.test', tokens) }; const enter = ${tty ? "'\\r'" : "'\\n'"}; const submit = text => input.write(text + enter); - const delay = ms => new Promise(r => setTimeout(r, ms)); + const until = async (predicate, label) => { + const deadline = Date.now() + 5000; + while (!predicate()) { + if (Date.now() >= deadline) throw new Error('waiting for ' + label + ': ' + observed.slice(-2000)); + await new Promise(resolve => setTimeout(resolve, 5)); + } + // Let the submit handler finish its continuation and clear busy before + // the next input; a terminal state event precedes that continuation. + await new Promise(resolve => setImmediate(resolve)); + }; + const completed = () => (observed.match(/\\| completed \\| exit \\d+ \\| session /g) ?? []).length; + const releaseTurn = async number => { + await until(() => calls === number && releaseModel !== null, 'model call ' + number); + const release = releaseModel; + releaseModel = null; + release(); + await until(() => (observed.match(/model response/g) ?? []).length >= number, 'model response ' + number); + }; const session = cmdChat(ctx, ''); - await delay(50); + await until(() => input.listenerCount('data') > 0 && observed.includes(${JSON.stringify(cwd)}), 'console input ready'); submit('!echo SHELL_ONLY'); - await delay(150); + await until(() => completed() === 1, 'first shell completion'); if (calls !== 0) throw new Error('shell made API call'); submit(${JSON.stringify(process.platform === 'win32' ? '!cd' : '!pwd')}); - await delay(100); + await until(() => completed() === 2, 'cwd command completion'); submit('hello'); - await delay(20); + await until(() => calls === 1 && releaseModel !== null, 'model turn running'); submit('!echo QUEUED_SHELL'); - await delay(250); + ${tty ? "await until(() => observed.includes('Local shell queued'), 'shell queued during model turn');" : ""} + await releaseTurn(1); + await until(() => completed() === 3, 'queued shell completion'); submit(${JSON.stringify(`!"${process.execPath}" -e "process.exit(7)"`)}); - await delay(100); + await until(() => completed() === 4 && observed.includes('completed | exit 7'), 'nonzero shell completion'); submit('!'); - await delay(50); - ${tty ? "input.write('\\x1b[200~!echo MULTILINE_BAD\\necho SECOND_BAD\\x1b[201~\\r'); await delay(50);" : ""} + await until(() => observed.includes('usage: !'), 'empty shell refusal'); + ${tty ? "input.write('\\x1b[200~!echo MULTILINE_BAD\\necho SECOND_BAD\\x1b[201~\\r'); await until(() => completed() === 5, 'multiline shell completion');" : ""} submit(${JSON.stringify(`!"${process.execPath}" -e "setTimeout(()=>{},30000)"`)}); - await delay(100); + await until(() => observed.includes(${JSON.stringify(`running] !"${process.execPath}" -e "setTimeout(()=>{},30000)"`)}), 'cancellable shell running'); ${tty ? "input.write('\\\\!literal'); input.write('\\x03');" : "process.emit('SIGINT');"} - await delay(350); + await until(() => observed.includes('cancelled | exit 130'), 'shell cancellation'); ${tty ? "input.write(enter);" : "submit('\\\\!literal');"} - await delay(200); + await releaseTurn(2); submit('/shell-reset'); - await delay(100); + await until(() => observed.includes('shell reset — cwd/environment/functions cleared'), 'shell reset'); submit('!echo SHARE_ALLOWED'); - await delay(100); + await until(() => completed() === ${tty ? 6 : 5}, 'shareable shell completion'); sharing = true; submit('/shell-result'); - await delay(200); + await releaseTurn(3); submit('/exit'); await session; if(calls !== 3) throw new Error('wrong model call count: ' + calls); @@ -125,7 +148,7 @@ for (const tty of [false, true]) { let output = ""; child.stdout.on("data", (chunk) => { output += chunk; }); child.stderr.on("data", (chunk) => { output += chunk; }); - const timer = setTimeout(() => { child.kill(); reject(new Error("console timed out: " + output)); }, 10000); + const timer = setTimeout(() => { child.kill(); reject(new Error("console timed out: " + output)); }, 30000); child.on("error", reject); child.on("close", (code) => { clearTimeout(timer); resolve({ code, output }); }); }); diff --git a/test/console_shell.test.ts b/test/console_shell.test.ts index 2d1d65f..3be759a 100644 --- a/test/console_shell.test.ts +++ b/test/console_shell.test.ts @@ -59,6 +59,21 @@ test("line console shell commands make zero model calls, keep output out of prom } }); +test("explicit shell-result sharing keeps a final summary after a long Unicode command", { skip: !supported }, async () => { + const root = mkdtempSync(join(tmpdir(), "aether-console-share-")); + const shell = new ConsoleShell(root, () => {}, true); + try { + await shell.run(`printf 'FINAL SUMMARY: 1 failed'; # ${"šŸ˜€".repeat(4000)}`); + const shared = shell.share(); + assert.equal(shared.kind, "chat"); + if (shared.kind !== "chat") return; + assert.ok(shared.text.endsWith("FINAL SUMMARY: 1 failed")); + assert.match(shared.text, /UTF-8 bytes elided/); + assert.doesNotMatch(shared.text, /\ufffd/); + assert.ok(Buffer.byteLength(shared.text) <= 8192 + 80); + } finally { shell.close(); rmSync(root, { recursive: true, force: true }); } +}); + class ShellBrain implements Brain { result: ToolResult | null = null; async *run(_task: TaskCommand): AsyncGenerator { diff --git a/test/shell_output_capture.test.ts b/test/shell_output_capture.test.ts new file mode 100644 index 0000000..343ceae --- /dev/null +++ b/test/shell_output_capture.test.ts @@ -0,0 +1,283 @@ +import { test } from "node:test"; +import assert from "node:assert/strict"; +import { rmSync, writeFileSync } from "node:fs"; +import { join } from "node:path"; +import { ToolExecutor, type RunOptions, type ToolResult } from "../src/core/tool_executor.js"; +import { ShellSession } from "../src/core/shell_session.js"; +import { tmpWorkspace } from "./tmp_workspace.js"; + +const supportedSession = process.platform === "linux" || process.platform === "darwin"; +const FIRST = "FIRST_STDOUT_247\n"; +const LAST_OUT = "FINAL_STDOUT_SUMMARY_247\n"; +const LAST_ERR = "FINAL_STDERR_SUMMARY_247\n"; +const READY = "READY_FOR_STOP_247"; +const FLOOD_BYTES = 66 * 1024 * 1024; +const quote = (value: string): string => process.platform === "win32" + ? `"${value.replaceAll('"', '""')}"` + : "'" + value.replaceAll("'", "'\\''") + "'"; + +// Generate the flood in the child and await each write callback. The parent +// never builds a giant fixture string or accumulates the live output stream. +const floodScript = ` +const write = (stream, data) => new Promise((resolve, reject) => { + stream.write(data, error => error ? reject(error) : resolve()); +}); +const delay = ms => new Promise(resolve => setTimeout(resolve, ms)); +(async () => { + await write(process.stdout, ${JSON.stringify(FIRST)}); + await write(process.stderr, "FIRST_STDERR_247\\n"); + await delay(20); + const chunk = Buffer.alloc(64 * 1024, 120); + for (let i = 0; i < 65 * 16; i++) await write(process.stdout, chunk); + for (let i = 0; i < 16; i++) await write(process.stderr, chunk); + await delay(20); + await write(process.stdout, ${JSON.stringify(LAST_OUT)}); + await write(process.stderr, ${JSON.stringify(LAST_ERR)}); + if (process.argv[2] === "wait") { + await write(process.stdout, ${JSON.stringify(`${READY}\n`.repeat(6))}); + setInterval(() => {}, 1000); + } else { + process.exitCode = 7; + } +})().catch(error => { console.error(error); process.exitCode = 1; }); +`; + +function fixture(persistent = false, script = floodScript): { + exec: ToolExecutor; + command: (mode?: string) => string; + session?: ShellSession; + close: () => void; +} { + const root = tmpWorkspace("aether-output-247-"); + const path = join(root, "flood.cjs"); + writeFileSync(path, script); + const session = persistent ? new ShellSession(root) : undefined; + const exec = new ToolExecutor(root, undefined, { mode: "coding", ...(session ? { shellSession: session } : {}) }); + return { + exec, + ...(session ? { session } : {}), + command: (mode = "exit") => `${quote(process.execPath)} ${quote(path)} ${quote(mode)}`, + close: () => { exec.close(); rmSync(root, { recursive: true, force: true }); }, + }; +} + +function streamed(): { bytes: number; tail: string; output: (chunk: string) => void } { + const result = { + bytes: 0, + tail: "", + output(chunk: string): void { + result.bytes += Buffer.byteLength(chunk, "utf8"); + result.tail = (result.tail + chunk).slice(-512); + }, + }; + return result; +} + +function assertSummary(result: ToolResult, exitCode: number): void { + assert.equal(result.exitCode, exitCode, result.output); + assert.ok(result.output.includes(FIRST), "the initial output must survive truncation"); + assert.ok(result.output.includes(LAST_OUT), "the final stdout summary must survive the 64 MiB boundary"); + assert.ok(result.output.includes(LAST_ERR), "the final stderr summary must survive the 64 MiB boundary"); + assert.ok(result.output.length <= 8300, "retained output must remain bounded, including status and truncation notices"); +} + +type Entrypoint = "run_shell" | "run_tests" | "runUserCommand"; +function run(exec: ToolExecutor, entrypoint: Entrypoint, command: string, options: RunOptions): Promise { + return entrypoint === "runUserCommand" + ? exec.runUserCommand(command, options) + : exec.executeAsync(entrypoint, { command }, options); +} + +for (const entrypoint of ["run_shell", "run_tests", "runUserCommand"] as const) { + test(`${entrypoint} retains first and final summaries after more than 64 MiB while streaming every byte`, { timeout: 30_000 }, async () => { + const w = fixture(); + const live = streamed(); + try { + const result = await run(w.exec, entrypoint, w.command(), { onOutput: live.output }); + assertSummary(result, 7); + assert.match(result.output, /^\[exit 7\]\n/); + assert.equal(live.bytes, FLOOD_BYTES + Buffer.byteLength(FIRST + "FIRST_STDERR_247\n" + LAST_OUT + LAST_ERR)); + assert.ok(live.tail.includes(LAST_OUT), "live stdout must continue after the retention cap"); + assert.ok(live.tail.includes(LAST_ERR), "live stderr must continue after the retention cap"); + } finally { w.close(); } + }); +} + +test("one-shot live output decodes split UTF-8 independently for stdout and stderr", { timeout: 10_000 }, async () => { + const root = tmpWorkspace("aether-output-utf8-247-"); + const path = join(root, "utf8.cjs"); + writeFileSync(path, ` +const pause = () => new Promise(resolve => setTimeout(resolve, 20)); +(async () => { + const out = Buffer.from("OUT:ā‚¬šŸ˜€\\n"); + const err = Buffer.from("ERR:ē•ŒšŸ¦€\\n"); + for (let i = 0; i < Math.max(out.length, err.length); i++) { + if (i < out.length) process.stdout.write(out.subarray(i, i + 1)); + if (i < err.length) process.stderr.write(err.subarray(i, i + 1)); + await pause(); + } +})(); +`); + const live = streamed(); + try { + const result = await new ToolExecutor(root).runUserCommand(`${quote(process.execPath)} ${quote(path)}`, { onOutput: live.output }); + assert.equal(result.exitCode, 0, result.output); + assert.doesNotMatch(result.output, /\ufffd/, "retained output must not replace a split multibyte character"); + assert.doesNotMatch(live.tail, /\ufffd/, "live output must not replace a split multibyte character"); + // The streams can interleave, so compare the emitted characters without + // imposing an artificial ordering across independent pipes. + const payload = "OUT:ā‚¬šŸ˜€\nERR:ē•ŒšŸ¦€\n"; + assert.equal(live.bytes, Buffer.byteLength(payload)); + assert.deepEqual([...live.tail].sort(), [...payload].sort()); + assert.deepEqual([...result.output.replace(/^\[exit 0\]\n/, "")].sort(), [...payload].sort()); + } finally { rmSync(root, { recursive: true, force: true }); } +}); + +test("one-shot cancellation after the output cap keeps recent summaries and reports 130", { timeout: 30_000 }, async () => { + const w = fixture(); + const live = streamed(); + const abort = new AbortController(); + let stopRequested = false; + try { + const result = await w.exec.runUserCommand(w.command("wait"), { + timeoutMs: 20_000, + signal: abort.signal, + onOutput(chunk) { + live.output(chunk); + if (!stopRequested && live.tail.includes(READY)) { + stopRequested = true; + setImmediate(() => abort.abort()); + } + }, + }); + assert.equal(stopRequested, true, "cancellation must happen only after the child emits its post-cap marker"); + assert.ok(live.bytes > 64 * 1024 * 1024); + assertSummary(result, 130); + assert.match(result.output, /^\[aborted\]/); + assert.ok(result.output.includes(READY)); + } finally { w.close(); } +}); + +test("one-shot timeout after the output cap keeps recent summaries and reports 124", { timeout: 30_000 }, async () => { + const w = fixture(); + const live = streamed(); + try { + const result = await w.exec.executeAsync("run_tests", { command: w.command("wait") }, { + timeoutMs: 10_000, + onOutput: live.output, + }); + assert.ok(live.bytes > 64 * 1024 * 1024); + assert.ok(live.tail.includes(READY), "the timeout must occur after the child emits its post-cap marker"); + assertSummary(result, 124); + assert.match(result.output, /^\[timeout after /); + assert.ok(result.output.includes(READY)); + } finally { w.close(); } +}); + +test("persistent model and user shell commands retain post-cap summaries, drain streams, and isolate the next command", { + skip: !supportedSession, + timeout: 30_000, +}, async () => { + const w = fixture(true); + try { + for (const entrypoint of ["run_shell", "runUserCommand"] as const) { + const live = streamed(); + const result = await run(w.exec, entrypoint, w.command(), { onOutput: live.output }); + assertSummary(result, 7); + assert.equal(live.bytes, FLOOD_BYTES + Buffer.byteLength(FIRST + "FIRST_STDERR_247\n" + LAST_OUT + LAST_ERR)); + assert.ok(live.tail.includes(LAST_OUT)); + assert.ok(live.tail.includes(LAST_ERR)); + assert.equal(w.session!.state, "ready"); + const nextLive = streamed(); + const next = await w.exec.runUserCommand("printf NEXT_COMMAND_247", { onOutput: nextLive.output }); + assert.deepEqual(next, { exitCode: 0, output: "[exit 0]\nNEXT_COMMAND_247" }); + assert.equal(nextLive.tail, "NEXT_COMMAND_247"); + assert.equal(nextLive.bytes, Buffer.byteLength("NEXT_COMMAND_247")); + } + } finally { w.close(); } +}); + +test("persistent live output preserves split UTF-8 through marker lookbehind without protocol leakage", { + skip: !supportedSession, + timeout: 10_000, +}, async () => { + // Follow the astral characters with enough individually written bytes to + // slide each surrogate pair across the protocol marker's lookbehind edge. + const out = "OUT:ā‚¬šŸ˜€" + ".".repeat(70) + "\n"; + const err = "ERR:ē•ŒšŸ¦€" + "-".repeat(70) + "\n"; + const w = fixture(true, ` +const pause = () => new Promise(resolve => setTimeout(resolve, 20)); +(async () => { + const out = Buffer.from(${JSON.stringify(out)}); + const err = Buffer.from(${JSON.stringify(err)}); + for (let i = 0; i < Math.max(out.length, err.length); i++) { + if (i < out.length) process.stdout.write(out.subarray(i, i + 1)); + if (i < err.length) process.stderr.write(err.subarray(i, i + 1)); + await pause(); + } +})(); +`); + const live = streamed(); + let completeUnicodeChunks = true; + try { + const result = await w.exec.executeAsync("run_shell", { command: w.command() }, { + onOutput(chunk) { + completeUnicodeChunks &&= Buffer.from(chunk, "utf8").toString("utf8") === chunk; + live.output(chunk); + }, + }); + assert.equal(result.exitCode, 0, result.output); + assert.equal(completeUnicodeChunks, true, "each live callback must contain complete Unicode characters"); + assert.doesNotMatch(result.output, /[\ufffd\x1e\x1f\x00]/); + assert.doesNotMatch(live.tail, /[\ufffd\x1e\x1f\x00]/); + assert.equal(live.bytes, Buffer.byteLength(out + err)); + assert.deepEqual([...live.tail].sort(), [...(out + err)].sort()); + assert.deepEqual([...result.output.replace(/^\[exit 0\]\n/, "")].sort(), [...(out + err)].sort()); + const next = await w.exec.runUserCommand("printf UTF8_NEXT_COMMAND_247"); + assert.deepEqual(next, { exitCode: 0, output: "[exit 0]\nUTF8_NEXT_COMMAND_247" }); + } finally { w.close(); } +}); + +for (const termination of ["cancellation", "timeout"] as const) { + test(`persistent ${termination} after the output cap keeps summaries, loses state visibly, and resets cleanly`, { + skip: !supportedSession, + timeout: 30_000, + }, async () => { + const w = fixture(true); + const live = streamed(); + const abort = new AbortController(); + let stopRequested = false; + try { + const result = await w.exec.runUserCommand(w.command("wait"), { + timeoutMs: termination === "timeout" ? 10_000 : 20_000, + signal: abort.signal, + onOutput(chunk) { + live.output(chunk); + if (termination === "cancellation" && !stopRequested && live.tail.includes(READY)) { + stopRequested = true; + setImmediate(() => abort.abort()); + } + }, + }); + assert.ok(live.bytes > 64 * 1024 * 1024); + assert.ok(live.tail.includes(READY), "termination must happen after the child emits its post-cap marker"); + assertSummary(result, termination === "cancellation" ? 130 : 124); + assert.ok(result.output.includes(READY)); + assert.match(result.output, termination === "cancellation" ? /^\[aborted;/ : /^\[shell command timed out;/); + assert.match(result.output, /shell state lost/); + assert.doesNotMatch(result.output, /[\x1e\x1f\x00]/, "command protocol must never appear in retained output"); + assert.doesNotMatch(live.tail, /[\x1e\x1f\x00]/, "command protocol must never appear in live output"); + if (termination === "cancellation") assert.equal(stopRequested, true); + assert.equal(w.session!.state, "lost"); + const refused = await w.exec.runUserCommand("printf SHOULD_NOT_RUN_247"); + assert.equal(refused.exitCode, 1); + assert.doesNotMatch(refused.output, /SHOULD_NOT_RUN_247/); + w.session!.reset(); + const nextLive = streamed(); + const next = await w.exec.runUserCommand("printf AFTER_RESET_247", { onOutput: nextLive.output }); + assert.deepEqual(next, { exitCode: 0, output: "[exit 0]\nAFTER_RESET_247" }); + assert.equal(nextLive.tail, "AFTER_RESET_247"); + assert.equal(nextLive.bytes, Buffer.byteLength("AFTER_RESET_247")); + } finally { w.close(); } + }); +}