Level 6
Streaming and events
- user
- assistant
- tool_result
EventBus
The problem
Calling a model and waiting for the whole answer works in a script. In front of a person it doesn’t: a long reply can take ten or twenty seconds, and during all of them the screen shows nothing. An agent makes it worse, because one request can run several turns and tools before the final answer exists.
And it isn’t only the person who wants to follow along. The interface wants the text, your logs want every step, billing wants the tokens of each turn. If each of those has to be written into the loop, the loop ends up knowing about every screen, file and database in your app.
The solution
Two pieces that go together. Streaming: ask the provider for the reply as a stream, and it sends each piece of text as soon as the model writes it. The total time is the same, but the first word shows up in a fraction of a second instead of at the end.
Events: the loop announces every step on a bus, and doesn’t care who is listening. Your app subscribes the parts that care: the screen listens for text, the log for everything, the meter for the end of each turn. Adding a listener doesn’t touch the loop. These are the events astorlm emits on a normal run:
-
turn_startA turn begins: the loop is about to call the model.
-
text_deltaA piece of the reply arrived. Many per turn.
-
assistant_messageThe model’s message is complete, text and tool calls included.
-
tool_execution_startA tool is about to run, with its input.
-
tool_execution_endThe tool finished: its output, whether it failed, how long it took.
-
turn_endThe turn is over, with the tokens its model call used.
-
session_endThe run is over: completed, aborted, or failed.
Tool calls stream from the provider too, a few characters of their arguments at a time. The loop puts them back
together before it runs them, so on the bus a tool call shows up whole, in assistant_message.
The cast
Same cast as always, on a rock stage.
- The Oracle the model
- The singer, up on the riser. Its reply comes down the highway as notes, one per word.
- The highway the stream
- Each text_delta is a handful of notes. Without streaming, nothing comes down until the end.
- The guitarist your app
- Plays each note as it reaches the frets, and the lyrics appear for the crowd.
- The cable the EventBus
- Every event runs along it. Only the listeners that subscribed to that event light up.
- The listeners agent.on(…)
- The lyrics screen (
'text'), a tape deck ('event') and a token meter (turn_end). - Astor the loop
- Runs the loop as usual, and announces each step on the bus as he goes.
- The stage the tools
-
The front edge, the lighting desk and the pyro box:
crowd_mood,set_lightsandpyro.
The code
With astorlm: the provider always streams, so there is nothing to switch on. Subscribe with
agent.on() before you call run(): 'text' for each piece of text,
'tool-start' and 'tool-end' for tools, 'event' for everything. Each call returns
a function that unsubscribes.
From scratch: the loop from level 2 with two changes. The request asks for
stream: true and reads the reply as server-sent events, gluing the pieces back together; and a list of
listeners gets every step through one emit().
import { OpenAIProvider, createLocalAgent, tool } from 'astorlm'
import { z } from 'zod'
// The stage's functions, wrapped as tools.
const crowdMood = tool({
name: 'crowd_mood',
description: 'Look at the crowd from the front of the stage: how they feel, and what they came to hear.',
schema: z.object({}),
execute: async () => stage.readCrowd(), // your code
})
const setLights = tool({
name: 'set_lights',
description: 'Set the color of every light on the stage.',
schema: z.object({ color: z.enum(['amber', 'red', 'blue', 'white']) }),
execute: async ({ color }) => stage.lights(color), // your code
})
const pyro = tool({
name: 'pyro',
description: 'Fire the pyrotechnics on both sides of the stage. Once per show.',
schema: z.object({}),
execute: async () => stage.firePyro(), // your code
})
const agent = await createLocalAgent({
// Any OpenAI-compatible endpoint: OpenAI, Ollama, LM Studio, vLLM, a proxy…
provider: new OpenAIProvider({
baseURL: 'http://localhost:11434/v1', // e.g. Ollama's default address
model: 'your-model', // e.g. 'llama3.1', 'gpt-4o-mini'
apiKey: 'YOUR_API_KEY', // local servers usually ignore it
}),
tools: [crowdMood, setLights, pyro],
maxTurns: 10,
})
// Streaming is always on: subscribe before run(). Each listener only hears what it asked for.
// The lyrics screen: every piece of text, the moment it arrives.
agent.on('text', (piece) => lyrics.append(piece))
// The tape deck: every event the loop emits, in order.
agent.on('event', (event) => tape.write(event))
// The meter: the tokens each turn reported.
agent.on('event', (event) => {
if (event.type === 'turn_end' && event.usage) meter.add(event.usage.inputTokens + event.usage.outputTokens)
})
// Each on() returns a function that unsubscribes.
const stopTape = agent.on('tool-start', ({ name }) => console.log('→', name))
const last = await agent.run('Open the show: read the crowd, set the mood, and greet them.')
stopTape()
console.log(last.content) // the same text, whole, once the run is over
// Streaming and events from scratch. Plain fetch, no SDK.
// Any OpenAI-compatible endpoint: OpenAI, Ollama, LM Studio, vLLM, a proxy…
const LLM = {
baseURL: 'http://localhost:11434/v1', // e.g. Ollama's default address
model: 'your-model', // e.g. 'llama3.1', 'gpt-4o-mini'
apiKey: 'YOUR_API_KEY', // local servers usually ignore it
}
type ToolCall = { id: string; type: 'function'; function: { name: string; arguments: string } }
type Message =
| { role: 'user'; content: string }
| { role: 'assistant'; content: string | null; tool_calls?: ToolCall[] }
| { role: 'tool'; tool_call_id: string; content: string }
const tools: Record<string, (args: Record<string, string>) => Promise<string>> = { crowd_mood: crowdMood, set_lights: setLights, pyro }
const toolSchemas = [/* one JSON Schema per tool */]
// 1. The bus: a list of listeners. The loop emits; it never knows who is listening.
type AgentEvent =
| { type: 'turn_start'; turn: number }
| { type: 'text_delta'; text: string }
| { type: 'tool_execution_start'; name: string }
| { type: 'tool_execution_end'; name: string; output: string }
| { type: 'turn_end'; turn: number; usage?: { prompt_tokens: number; completion_tokens: number } }
const listeners: ((event: AgentEvent) => void)[] = []
export const on = (listener: (event: AgentEvent) => void) => {
listeners.push(listener)
return () => listeners.splice(listeners.indexOf(listener), 1)
}
const emit = (event: AgentEvent) => listeners.forEach((listener) => listener(event))
// 2. One streamed model call: read the server-sent events and rebuild the message piece by piece.
async function streamReply(messages: Message[]) {
const res = await fetch(`${LLM.baseURL}/chat/completions`, {
method: 'POST',
headers: { 'content-type': 'application/json', authorization: `Bearer ${LLM.apiKey}` },
body: JSON.stringify({ model: LLM.model, messages, tools: toolSchemas, stream: true, stream_options: { include_usage: true } }),
})
const reader = res.body!.pipeThrough(new TextDecoderStream()).getReader()
let text = ''
let buffer = ''
let usage
const calls: ToolCall[] = []
for (let read = await reader.read(); !read.done; read = await reader.read()) {
buffer += read.value
const lines = buffer.split('\n')
buffer = lines.pop()! // a line can arrive cut in half: keep it for the next chunk
for (const line of lines) {
if (!line.startsWith('data: ') || line === 'data: [DONE]') continue
const chunk = JSON.parse(line.slice(6))
usage = chunk.usage ?? usage
const delta = chunk.choices?.[0]?.delta
if (delta?.content) {
text += delta.content
emit({ type: 'text_delta', text: delta.content }) // show it now, not at the end
}
// Tool calls arrive in pieces too: glue each one's arguments together by index.
for (const part of delta?.tool_calls ?? []) {
const call = (calls[part.index] ??= { id: '', type: 'function', function: { name: '', arguments: '' } })
call.id ||= part.id ?? ''
call.function.name += part.function?.name ?? ''
call.function.arguments += part.function?.arguments ?? ''
}
}
}
const reply: Message = { role: 'assistant', content: text || null, ...(calls.length ? { tool_calls: calls } : {}) }
return { reply, usage }
}
// 3. The loop from level 2, emitting as it goes.
export async function runAgent(prompt: string, maxTurns = 10): Promise<string> {
const messages: Message[] = [{ role: 'user', content: prompt }]
for (let turn = 1; turn <= maxTurns; turn++) {
emit({ type: 'turn_start', turn })
const { reply, usage } = await streamReply(messages)
messages.push(reply)
if (reply.role !== 'assistant' || !reply.tool_calls) {
emit({ type: 'turn_end', turn, usage })
return reply.content ?? ''
}
for (const call of reply.tool_calls) {
emit({ type: 'tool_execution_start', name: call.function.name })
const output = await tools[call.function.name]!(JSON.parse(call.function.arguments || '{}'))
emit({ type: 'tool_execution_end', name: call.function.name, output })
messages.push({ role: 'tool', tool_call_id: call.id, content: output })
}
emit({ type: 'turn_end', turn, usage })
}
throw new Error(`No answer after ${maxTurns} turns`)
}
// 4. Your app subscribes. The loop above doesn't change to add one.
on((e) => e.type === 'text_delta' && lyrics.append(e.text))
on((e) => tape.write(e))
on((e) => e.type === 'turn_end' && e.usage && meter.add(e.usage.prompt_tokens + e.usage.completion_tokens))
# Streaming and events from scratch. Standard library only, no SDK.
import json
import urllib.request
# Any OpenAI-compatible endpoint: OpenAI, Ollama, LM Studio, vLLM, a proxy...
LLM = {
"base_url": "http://localhost:11434/v1", # e.g. Ollama's default address
"model": "your-model", # e.g. "llama3.1", "gpt-4o-mini"
"api_key": "YOUR_API_KEY", # local servers usually ignore it
}
TOOLS = {"crowd_mood": crowd_mood, "set_lights": set_lights, "pyro": pyro}
TOOL_SCHEMAS = [...] # one JSON Schema per tool
# 1. The bus: a list of listeners. The loop emits; it never knows who is listening.
listeners = []
def on(listener):
listeners.append(listener)
return lambda: listeners.remove(listener)
def emit(event):
for listener in listeners:
listener(event)
# 2. One streamed model call: read the server-sent events and rebuild the message piece by piece.
def stream_reply(messages):
body = {
"model": LLM["model"],
"messages": messages,
"tools": TOOL_SCHEMAS,
"stream": True,
"stream_options": {"include_usage": True},
}
request = urllib.request.Request(
f"{LLM['base_url']}/chat/completions",
data=json.dumps(body).encode(),
headers={"Content-Type": "application/json", "Authorization": f"Bearer {LLM['api_key']}"},
)
text, calls, usage = "", {}, None
with urllib.request.urlopen(request) as response:
for raw in response: # one server-sent event per line
line = raw.decode().strip()
if not line.startswith("data: ") or line == "data: [DONE]":
continue
chunk = json.loads(line[6:])
usage = chunk.get("usage") or usage
delta = (chunk.get("choices") or [{}])[0].get("delta", {})
if delta.get("content"):
text += delta["content"]
emit({"type": "text_delta", "text": delta["content"]}) # show it now, not at the end
# Tool calls arrive in pieces too: glue each one's arguments together by index.
for part in delta.get("tool_calls") or []:
call = calls.setdefault(part["index"], {"id": "", "type": "function", "function": {"name": "", "arguments": ""}})
call["id"] = call["id"] or part.get("id") or ""
call["function"]["name"] += part.get("function", {}).get("name") or ""
call["function"]["arguments"] += part.get("function", {}).get("arguments") or ""
reply = {"role": "assistant", "content": text or None}
if calls:
reply["tool_calls"] = [calls[i] for i in sorted(calls)]
return reply, usage
# 3. The loop from level 2, emitting as it goes.
def run_agent(prompt, max_turns=10):
messages = [{"role": "user", "content": prompt}]
for turn in range(1, max_turns + 1):
emit({"type": "turn_start", "turn": turn})
reply, usage = stream_reply(messages)
messages.append(reply)
if not reply.get("tool_calls"):
emit({"type": "turn_end", "turn": turn, "usage": usage})
return reply["content"] or ""
for call in reply["tool_calls"]:
name = call["function"]["name"]
emit({"type": "tool_execution_start", "name": name})
output = TOOLS[name](**json.loads(call["function"]["arguments"] or "{}"))
emit({"type": "tool_execution_end", "name": name, "output": output})
messages.append({"role": "tool", "tool_call_id": call["id"], "content": output})
emit({"type": "turn_end", "turn": turn, "usage": usage})
raise RuntimeError(f"No answer after {max_turns} turns")
# 4. Your app subscribes. The loop above doesn't change to add one.
on(lambda e: e["type"] == "text_delta" and print(e["text"], end="", flush=True))
on(lambda e: tape.write(e))
on(lambda e: e["type"] == "turn_end" and e["usage"] and meter.add(e["usage"]["total_tokens"]))
What to watch
- Streaming doesn’t make the model faster. The reply takes as long as before; what changes is when the person sees the first word. Measure time to first token, not only total time.
- Listeners watch; they don’t steer. Whatever a listener returns is ignored, and the loop doesn’t wait for it. To block a tool or change what the model reads, you need a hook: the next level.
- Keep listeners quick. astorlm’s bus calls them one after the other, inside the loop. Slow work, like writing to a database, belongs in a queue the listener only pushes to.
- A broken listener fails quietly. astorlm catches its error so the run goes on, which also means nobody hears about it. Log inside your listeners.
- Pieces don’t respect words. A text_delta can end in the middle of a word or of a Markdown table. Render what you have so far, and redraw as more arrives.
-
Let the person stop it. Once they can watch the reply arrive, they will want to cut a wrong one
short. Wire a stop button to
agent.abort().