第 6 关
流式输出与事件
- user
- assistant
- tool_result
EventBus
问题
调用模型、等整个回复回来,在脚本里没问题。可当有人看着屏幕时就不行了:一段长回复可能要十秒、二十秒,这期间屏幕上什么都没有。智能体让情况更糟,因为一次请求在最终回复出现之前,可能要跑好几轮、调用好几个工具。
而且想跟进进度的不只是人。界面要文本,日志要每一步,计费要每一轮的 token。如果这些都得写进循环里,循环最后就得认识你应用里的每一块屏幕、每一个文件、每一个数据库。
解决方案
两块配套使用的东西。流式输出:让服务商以流的形式返回回复,模型每写出一小段文本就立刻发过来。总时间不变,但第一个字会在零点几秒内出现,而不是等到最后。
事件:循环在一条总线上广播每一步,并不在乎谁在听。你的应用让关心的部分去订阅:屏幕听文本,日志听所有事件,计量表听每一轮的结束。加一个监听者不用改动循环。下面是 astorlm 在一次正常运行中发出的事件:
-
turn_start一轮开始:循环马上要调用模型。
-
text_delta回复的一小段到了。每一轮有很多个。
-
assistant_message模型的消息完整了,包括文本和工具调用。
-
tool_execution_start一个工具马上要运行,带着它的输入。
-
tool_execution_end工具跑完了:它的输出、有没有失败、花了多久。
-
turn_end这一轮结束了,带着这次模型调用用掉的 token。
-
session_end这次运行结束了:完成、取消或出错。
工具调用从服务商那里也是流式过来的,参数一次来几个字符。循环会先把它们拼完整再运行,所以在总线上,工具调用是整个出现的,在 assistant_message 里。
角色
还是那群熟悉的角色,这次站上了摇滚舞台。
- 神谕者 模型
- 主唱,站在升降台上。它的回复变成音符顺着轨道滑下来,一个词一个音符。
- 轨道 流
- 每个 text_delta 是一小把音符。不用流式输出,直到最后什么都不会滑下来。
- 弹吉他的人 你的应用
- 音符一到音键就弹出来,歌词随之出现在观众面前。
- 线缆 EventBus
- 所有事件都从这里跑过。只有订阅了这个事件的监听者才会亮。
- 监听者 agent.on(…)
- 歌词屏(
'text')、一台录音机('event')和一个 token 计量表(turn_end)。 - Astor 循环
- 照常跑循环,一边走一边在总线上广播每一步。
- 舞台 工具
-
台口、灯光控制台和烟火箱:
crowd_mood、set_lights和pyro。
代码
使用 astorlm:服务商总是以流式返回,所以没有什么需要打开的。在调用 run() 之前用 agent.on() 订阅:'text' 接收每一小段文本,'tool-start' 和 'tool-end' 接收工具事件,'event' 接收所有事件。每次调用都会返回一个用来取消订阅的函数。
从零手写:第 2 关的循环,改两处。请求带上 stream: true,把回复当作 server-sent events 来读,再把一小段一小段拼起来;另外有一个监听者列表,通过同一个 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"]))
注意事项
- 流式输出不会让模型变快。回复花的时间和以前一样;变的是人看到第一个字的时间。要测首个 token 的时间,而不只是总时间。
- 监听者只看,不掌舵。监听者的返回值会被忽略,循环也不会等它。想拦下一个工具,或者改变模型读到的内容,你需要的是 hook:下一关。
- 让监听者跑得快。astorlm 的总线在循环内部一个接一个地调用它们。慢活儿,比如写数据库,应该放进一个队列,监听者只负责往里推。
- 坏掉的监听者会悄无声息地失败。astorlm 会接住它的错误,让运行继续,这也意味着没人会知道。在你的监听者里自己记日志。
- 一小段不一定是一个完整的词。一个 text_delta 可能停在一个词中间,或一张 Markdown 表格中间。先显示已有的内容,等更多内容到了再重画。
-
让人能叫停。人一旦能看着回复一点点到来,就会想打断一个跑偏的回复。把停止按钮接到
agent.abort()上。