From 8185fdbfa9503263c41dce4d3ce6ea4a8a1931d2 Mon Sep 17 00:00:00 2001 From: can1357 Date: Sun, 28 Jun 2026 17:11:43 +0200 Subject: [PATCH] feat(coding-agent): queued tool argument updates to prevent edit loss - Implemented a queueing mechanism in `ToolExecutionComponent` to prevent starvation of edit previews during high-frequency argument updates. - Replaced eager cancellation of in-flight diff computations with a drain loop that ensures every update is processed once the current compute settles. - Added `partialJsonOf` helper to safely narrow streamed JSON buffers from tool arguments. - Added regression test to verify that slow diff computations are not aborted by incoming stream chunks and instead queue a subsequent re-run. --- .../src/modes/components/tool-execution.ts | 53 +++++++- .../tool-execution-preview-coalesce.test.ts | 122 ++++++++++++++++++ 2 files changed, 169 insertions(+), 6 deletions(-) create mode 100644 packages/coding-agent/test/tool-execution-preview-coalesce.test.ts diff --git a/packages/coding-agent/src/modes/components/tool-execution.ts b/packages/coding-agent/src/modes/components/tool-execution.ts index b2dee36a2..be7e85a0c 100644 --- a/packages/coding-agent/src/modes/components/tool-execution.ts +++ b/packages/coding-agent/src/modes/components/tool-execution.ts @@ -132,6 +132,14 @@ function rawTextInputFromPartialJson(partialJson: unknown): string | undefined { return partialJson; } +/** Read the streamed raw-JSON buffer a tool block stashes on its args, narrowed + * rather than cast: a missing or non-string `__partialJson` yields `undefined`. */ +function partialJsonOf(args: unknown): string | undefined { + if (args == null || typeof args !== "object" || !("__partialJson" in args)) return undefined; + const value = args.__partialJson; + return typeof value === "string" ? value : undefined; +} + function getArgsWithStreamedTextInput(args: unknown): unknown { if (args == null || typeof args !== "object") return args; const record = args as Record; @@ -242,6 +250,10 @@ export class ToolExecutionComponent extends Container implements NativeScrollbac #editDiffLastArgsKey?: string; // Latest in-flight streaming diff recompute, captured so it can be awaited. #editDiffInFlight?: Promise; + /** Set when newer args arrived while a preview compute was in flight; the + * drain loop re-runs once the current compute settles, so a slow diff + * coalesces streamed ticks instead of being aborted by each one. */ + #editDiffDirty = false; // Cached converted images for Kitty protocol (which requires PNG), keyed by index #convertedImages: Map = new Map(); // Spinner animation for partial task results @@ -323,7 +335,7 @@ export class ToolExecutionComponent extends Container implements NativeScrollbac this.setIgnoreTight(true); this.#updateDisplay(); - this.#editDiffInFlight = this.#runPreviewDiff(); + this.#schedulePreviewDiff(); } updateArgs(args: any, _toolCallId?: string): void { @@ -335,7 +347,7 @@ export class ToolExecutionComponent extends Container implements NativeScrollbac this.#args = args; this.#displayInputVersion++; this.#updateSpinnerAnimation(); - this.#editDiffInFlight = this.#runPreviewDiff(); + this.#schedulePreviewDiff(); this.#updateDisplay(); } @@ -346,7 +358,7 @@ export class ToolExecutionComponent extends Container implements NativeScrollbac setArgsComplete(_toolCallId?: string): void { this.#argsComplete = true; this.#updateSpinnerAnimation(); - this.#editDiffInFlight = this.#runPreviewDiff(); + this.#schedulePreviewDiff(); } /** @@ -360,7 +372,32 @@ export class ToolExecutionComponent extends Container implements NativeScrollbac await this.#editDiffInFlight; } - async #runPreviewDiff(): Promise { + /** + * Schedule a streaming diff preview recompute, coalescing bursts of + * `updateArgs` into one compute at a time: run the current compute to + * completion and re-run only after it settles when newer args arrived, never + * cancelling an in-flight compute on a fresh tick. The reveal controller pushes + * args ~30fps and a whole-file hashline/large-file diff can outlast a frame, so + * cancel-per-tick would starve every compute and no preview would land until + * args complete. Coalescing lets each diff land, so the preview tracks the + * stream at the rate the diffs can sustain. + */ + #schedulePreviewDiff(): void { + this.#editDiffDirty = true; + if (this.#editDiffInFlight) return; + this.#editDiffInFlight = this.#drainPreviewDiff().finally(() => { + this.#editDiffInFlight = undefined; + }); + } + + async #drainPreviewDiff(): Promise { + while (this.#editDiffDirty) { + this.#editDiffDirty = false; + await this.#computePreviewDiff(); + } + } + + async #computePreviewDiff(): Promise { const editMode = this.#editMode; if (!editMode) return; const strategy = EDIT_MODE_STRATEGIES[editMode]; @@ -370,7 +407,7 @@ export class ToolExecutionComponent extends Container implements NativeScrollbac if (args == null || typeof args !== "object") return; const previewArgs = getArgsWithStreamedTextInput(args); - const partialJson = (previewArgs as { __partialJson?: string }).__partialJson; + const partialJson = partialJsonOf(previewArgs); let effectiveArgs: unknown; try { effectiveArgs = strategy.extractCompleteEdits(previewArgs, partialJson); @@ -400,7 +437,8 @@ export class ToolExecutionComponent extends Container implements NativeScrollbac if (argsKey === this.#editDiffLastArgsKey) return; this.#editDiffLastArgsKey = argsKey; - this.#editDiffAbort?.abort(); + // Single-flight (the drain loop never overlaps computes), so this controller + // only ever cancels the live compute on teardown via `stopAnimation`. const controller = new AbortController(); this.#editDiffAbort = controller; @@ -732,6 +770,9 @@ export class ToolExecutionComponent extends Container implements NativeScrollbac this.#stopTodoStrikeAnimation(); this.#editDiffAbort?.abort(); this.#editDiffAbort = undefined; + // Drop any queued rerun so the drain loop exits instead of recomputing a + // preview for a torn-down block after its in-flight compute is aborted. + this.#editDiffDirty = false; } setExpanded(expanded: boolean): void { diff --git a/packages/coding-agent/test/tool-execution-preview-coalesce.test.ts b/packages/coding-agent/test/tool-execution-preview-coalesce.test.ts new file mode 100644 index 000000000..48ff6e13a --- /dev/null +++ b/packages/coding-agent/test/tool-execution-preview-coalesce.test.ts @@ -0,0 +1,122 @@ +import { afterEach, beforeEach, describe, expect, spyOn, test } from "bun:test"; +import * as fs from "node:fs/promises"; +import * as os from "node:os"; +import * as path from "node:path"; +import type { AgentTool } from "@oh-my-pi/pi-agent-core"; +import { EDIT_MODE_STRATEGIES, type PerFileDiffPreview } from "@oh-my-pi/pi-coding-agent/edit"; +import { ToolExecutionComponent } from "@oh-my-pi/pi-coding-agent/modes/components/tool-execution"; +import { initTheme } from "@oh-my-pi/pi-coding-agent/modes/theme/theme"; +import type { TUI } from "@oh-my-pi/pi-tui"; +import { removeWithRetries } from "@oh-my-pi/pi-utils"; + +// The reveal controller pushes streamed args at ~30fps; a whole-file diff can +// outlast a frame. The component must coalesce those ticks into one compute at a +// time — running the current compute to completion and re-running with the latest +// args once it settles — rather than aborting the in-flight compute on every +// tick, which starved the diff so no preview ever landed until args completed +// (the "blank edit box for the whole stream" regression). +describe("streaming edit preview coalescing", () => { + let tmpDir: string; + let file: string; + let themed = false; + let restore: (() => void) | undefined; + + beforeEach(async () => { + if (!themed) { + await initTheme(); + themed = true; + } + tmpDir = await fs.mkdtemp(path.join(os.tmpdir(), "preview-coalesce-")); + file = path.join(tmpDir, "mod.ts"); + await fs.writeFile(file, "const a = 1;\n"); + }); + + afterEach(async () => { + restore?.(); + restore = undefined; + await removeWithRetries(tmpDir); + }); + + // Read `edits[0].new_text` by narrowing rather than asserting an inline shape, + // so the captured args identity stays type-checked. + function firstNewText(args: unknown): unknown { + if (!args || typeof args !== "object" || !("edits" in args)) return undefined; + const edits = args.edits; + if (!Array.isArray(edits) || edits.length === 0) return undefined; + const first: unknown = edits[0]; + if (!first || typeof first !== "object" || !("new_text" in first)) return undefined; + return first.new_text; + } + + test("a slow compute is not aborted by a newer chunk; it lands, then re-runs with the latest args", async () => { + const deferreds: Array> = []; + const calls: Array<{ newText: unknown; signal: AbortSignal }> = []; + // One gate per compute invocation, resolved by the mock as each call + // starts, so the test awaits the real "compute N began" signal instead of a + // wall-clock delay. + const gates: Array> = []; + const gateFor = (index: number): PromiseWithResolvers => { + while (gates.length <= index) gates.push(Promise.withResolvers()); + return gates[index]!; + }; + const spy = spyOn(EDIT_MODE_STRATEGIES.replace, "computeDiffPreview").mockImplementation(async (args, ctx) => { + calls.push({ newText: firstNewText(args), signal: ctx.signal }); + const deferred = Promise.withResolvers(); + deferreds.push(deferred); + gateFor(calls.length - 1).resolve(); + return deferred.promise; + }); + restore = () => spy.mockRestore(); + + let renders = 0; + const ui = { + requestRender() { + renders++; + }, + } as unknown as TUI; + const tool = { mode: "replace" } as unknown as AgentTool; + + // Construction kicks off compute #0 for the first chunk; it stays in flight + // (mock returns an unresolved promise) so we can race a newer chunk against it. + const component = new ToolExecutionComponent( + "edit", + { path: file, edits: [{ old_text: "const a = 1;", new_text: "a" }] }, + {}, + tool, + ui, + tmpDir, + ); + try { + await gateFor(0).promise; + expect(calls.length).toBe(1); + expect(calls[0]!.newText).toBe("a"); + + // A newer chunk arrives mid-compute. Coalescing must NOT cancel #0 and + // must NOT launch a second concurrent compute — only mark a rerun pending. + component.updateArgs({ path: file, edits: [{ old_text: "const a = 1;", new_text: "ab" }] }); + expect(calls.length).toBe(1); + expect(calls[0]!.signal.aborted).toBe(false); + + // Resolving #0 lands its (now slightly stale) preview mid-stream — the + // behavior the starvation bug suppressed — and only then drives the rerun. + const rendersBeforeLanding = renders; + deferreds[0]!.resolve([{ path: file, diff: "@@ -1 +1 @@\n-const a = 1;\n+a", firstChangedLine: 1 }]); + + // Awaiting compute #1's start proves the rerun fired off the back of #0 + // settling; #0's landing (requestRender) runs synchronously before it. + await gateFor(1).promise; + expect(renders).toBeGreaterThan(rendersBeforeLanding); + expect(calls.length).toBe(2); + expect(calls[1]!.newText).toBe("ab"); + expect(calls[1]!.signal.aborted).toBe(false); + + // Settle the rerun so the drain loop exits cleanly. + deferreds[1]!.resolve([{ path: file, diff: "@@ -1 +1 @@\n-const a = 1;\n+ab", firstChangedLine: 1 }]); + await component.whenPreviewSettled(); + } finally { + // updateArgs starts the edit spinner interval; clear it so the timer + // never leaks into later tests. + component.stopAnimation(); + } + }); +});