Nível 6
Streaming e eventos
- user
- assistant
- tool_result
EventBus
O problema
Chamar um modelo e esperar a resposta inteira funciona num script. Na frente de uma pessoa, não: uma resposta longa pode levar dez ou vinte segundos, e durante todos eles a tela não mostra nada. Um agente piora isso, porque uma requisição pode rodar vários turnos e ferramentas antes de a resposta final existir.
E não é só a pessoa que quer acompanhar. A interface quer o texto, seus logs querem cada passo, o faturamento quer os tokens de cada turno. Se cada uma dessas coisas tiver que ser escrita dentro do loop, o loop acaba conhecendo cada tela, arquivo e banco de dados do seu app.
A solução
Duas peças que andam juntas. Streaming: peça ao provedor a resposta como um fluxo, e ele manda cada pedacinho de texto assim que o modelo o escreve. O tempo total é o mesmo, mas a primeira palavra aparece numa fração de segundo em vez de no final.
Eventos: o loop anuncia cada passo num bus, e não se importa com quem está ouvindo. Seu app inscreve as partes que se interessam: a tela escuta o texto, o log escuta tudo, o medidor escuta o fim de cada turno. Adicionar um ouvinte não mexe no loop. Estes são os eventos que o astorlm emite numa execução normal:
-
turn_startComeça um turno: o loop está prestes a chamar o modelo.
-
text_deltaChegou um pedacinho da resposta. São muitos por turno.
-
assistant_messageA mensagem do modelo está completa, com o texto e as chamadas de ferramenta.
-
tool_execution_startUma ferramenta está prestes a rodar, com sua entrada.
-
tool_execution_endA ferramenta terminou: sua saída, se falhou e quanto tempo levou.
-
turn_endO turno terminou, com os tokens que sua chamada ao modelo usou.
-
session_endA execução terminou: concluída, cancelada ou com erro.
As chamadas de ferramenta também chegam do provedor em streaming, alguns caracteres dos argumentos por vez. O loop as
remonta antes de rodá-las, então no bus uma chamada de ferramenta aparece inteira, em assistant_message.
O elenco
O mesmo elenco de sempre, num palco de rock.
- O Oráculo o modelo
- O vocalista, lá no praticável. Sua resposta desce pela pista como notas, uma por palavra.
- A pista o stream
- Cada text_delta é um punhado de notas. Sem streaming, nada desce até o final.
- A pessoa da guitarra seu app
- Toca cada nota quando ela chega aos trastes, e as letras aparecem para o público.
- O cabo o EventBus
- Todos os eventos correm por ele. Só acendem os ouvintes que se inscreveram naquele evento.
- Os ouvintes agent.on(…)
- A tela de letras (
'text'), um gravador de fita ('event') e um medidor de tokens (turn_end). - Astor o loop
- Roda o loop como sempre, e anuncia cada passo no bus conforme avança.
- O palco as ferramentas
-
A beira do palco, a mesa de luz e a caixa de pirotecnia:
crowd_mood,set_lightsepyro.
O código
Com astorlm: o provedor sempre faz streaming, então não há nada para ligar. Inscreva-se com
agent.on() antes de chamar run(): 'text' para cada pedacinho de texto,
'tool-start' e 'tool-end' para as ferramentas, 'event' para tudo. Cada chamada
devolve uma função que cancela a inscrição.
Do zero: o loop do nível 2 com duas mudanças. A requisição usa stream: true e lê a
resposta como server-sent events, colando os pedacinhos; e uma lista de ouvintes recebe cada passo por um único
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"]))
O que observar
- O streaming não deixa o modelo mais rápido. A resposta leva o mesmo tempo de antes; o que muda é quando a pessoa vê a primeira palavra. Meça o tempo até o primeiro token, não só o tempo total.
- Ouvintes observam; não dirigem. O que um ouvinte devolve é ignorado, e o loop não espera por ele. Para bloquear uma ferramenta ou mudar o que o modelo lê, você precisa de um hook: o próximo nível.
- Mantenha os ouvintes rápidos. O bus do astorlm os chama um depois do outro, dentro do loop. Trabalho lento, como gravar num banco de dados, vai numa fila para a qual o ouvinte só empurra.
- Um ouvinte quebrado falha em silêncio. O astorlm captura o erro para a execução seguir, o que também significa que ninguém fica sabendo. Registre os erros dentro dos seus ouvintes.
- Os pedacinhos não respeitam as palavras. Um text_delta pode terminar no meio de uma palavra ou de uma tabela em Markdown. Mostre o que você tem até agora e redesenhe conforme chega mais.
-
Deixe a pessoa parar. Quando ela pode ver a resposta chegando, vai querer cortar uma que está indo
mal. Ligue um botão de parar a
agent.abort().