From 5dab80f7643f50764222d88fcc6e9d562794cfe3 Mon Sep 17 00:00:00 2001 From: Pawan Osman Date: Wed, 15 Jul 2026 16:45:46 +0300 Subject: [PATCH] Implement coalescing for high-frequency UI events and enhance subagent budget management - Introduced a coalescing mechanism for UI events in the sidebar and agent loop to prevent performance issues caused by high-frequency updates. - Added budget management for foreground and background subagents, ensuring timely execution and preventing timeouts. - Enhanced the `runAgent` function to handle budget hits and provide appropriate feedback on subagent execution status. - Updated the UI to batch process events, improving responsiveness and reducing the likelihood of UI freezes during agent operations. --- src/agent/loop.ts | 193 ++++++++++++++++++++++++++++++++----- src/agent/tools/shared.ts | 10 +- src/ui/sidebarProvider.ts | 93 ++++++++++++++++++ webview-ui/sidebar/App.tsx | 25 ++++- 4 files changed, 295 insertions(+), 26 deletions(-) diff --git a/src/agent/loop.ts b/src/agent/loop.ts index dff86fa..6702e16 100644 --- a/src/agent/loop.ts +++ b/src/agent/loop.ts @@ -30,9 +30,98 @@ const MULTITASK_REMINDER = "(run_in_background=true), launching multiple subagents AT THE SAME TIME in a single turn.\n"; const MAX_STEPS = 50; +/** Foreground subagent hard budget (abort + settle). */ +const SUBAGENT_MAX_MS = 6 * 60_000; +/** Background subagent hard budget. */ +const BG_SUBAGENT_MAX_MS = 10 * 60_000; +/** Coalesce high-frequency stream UI events (ms). */ +const STREAM_COALESCE_MS = 40; + +/** + * Batch text/thinking/tool-args deltas so streaming cannot flood the host + * reducer + webview (main cause of UI freezes that look like "stuck" tools). + * Terminal events flush pending deltas first to preserve order. + */ +function coalesceEmit(raw: (e: AgentEvent) => void): (e: AgentEvent) => void { + const pending = new Map(); + let timer: ReturnType | undefined; + const flush = () => { + timer = undefined; + if (!pending.size) return; + const batch = [...pending.values()]; + pending.clear(); + for (const e of batch) { + try { raw(e); } catch { /* ignore */ } + } + }; + const schedule = () => { + if (!timer) timer = setTimeout(flush, STREAM_COALESCE_MS); + }; + return (event: AgentEvent) => { + if (event.type === "text-delta") { + const prev = pending.get("text"); + if (prev && prev.type === "text-delta") { + pending.set("text", { type: "text-delta", text: prev.text + event.text }); + } else { + pending.set("text", event); + } + schedule(); + return; + } + if (event.type === "thinking-delta") { + const prev = pending.get("think"); + if (prev && prev.type === "thinking-delta") { + pending.set("think", { type: "thinking-delta", text: prev.text + event.text }); + } else { + pending.set("think", event); + } + schedule(); + return; + } + if (event.type === "tool-call-args") { + // Latest full argsText wins (provider sends cumulative chunks). + pending.set(`args:${event.callId}`, event); + schedule(); + return; + } + if (event.type === "subagent-event") { + const child = event.event; + // Coalesce nested high-freq child stream events per parent call. + if (child.type === "text-delta" || child.type === "thinking-delta" || child.type === "tool-call-args") { + const key = + child.type === "tool-call-args" + ? `sub:${event.callId}:args:${child.callId}` + : `sub:${event.callId}:${child.type}`; + if (child.type === "text-delta" || child.type === "thinking-delta") { + const prev = pending.get(key); + if (prev && prev.type === "subagent-event" && prev.event.type === child.type) { + pending.set(key, { + type: "subagent-event", + callId: event.callId, + event: { type: child.type, text: (prev.event as { text: string }).text + child.text }, + }); + } else { + pending.set(key, event); + } + } else { + pending.set(key, event); + } + schedule(); + return; + } + } + // Ordering: flush coalesced deltas before discrete events. + if (pending.size) { + if (timer) { clearTimeout(timer); timer = undefined; } + flush(); + } + try { raw(event); } catch { /* ignore */ } + }; +} export async function runAgent(opts: RunAgentOptions): Promise { - const { apiBaseUrl, apiKey, model, prompt, attachments, history, maxTokens, maxSteps, autoContinue, contextTokens, sampling, modelParams, anthropic, oauthKind, systemPromptOverride, extraInstructions, enableFileReading, enableTerminalSuggestions, enableWorkspaceContext, approve, isSubagent, customSubagents, subagentModel, registerSubagentAbort, askUser, onAfterRun, onBeforeShell, onAfterEdit, onHook, signal, emit } = opts; + const { apiBaseUrl, apiKey, model, prompt, attachments, history, maxTokens, maxSteps, autoContinue, contextTokens, sampling, modelParams, anthropic, oauthKind, systemPromptOverride, extraInstructions, enableFileReading, enableTerminalSuggestions, enableWorkspaceContext, approve, isSubagent, customSubagents, subagentModel, registerSubagentAbort, askUser, onAfterRun, onBeforeShell, onAfterEdit, onHook, signal, emit: rawEmit } = opts; + const emit = coalesceEmit(rawEmit); // Mutable so the SwitchMode tool can change it mid-run. let mode = opts.mode; // multitask is agentic (full tool access); treat it like agent for gating. @@ -96,6 +185,12 @@ export async function runAgent(opts: RunAgentOptions): Promise { }); } let finalText = ""; + const budgetMs = opts?.runInBackground ? BG_SUBAGENT_MAX_MS : SUBAGENT_MAX_MS; + let budgetHit = false; + const budgetTimer = setTimeout(() => { + budgetHit = true; + try { childAC.abort(); } catch { /* ignore */ } + }, budgetMs); // Isolated chat: empty history, own context budget. Parent only gets // the final summary string (tool result / run-result) — never sub steps. const runP = runAgent({ @@ -120,9 +215,11 @@ export async function runAgent(opts: RunAgentOptions): Promise { signal: childAC.signal, emit: (e) => { if (e.type === "run-result") finalText = e.text; - // UI stream only — not parent history. + // UI stream only — not parent history. Coalesced via parent emit. if (callId) emit({ type: "subagent-event", callId, event: e }); }, + }).finally(() => { + clearTimeout(budgetTimer); }); // Background subagents return immediately; they keep streaming via emit. if (opts?.runInBackground) { @@ -131,8 +228,18 @@ export async function runAgent(opts: RunAgentOptions): Promise { // and capture its summary so it can be fed back into the loop on completion. const idx = bgSubagents.length; const tracked = runP - .then(() => ({ title, text: finalText || "(subagent finished with no summary)" })) - .catch((e) => ({ title, text: `(subagent failed: ${e instanceof Error ? e.message : String(e)})` })) + .then(() => ({ + title, + text: budgetHit + ? `(subagent timed out after ${Math.round(budgetMs / 1000)}s)` + : (finalText || "(subagent finished with no summary)"), + })) + .catch((e) => ({ + title, + text: budgetHit + ? `(subagent timed out after ${Math.round(budgetMs / 1000)}s)` + : `(subagent failed: ${e instanceof Error ? e.message : String(e)})`, + })) .finally(() => { parentSig.removeEventListener("abort", onParentAbort); onHook?.("subagentStop", { subagent: title }); @@ -144,14 +251,17 @@ export async function runAgent(opts: RunAgentOptions): Promise { try { await runP; } catch (e) { + if (budgetHit) return `(subagent timed out after ${Math.round(budgetMs / 1000)}s)`; if (childAC.signal.aborted || parentSig.aborted) { return "(subagent cancelled)"; } return `(subagent failed: ${e instanceof Error ? e.message : String(e)})`; } finally { + clearTimeout(budgetTimer); parentSig.removeEventListener("abort", onParentAbort); onHook?.("subagentStop", { subagent: subagentName || "subagent" }); } + if (budgetHit) return `(subagent timed out after ${Math.round(budgetMs / 1000)}s)`; if (childAC.signal.aborted || parentSig.aborted) return "(subagent cancelled)"; return finalText || "(subagent finished with no summary)"; }; @@ -537,7 +647,25 @@ export async function runAgent(opts: RunAgentOptions): Promise { }); const results = new Array<{ status: "completed" | "error"; output: string; diff?: string; startLine?: number; endLine?: number; image?: { mime: string; base64: string } }>(parsed.length); - const ro: Promise[] = []; + const completedUi = new Set(); + const finishUi = (i: number) => { + if (completedUi.has(i) || !results[i]) return; + completedUi.add(i); + const { call } = parsed[i]; + const r = results[i]; + // Surface completion as soon as the tool settles — don't wait for + // siblings. Prevents one slow tool from freezing the whole card strip. + emit({ + type: "tool-call-completed", + callId: call.id, + name: call.name, + status: r.status, + result: r.output, + diff: r.diff, + startLine: r.startLine, + endLine: r.endLine, + }); + }; const exec = async (i: number) => { const { call, input, badArgs } = parsed[i]; @@ -546,6 +674,7 @@ export async function runAgent(opts: RunAgentOptions): Promise { status: "error", output: `error: tool arguments were not valid JSON (likely truncated — the payload was too large). Retry with a smaller edit: split the change into multiple smaller ${call.name} calls.`, }; + finishUi(i); return; } // MCP tool dispatch. @@ -653,44 +782,64 @@ export async function runAgent(opts: RunAgentOptions): Promise { if (status === "completed" && isEditTool && onAfterEdit) { onAfterEdit(String(input?.path ?? "")); } + finishUi(i); } catch (e) { results[i] = { status: "error", output: `error: ${e instanceof Error ? e.message : String(e)}` }; + finishUi(i); } }; + // Early-exit paths inside exec that set results without finishUi. + const wrapExec = async (i: number) => { + try { + await exec(i); + } finally { + // Guarantee UI settles even if a branch forgot finishUi. + if (results[i]) finishUi(i); + else { + results[i] = { status: "error", output: "error: tool produced no result" }; + finishUi(i); + } + } + }; + + // Cap parallel RO tools so a burst of Grep/Glob/Task can't thrash CPU/IO. + // Worker pool (not batch-wait): a long Task doesn't block the next free slot. + const RO_CONCURRENCY = 8; + const roIdx: number[] = []; for (let i = 0; i < parsed.length; i++) { const name = parsed[i].call.name; const tool = TOOLS[name]; - // Run read-only built-in tools in parallel; MCP + mutating tools serialize below. - if (tool && !tool.mutating && !name.startsWith("mcp__")) { - ro.push(exec(i)); - } + if (tool && !tool.mutating && !name.startsWith("mcp__")) roIdx.push(i); + } + if (roIdx.length) { + let cursor = 0; + const workers = Array.from( + { length: Math.min(RO_CONCURRENCY, roIdx.length) }, + async () => { + while (cursor < roIdx.length) { + const i = roIdx[cursor++]; + await wrapExec(i); + } + }, + ); + await Promise.all(workers); } - await Promise.all(ro); for (let i = 0; i < parsed.length; i++) { const name = parsed[i].call.name; const tool = TOOLS[name]; if (!tool || tool.mutating || name.startsWith("mcp__")) { - await exec(i); + await wrapExec(i); } } for (let i = 0; i < parsed.length; i++) { const { call } = parsed[i]; - const r = results[i]; + const r = results[i] ?? { status: "error" as const, output: "error: tool produced no result" }; if (call.name === "WritePlan" && r.status === "completed") { planWritten = true; } - emit({ - type: "tool-call-completed", - callId: call.id, - name: call.name, - status: r.status, - result: r.output, - diff: r.diff, - startLine: r.startLine, - endLine: r.endLine, - }); + // History in call order (model expects stable tool-result sequencing). history.push({ kind: "tool-result", callId: call.id, name: call.name, output: r.output, status: r.status, image: r.image }); } // After launching background Task(s), wait for that wave before calling the diff --git a/src/agent/tools/shared.ts b/src/agent/tools/shared.ts index 1dd0ffb..a05d19b 100644 --- a/src/agent/tools/shared.ts +++ b/src/agent/tools/shared.ts @@ -110,8 +110,14 @@ export function setToolTimeoutOverrides(sec: Record | undefined) * Always settles (never hangs) even if `p` never resolves. */ export function withToolTimeout(p: Promise, ms: number, label: string): Promise { - // Guard against a never-settling promise with no explicit limit. - const limit = !ms || ms <= 0 ? 120_000 : ms; + // ms <= 0: no outer race (Task/AskQuestion manage their own lifetime). + // Still attach a settle handler so a late rejection cannot become unhandled. + if (!ms || ms <= 0) { + return Promise.resolve(p).catch((e) => { + throw e instanceof Error ? e : new Error(String(e)); + }); + } + const limit = ms; return new Promise((resolve, reject) => { let settled = false; const timer = setTimeout(() => { diff --git a/src/ui/sidebarProvider.ts b/src/ui/sidebarProvider.ts index 348a989..970037d 100644 --- a/src/ui/sidebarProvider.ts +++ b/src/ui/sidebarProvider.ts @@ -1275,6 +1275,27 @@ export class SidebarProvider implements vscode.WebviewViewProvider { } // Only deliver events to the webview when this conversation is on screen. + // High-frequency stream events are coalesced so postMessage + applyEvent + // cannot stall the extension host (looks like stuck tools/subagents). + type PendingUi = { event: AgentEvent; apply: boolean }; + const uiPending = new Map(); + let uiTimer: ReturnType | undefined; + const flushUi = () => { + uiTimer = undefined; + if (!uiPending.size) return; + const batch = [...uiPending.values()]; + uiPending.clear(); + for (const { event, apply } of batch) { + if (apply) { + session.turns = applyEvent(session.turns, event as unknown as SharedAgentEvent); + this._schedulePersistTurns(convId, session); + } + this._view?.webview.postMessage({ type: "agentEvent", convId, event }); + } + }; + const scheduleUi = () => { + if (!uiTimer) uiTimer = setTimeout(flushUi, 48); + }; const emit = (event: AgentEvent) => { if (event.type === "error") { SidebarProvider.log.appendLine(`[${new Date().toISOString()}] [agent] ${event.message}`); @@ -1284,6 +1305,78 @@ export class SidebarProvider implements vscode.WebviewViewProvider { // Maintain authoritative host turns from the same reducer the UI uses, then // persist (throttled) so any webview reload restores the live state exactly. const ev = event as unknown as SharedAgentEvent; + + // Coalesce stream deltas before applying / posting. + if (ev.type === "text-delta") { + const prev = uiPending.get("text"); + if (prev && (prev.event as SharedAgentEvent).type === "text-delta") { + const p = prev.event as Extract; + uiPending.set("text", { + event: { type: "text-delta", text: p.text + ev.text }, + apply: true, + }); + } else { + uiPending.set("text", { event, apply: true }); + } + scheduleUi(); + return; + } + if (ev.type === "thinking-delta") { + const prev = uiPending.get("think"); + if (prev && (prev.event as SharedAgentEvent).type === "thinking-delta") { + const p = prev.event as Extract; + uiPending.set("think", { + event: { type: "thinking-delta", text: p.text + ev.text }, + apply: true, + }); + } else { + uiPending.set("think", { event, apply: true }); + } + scheduleUi(); + return; + } + if (ev.type === "tool-call-args") { + uiPending.set(`args:${ev.callId}`, { event, apply: true }); + scheduleUi(); + return; + } + if (ev.type === "subagent-event") { + const child = ev.event as SharedAgentEvent; + if (child.type === "text-delta" || child.type === "thinking-delta" || child.type === "tool-call-args") { + const key = + child.type === "tool-call-args" + ? `sub:${ev.callId}:args:${(child as { callId: string }).callId}` + : `sub:${ev.callId}:${child.type}`; + if (child.type === "text-delta" || child.type === "thinking-delta") { + const prev = uiPending.get(key); + if (prev && (prev.event as SharedAgentEvent).type === "subagent-event") { + const pe = (prev.event as Extract).event; + if (pe.type === child.type) { + uiPending.set(key, { + event: { + type: "subagent-event", + callId: ev.callId, + event: { type: child.type, text: (pe as { text: string }).text + child.text }, + }, + apply: true, + }); + scheduleUi(); + return; + } + } + } + uiPending.set(key, { event, apply: true }); + scheduleUi(); + return; + } + } + + // Discrete events: flush coalesced stream first (ordering). + if (uiPending.size) { + if (uiTimer) { clearTimeout(uiTimer); uiTimer = undefined; } + flushUi(); + } + if (ev.type === "run-status") { if (ev.status === "finished" || ev.status === "cancelled" || ev.status === "error") { session.turns = forceSettleOpenWork( diff --git a/webview-ui/sidebar/App.tsx b/webview-ui/sidebar/App.tsx index d0c25d2..55a5c2c 100644 --- a/webview-ui/sidebar/App.tsx +++ b/webview-ui/sidebar/App.tsx @@ -882,6 +882,15 @@ export function App() { }; React.useEffect(() => { + // rAF-batch stream-driven re-renders so tool/args/text deltas don't thrash React. + let raf = 0; + const scheduleForce = () => { + if (raf) return; + raf = requestAnimationFrame(() => { + raf = 0; + force(); + }); + }; const handler = (event: MessageEvent) => { const msg = event.data; switch (msg.type) { @@ -1006,12 +1015,23 @@ export function App() { // Keep the rendered error block; persist so the chat survives reloads. s.status = { text: "Error: " + ev.message, error: true }; post({ type: "persistTurns", convId: msg.convId, turns: s.turns }); - } else if (ev.type === "mode-changed") { + } else if (ev.type === "mode-changed") { setMode(ev.mode); post({ type: "setMode", mode: ev.mode }); } } - force(); + // Immediate paint on settle / tool complete; coalesce stream deltas. + if ( + settled || + ev.type === "tool-call-completed" || + ev.type === "tool-call-started" || + ev.type === "error" || + ev.type === "run-result" + ) { + force(); + } else { + scheduleForce(); + } break; } } @@ -1029,6 +1049,7 @@ export function App() { window.addEventListener("beforeunload", flush); post({ type: "ready" }); return () => { + if (raf) cancelAnimationFrame(raf); flush(); window.removeEventListener("message", handler); window.removeEventListener("pagehide", flush);