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); + } +});