refactor: fold agent stream events incrementally per poll

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
This commit is contained in:
hugocasa
2026-09-17 10:07:08 +02:00
co-authored by Claude Opus 5
parent 9819b16186
commit fae7c8bd2c
4 changed files with 163 additions and 92 deletions
@@ -1,20 +1,50 @@
<script lang="ts">
import { Loader2 } from 'lucide-svelte'
import { CircleX, Loader2 } from 'lucide-svelte'
import { untrack } from 'svelte'
import GfmMarkdown from './GfmMarkdown.svelte'
import type { AgentStream } from './aiAgentResult'
import {
advanceAgentStream,
emptyAgentStreamProgress,
type AgentStreamProgress
} from './aiAgentResult'
interface Props {
stream: AgentStream
/** The raw `result_stream` buffer, which only ever grows for a given run. */
raw: string
/** Identifies the run, so a viewer reused for another one starts over. */
streamKey?: string
}
let { stream }: Props = $props()
let { raw, streamKey }: Props = $props()
let progress: AgentStreamProgress & { key?: string } = $state(emptyAgentStreamProgress())
$effect(() => {
raw
streamKey
untrack(() => {
// A shorter buffer than what was already folded in cannot be a longer
// version of the same stream, so treat it as a different one.
const continues = progress.key === streamKey && raw.length >= progress.consumed
const base = continues ? progress : emptyAgentStreamProgress()
progress = { ...advanceAgentStream(raw, base), key: streamKey }
})
})
let stream = $derived(progress.stream)
</script>
<div class="flex flex-col gap-2 w-full">
{#if stream.tool}
<div class="flex items-center gap-2 text-secondary text-xs">
<div
class="flex items-center gap-2 text-xs {stream.tool.success === false
? 'text-red-500'
: 'text-secondary'}"
>
{#if stream.tool.running}
<Loader2 class="animate-spin shrink-0" size={14} />
{:else if stream.tool.success === false}
<CircleX class="shrink-0" size={14} />
{/if}
<span class="font-mono truncate">{stream.tool.name}</span>
</div>
@@ -57,7 +57,7 @@
import type { MarkupTrust } from './apps/markupTrust'
import AgentResultDisplay from './AgentResultDisplay.svelte'
import AgentStreamDisplay from './AgentStreamDisplay.svelte'
import { parseAgentResult, parseAgentStream } from './aiAgentResult'
import { isAgentStream, parseAgentResult } from './aiAgentResult'
const TABLE_MAX_SIZE = 5000000
const DISPLAY_MAX_SIZE = 100000
@@ -99,7 +99,10 @@
const REPLAY_INERT_KINDS: ResultKind[] = ['s3object', 's3object-list', 'materialized', 'approval']
/** Kinds whose markup pulls subresources: DOMPurify stops scripting but keeps
* `<img src>` and SVG `<image href>`, and `map` tiles are requests by
* construction. Kinds absent here carry their bytes as `data:` and reach nothing.
* construction. Kinds absent here carry their bytes as `data:` and reach nothing,
* or render through a component that is itself inert on the public page —
* `aiagent` is the second case, via `GfmMarkdown`, which is why it renders
* markdown yet is not listed while `markdown` still is.
* Inert only on the public page, which promises to issue no requests. */
const OFFLINE_INERT_KINDS: ResultKind[] = ['markdown', 'html', 'svg', 'map']
let length = $state(1)
@@ -743,15 +746,14 @@
{/if}
{#if result_stream && result == undefined}
{@const agentStream = parseAgentStream(result_stream)}
<div class="flex flex-col w-full gap-2">
<div class="flex items-center gap-2 text-secondary text-xs">
<Loader2 class="animate-spin" size={14} /> Streaming result
</div>
{#if agentStream}
{#if isAgentStream(result_stream)}
<!-- An agent streams one JSON event per line, so the raw stream is a wall of
event objects rather than the answer being written. -->
<AgentStreamDisplay stream={agentStream} />
<AgentStreamDisplay raw={result_stream} streamKey={jobId} />
{:else}
<ResultStreamDisplay {result_stream} />
{/if}
@@ -1,9 +1,10 @@
import { describe, expect, it } from 'vitest'
import {
advanceAgentStream,
emptyAgentStreamProgress,
formatTokenCount,
parseAgentErrorMessages,
isAgentStream,
parseAgentResult,
parseAgentStream,
summarizeAgentResult
} from './aiAgentResult'
@@ -38,25 +39,6 @@ describe('parseAgentResult', () => {
})
})
describe('parseAgentErrorMessages', () => {
it('reads the partial transcript a max-iterations failure carries', () => {
const messages = parseAgentErrorMessages({
error: {
name: 'ExecutionErr',
message: 'AI agent reached max iterations (10)',
result: { messages: [{ role: 'user', content: 'ask' }] }
}
})
expect(messages).toHaveLength(1)
})
it('ignores an error that carries no transcript', () => {
expect(
parseAgentErrorMessages({ error: { name: 'ExecutionErr', message: 'boom' } })
).toBeUndefined()
})
})
describe('summarizeAgentResult', () => {
it('counts the actions and falls back to the parts when no total is reported', () => {
const summary = summarizeAgentResult({
@@ -79,37 +61,58 @@ describe('summarizeAgentResult', () => {
})
})
describe('parseAgentStream', () => {
const events = [
describe('agent stream', () => {
const lines = [
'{"type":"tool_call","call_id":"c1","function_name":"query_metrics"}',
'{"type":"tool_result","call_id":"c1","function_name":"query_metrics","result":"{}","success":true}',
'{"type":"reasoning_token_delta","content":"checking"}',
'{"type":"token_delta","content":"eu-central-1"}',
'{"type":"token_delta","content":" is down"}'
].join('\n')
]
const events = lines.join('\n') + '\n'
it('joins the token deltas into the answer so far', () => {
const stream = parseAgentStream(events)
expect(stream?.answer).toBe('eu-central-1 is down')
expect(stream?.reasoning).toBe('checking')
expect(stream?.tool).toEqual({ name: 'query_metrics', running: false, success: true })
it('recognises an agent stream from its first line only', () => {
expect(isAgentStream(events)).toBe(true)
expect(isAgentStream('processing row 1\nprocessing row 2\n')).toBe(false)
expect(isAgentStream('{"level":"info","msg":"hello"}\n')).toBe(false)
// No newline yet, so the first line may still be half-written.
expect(isAgentStream('{"type":"token_delta","content":"a"}')).toBe(false)
})
it('marks a tool still running', () => {
expect(
parseAgentStream('{"type":"tool_execution","call_id":"c1","function_name":"fetch"}')?.tool
).toEqual({ name: 'fetch', running: true, success: undefined })
it('folds the token deltas into the answer so far', () => {
const { stream } = advanceAgentStream(events, emptyAgentStreamProgress())
expect(stream.answer).toBe('eu-central-1 is down')
expect(stream.reasoning).toBe('checking')
expect(stream.tool).toEqual({ name: 'query_metrics', running: false, success: true })
})
// A poll can cut the last event in half, and any other job may stream something
// that is not an agent's events at all.
it('survives a truncated trailing line', () => {
expect(parseAgentStream(`${events}\n{"type":"token_de`)?.answer).toBe('eu-central-1 is down')
// The stream only grows, so each poll must fold in the new lines and re-read
// none of the old ones — the reason this is incremental at all.
it('resumes where the previous poll stopped', () => {
const firstPoll = advanceAgentStream(lines.slice(0, 3).join('\n') + '\n', emptyAgentStreamProgress())
const secondPoll = advanceAgentStream(events, firstPoll)
expect(secondPoll.consumed).toBe(events.length)
expect(secondPoll.stream.answer).toBe('eu-central-1 is down')
expect(secondPoll.stream.reasoning).toBe('checking')
})
it('ignores a stream that carries no agent events', () => {
expect(parseAgentStream('processing row 1\nprocessing row 2')).toBeUndefined()
expect(parseAgentStream('{"level":"info","msg":"hello"}')).toBeUndefined()
it('leaves a half-written trailing line for the next poll', () => {
const partial = advanceAgentStream(`${events}{"type":"token_de`, emptyAgentStreamProgress())
expect(partial.stream.answer).toBe('eu-central-1 is down')
const completed = advanceAgentStream(`${events}{"type":"token_delta","content":"!"}\n`, partial)
expect(completed.stream.answer).toBe('eu-central-1 is down!')
})
it('marks a tool still running, and a failed one', () => {
const started = '{"type":"tool_execution","call_id":"c1","function_name":"fetch"}\n'
const running = advanceAgentStream(started, emptyAgentStreamProgress())
expect(running.stream.tool).toEqual({ name: 'fetch', running: true, success: undefined })
const failed = advanceAgentStream(
started +
'{"type":"tool_result","call_id":"c1","function_name":"fetch","result":"boom","success":false}\n',
running
)
expect(failed.stream.tool).toEqual({ name: 'fetch', running: false, success: false })
})
})
+78 -42
View File
@@ -52,6 +52,14 @@ function hasRole(message: unknown): boolean {
* The signature is deliberately closed — no key outside `ENVELOPE_KEYS`, and
* every message carrying a `role` — so an ordinary result that happens to have
* an `output` field cannot claim it.
*
* It stops short of also requiring a recognised `agent_action`, which would rule
* out a hand-written script returning this same shape. Not every completed run
* is guaranteed to tag a message (a run whose provider returns its answer
* through a structured-output tool leaves the final assistant message untagged),
* and the two failures are not symmetric: claiming a lookalike costs a viewer
* one click on the JSON toggle, while rejecting a real agent hides its answer
* with nothing on screen to say why.
*/
export function parseAgentResult(result: unknown): AgentResult | undefined {
if (!isRecord(result)) {
@@ -77,25 +85,6 @@ export function parseAgentResult(result: unknown): AgentResult | undefined {
}
}
/**
* A run stopped by `max_iterations` fails, so it returns an error rather than an
* envelope — but the worker attaches the conversation so far to it. That partial
* transcript is the whole reason to look at a run that hit the cap.
*/
export function parseAgentErrorMessages(result: unknown): AgentMessage[] | undefined {
if (!isRecord(result) || !isRecord(result.error)) {
return undefined
}
const inner = result.error.result
if (!isRecord(inner) || !Array.isArray(inner.messages)) {
return undefined
}
if (inner.messages.length === 0 || !inner.messages.every(hasRole)) {
return undefined
}
return inner.messages as AgentMessage[]
}
export type AgentResultSummary = {
toolCalls: number
webSearches: number
@@ -141,6 +130,13 @@ export type AgentStream = {
tool?: { name: string; running: boolean; success?: boolean }
}
/** How much of the stream has been folded in, so the next poll starts there. */
export type AgentStreamProgress = { consumed: number; stream: AgentStream }
export function emptyAgentStreamProgress(): AgentStreamProgress {
return { consumed: 0, stream: { answer: '', reasoning: '' } }
}
const STREAM_EVENT_TYPES = [
'token_delta',
'reasoning_token_delta',
@@ -150,32 +146,68 @@ const STREAM_EVENT_TYPES = [
'tool_result'
]
function parseStreamEvent(line: string): Record<string, unknown> | undefined {
let event: unknown
try {
event = JSON.parse(line)
} catch {
return undefined
}
if (!isRecord(event) || typeof event.type !== 'string') {
return undefined
}
return STREAM_EVENT_TYPES.includes(event.type) ? event : undefined
}
/**
* `result_stream` carries one `StreamingEvent` per line while an agent runs.
* Returns undefined for a stream that is not an agent's, so any other streaming
* result keeps being shown verbatim.
* Whether `result_stream` is an agent's event stream rather than something a
* script printed. Reads only the first complete line, because it runs on every
* poll of a running job.
*/
export function parseAgentStream(raw: string): AgentStream | undefined {
let sawEvent = false
const stream: AgentStream = { answer: '', reasoning: '' }
for (const line of raw.split('\n')) {
export function isAgentStream(raw: string): boolean {
let start = 0
while (start < raw.length) {
const end = raw.indexOf('\n', start)
if (end === -1) {
// Only a partial first line so far; wait for the poll that completes it.
return false
}
const line = raw.slice(start, end)
if (line.trim() !== '') {
return parseStreamEvent(line) !== undefined
}
start = end + 1
}
return false
}
/**
* Fold the events that arrived since `previous` into the answer so far.
*
* Incremental rather than a parse of the whole buffer: the stream only ever
* grows, a poll can arrive every 50ms, and a `tool_result` event carries the
* tool's entire output — so re-reading everything each time is quadratic in the
* number of events with a large constant.
*/
export function advanceAgentStream(
raw: string,
previous: AgentStreamProgress
): AgentStreamProgress {
// A trailing line with no newline yet is still being written, so it stays
// unconsumed until the poll that completes it.
const complete = raw.lastIndexOf('\n') + 1
if (complete <= previous.consumed) {
return previous
}
const stream: AgentStream = { ...previous.stream }
for (const line of raw.slice(previous.consumed, complete).split('\n')) {
if (line.trim() === '') {
continue
}
let event: unknown
try {
event = JSON.parse(line)
} catch {
// A trailing partial line is normal: the poll can cut an event in half.
const event = parseStreamEvent(line)
if (!event) {
continue
}
if (!isRecord(event) || typeof event.type !== 'string') {
continue
}
if (!STREAM_EVENT_TYPES.includes(event.type)) {
continue
}
sawEvent = true
if (event.type === 'token_delta' && typeof event.content === 'string') {
stream.answer += event.content
} else if (event.type === 'reasoning_token_delta' && typeof event.content === 'string') {
@@ -188,14 +220,18 @@ export function parseAgentStream(raw: string): AgentStream | undefined {
}
}
}
return sawEvent ? stream : undefined
return { consumed: complete, stream }
}
/** Token counts run to five and six figures, where the exact digit is noise. */
/**
* Token counts run to five and six figures, where the exact digit is noise. The
* millions branch is not decoration: usage accumulates over every loop
* iteration, and each one re-sends the whole context.
*/
export function formatTokenCount(count: number): string {
if (count < 1000) {
return String(count)
}
const thousands = count / 1000
return `${thousands < 10 ? thousands.toFixed(1) : Math.round(thousands)}k`
const [scaled, unit] = count < 1_000_000 ? [count / 1000, 'k'] : [count / 1_000_000, 'M']
return `${scaled < 10 ? scaled.toFixed(1) : Math.round(scaled)}${unit}`
}