From 71aaad50e6a5a0300415c07b1913cc34eb9af94a Mon Sep 17 00:00:00 2001 From: Wolfgang Schoenberger <221313372+wolfiesch@users.noreply.github.com> Date: Fri, 10 Jul 2026 18:43:31 -0700 Subject: [PATCH] fix(coding-agent): flush throttled output tails --- packages/coding-agent/CHANGELOG.md | 4 ++ .../src/session/streaming-output.ts | 52 ++++++++++++++----- .../test/streaming-output.test.ts | 45 +++++++++++++++- 3 files changed, 88 insertions(+), 13 deletions(-) diff --git a/packages/coding-agent/CHANGELOG.md b/packages/coding-agent/CHANGELOG.md index d8041c95c..f5e32d613 100644 --- a/packages/coding-agent/CHANGELOG.md +++ b/packages/coding-agent/CHANGELOG.md @@ -2,6 +2,10 @@ ## [Unreleased] +### Fixed + +- Fixed throttled live command output holding a quiet final chunk until the command exited. + ## [16.4.2] - 2026-07-10 ### Fixed diff --git a/packages/coding-agent/src/session/streaming-output.ts b/packages/coding-agent/src/session/streaming-output.ts index 5eafeb528..449b57f6b 100644 --- a/packages/coding-agent/src/session/streaming-output.ts +++ b/packages/coding-agent/src/session/streaming-output.ts @@ -735,6 +735,7 @@ export class OutputSink { #truncated = false; #lastChunkTime = 0; #pendingChunk = ""; + #pendingChunkTimer: Timer | undefined; // Per-line column cap streaming state (persists across `push` calls so a // long line split across chunks still trips the same trigger). @@ -806,20 +807,18 @@ export class OutputSink { push(chunk: string): void { chunk = sanitizeWithOptionalSixelPassthrough(chunk, sanitizeText); - // Throttled onChunk: coalesce chunks arriving inside the throttle window - // and flush the buffered concatenation on the next eligible tick (plus a - // final flush in dump()) so the preview never has silent gaps. + // Throttled onChunk: coalesce chunks arriving inside the throttle window. + // A timer flushes quiet tails at the throttle boundary; dump() catches a + // final pending chunk when the process exits before that timer fires. // Live preview gets the raw (pre-cap) chunk so the TUI never lags behind // what reached the sink — the column cap is for the persisted LLM view. if (this.#onChunk) { const now = Date.now(); if (now - this.#lastChunkTime >= this.#chunkThrottleMs) { - this.#lastChunkTime = now; - const merged = this.#pendingChunk + chunk; - this.#pendingChunk = ""; - this.#onChunk(merged); + this.#emitPendingChunkWith(chunk, now); } else { this.#pendingChunk += chunk; + this.#schedulePendingChunkFlush(); } } @@ -1121,6 +1120,7 @@ export class OutputSink { * branch in `dump()` against stale totals. */ replace(text: string): void { + this.#clearPendingChunkTimer(); this.#buffer = text; this.#bufferBytes = Buffer.byteLength(text, "utf-8"); this.#head = ""; @@ -1138,6 +1138,38 @@ export class OutputSink { this.#pendingChunk = ""; } + #clearPendingChunkTimer(): void { + if (!this.#pendingChunkTimer) return; + clearTimeout(this.#pendingChunkTimer); + this.#pendingChunkTimer = undefined; + } + + #emitPendingChunkWith(chunk: string, now: number): void { + this.#clearPendingChunkTimer(); + this.#lastChunkTime = now; + const merged = this.#pendingChunk + chunk; + this.#pendingChunk = ""; + this.#onChunk?.(merged); + } + + #flushPendingChunk(): void { + if (this.#pendingChunk.length === 0) { + this.#clearPendingChunkTimer(); + return; + } + this.#emitPendingChunkWith("", Date.now()); + } + + #schedulePendingChunkFlush(): void { + if (this.#chunkThrottleMs <= 0 || this.#pendingChunkTimer) return; + const elapsed = Date.now() - this.#lastChunkTime; + const delay = Math.max(0, this.#chunkThrottleMs - elapsed); + this.#pendingChunkTimer = setTimeout(() => { + this.#pendingChunkTimer = undefined; + this.#flushPendingChunk(); + }, delay); + } + /** * Replay the rolling tail ring back into the artifact sink. When bytes * were actually dropped from the middle (the head budget was exhausted @@ -1179,11 +1211,7 @@ export class OutputSink { // Flush any chunk still held back by the throttle so the live preview // ends with the complete stream. - if (this.#onChunk && this.#pendingChunk.length > 0) { - const pending = this.#pendingChunk; - this.#pendingChunk = ""; - this.#onChunk(pending); - } + this.#flushPendingChunk(); const totalLines = this.#sawData ? this.#totalLines + 1 : 0; if (this.#file) { diff --git a/packages/coding-agent/test/streaming-output.test.ts b/packages/coding-agent/test/streaming-output.test.ts index f18d05fa9..4b5ec3bff 100644 --- a/packages/coding-agent/test/streaming-output.test.ts +++ b/packages/coding-agent/test/streaming-output.test.ts @@ -1,4 +1,4 @@ -import { afterEach, describe, expect, test } from "bun:test"; +import { afterEach, describe, expect, test, vi } from "bun:test"; import * as fs from "node:fs/promises"; import * as os from "node:os"; import * as path from "node:path"; @@ -32,6 +32,7 @@ function byteLength(text: string): number { } afterEach(async () => { + vi.useRealTimers(); for (const dir of createdTempDirs.splice(0)) { await removeWithRetries(dir); } @@ -317,6 +318,48 @@ describe("OutputSink", () => { expect(dumped.output).toBe("abc"); }); + test("throttled onChunk emits a quiet tail at the throttle boundary", () => { + vi.useFakeTimers(); + const chunks: string[] = []; + const sink = new OutputSink({ onChunk: chunk => chunks.push(chunk), chunkThrottleMs: 20 }); + + sink.push("a"); + sink.push("b"); + expect(chunks).toEqual(["a"]); + + vi.advanceTimersByTime(20); + + expect(chunks).toEqual(["a", "b"]); + }); + + test("dump flushes a throttled tail once and cancels its timer", async () => { + vi.useFakeTimers(); + const chunks: string[] = []; + const sink = new OutputSink({ onChunk: chunk => chunks.push(chunk), chunkThrottleMs: 20 }); + + sink.push("a"); + sink.push("b"); + expect((await sink.dump()).output).toBe("ab"); + expect(chunks).toEqual(["a", "b"]); + + vi.advanceTimersByTime(20); + + expect(chunks).toEqual(["a", "b"]); + }); + + test("replace cancels a throttled tail and discards its pending preview", () => { + vi.useFakeTimers(); + const chunks: string[] = []; + const sink = new OutputSink({ onChunk: chunk => chunks.push(chunk), chunkThrottleMs: 20 }); + + sink.push("a"); + sink.push("superseded"); + sink.replace("replacement"); + vi.advanceTimersByTime(20); + + expect(chunks).toEqual(["a"]); + }); + test("caps artifact-on-disk size: head + notice + tail when stream exceeds cap", async () => { const dir = await createTempDir(); const artifactPath = path.join(dir, "capped.log");