Nivel 6
Streaming y eventos
- user
- assistant
- tool_result
EventBus
El problema
Llamar a un modelo y esperar la respuesta entera funciona en un script. Frente a una persona, no: una respuesta larga puede tardar diez o veinte segundos, y durante todos ellos la pantalla no muestra nada. Un agente lo empeora, porque un pedido puede correr varios turnos y herramientas antes de que exista la respuesta final.
Y no es solo la persona la que quiere seguir lo que pasa. La interfaz quiere el texto, tus logs quieren cada paso, la facturación quiere los tokens de cada turno. Si cada una de esas cosas hay que escribirla dentro del bucle, el bucle termina conociendo cada pantalla, archivo y base de datos de tu app.
La solución
Dos piezas que van juntas. Streaming: pídele al proveedor la respuesta como un flujo, y te manda cada pedacito de texto apenas el modelo lo escribe. El tiempo total es el mismo, pero la primera palabra aparece en una fracción de segundo en lugar de al final.
Eventos: el bucle anuncia cada paso en un bus, y no le importa quién escucha. Tu app suscribe las partes a las que les interesa: la pantalla escucha el texto, el log escucha todo, el medidor escucha el final de cada turno. Agregar un oyente no toca el bucle. Estos son los eventos que astorlm emite en una ejecución normal:
-
turn_startEmpieza un turno: el bucle está por llamar al modelo.
-
text_deltaLlegó un pedacito de la respuesta. Hay muchos por turno.
-
assistant_messageEl mensaje del modelo está completo, con el texto y las llamadas a herramientas.
-
tool_execution_startUna herramienta está por ejecutarse, con su entrada.
-
tool_execution_endLa herramienta terminó: su salida, si falló y cuánto tardó.
-
turn_endEl turno terminó, con los tokens que usó su llamada al modelo.
-
session_endLa ejecución terminó: completada, cancelada o con error.
Las llamadas a herramientas también llegan del proveedor en streaming, de a unos pocos caracteres de sus argumentos.
El bucle las vuelve a armar antes de ejecutarlas, así que en el bus una llamada a una herramienta aparece entera, en
assistant_message.
El elenco
El mismo elenco de siempre, sobre un escenario de rock.
- El Oráculo el modelo
- El cantante, arriba en la tarima. Su respuesta baja por el mástil como notas, una por palabra.
- El mástil el stream
- Cada text_delta es un puñado de notas. Sin streaming, no baja nada hasta el final.
- La persona de la guitarra tu app
- Toca cada nota cuando llega a los trastes, y las letras aparecen para el público.
- El cable el EventBus
- Todos los eventos corren por él. Solo se encienden los oyentes que se suscribieron a ese evento.
- Los oyentes agent.on(…)
- La pantalla de letras (
'text'), una casetera ('event') y un medidor de tokens (turn_end). - Astor el bucle
- Corre el bucle como siempre, y anuncia cada paso en el bus a medida que avanza.
- El escenario las herramientas
-
El borde del escenario, la consola de luces y la caja de pirotecnia:
crowd_mood,set_lightsypyro.
El código
Con astorlm: el proveedor siempre hace streaming, así que no hay nada que activar. Suscríbete con
agent.on() antes de llamar a run(): 'text' para cada pedacito de texto,
'tool-start' y 'tool-end' para las herramientas, 'event' para todo. Cada llamada
devuelve una función que cancela la suscripción.
Desde cero: el bucle del nivel 2 con dos cambios. El pedido usa stream: true y lee la
respuesta como server-sent events, pegando los pedacitos; y una lista de oyentes recibe cada paso a través de un ú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"]))
Qué vigilar
- El streaming no hace más rápido al modelo. La respuesta tarda lo mismo que antes; lo que cambia es cuándo ve la persona la primera palabra. Mide el tiempo hasta el primer token, no solo el tiempo total.
- Los oyentes miran; no manejan. Lo que devuelve un oyente se ignora, y el bucle no lo espera. Para bloquear una herramienta o cambiar lo que lee el modelo, necesitas un hook: el próximo nivel.
- Que los oyentes sean rápidos. El bus de astorlm los llama uno detrás del otro, dentro del bucle. El trabajo lento, como escribir en una base de datos, va en una cola a la que el oyente solo empuja.
- Un oyente roto falla en silencio. astorlm atrapa su error para que la ejecución siga, lo que también significa que nadie se entera. Registra los errores dentro de tus oyentes.
- Los pedacitos no respetan las palabras. Un text_delta puede terminar en mitad de una palabra o de una tabla en Markdown. Muestra lo que tienes hasta ahora y vuelve a dibujar a medida que llega más.
-
Deja que la persona lo frene. Cuando puede ver llegar la respuesta, va a querer cortar una que va
mal. Conecta un botón de stop a
agent.abort().