From 865d7de858906efe41e3cb2041ad83fb542bc688 Mon Sep 17 00:00:00 2001 From: Tianyi Cui <53024+tianyicui@users.noreply.github.com> Date: Wed, 22 Jul 2026 01:27:03 +0800 Subject: [PATCH] fix(code-runtime): drain late worker pipe output --- ...-20-code-mode-typed-tool-returns.i18n.yaml | 4 +- ...2026-07-20-code-mode-typed-tool-returns.md | 2 +- ...6-07-20-code-mode-typed-tool-returns.zh.md | 2 +- .../code-runtime-worker/README.md | 2 +- .../code-runtime-worker/src/index.ts | 79 ++++++++++++++----- .../code-runtime-worker/tests/runtime.spec.ts | 18 +++++ 6 files changed, 82 insertions(+), 25 deletions(-) diff --git a/.agents/notes/implemented/feature/2026-07-20-code-mode-typed-tool-returns.i18n.yaml b/.agents/notes/implemented/feature/2026-07-20-code-mode-typed-tool-returns.i18n.yaml index 3b3e9544f2..7981bedde0 100644 --- a/.agents/notes/implemented/feature/2026-07-20-code-mode-typed-tool-returns.i18n.yaml +++ b/.agents/notes/implemented/feature/2026-07-20-code-mode-typed-tool-returns.i18n.yaml @@ -2,5 +2,5 @@ # side as of the last confirmed-consistent state. Both languages carry equal authority; # after editing either side, bring the other along and re-record with: # pnpm run verify-translation-pairing --write -2026-07-20-code-mode-typed-tool-returns.md: 1beecbc9e5f61ac5dce50aeba508ce75cb0ce507 -2026-07-20-code-mode-typed-tool-returns.zh.md: 0dbad69f6b8120e0904f026961b004855d23f3ab +2026-07-20-code-mode-typed-tool-returns.md: 42fe44c1b6f6a0d0debc9d8e1f742b436cd22047 +2026-07-20-code-mode-typed-tool-returns.zh.md: c567fa86bbe381b9bd54ed37d83fbc0e18252b51 diff --git a/.agents/notes/implemented/feature/2026-07-20-code-mode-typed-tool-returns.md b/.agents/notes/implemented/feature/2026-07-20-code-mode-typed-tool-returns.md index 1beecbc9e5..42fe44c1b6 100644 --- a/.agents/notes/implemented/feature/2026-07-20-code-mode-typed-tool-returns.md +++ b/.agents/notes/implemented/feature/2026-07-20-code-mode-typed-tool-returns.md @@ -61,7 +61,7 @@ The runtime accepts an exact lossless JSON completion of any root. Returning `un `WorkerCodeRuntime` replaces the former independent log and value caps with configurable `maxOutputBytes`, defaulting to `67_108_864` bytes. One host-side hostile-peer ledger accounts the JSON serialization of the outer logs array plus either the completion value or failure diagnostic. A result at or below the cap is exact. A completion that cannot survive lossless JSON snapshotting fails as `invalid-output`; a value or combined logs/value outcome over the cap fails as `output-limit` rather than becoming inspected or truncated text. -Logs stream eagerly so a terminated run can retain output already admitted. When the cap is crossed, the runtime returns an explicit bounded failure with the fitting captured prefix. That outer result then traverses the ordinary `run_code` rendering and spill policy, which may save the captured text and expose its configured head/tail preview. The spill layer cannot recover bytes the runtime rejected beyond the hard cap. +Logs stream eagerly so a terminated run can retain output already admitted. Native stdout and stderr writes that bypass the worker's patched stream slots use independent pipes, so terminal settlement continues bounded capture until worker termination completes before materializing the result. When the cap is crossed, the runtime returns an explicit bounded failure with the fitting captured prefix. That outer result then traverses the ordinary `run_code` rendering and spill policy, which may save the captured text and expose its configured head/tail preview. The spill layer cannot recover bytes the runtime rejected beyond the hard cap. Compute time, wall time, worker heap, cancellation, and fresh-worker isolation remain independent limits. The outer ledger never charges intermediate bindings, so structured-clone cost and available process or worker memory are their practical bounds. diff --git a/.agents/notes/implemented/feature/2026-07-20-code-mode-typed-tool-returns.zh.md b/.agents/notes/implemented/feature/2026-07-20-code-mode-typed-tool-returns.zh.md index 0dbad69f6b..c567fa86bb 100644 --- a/.agents/notes/implemented/feature/2026-07-20-code-mode-typed-tool-returns.zh.md +++ b/.agents/notes/implemented/feature/2026-07-20-code-mode-typed-tool-returns.zh.md @@ -61,7 +61,7 @@ worker 暴露的是真正用于 `tools` 绑定失败的 `ToolCallError` 构造 `WorkerCodeRuntime` 以可配置的 `maxOutputBytes` 取代彼此独立的日志与值上限,默认值为 `67_108_864` 字节。宿主侧为不可信对端维护一份统一账本,计入外层日志数组以及完成值或失败诊断的 JSON 序列化大小。结果不超过上限时会保持精确。完成值无法通过无损 JSON 快照时,以 `invalid-output` 失败;值本身或日志与值的组合超过上限时,以 `output-limit` 失败,而不会变成检查格式化后或截断的文本。 -日志会在产生时立即流出,因此运行被终止时仍可保留已经纳入额度的输出。超过上限后,运行时会返回一个显式的有界失败,并携带可容纳的已捕获前缀。该外层结果随后通过普通的 `run_code` 渲染与输出落盘策略;策略可以保存已捕获的文本,并暴露其配置指定的头尾预览。输出落盘层无法恢复运行时在硬上限之外拒绝的字节。 +日志会在产生时立即流出,因此运行被终止时仍可保留已经纳入额度的输出。绕过 worker 中已改写流写入入口的原生 stdout 和 stderr 写入会经由彼此独立的管道传输,因此运行时在终态结算期间仍会继续在上限内捕获输出,直至 worker 完全终止,然后才组装结果。超过上限后,运行时会返回一个显式的有界失败,并携带可容纳的已捕获前缀。该外层结果随后通过普通的 `run_code` 渲染与输出落盘策略;策略可以保存已捕获的文本,并暴露其配置指定的头尾预览。输出落盘层无法恢复运行时在硬上限之外拒绝的字节。 计算时间、墙钟时间、worker 堆内存、取消和每次运行使用全新 worker 的隔离仍是互相独立的限制。外层账本从不计入中间绑定值,因此这些值实际受结构化克隆开销以及进程或 worker 可用内存限制。 diff --git a/packages/code-runtime/code-runtime-worker/README.md b/packages/code-runtime/code-runtime-worker/README.md index bfa19b1189..6ff75e9fe5 100644 --- a/packages/code-runtime/code-runtime-worker/README.md +++ b/packages/code-runtime/code-runtime-worker/README.md @@ -23,7 +23,7 @@ Every field is validated and defaulted; `maxOutputBytes` is a safe integer of at - **The port assumes a hostile peer** — model code can reach `parentPort` and forge traffic, so every inbound message is shape-validated and REBUILT before anything reads it (`null`, primitives, junk types, and malformed payloads drop without a throw; forged extra fields never ride along), the host answers each call id at most once, resolves binding names as OWN properties only (a forged `constructor` cannot walk a prototype chain), drops post-settlement replies, and validates every binding resolution and completion as lossless JSON. Forged `log`/`done` messages cannot bypass the outer cap: the host repeats validation and accounts every admitted log plus the completion or diagnostic. Worker-side namespaces are null-prototype with `defineProperty`, so `__proto__`-shaped binding names are ordinary keys. - **Two independent budgets, because the peer is hostile** — `computeMs` meters the worker's MEASURED busy time (`worker.performance.eventLoopUtilization()` polling): a hot loop cannot hide behind a pending decoy dispatch, and a program awaiting a slow tool accrues nothing. `maxWallMs` backstops what busy time cannot see (awaiting a promise nobody resolves). Both funnel into `worker.terminate()`, which ends hot synchronous loops too; heap overflow surfaces as the worker's OOM exit (`kind: 'worker-exit'`). - **Intermediate binding values are complete JSON** — binding arguments and resolutions cross by structured clone after lossless-JSON validation and have no byte cap. They never enter the outer-output ledger or model context; provider/executor acquisition bounds and process/worker memory remain the limits. -- **Logs stream eagerly into one outer ledger** — console/stdout/stderr text crosses the port in emission order, so a timed-out or killed program still shows what it printed. `maxOutputBytes` accounts the JSON serialization of the outer `logs` array plus the completion value or failure diagnostic. At or below the cap the exact value returns; a lossy completion is `invalid-output`, and a combined overflow is `output-limit` rather than a substituted inspected string. The failure retains the fitting captured prefix and later follows the normal outer `run_code` spill policy. +- **Logs stream eagerly into one outer ledger** — console/stdout/stderr text crosses the port in emission order, so a timed-out or killed program still shows what it printed. Native writes that bypass the patched stream slots arrive on pipes independent of the completion port; settlement therefore continues bounded pipe capture until worker termination completes before materializing the result. `maxOutputBytes` accounts the JSON serialization of the outer `logs` array plus the completion value or failure diagnostic. At or below the cap the exact value returns; a lossy completion is `invalid-output`, and a combined overflow is `output-limit` rather than a substituted inspected string. The failure retains the fitting captured prefix and later follows the normal outer `run_code` spill policy. - **Empty environment** — the worker gets `env: {}` and `execArgv: []`: no ambient credentials (stronger than the scrubbed-env rule for spawned commands) and no inherited loader flags. - **Dispose to quiescence** — teardown fails in-flight runs as `abort` and AWAITS each worker's exit before resolving. diff --git a/packages/code-runtime/code-runtime-worker/src/index.ts b/packages/code-runtime/code-runtime-worker/src/index.ts index 198bec694b..121b1a47d4 100644 --- a/packages/code-runtime/code-runtime-worker/src/index.ts +++ b/packages/code-runtime/code-runtime-worker/src/index.ts @@ -8,6 +8,7 @@ import { Worker } from 'node:worker_threads' import { stripTypeScriptTypes } from 'node:module' +import type { Readable } from 'node:stream' import { fileURLToPath } from 'node:url' import { Context } from 'cordis' import z from 'schemastery' @@ -106,6 +107,25 @@ function messageOf(error: unknown): string { return error instanceof Error ? error.message : String(error) } +/** Resolve after a worker pipe emits all queued data, or closes/errors during termination. */ +function waitForPipeDrain(stream: Readable): Promise { + if (stream.readableEnded || stream.destroyed) return Promise.resolve() + return new Promise((resolve) => { + const done = (): void => { + stream.off('end', done) + stream.off('close', done) + stream.off('error', done) + resolve() + } + stream.once('end', done) + stream.once('close', done) + stream.once('error', done) + // Close the event-registration race if termination finished between the + // initial state check and the listeners above. + if (stream.readableEnded || stream.destroyed) done() + }) +} + /** * Runtime shape gate for inbound port traffic. The peer runs MODEL CODE and * can post anything — `null`, primitives, objects with poisoned fields — so @@ -335,13 +355,20 @@ export class WorkerCodeRuntime extends CodeRuntime { const logs: string[] = [] const strayLogs: string[] = [] const output = new OutputLedger(this.config.maxOutputBytes) + let terminalOverride: CodeRunResult | undefined - // No settled guard: `finish` snapshots the arrays when it resolves, so - // a chunk flushing after settlement mutates only the discarded buffers, - // and the ledger bounds that growth until the pipes close. + // Pipe and message-port delivery are independent. Continue bounded pipe + // capture after a terminal message while worker termination drains bytes + // that were already queued; `finish` materializes the result only after + // termination completes. const captureStray = (chunk: Buffer): void => { + if (terminalOverride !== undefined) return const text = chunk.toString('utf8') - if (!settled && !output.admit(text, strayLogs)) finish(output.limit([...logs, ...strayLogs, text])) + if (!output.admit(text, strayLogs)) { + const limited = output.limit([...logs, ...strayLogs, text]) + terminalOverride = limited + finish(() => limited) + } } worker.stdout.on('data', captureStray) worker.stderr.on('data', captureStray) @@ -350,14 +377,20 @@ export class WorkerCodeRuntime extends CodeRuntime { // logs captured before timeout, abort, or failure remain in the result. let finishResolve!: () => void const finished = new Promise((done) => { finishResolve = done }) - const finish = (result: CodeRunResult): void => { + const finish = (finalize: () => CodeRunResult): void => { if (settled) return settled = true clearInterval(eluTimer) clearTimeout(wallTimer) request.signal?.removeEventListener('abort', onAbort) this.live.delete(live) - void worker.terminate().then(() => { + // Let the poll phase deliver pipe bytes already queued independently + // of the terminal port message before termination closes the streams. + void new Promise((resume) => { setImmediate(resume) }).then(async () => { + const stdoutDrained = waitForPipeDrain(worker.stdout) + const stderrDrained = waitForPipeDrain(worker.stderr) + await Promise.all([worker.terminate(), stdoutDrained, stderrDrained]) + const result = terminalOverride ?? finalize() finishResolve() resolve(result) }) @@ -365,22 +398,24 @@ export class WorkerCodeRuntime extends CodeRuntime { const onDone = (message: WorkerToHost): void => { if (message.type !== 'done') return - const captured = [...logs, ...strayLogs] if (message.error) { - finish(output.failure(captured, message.error)) + const error = message.error + finish(() => output.failure([...logs, ...strayLogs], error)) return } if (message.value === undefined) { - finish(output.success(captured)) + finish(() => output.success([...logs, ...strayLogs])) return } // The worker-thread boundary has already structured-cloned this // hostile value, so accessors and proxies cannot survive to throw // during the lossless-JSON snapshot. const value = snapshotJsonValue(message.value) as CodeJsonValue | undefined - finish(value === undefined - ? output.failure(captured, { kind: 'invalid-output', message: 'program completion must be lossless JSON' }) - : output.success(captured, value)) + if (value === undefined) { + finish(() => output.failure([...logs, ...strayLogs], { kind: 'invalid-output', message: 'program completion must be lossless JSON' })) + } else { + finish(() => output.success([...logs, ...strayLogs], value)) + } } const onCall = (message: WorkerToHost): void => { @@ -438,21 +473,25 @@ export class WorkerCodeRuntime extends CodeRuntime { const message = parseWorkerMessage(raw) if (!message) return if (message.type === 'log' && !settled && !output.admit(message.text, logs)) { - finish(output.limit([...logs, ...strayLogs, message.text])) + const limited = output.limit([...logs, ...strayLogs, message.text]) + terminalOverride = limited + finish(() => limited) return } if (message.type === 'output-limit' && !settled) { - finish(output.limit([...logs, ...strayLogs])) + const limited = output.limit([...logs, ...strayLogs]) + terminalOverride = limited + finish(() => limited) return } onCall(message) onDone(message) }) worker.on('error', (error: Error) => { - finish(output.failure([...logs, ...strayLogs], { kind: 'worker-exit', message: `worker error: ${error.message}` })) + finish(() => output.failure([...logs, ...strayLogs], { kind: 'worker-exit', message: `worker error: ${error.message}` })) }) worker.on('exit', (exitCode: number) => { - finish(output.failure([...logs, ...strayLogs], { kind: 'worker-exit', message: `worker exited with code ${exitCode} before completing` })) + finish(() => output.failure([...logs, ...strayLogs], { kind: 'worker-exit', message: `worker exited with code ${exitCode} before completing` })) }) // The compute budget reads the worker's own measured busy time, so a @@ -461,21 +500,21 @@ export class WorkerCodeRuntime extends CodeRuntime { const eluTimer = setInterval(() => { const elu = worker.performance.eventLoopUtilization() if (elu.active > this.config.computeMs) { - finish(output.failure([...logs, ...strayLogs], { kind: 'timeout', message: `compute budget exhausted (${this.config.computeMs}ms busy)` })) + finish(() => output.failure([...logs, ...strayLogs], { kind: 'timeout', message: `compute budget exhausted (${this.config.computeMs}ms busy)` })) } }, ELU_POLL_INTERVAL_MS) const wallTimer = setTimeout(() => { - finish(output.failure([...logs, ...strayLogs], { kind: 'timeout', message: `wall-clock ceiling reached (${this.config.maxWallMs}ms)` })) + finish(() => output.failure([...logs, ...strayLogs], { kind: 'timeout', message: `wall-clock ceiling reached (${this.config.maxWallMs}ms)` })) }, this.config.maxWallMs) const onAbort = (): void => { - finish(output.failure([...logs, ...strayLogs], { kind: 'abort', message: String(request.signal?.reason) })) + finish(() => output.failure([...logs, ...strayLogs], { kind: 'abort', message: String(request.signal?.reason) })) } request.signal?.addEventListener('abort', onAbort, { once: true }) const live: LiveRun = { worker, finished, - settle: (failure: CodeRunFailure) => { finish(output.failure([...logs, ...strayLogs], failure)) }, + settle: (failure: CodeRunFailure) => { finish(() => output.failure([...logs, ...strayLogs], failure)) }, } this.live.add(live) }) diff --git a/packages/code-runtime/code-runtime-worker/tests/runtime.spec.ts b/packages/code-runtime/code-runtime-worker/tests/runtime.spec.ts index c38ee7f492..4d73c859ce 100644 --- a/packages/code-runtime/code-runtime-worker/tests/runtime.spec.ts +++ b/packages/code-runtime/code-runtime-worker/tests/runtime.spec.ts @@ -333,6 +333,24 @@ describe('WorkerCodeRuntime — budgets and containment (real workers)', () => { expect(result.logs[1]?.length).toBeGreaterThan(0) expect('b'.repeat(100).startsWith(result.logs[1] ?? '')).toBe(true) }, 15_000) + + it('drains pipe output queued before terminal worker teardown completes', async () => { + const { runtime } = await setup({ maxOutputBytes: 200_000 }) + const payload = `late-pipe-${'x'.repeat(100_000)}` + const result = await runtime.run({ + program: ` + const { parentPort } = await import('node:worker_threads'); + const write = (text) => Object.getPrototypeOf(process.stdout).write.call(process.stdout, text); + write('late-pipe-' + 'x'.repeat(100_000)); + parentPort.postMessage({ type: 'done', value: 'done' }); + for (;;) {} + `, + bindings: [], + }) + expect(result.error).toBeUndefined() + expect(result.value).toBe('done') + expect(result.logs.join('') === payload).toBe(true) + }, 15_000) }) describe('WorkerCodeRuntime — hostile programs (real workers)', () => {