青蛙小白
博客 / 2026/08

Pi SDK 学习笔记(二):事件流与 SSE

2026/08/04 · — 字 · 阅读约 — 分钟 ·
目录

上一篇把最小会话跑起来了:ModelRuntime + createAgentSession 在进程内订阅 text_deltafinallydispose()。这一篇把它搬到 HTTP 上——用一个稳定的公开事件协议把内部事件流喂给浏览器,并处理好 SSE 的建流、心跳、断连和错误帧。

这一篇是把「进程内的 AgentSession」变成「跨信任边界的服务」的第一步,核心是事件映射和安全边界。

事件层次:一次 prompt 会经历什么

先理清 Pi 内部的事件层次。一次 session.prompt(text) 可能经历多个 turn:

agent_start
  turn_start
    message_start            ← 模型开始输出
    message_update(text_delta) …   ← 流式增量
    tool_execution_start     ← 模型决定调工具
    tool_execution_update
    tool_execution_end
    message_end
  turn_end                   ← 工具结果送回模型,可能再来一个 turn
  turn_start …
agent_end

这里有几条必须记住的事实,它们决定了后面所有设计:

  • message_update 里的 assistantMessageEvent 既可能是 text_delta(可展示),也可能是 thinking_delta不应默认公开)。
  • tool_execution_* 事件里带 argspartialResultresult,可能含绝对路径、文件内容、命令参数,不能原样跨越 HTTP 信任边界
  • queue_updatesteering / followUp 是排队的用户消息正文,同样是 payload。
  • agent_endmessageswillRetry,对客户端只需要一个「done」信号。
  • 一次正常 run 只应产生一个业务完成信号。

Pi 的内部事件类型是 AgentSessionEvent(在 @earendil-works/pi-coding-agent 里),它是 AgentEvent 的扩展联合。直接把它 JSON.stringify 丢给浏览器是错的——上游类型会随版本演进,这会把内部细节泄漏给了客户端。

映射:把内部事件变成公开协议

这是这一篇最关键的一个小文件。思路是:写一个穷举 switch,只挑允许的字段,未知事件默认丢弃,而不是透传。

// lessons/02-events-and-sse/event-mapper.ts
export type PublicEvent =
  | { type: "text-delta"; delta: string }
  | { type: "tool-start"; toolCallId: string; toolName: string }
  | { type: "tool-update"; toolCallId: string; toolName: string }
  | { type: "tool-end"; toolCallId: string; toolName: string; isError: boolean }
  | { type: "queue"; steering: number; followUp: number }
  | { type: "done" }
  | { type: "error"; message: string };

export function toPublicEvent(event: AgentSessionEvent): PublicEvent | undefined {
  switch (event.type) {
    case "message_update":
      if (event.assistantMessageEvent.type !== "text_delta") return undefined;
      return { type: "text-delta", delta: event.assistantMessageEvent.delta };
    case "tool_execution_start":
      return { type: "tool-start", toolCallId: event.toolCallId, toolName: event.toolName };
    case "tool_execution_update":
      return { type: "tool-update", toolCallId: event.toolCallId, toolName: event.toolName };
    case "tool_execution_end":
      return { type: "tool-end", toolCallId: event.toolCallId, toolName: event.toolName, isError: event.isError };
    case "queue_update":
      return { type: "queue", steering: event.steering.length, followUp: event.followUp.length };
    case "agent_end":
      return { type: "done" };
    default:
      return undefined;
  }
}

PublicEvent 是一个我们定义的全新的联合类型,不是 AgentSessionEvent 的子集别名。 这意味着上游改名、加字段,我们的协议不变;客户端可以放心地按这个类型写解析器。

thinking_deltareturn undefined,被显式过滤——message_update 分支只放行 text_delta。工具事件只保留 toolCallId + toolName(+ isError),argspartialResultresult 全部丢弃。这就是「工具状态可见,工具内容不可见」:toolCallId + toolName 给前端做 loading 动画足够了,参数和结果不能过边界。queue_update 只给数量(steering.length / followUp.length),不给排队消息的文本。agent_enddone,不带 messages、不带 willRetry

最后那条 defaultundefined 是白名单语义:上游新增的事件类型,在我们没有显式映射之前,默认不透传。这比「黑名单」安全得多。

测试里有一组断言专门盯这条边界,给 tool_execution_start{ path: "/secret/.env" },断言序列化结果里不含敏感信息:

const serialized = JSON.stringify([toPublicEvent(start), toPublicEvent(update), toPublicEvent(end)]);
expect(serialized).not.toContain("/secret/.env");
expect(serialized).not.toContain("secret contents");

SSE 帧的编码也在这里:

export function encodeSse(event: PublicEvent): string {
  return `event: ${event.type}\ndata: ${JSON.stringify(event)}\n\n`;
}

用命名事件(event: text-delta),客户端可以用 EventSource 或手写解析按类型分发。data 是单行 JSON——注意 JSON.stringify 会把 \n 转义成 \\n,所以文本里的换行不会破坏 SSE 帧结构,一条 data: 行就是一个完整事件。帧以空行 \n\n 结束,这是 SSE 规范要求的帧分隔符。

SSE 路由的四个要点

buildLessonServer(session) 起一个最小 Fastify,只有 /health/prompt。`它要做对四件事:建流、心跳、断连中止、错误帧。

// lessons/02-events-and-sse/app.ts(节选)
app.post<{ Body: PromptBody }>("/prompt", { schema: { body: promptBodySchema } },
async (request, reply) => {
  const text = request.body.text.trim();
  if (!text) return reply.code(400).send({ error: "text must not be blank" });

  reply.hijack();                       // ① 接管原生响应
  reply.raw.writeHead(200, {
    "content-type": "text/event-stream; charset=utf-8",
    "cache-control": "no-cache, no-transform",
    connection: "keep-alive",
    "x-accel-buffering": "no",          // ② 关掉 Nginx 等反向代理的缓冲
  });

  let finished = false;
  const heartbeat = setInterval(() => {
    if (!reply.raw.destroyed) reply.raw.write(": heartbeat\n\n");  // ③ 注释心跳
  }, 15_000);

  const onClose = () => {               // ④ 客户端断开 → 中止 Agent
    clearInterval(heartbeat);
    if (!finished) void session.abort().catch(() => undefined);
  };
  reply.raw.on("close", onClose);

  const send = (event: PublicEvent) => {
    if (!reply.raw.destroyed) reply.raw.write(encodeSse(event));
  };
  const unsubscribe = session.subscribe((event) => {
    const mapped = toPublicEvent(event);
    if (mapped) send(mapped);
  });

  try {
    await session.prompt(text);
  } catch {
    send({ type: "error", message: "Agent request failed" });  // ⑤ 通用错误
  } finally {
    finished = true;
    clearInterval(heartbeat);
    unsubscribe();
    reply.raw.off("close", onClose);
    if (!reply.raw.destroyed) reply.raw.end();
  }
});

逐条解释。

reply.hijack():SSE 是长连接、手动写帧,Fastify 默认的「序列化 + 发完就结束」模型不适用。hijack 把底层 reply.raw 交出来,后面所有写入都直接操作 Node 原生 ServerResponse

x-accel-buffering: no:开发时直连没问题,但生产常挂在 Nginx 后面。Nginx 默认会缓冲响应,这会让 SSE 帧攒一批才发出去,破坏「逐字到达」的体验。这个 header 是给 Nginx 看的信号:别缓冲我。cache-control: no-cache, no-transform 同理,防止中间代理改写或缓存流。

③ 心跳用 : heartbeat\n\n:SSE 里以 : 开头的行是注释,客户端会忽略,但它能保持连接活跃、防止空闲超时被代理掐断。15 秒是一个常见值。

④ 断连必须 abort():浏览器关了标签页,服务端还在 await session.prompt(text) 的话,模型会继续跑完、继续烧 token——这种没人再看的 run 叫孤儿 run。reply.rawclose 事件是客户端断开的信号,接到它就调 session.abort() 让 Agent 停下来。finished 标志位避免在正常结束后重复 abort。

⑤ 错误帧,而不是改 status:这是 SSE 一个容易忽略的约束——一旦 writeHead(200) 把状态行发出去了,就再也改不了 HTTP status 了。所以建流之后的任何错误(provider 报错、工具抛异常)都只能通过事件帧告诉客户端:

} catch {
  send({ type: "error", message: "Agent request failed" });
}

而且消息是通用的,不是 error.message。原始错误里可能有 provider 名、文件路径、prompt 片段——这些都是不该跨信任边界的东西。测试里有专门一条:

// session.prompt 抛 "provider response contains a secret path"
expect(response.body).toContain('event: error\ndata: {"type":"error","message":"Agent request failed"}\n\n');
expect(response.body).not.toContain("secret path");

输入校验在 hijack 之前:空文本直接 400,不建流、不订阅、不 prompt。这也是「能在建流前拒绝就在建流前拒绝」的原则——建流后想改 status 已经不可能了。

用真模型跑一次 smoke

app.test.ts 用 mock session 测的是协议和路由行为,不消耗模型额度。但 SSE 这东西,「真的能逐帧到达」这件事得用真模型 + curl -N 亲自跑一遍才算数——这就是 smoke.ts 的职责。

这里关键是资源隔离:专用 agentDir.data/pi-agent,权限 0o700)跟真实 ~/.pi/agent 隔开;models.jsonapiKey$COURSE_OPENAI_API_KEY 引用环境变量,不落 key 值;然后还是第 1 课学过的三道闸门——ModelRuntime.creategetModelcheckAuth,都过了才 createAgentSession,并显式用 SessionManager.inMemory() 避免第一个实验就继承真实会话状态。

验证方式:

npm run lesson:02:smoke
# 另一个终端
curl -N -X POST http://127.0.0.1:3001/prompt \
  -H 'content-type: application/json' \
  -d '{"text":"用一句话说明当前工作目录有什么文件"}'

-N 关掉 curl 的缓冲,能看到 event: text-delta 一帧一帧地到,最后 event: done。这才是 SSE 逐帧到达该有的样子。

作业:事件时间线记录器

课程作业是写一个事件时间线记录器,统计每类事件出现次数和 turn 数量,只保留类型、时间和耗时,不保存 message 正文。

我把「类型、时间、耗时」落到一个极简的条目类型上:

export interface TimelineEntry {
  type: string;       // 事件类型
  time: number;       // 记录时的 epoch 毫秒
  durationMs?: number;// 仅对「结束事件」:从匹配的「开始事件」起的耗时
}

TimelineEntry 只有这三个字段,没有任何 payload 容器。这是安全约束的物理体现——结构上就不可能装下 message 正文或工具参数。把安全约束落到类型定义里,比靠注释约束更可靠。

核心 record(event) 也是一个穷举 switch,和 toPublicEvent 同构,但目的不同:这里每个事件都记一条 TimelineEntry(包括被映射器丢弃的 thinking_deltaqueue_update 等),只是不记它们的 payload:

record(event) {
  const t = now();
  switch (event.type) {
    case "turn_start": turnStart = t; turnCount += 1; push("turn_start", t); break;
    case "turn_end":
      completedTurns += 1;
      if (turnStart !== undefined) turnDurationsMs.push(t - turnStart);
      push("turn_end", t, turnStart === undefined ? undefined : t - turnStart);
      turnStart = undefined;                     // 不读 event.message / toolResults
      break;
    case "tool_execution_start":
      toolStarts.set(event.toolCallId, t);       // 按 id 索引,支持并发
      push("tool_execution_start", t);           // 不读 event.args
      break;
    case "tool_execution_end": {
      const start = toolStarts.get(event.toolCallId);
      toolStarts.delete(event.toolCallId);
      const dur = start === undefined ? undefined : t - start;
      completedToolCalls += 1;
      if (dur !== undefined) toolDurationsMs.push(dur);
      push("tool_execution_end", t, dur);        // 不读 event.result
      break;
    }
    default:
      // message_update(含 thinking_delta)、queue_update、bash_execution_update …
      // 一律只记 type + time,绝不保留 payload。
      push(event.type, t); break;
  }
}

几个细节值得记一笔。

工具调用按 toolCallId 索引(Map),不是单变量。 同一个 turn 里模型可能并发发起多个工具调用:tool_execution_start A → tool_execution_start B → tool_execution_end A → tool_execution_end B,用单变量会错配。测试里专门有一例「overlapping tool calls by toolCallId」断言 A=40ms、B=50ms。消息则用栈(messageStartStack),理论上 message 一般不嵌套,但用栈比单变量更鲁棒,pop 自动配对最近的 message_start

default 仍然只记 type + time,和映射器一样是白名单语义。 比如 bash_execution_updatedelta(命令输出)是 payload,记了就泄漏了,所以只记类型。

now 可注入。 默认 Date.now,测试里用一个手动 tick(ms) 的确定性时钟,这样耗时可精确断言(turnDurationsMs === [60] 这种)。

插曲:smoke 为什么会打印 ~/.agents/skills

跑 smoke 时,控制台里除了课程自己的输出,还混进了几个跟本课无关的 skill——比如本机装的一个 ticket 工具和一个终端会话工具——它们来自 ~/.agents/skills。排查过程是一次很好的「读源码理解资源发现」的练习。

链路是这样的。smoke.tscreateAgentSession({...})没有传 resourceLoader。看 SDK 的 core/sdk.js

if (!resourceLoader) {
  resourceLoader = new DefaultResourceLoader({ cwd, agentDir, settingsManager });
  await resourceLoader.reload();
}

于是 DefaultResourceLoader.reload() 跑起来,里面调用 packageManager.resolve(),而 core/package-manager.js 里有这一段:

const userAgentsSkillsDir = join(getHomeDir(), ".agents", "skills");
const userAgentsBaseDir = dirname(userAgentsSkillsDir);
addResources("skills",
  collectAutoSkillEntries(userAgentsSkillsDir, "agents"),
  userAgentsMetadata, userOverrides.skills, userAgentsBaseDir);

也就是说,~/.agents/skills 是一个「用户/全局」作用域的 skill 目录,无论项目是否受信任都会被扫描。对比之下,项目里的 .agents/skills(沿 cwd 向上找)要求 projectTrusted === true 才扫。

为什么会「打印出来」?因为 smoke.ts 里有这一行:

console.log("[lesson-02] system prompt:\n" + session.systemPrompt);

AgentSession.systemPrompt 是由 resource loader 的 skills 列表拼进来的。于是 ~/.agents/skills 下的 skill 就顺着「loader 扫描 → skills 列表 → 系统提示词」这条路,出现在了 smoke 的日志里。

这说明什么:smoke.ts 传了隔离的 agentDir,这控制了 auth.json / models.json 的位置,但没传 resourceLoader / settingsManager,所以资源发现仍然会落到全局用户目录。想要一个「干净」的 smoke,需要自己构造一个 DefaultResourceLoader(比如开 noSkills: true)。

这其实是后边要学习的内容。第 2 课的 smoke 暂时不必处理,但要看得懂这个现象从哪来。

小结

这一篇最重要的不是「跑通了一个 SSE 接口」,而是事件映射、SSE 的几个工程要点、以及把安全约束落到类型结构里这三块。事件映射:一个由自定义的新公开类型,穷举 switch,未知事件默认丢弃。SSE 的工程要点:hijack、注释心跳、断连 abort、建流后如果出错只能用错误帧。

代码放在本地仓库 lessons/02-events-and-sse/ 下:event-mapper.ts(映射 + 编码)、app.ts(SSE 路由)、smoke.ts(真模型验证)、timeline-recorder.tshomework.ts(作业)。

参考

代码固定使用 @earendil-works/[email protected],避免上游快速变化破坏可复现性。

评论