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.
This commit is contained in:
@@ -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<string, unknown>;
|
||||
@@ -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<void>;
|
||||
/** 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<number, { data: string; mimeType: string }> = 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<void> {
|
||||
/**
|
||||
* 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<void> {
|
||||
while (this.#editDiffDirty) {
|
||||
this.#editDiffDirty = false;
|
||||
await this.#computePreviewDiff();
|
||||
}
|
||||
}
|
||||
|
||||
async #computePreviewDiff(): Promise<void> {
|
||||
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 {
|
||||
|
||||
@@ -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<PromiseWithResolvers<PerFileDiffPreview[] | null>> = [];
|
||||
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<PromiseWithResolvers<void>> = [];
|
||||
const gateFor = (index: number): PromiseWithResolvers<void> => {
|
||||
while (gates.length <= index) gates.push(Promise.withResolvers<void>());
|
||||
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<PerFileDiffPreview[] | null>();
|
||||
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();
|
||||
}
|
||||
});
|
||||
});
|
||||
Reference in New Issue
Block a user