> Agent Harness Patterns 第 6 关，一条讲解 AI 智能体工作原理的模式路线。网页版：https://harnesspatterns.dev/zh/patterns/streaming · 全部模式（英文）：https://harnesspatterns.dev/llms.txt

# 流式输出与事件

模型写回复是几个 token 几个 token 地写。流式输出把每一小段一写出来就交给你，而循环发出的事件会告诉应用的其他部分：智能体正在做什么，就在它做的时候。

## 问题

调用模型、等整个回复回来，在脚本里没问题。可当有人看着屏幕时就不行了：一段长回复可能要十秒、二十秒，这期间屏幕上什么都没有。智能体让情况更糟，因为一次请求在最终回复出现之前，可能要跑好几轮、调用好几个工具。

而且想跟进进度的不只是人。界面要文本，日志要每一步，计费要每一轮的 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()` 收到每一步。

**使用 astorlm**

```ts
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
```

**TypeScript**

```ts
// 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))
```

**Python**

```python
# 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()` 上。

## 相关模式

- [2 · 智能体循环](https://harnesspatterns.dev/zh/patterns/agent-loop.md)
- [7 · 钩子](https://harnesspatterns.dev/zh/patterns/hooks.md)
- [15 · 可观测性与评估](https://harnesspatterns.dev/zh/patterns/observability.md)
- [17 · 人在回路](https://harnesspatterns.dev/zh/patterns/human-in-the-loop.md)
