From ec783437c8ae6b9c018f488273419b993c8dcb07 Mon Sep 17 00:00:00 2001 From: Kieran Klukas Date: Wed, 26 Aug 2026 00:08:17 -0400 Subject: [PATCH] feat: watch a tool call's arguments being written A tool call only exists once its whole argument object has been generated, so a long shell command or a written-out file is a silent stretch that reads as a hang. The provider streams the argument JSON as it comes; the actor batches it into tool-input-delta events like any other delta, and the client repairs the partial JSON on each one to keep the step's own label current, showing the text in full once it grows past a line. --- src/actor.ts | 46 ++++++++++++++ src/client/app.css | 19 ++++++ src/client/app.js | 131 +++++++++++++++++++++++++++++++++++++++- src/events.ts | 23 +++++++ src/inference.ts | 17 ++++++ tests/actor.test.ts | 78 ++++++++++++++++++++++++ tests/inference.test.ts | 88 +++++++++++++++++++++++++++ 7 files changed, 401 insertions(+), 1 deletion(-) diff --git a/src/actor.ts b/src/actor.ts index d237cbe..94a6c03 100644 --- a/src/actor.ts +++ b/src/actor.ts @@ -142,6 +142,17 @@ export interface ReasoningSignatureStep { kind: "reasoning-signature"; providerOptions: Record>; } +/** + * A fragment of a tool call's arguments, as the model writes them — raw JSON, + * not a JSON document. See ToolInputDeltaData: this is what lets a client raise + * the step while a long command is still being typed. + */ +export interface ToolInputStep { + kind: "tool-input"; + toolCallId: string; + toolName: string; + delta: string; +} /** The model invoked a tool. */ export interface ToolCallStep { kind: "tool-call"; @@ -166,6 +177,7 @@ export type RunStep = | TextStep | ReasoningStep | ReasoningSignatureStep + | ToolInputStep | ToolCallStep | ToolResultStep | UsageStep; @@ -641,6 +653,32 @@ export class ConversationActor { flush(Event.ReasoningDelta, reasoning); reasoning = ""; }; + // Arguments-in-progress, batched like the other two streams but keyed by the + // call they belong to: the tail has to be flushed before a delta for a + // DIFFERENT call joins it, or two calls' JSON would be spliced together. + let input: { toolCallId: string; toolName: string; text: string } | null = null; + let ibatch = 0; + const flushInput = () => { + if (!input || input.text.length === 0) { + input = null; + ibatch = 0; + return; + } + const data = { + runId, + threadId: this.conversationId, + messageId, + toolCallId: input.toolCallId, + toolName: input.toolName, + delta: truncateUtf8(input.text), + }; + const seq = this.store.appendAndBump(this.conversationId, Event.ToolInput, data); + this.emit({ id: makeId(this.conversationId, seq), event: Event.ToolInput, data }); + this.lastActivity = Date.now(); + onProgress?.(seq); + input = null; + ibatch = 0; + }; const tick = setInterval(() => { if (this.store.isConversationDeleted(this.conversationId)) { @@ -650,6 +688,7 @@ export class ConversationActor { } flushReasoning(); flushDelta(); + flushInput(); }, BATCH_FLUSH_MS); // The HTTP server and the worker can be separate processes. Poll the // durable request while a provider is quiet so a cancel reaches whichever @@ -707,6 +746,11 @@ export class ConversationActor { messageId, providerOptions: step.providerOptions, }); + } else if (step.kind === "tool-input") { + if (input && input.toolCallId !== step.toolCallId) flushInput(); + input ??= { toolCallId: step.toolCallId, toolName: step.toolName, text: "" }; + input.text += step.delta; + if (++ibatch >= BATCH_MAX_DELTAS) flushInput(); } else if (step.kind === "tool-call") { // Flush pending text/reasoning first so the durable log stays ordered // (a tool call sits after the text that preceded it). Tool events are @@ -718,6 +762,7 @@ export class ConversationActor { // `persist` calls direct — batching them would reopen that window. flushReasoning(); flushDelta(); + flushInput(); // the arguments finished being typed; the call follows them this.persist(Event.ToolCall, { runId, threadId: this.conversationId, @@ -782,6 +827,7 @@ export class ConversationActor { } flushReasoning(); flushDelta(); + flushInput(); const finish = errored ? "error" : this.isCancelled() ? "aborted" : "stop"; if (finish === "aborted") { this.persist(Event.Cancelled, { runId, threadId: this.conversationId }); diff --git a/src/client/app.css b/src/client/app.css index 7edc66d..8f0e280 100644 --- a/src/client/app.css +++ b/src/client/app.css @@ -1705,6 +1705,25 @@ body.selecting #selectBtn { } /* Artifacts collect at the bottom of the message that produced them. */ +/* A tool call's arguments while they're still being written: the terminal + vocabulary (mono, sunk, tail-following), capped so a 400-line file being + written doesn't push the conversation off screen. Replaced by the real + result once the call lands. */ +.tprev { + max-height: 190px; + overflow: auto; + margin: 4px 0 2px; + padding: 8px 10px; + border-radius: var(--radius-sm); + background: var(--bg-sunk); + font-family: var(--mono); + font-size: 12px; + line-height: 1.5; + color: var(--ink-muted); + white-space: pre-wrap; + word-break: break-word; +} + .artifacts { display: flex; flex-direction: column; diff --git a/src/client/app.js b/src/client/app.js index 56111f6..87cdeca 100644 --- a/src/client/app.js +++ b/src/client/app.js @@ -1276,7 +1276,14 @@ import { mountSidebar } from "./sidebar.js"; // tool name. Keyed by toolCallId so the result attaches. Shimmers while running. function toolStep(rec, data) { rec.toolSteps = rec.toolSteps || Object.create(null); - if (rec.toolSteps[data.toolCallId]) return rec.toolSteps[data.toolCallId]; + if (rec.toolSteps[data.toolCallId]) { + // The step may already exist because the ARGUMENTS raised it (see + // toolInputDelta) with only as much as had been typed. The call event is + // the authority on what was actually called, so it settles the step. + var had = rec.toolSteps[data.toolCallId]; + if (data.input !== undefined) settleToolInput(had, data.input); + return had; + } var a = openActivity(rec); blockEndReasoning(a); // the thinking that led to this call is done var ui = toolUI(data.toolName); @@ -1308,6 +1315,7 @@ import { mountSidebar } from "./sidebar.js"; row: step.row, label: step.label, body: step.body, + host: step.host, // a custom summary is rebuilt against it when the input settles toolName: data.toolName, input: data.input, block: a, @@ -1324,6 +1332,123 @@ import { mountSidebar } from "./sidebar.js"; autoScroll(); return t; } + /* + * A tool call's arguments, watched as the model writes them. + * + * `tool-call` only exists once the whole argument object has been generated, + * so a 60-line shell heredoc or a written-out file is many seconds in which + * nothing moves and the app looks hung. These deltas are raw JSON fragments — + * only their accumulation parses — so the buffer is repaired and re-parsed on + * every one, and whatever survives feeds the SAME per-tool row renderer the + * finished call uses. Nothing here is authoritative: the call event overwrites + * all of it. + */ + function toolInputDelta(rec, data) { + var t = + (rec.toolSteps && rec.toolSteps[data.toolCallId]) || + toolStep(rec, { toolCallId: data.toolCallId, toolName: data.toolName, input: {} }); + t.raw = (t.raw || "") + data.delta; + // Some fragments land mid-key or mid-escape and can't be read at all. The + // last one that could is still true, and holding it steady beats flashing + // raw JSON for one frame. + t.partial = parsePartialJson(t.raw) || t.partial; + var ui = toolUI(t.toolName); + if (t.partial && !ui.summaryDom) { + t.input = t.partial; + if (t.entry) t.entry.input = t.partial; + t.label.textContent = ui.row(t.partial, t.toolName); + blockUpdateHead(t.block); + } + // The row label carries a short command on its own. A long one is exactly + // the case that reads as a hang, so it gets shown in full, being typed. + var text = longestString(t.partial) || t.raw; + if (text.length > 120 || text.indexOf("\n") >= 0) { + if (!t.preview) { + t.preview = document.createElement("pre"); + t.preview.className = "tprev"; + t.body.appendChild(t.preview); + if (t.row.tagName === "DETAILS" && !t.row.open) { + t.row.open = true; + t.autoOpened = true; // ours to close again when the arguments settle + } + } + t.preview.textContent = text; + // Follow the tail, like a terminal. Not while a whole history is being + // replayed at once: reading scrollHeight forces a layout, and there is + // nobody watching a preview that is about to be settled anyway. + if (!bulkLoading) t.preview.scrollTop = t.preview.scrollHeight; + } + autoScroll(); + return t; + } + /* The finished arguments land: the live preview has served its purpose. */ + function settleToolInput(t, input) { + t.raw = null; + t.partial = null; + t.input = input; + if (t.entry) t.entry.input = input; + var ui = toolUI(t.toolName); + // A tool that owns its summary (fetch_url's favicon + hostname) built it + // from whatever had been typed at the time; now there's the real thing. + if (ui.summaryDom) t.label = ui.summaryDom(t.row.querySelector("summary"), input, t.host).label; + else t.label.textContent = ui.row(input, t.toolName); + if (t.preview) { + t.preview.remove(); + t.preview = null; + } + if (t.autoOpened) { + t.row.open = false; // back to the resting state; the result renders inside + t.autoOpened = false; + } + blockUpdateHead(t.block); + } + /* + * The longest string value in an object — for a shell call that's the command, + * for a file write it's the contents. Which is to say: the part worth watching. + */ + function longestString(obj) { + if (!obj || typeof obj !== "object") return ""; + var best = ""; + for (var k in obj) if (typeof obj[k] === "string" && obj[k].length > best.length) best = obj[k]; + return best; + } + /* + * Parse a JSON document that isn't finished yet: close whatever is still open + * and try again. `{"command":"ls -la` becomes `{"command":"ls -la"}`, which is + * a truthful reading of what has arrived so far. Returns null while the + * fragment can't be made sense of (a half-written key, say) — the next delta + * is a few milliseconds away. + */ + function parsePartialJson(s) { + try { + return JSON.parse(s); + } catch (_) {} + var close = [], + inStr = false, + esc = false; + for (var i = 0; i < s.length; i++) { + var c = s.charAt(i); + if (inStr) { + if (esc) esc = false; + else if (c === "\\") esc = true; + else if (c === '"') inStr = false; + continue; + } + if (c === '"') inStr = true; + else if (c === "{") close.push("}"); + else if (c === "[") close.push("]"); + else if (c === "}" || c === "]") close.pop(); + } + var t = esc ? s.slice(0, -1) : s; // a dangling escape has nothing to escape + if (inStr) t += '"'; + t = t.replace(/[,:]\s*$/, ""); // a separator with nothing after it yet + for (var j = close.length - 1; j >= 0; j--) t += close[j]; + try { + return JSON.parse(t); + } catch (_) { + return null; + } + } // ---- live tool progress (deep_research) -------------------------------- // A tool that runs for minutes reports as it goes (see ToolProgressData). The // panel is built once per tool call and mutated in place by each phase, so a @@ -2985,6 +3110,9 @@ import { mountSidebar } from "./sidebar.js"; scheduleFlush(); break; } + case "tool-input-delta": + toolInputDelta(assistantTurn(data.messageId), data); + break; case "tool-call": toolStep(assistantTurn(data.messageId), data); // The question tool's whole UI is the composer form, and the call is @@ -3072,6 +3200,7 @@ import { mountSidebar } from "./sidebar.js"; "run-started", "message-start", "reasoning-delta", + "tool-input-delta", "tool-call", "tool-progress", "tool-result", diff --git a/src/events.ts b/src/events.ts index e5359f3..9e08ead 100644 --- a/src/events.ts +++ b/src/events.ts @@ -13,6 +13,7 @@ export const Event = { TextDelta: "text-delta", ReasoningDelta: "reasoning-delta", ReasoningSig: "reasoning-signature", + ToolInput: "tool-input-delta", ToolCall: "tool-call", ToolProgress: "tool-progress", ToolResult: "tool-result", @@ -122,6 +123,27 @@ export interface ReasoningSignatureData { providerOptions: Record>; } +/** + * A chunk of a tool call's arguments, as the model writes them. + * + * The call itself (`tool-call`) only exists once the whole argument object has + * been generated, which for a long shell command or a written-out file is many + * seconds of nothing at all — indistinguishable, from the outside, from a hang. + * These carry the raw JSON as it comes off the provider so the client can raise + * the step immediately and show the arguments being typed. `delta` is a + * fragment of a JSON document, not a JSON document: only the accumulation of a + * call's deltas parses, and the `tool-call` event remains the authority on what + * was actually called. + */ +export interface ToolInputDeltaData { + threadId: string; + runId: string; + messageId: string; + toolCallId: string; + toolName: string; + delta: string; +} + export interface ToolCallData { threadId: string; runId: string; @@ -275,6 +297,7 @@ export type EventData = | TextDeltaData | ReasoningDeltaData | ReasoningSignatureData + | ToolInputDeltaData | ToolCallData | ToolProgressData | ToolResultData diff --git a/src/inference.ts b/src/inference.ts index d444190..4cfc852 100644 --- a/src/inference.ts +++ b/src/inference.ts @@ -570,6 +570,10 @@ export async function* run(messages: ModelMessage[], opts: RunOptions): AsyncGen // could possibly know the answer. It goes nowhere: the turn ended at the // question. let asked = false; + // Which tool each in-flight argument stream belongs to. The provider names the + // tool when the arguments START and never again, so the name has to be carried + // forward to tag the deltas that follow. + const inputNames = new Map(); for await (const part of result.fullStream) { if (part.type === "text-delta") { if (asked) continue; @@ -591,7 +595,20 @@ export async function* run(messages: ModelMessage[], opts: RunOptions): AsyncGen (part.providerMetadata as Record>) ?? reasoningMeta.get(part.id); if (meta && hasSignature(meta)) yield { kind: "reasoning-signature", providerOptions: meta }; + } else if (part.type === "tool-input-start") { + inputNames.set(part.id, part.toolName); + } else if (part.type === "tool-input-delta") { + // The arguments as they're written. Not counted toward `grownChars` — the + // tool-call below counts the finished object, and these are the same bytes + // arriving early. + yield { + kind: "tool-input", + toolCallId: part.id, + toolName: inputNames.get(part.id) ?? "", + delta: part.delta, + }; } else if (part.type === "tool-call") { + inputNames.delete(part.toolCallId); if (part.toolName === ASK_TOOL) asked = true; grownChars += JSON.stringify(part.input ?? "").length; yield { diff --git a/tests/actor.test.ts b/tests/actor.test.ts index 017c545..6480b22 100644 --- a/tests/actor.test.ts +++ b/tests/actor.test.ts @@ -482,6 +482,84 @@ test("tool-call and tool-result steps persist as durable, ordered events", async expect(res.output.results.length).toBe(1); }); +/** Just enough of a message part to name it. */ +interface PartType { + type: string; +} + +test("streamed tool arguments persist as batched deltas, ahead of the call", async () => { + const a = new ConversationActor("t-input", store); + const events: WireEvent[] = []; + a.follow({ push: (e) => events.push(e), closed: false }); + const raw = '{"command":"echo hello"}'; + await a.runText("ri", "mi", async function* (_signal) { + for (const ch of raw) { + yield { kind: "tool-input", toolCallId: "c1", toolName: "run_shell", delta: ch }; + } + yield { + kind: "tool-call", + toolCallId: "c1", + toolName: "run_shell", + input: { command: "echo hello" }, + }; + }); + const inputs = events.filter((e) => e.event === Event.ToolInput); + expect(inputs.length).toBeGreaterThan(0); + // Batched, not one event per character. + expect(inputs.length).toBeLessThan(raw.length); + // The fragments reassemble into exactly what the model wrote. + expect(inputs.map((e) => (e.data as { delta: string }).delta).join("")).toBe(raw); + const first = inputs[0]!.data as { toolCallId: string; toolName: string; messageId: string }; + expect(first.toolCallId).toBe("c1"); + expect(first.toolName).toBe("run_shell"); + expect(first.messageId).toBe("mi"); + // Everything typed lands before the call it belongs to. + const names = events.map((e) => e.event); + expect(names.lastIndexOf(Event.ToolInput)).toBeLessThan(names.indexOf(Event.ToolCall)); +}); + +test("streamed arguments for two calls never splice together", async () => { + const a = new ConversationActor("t-input2", store); + const events: WireEvent[] = []; + a.follow({ push: (e) => events.push(e), closed: false }); + await a.runText("ri2", "mi2", async function* (_signal) { + yield { kind: "tool-input", toolCallId: "c1", toolName: "run_shell", delta: '{"a":1}' }; + yield { kind: "tool-input", toolCallId: "c2", toolName: "view_file", delta: '{"b":2}' }; + }); + const byCall: Record = {}; + for (const e of events.filter((x) => x.event === Event.ToolInput)) { + const d = e.data as { toolCallId: string; delta: string }; + byCall[d.toolCallId] = (byCall[d.toolCallId] ?? "") + d.delta; + } + expect(byCall.c1).toBe('{"a":1}'); + expect(byCall.c2).toBe('{"b":2}'); +}); + +test("streamed tool arguments are not replayed to the model as history", async () => { + const a = new ConversationActor("t-input3", store); + a.appendUser("run it", "r0"); + await a.runText("ri3", "mi3", async function* (_signal) { + yield { + kind: "tool-input", + toolCallId: "c1", + toolName: "run_shell", + delta: '{"command":"ls"}', + }; + yield { kind: "tool-call", toolCallId: "c1", toolName: "run_shell", input: { command: "ls" } }; + yield { kind: "tool-result", toolCallId: "c1", toolName: "run_shell", output: "a.txt" }; + }); + // The call carries the arguments; the fragments that built it are UI, and + // replaying them would show the model its own typing a second time — as prose. + const history = await a.history(); + const kinds: string[] = []; + for (const m of history) { + const content = m.content as unknown; + if (Array.isArray(content)) for (const part of content) kinds.push((part as PartType).type); + } + expect(kinds).toContain("tool-call"); + expect(kinds).not.toContain("text"); +}); + test("history folds a tool turn into paired assistant/tool messages", async () => { const a = new ConversationActor("t-tool-hist", store); a.appendUser("what's new"); diff --git a/tests/inference.test.ts b/tests/inference.test.ts index ff91224..4e2cbfd 100644 --- a/tests/inference.test.ts +++ b/tests/inference.test.ts @@ -1,10 +1,12 @@ import { expect, test } from "bun:test"; import type { ModelMessage } from "ai"; +import type { RunStep, ToolInputStep } from "../src/actor"; import { Catalog } from "../src/catalog"; import { dropForeignReasoning, effortFor, promptChars, + run, setRegistry, usageFor, } from "../src/inference"; @@ -174,3 +176,89 @@ test("an assistant turn that was only a foreign thinking block is dropped whole" // An empty assistant message is itself a 400, so the turn goes with it. expect(dropForeignReasoning(messages, "openai")).toEqual([messages[0]!]); }); + +/** + * Arguments, as they are written. + * + * A tool call only exists once its whole argument object has been generated, so + * a long one is a silent stretch that reads as a hang. The adapter streams the + * fragments; this is the check that they survive the trip out of `run()` and + * still add up to what the model wrote. + */ +test("a tool call's arguments stream as fragments before the call itself", async () => { + const args = '{"questions":[{"type":"free_text","question":"which one?"}]}'; + const sse = (o: unknown) => `data: ${JSON.stringify(o)}\n\n`; + // The arguments arrive in three pieces, the way a provider actually sends them. + const pieces = [args.slice(0, 12), args.slice(12, 40), args.slice(40)]; + const endpoint = Bun.serve({ + port: 0, + fetch: () => + new Response( + sse({ + id: "1", + choices: [ + { + index: 0, + delta: { + role: "assistant", + tool_calls: [ + { index: 0, id: "call_1", type: "function", function: { name: "ask_user" } }, + ], + }, + }, + ], + }) + + pieces + .map((p) => + sse({ + id: "1", + choices: [ + { + index: 0, + delta: { tool_calls: [{ index: 0, function: { arguments: p } }] }, + }, + ], + }), + ) + .join("") + + sse({ id: "1", choices: [{ index: 0, delta: {}, finish_reason: "tool_calls" }] }) + + "data: [DONE]\n\n", + { headers: { "content-type": "text/event-stream" } }, + ), + }); + try { + setRegistry( + new ProviderRegistry( + Catalog.fromRaw([ + { + id: "acme", + name: "Acme", + type: "openai-compat", + api_endpoint: `${endpoint.url.origin}/v1`, + models: [{ id: "cheap", name: "Cheap", context_window: 8000 }], + }, + ]), + { config: { providers: [{ id: "acme", apiKey: "sk-test" }] } }, + ), + ); + const steps: RunStep[] = []; + for await (const s of run([{ role: "user", content: "hi" }], { + model: "acme/cheap", + runId: "r1", + canAsk: true, // so ask_user is offered at all + })) { + steps.push(s); + } + const typed = steps.filter((s): s is ToolInputStep => s.kind === "tool-input"); + expect(typed.length).toBeGreaterThan(0); + expect(typed.map((s) => s.delta).join("")).toBe(args); + // Every fragment knows which call it belongs to, and to which tool. + expect(typed.every((s) => s.toolCallId === "call_1")).toBe(true); + expect(typed.every((s) => s.toolName === "ask_user")).toBe(true); + // And they all land before the call they were building. + const kinds = steps.map((s) => s.kind); + expect(kinds.lastIndexOf("tool-input")).toBeLessThan(kinds.indexOf("tool-call")); + } finally { + endpoint.stop(true); + } +}); -- 2.51.2