服务化 v02:事件协议设计

服务化 v02:事件协议设计

第 14 章我们组装出了产品的 v01:一个能查天气、查时间、算数的 CLI Agent。它工作得很好,只要使用者愿意坐在终端前、只有一个人用、并且能忍受「回答要等全部生成完才一次性出现」。

这三条限制,正是本章要拆掉的东西。我们把产品从「进程内对话」升级为「服务」:Agent Core 与传输层分离,主循环不再对着终端说话,每一步决策都发成结构化的事件流;新加一个 HTTP/SSE 服务层,让任意多个客户端通过网络订阅这个事件流。这就是产品的 v02。

本章定义的事件协议跨章契约:ch16 的断线恢复、ch17 的 Web 前端、ch18 的权限审批,都会在这套协议上扩展。所以这一章的代码量不大,但每一行设计都值得仔细读。

本章目标

读完本章并做完配套练习后,你应该能够:

概念与动机:从进程内对话到服务

先看一个「错误答案」:把 CLI 挂到 HTTP 上

假设我们不学新东西,直接把 ch14 的 CLI 逻辑搬进一个 HTTP 处理器:

// 反面教材:把一次对话当作一次 HTTP 请求
server.on('request', async (req, res) => {
  const question = readBody(req);                 // 1. 读问题
  const result = await runAgent([system, question], { /* ... */ });
  res.end(JSON.stringify({ answer: result.finalText })); // 3. 一次性返回
});

能跑。但三分钟后就崩了:

  1. 没有流式:模型回答要 5–15 秒,客户端拿到的第一个字节是 15 秒后,体验像「卡死」。用户对着空白屏幕,不知道 Agent 是在思考、在调工具,还是真的死了。
  2. 没有会话:HTTP 请求是无状态的。同一个用户问完「上海的天气」再问「那北京呢?」,第二个请求不知道上下文里已经查过上海。每次对话都从零开始。
  3. 没有并发await runAgent 一个请求跑完才处理下一个。两个用户同时问,第二个排队;更糟的是,如果 runAgent 里还有共享可变状态,两个请求会互相污染。
  4. 客户端被绑定成「终端」:curl 能拿到最终答案,但拿不到「正在调用 get_weather」的中间过程。任何想做打字机效果、工具中间态渲染的界面都无从下手。

服务化需要回答的三件事

多客户端:一个服务同时服务 N 个客户端,每个客户端有自己的对话上下文,互不串流。→ 需要**会话**(session)概念,每个会话独立的状态与事件流。

流式:事件的到达是增量的:先看到「它开始调用工具了」,再看到「结果回来了」,最后看到「回答在逐字出现」。→ 需要一种从服务器单向推送给客户端的传输通道。

会话隔离:并发客户端之间不能互相看到对方的中间过程。→ 会话是事件流的分桶键,订阅按会话隔离。

三条加起来,指向一个经典组合:结构化事件流 + 会话 + SSE 传输

为什么选 SSE 而不是 WebSocket

流式推送有两个主流方案。WebSocket 是全双工的(双向都可以任意时刻发),但代价是独立的连接协议、需要专门处理握手、心跳、重连。而我们的场景有个显著的不对称:客户端发给服务器的是低频请求(提交一句话),服务器发给客户端的是高频推送(一堆事件):大部分流量是单向的。

Server-Sent Events(SSE) 正是为「单向推送」设计的:

WebSocket 的价值在「服务器也要主动给客户端发消息之外的交互」,比如在线游戏、协同编辑。我们的 Agent 服务不需要:客户端说话、服务器流式回答,一轮结束,仅此而已。选 SSE。

参照:业界怎么做这件事

「Agent 内核与传输层分离、事件流是中间产物」不是我们拍脑袋发明的。OpenAI Agents SDK 把「运行 agent」(agent loop、流式输出、会话延续)作为独立概念,与具体的前端/传输实现解耦:运行中的每一步都可以作为流式事件(run item 的增量)被上层消费(OpenAI Agents SDK · Running agents)。《从 0 开始构建 AI 智能体》的 FunHarness 项目也把差异化的流式事件(文本流、工具调用流、审批流各自成事件)作为产品演进的骨架(hyyhf/agent-book-code)。我们借鉴的是「事件驱动分层」这个机制,事件的具体形状是我们自己的(下面逐条定义)。

事件模型设计:为什么需要「统一事件」

从 trace hook 到正式协议

ch14 的 onEvent 是这样的:类型是 think / act / observe / done,字段为 CLI 打印而定制(finalTextiterationstruncated)。CLI 用它打日志,仅此而已。

服务化要求同一份 core 输出被多个消费方使用:终端打印、SSE 推送、Web 前端渲染、未来的存储与回放。日志用的 trace 不够用了,我们需要一个正式的、自描述的、可序列化的事件协议

七种事件,逐条拆

{ type: 'text_delta',       sessionId, delta }
{ type: 'tool_call_start',  sessionId, callId, name, args }
{ type: 'tool_result',      sessionId, callId, ok, result? | error? }
{ type: 'approval_request', sessionId, requestId, toolName, args, risk }
{ type: 'approval_result',  sessionId, requestId, approved, reason? }
{ type: 'done',             sessionId, usage? }
{ type: 'error',            sessionId, message }
类型核心字段含义谁发出
text_deltadelta一段助手文本(流式输出的一块)loop(停止轮)
tool_call_startcallId, name, args循环决定执行某个工具loop(执行前)
tool_resultcallId, ok, result?/error?工具执行结果;失败也是数据loop(执行后)
approval_requestrequestId, toolName, args, risk请求用户审批某个工具(risk: auto/suggest/approve 三档)ch18 权限层
approval_resultrequestId, approved, reason?审批结果回传ch18 权限层
doneusage?一轮运行正常结束(可选 token 用量)loop
errormessage出错(如 provider 调用失败)loop

几个值得停下来的设计点:

传输封装:SSE 帧格式

事件本身是 JSON 对象;上线时套一层 SSE 帧:

event: text_delta
data: {"type":"text_delta","sessionId":"s1","delta":"上海"}

(帧格式细节,字段、多行 data 拼接、注释行保活、断线重连,详见 MDN · Using server-sent events。ch03 你已经手写过 SSE 客户端解析器,这一章补上服务端那一半。)

手写实现:把事件协议装进产品

一、升级主循环:lib/loop.mjs

ch14 的循环骨架一行没动(think → act → observe、两个停止条件),变的只有它发什么。原来发 trace 日志,现在发协议事件。关键是一处收敛。所有事件从同一个 emit 出去,统一盖章:

export async function runAgent(initialMessages, { callLlm, tools, sessionId = 'default', maxIterations = DEFAULT_MAX_ITERATIONS, onEvent }) {
  const messages = [...initialMessages];
  let iterations = 0;

  // The single chokepoint through which every protocol event leaves the loop.
  const emit = (event) => onEvent?.({ sessionId, ...event });

  while (iterations < maxIterations) {
    const toolList = modelTools(tools);
    let response;
    try {
      response = await callLlm(messages, toolList);
    } catch (err) {
      const message = err instanceof Error ? err.message : String(err);
      emit({ type: 'error', message });   // provider 挂了:发 error 事件…
      throw err;                          // …然后交给传输层决定怎么办
    }

    const assistant = { role: 'assistant', content: response.content };
    if (response.toolCalls.length > 0) assistant.tool_calls = response.toolCalls;
    messages.push(assistant);

    if (response.toolCalls.length === 0) {
      // 停止轮:文本流出去,done 收尾(带上 provider 报的 usage)
      if (response.content) emit({ type: 'text_delta', delta: response.content });
      emit({ type: 'done', usage: response.usage });
      return { messages, finalText: response.content, iterations: iterations + 1 };
    }

    for (const call of response.toolCalls) {
      // act:执行前发 tool_call_start
      emit({ type: 'tool_call_start', callId: call.id, name: call.name, args: call.arguments });
      const outcome = executeTool(tools, call.name, call.arguments);
      if (outcome.ok) {
        const content = stringifyToolResult(outcome.result);
        messages.push({ role: 'tool', tool_call_id: call.id, content });
        emit({ type: 'tool_result', callId: call.id, ok: true, result: content });
      } else {
        // observe:失败也是数据
        messages.push({ role: 'tool', tool_call_id: call.id, content: outcome.error });
        emit({ type: 'tool_result', callId: call.id, ok: false, error: outcome.error });
      }
    }
    iterations += 1;
  }
  // 保险丝:最后一段文本仍以 text_delta 流出,绝不让用户面对空答案
  const lastAssistant = [...messages].reverse().find((m) => m.role === 'assistant');
  const finalText = lastAssistant?.content ?? '';
  if (finalText) emit({ type: 'text_delta', delta: finalText });
  emit({ type: 'done' });
  return { messages, finalText, iterations };
}

注意两个「升级点」:

  1. emit 是唯一出口{ sessionId, ...event } 一行把 sessionId 盖到所有事件上。循环本身不知道「会话」是什么,是调用方把 sessionId 传进来当选项的。core 保持无状态、无感知,这正是 ch14 分层红线的兑现。
  2. error 事件的取舍:provider 抛异常时,先发 error 事件、再 rethrow。为什么不是「吞掉错误继续循环」?因为 provider 挂了重试多少次都大概率还是挂,继续循环只会空转烧钱;为什么不是「不发事件直接抛」?因为客户端需要知道发生了什么:error 事件让事件流完整(以 error 收尾,不会半截断掉)。rethrow 之后由传输层(server)决定怎么收尾。这是「core 发事件、界面做决策」的又一例。

二、事件协议模块:lib/events.mjs

新的 core 模块,三件事:类型定义、SSE 序列化、会话缓冲。全是纯函数或纯数据结构,浏览器沙箱里也能跑,它就是练习的判题对象

序列化与反序列化,互为逆运算:

// One event -> one SSE frame:
//   event: <type>\n
//   data:   <json>\n
//   \n
export function encodeSse(event) {
  const data = JSON.stringify(event);
  return `event: ${event.type}\ndata: ${data}\n\n`;
}

// One SSE frame -> the event object carried in the data field. The `event:`
// line names the type for the browser's EventSource, but the data JSON is the
// source of truth (it carries type too) — so parseSse simply returns the JSON.
// Multi-line data: SSE concatenates consecutive `data:` lines with a newline.
export function parseSse(frame) {
  const dataLines = [];
  for (const line of frame.split('\n')) {
    if (line.startsWith('data:')) dataLines.push(line.slice(5).replace(/^\s+/, ''));
  }
  return JSON.parse(dataLines.join('\n'));
}

parseSse 只认 data: 行,这正是「data.type 是事实来源」的落地:event: 行写错或漏掉,解析结果不受影响。

会话缓冲是服务端事件分发的核心数据结构:

export function createSessionBuffer() {
  const sessions = new Map(); // sessionId -> { events: [], listeners: Set<fn> }

  function bucket(sessionId) {
    let s = sessions.get(sessionId);
    if (!s) { s = { events: [], listeners: new Set() }; sessions.set(sessionId, s); }
    return s;
  }

  return {
    append(sessionId, event) {
      const s = bucket(sessionId);
      s.events.push(event);
      for (const listener of s.listeners) listener(event);
      return event;
    },
    events(sessionId) {                    // 回放:该会话的全部历史(副本)
      return (sessions.get(sessionId)?.events ?? []).slice();
    },
    subscribe(sessionId, listener) {       // 实时:注册监听,返回退订函数
      bucket(sessionId).listeners.add(listener);
      return () => bucket(sessionId).listeners.delete(listener);
    },
    sessions() { /* ... { sessionId, eventCount, lastEvent } ... */ },
  };
}

一个缓冲同时干两件事:回放events(),迟到的订阅者从第一帧补起)和实时分发subscribe(),新事件推给所有在线订阅者)。这是 GET /events 的全部后台。

三、HTTP/SSE 服务:lib/server.mjs

服务层是新的界面层(ch14 的 CLI 是另一个界面)。三个端点:

方法路径做什么
POST/chat提交一轮对话;立即返回,Agent 在后台跑
GET/events?sessionId=xSSE 事件流:先回放、再实时
GET/sessions会话索引(id + 事件数 + 最后事件)

POST /chat 的核心决策是不等待:提交后立刻回 { sessionId, status: 'started' },Agent 在后台跑,进度全部走事件流。如果像反面教材那样 await 整个 runAgent 才返回,客户端就退回了「15 秒空白」,流式的意义全没了:

async function handleChat(req, res) {
  const { sessionId = 'default', text } = await readJsonBody(req);
  // ...校验、取会话、把用户消息追加进会话历史...
  if (session.status === 'running') {          // 同一会话并发第二轮:拒绝
    res.writeHead(409, ...); return;
  }
  startRun(session, sessionId);                // fire-and-forget:立即返回
  res.writeHead(200, { 'content-type': 'application/json' });
  res.end(JSON.stringify({ sessionId, status: 'started' }));
}

startRun 把 loop 的 onEvent 接到缓冲上,事件自然流向所有订阅者:

function startRun(session, sessionId) {
  session.status = 'running';
  runAgent(session.messages, {
    callLlm, tools, sessionId,
    maxIterations: config.maxIterations,
    onEvent: (event) => buffer.append(sessionId, event),   // 事件进缓冲 = 广播
  })
    .then((result) => { session.messages = result.messages; session.status = 'done'; })
    .catch(() => { session.status = 'error'; });           // loop 已发过 error 事件
}

GET /events 处理一个关键竞态:订阅者来晚了怎么办

function handleEvents(req, res) {
  const sessionId = url.searchParams.get('sessionId') || 'default';
  res.writeHead(200, { 'content-type': 'text/event-stream', /* no-cache */ });

  // Subscribe FIRST, then replay: both happen in this same synchronous tick,
  // so no event can be appended in between — a late subscriber never misses
  // the beginning (and never sees a duplicate).
  let done = false;
  const close = () => { if (done) return; done = true; clearInterval(keepAlive); unsubscribe(); res.end(); };
  const unsubscribe = buffer.subscribe(sessionId, (event) => {
    res.write(encodeSse(event));
    if (event.type === 'done' || event.type === 'error') close();
  });
  const keepAlive = setInterval(() => res.write(': keep-alive\n\n'), 15_000); // 注释行保活
  req.on('close', close);

  for (const event of buffer.events(sessionId)) {          // 回放已发生的历史
    res.write(encodeSse(event));
    if (event.type === 'done' || event.type === 'error') { close(); return; }
  }
}

三行关键代码:

多会话是怎么来的?服务端一个 Map<sessionId, Session> 存每会话的消息历史,缓冲按 sessionId 分桶,客户端按自己的 sessionId 订阅,三个层次各管一截,互不越位。注意会话状态在界面层(server),不在 core:loop 每次运行拿到的是完整历史(ch14 的契约),它自己不保存任何东西。

四、跑起来看事件流

一次完整的两轮对话(mock 模型:第一轮并行要两个工具,第二轮给最终答案),订阅者看到的事件流长这样:

--- session s1: GET /events?sessionId=s1 (SSE) ---
[1] event: tool_call_start
    data: {"sessionId":"s1","type":"tool_call_start","callId":"call_time","name":"get_time","args":{}}
[2] event: tool_result
    data: {"sessionId":"s1","type":"tool_result","callId":"call_time","ok":true,"result":"2026-08-11T13:12:07.601Z"}
[3] event: tool_call_start
    data: {"sessionId":"s1","type":"tool_call_start","callId":"call_weather","name":"get_weather","args":{"location":"Shanghai"}}
[4] event: tool_result
    data: {"sessionId":"s1","type":"tool_result","callId":"call_weather","ok":true,"result":"上海:晴,27°C,微风,紫外线中等"}
[5] event: text_delta
    data: {"sessionId":"s1","type":"text_delta","delta":"上海今天晴,27°C,微风,适合出门,记得防晒。"}
[6] event: done
    data: {"sessionId":"s1","type":"done","usage":{"inputTokens":42,"outputTokens":24}}
← s1 stream closed

整个序列是确定的:工具调用成对出现(callId 配对)、失败是数据、停止轮才发文本、done 带用量。第二个会话(问北京)走完全一样的流程,但事件全部落在自己的 sessionId 桶里,两个客户端同时在线,互不串流。GET /sessions 能看到两个会话各自 6 条事件:

GET /sessions →
{
  "sessions": [
    { "sessionId": "s1", "eventCount": 6, "lastEvent": "done" },
    { "sessionId": "s2", "eventCount": 6, "lastEvent": "done" }
  ]
}

这张图把 v02 的全貌画出来。core 在中间,两个界面(终端 CLI 与 HTTP/SSE 服务)各占一边:

flowchart TD
  subgraph clients["客户端(任意多个)"]
    WEB["浏览器 / Web 前端(ch17)"]
    CLI["终端 CLI(ch14 交互面,进程内直连 core)"]
    CURL["curl / 脚本(HTTP 客户端)"]
  end

  subgraph interface["界面层 · interface"]
    SRV["lib/server.mjs<br/>HTTP/SSE 服务:POST /chat · GET /events · GET /sessions"]
    CLIAPP["lib/cli.mjs<br/>进程内组合根 + 事件打印"]
  end

  subgraph core["核心层 · core(无状态 / 依赖注入)"]
    EV["lib/events.mjs<br/>事件协议:SSE 序列化 + 会话缓冲"]
    LOOP["lib/loop.mjs<br/>主循环:发出协议事件(含 sessionId)"]
    PROV["lib/provider.mjs<br/>传输:内部消息 ↔ 厂商 wire 格式"]
    TOOLS["lib/tools.mjs<br/>工具:注册表 + error-as-data"]
  end

  MOCK["mock LLM / 真实供应商(HTTP)"]

  CURL -->|"POST /chat / GET /events"| SRV
  WEB -. 订阅同一事件流 .-> SRV
  SRV -->|onEvent 注入| LOOP
  SRV -->|事件进缓冲| EV
  CLIAPP -->|"runAgent(..., onEvent)"| LOOP
  CLIAPP --> EV
  LOOP -->|callLlm 注入| PROV
  LOOP -->|registry 注入| TOOLS
  PROV -->|fetch| MOCK

再放大一次对话,看事件如何在各个模块间流转:

sequenceDiagram
  participant C as 客户端
  participant SRV as lib/server.mjs
  participant EV as events.mjs 缓冲
  participant LOOP as lib/loop.mjs
  participant PROV as lib/provider.mjs
  participant LLM as mock LLM
  participant TOOL as 工具注册表

  C->>SRV: POST /chat { sessionId, text }
  SRV->>SRV: 追加 user 消息 → startRun(fire-and-forget)
  SRV-->>C: 200 { sessionId, status: "started" }

  C->>SRV: GET /events?sessionId=s1(SSE)
  SRV->>EV: subscribe(s1, 写入响应流)
  SRV->>EV: 回放 events(s1)

  LOOP->>PROV: callLlm(历史, 工具清单)
  PROV->>LLM: POST /v1/chat/completions
  LLM-->>PROV: tool_calls ×2(并行)
  PROV-->>LOOP: { content, toolCalls }
  loop 每个 tool_call
    LOOP->>EV: tool_call_start(callId, name, args)
    EV-->>C: 实时推送(SSE 帧)
    LOOP->>TOOL: executeTool(name, args)
    TOOL-->>LOOP: { ok, result } | { ok, error }
    LOOP->>EV: tool_result(callId, ok, result|error)
    EV-->>C: 实时推送
  end
  LOOP->>PROV: callLlm(带 tool 结果的历史)
  PROV-->>LOOP: { content: 最终答案 }
  LOOP->>EV: text_delta(delta) → done(usage)
  EV-->>C: 实时推送 ×2
  EV-->>SRV: done 到达 → 关闭流
  SRV-->>C: SSE 流结束

时序图里有两点值得再强调:客户端早在事件产生之前就完成了订阅(POST 先回 200、GET 后订阅、再回放,所以第一帧 tool_call_start 不会丢);所有事件只从 loop 经缓冲单向流出,传输层不改写事件,只是把缓冲里的事件原样 encodeSse 转发。

常见坑与失败模式

坑一:SSE 输出被缓冲,事件「攒着」不出去。 服务器写了 res.write(),但客户端半天收不到,很可能是中间层(反向代理/网关)在缓冲响应。SSE 要求流式输出:响应头必须是 Content-Type: text/event-stream + Cache-Control: no-cache,代理类中间件常需要 X-Accel-Buffering: no 这类开关;没事件时还要靠注释行: 开头的一行 + 空行)保活,防止连接被空闲超时掐断(MDN · Sending events from the server)。

坑二:连接断开后资源泄漏。 客户端关了页面,订阅监听器还挂在缓冲上、保活定时器还在跑,时间一长,服务器上堆满死连接。正确姿势是 req.on('close') 里统一清理:清定时器、退订、标记关闭。本章 close() 里三件事一口气做完,就是为了防这个。

坑三:订阅与回放的顺序搞反。 先回放、再订阅,中间产生的事件就丢了;回放和订阅不在同一个同步 tick 里完成,事件可能在两者之间插队导致重复。正确姿势是先订阅、再回放,且两步连续执行(JS 单线程保证中间插不进 append),迟到订阅者既不错过开头、也不看到重复。

坑四:把会话状态塞进 core。 loop 里自己存历史、自己管理 sessionId,core 从此带着状态,多会话直接互相污染,ch16 换持久化时整个 core 要重写。正确分层:core 无状态(每次运行喂完整历史),会话状态是界面层的事。判断标准一句话:把 server.mjs 删掉,loop/events 还能不能单独跑、单独被测试?

坑五:POST /chat 等回答跑完才返回。 那就退回了「15 秒空白」,流式的意义全无。正确姿势:提交立即返回,进度走事件流;同一会话正在跑时再提交,回 409 session busy(顺序保证:不让同一会话的两轮事件交织)。

坑六:文本事件乱发。 模型「一边给文本一边要工具」的轮次也发 text_delta,前端会把中间草稿和最终回答拼在一起渲染。文本只在停止轮流出(ch14 里这叫「正常停止条件」),有 tool_calls 的轮次文本只留在历史里。

小结

下一章(ch16)我们做断线恢复与后台运行 v03:给事件加序号、接上 ch10 的事件存储,让「断开连接后 Agent 继续跑、重连后从断点续传」成为可能。你会发现,因为本章把事件定义成了可序列化、有序、带会话键的协议,ch16 几乎不用改任何事件形状,只需在流的外面套一层序号。再下一章(ch17)我们写 Web 前端 v04:浏览器消费这套 SSE 事件流,把 text_delta 渲染成打字机、把 tool_call_start/tool_result 渲染成工具调用卡片。事件协议越稳,前端越简单。先把本章练习做完:亲手实现 SSE 序列化、事件版主循环、会话缓冲这三层。

延伸阅读

完成阅读,去做练习 →