上一篇把最小会话跑起来了:ModelRuntime + createAgentSession 在进程内订阅 text_delta、finally 里 dispose()。这一篇把它搬到 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_*事件里带args、partialResult、result,可能含绝对路径、文件内容、命令参数,不能原样跨越 HTTP 信任边界。queue_update的steering/followUp是排队的用户消息正文,同样是 payload。agent_end带messages和willRetry,对客户端只需要一个「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_delta 走 return undefined,被显式过滤——message_update 分支只放行 text_delta。工具事件只保留 toolCallId + toolName(+ isError),args、partialResult、result 全部丢弃。这就是「工具状态可见,工具内容不可见」:toolCallId + toolName 给前端做 loading 动画足够了,参数和结果不能过边界。queue_update 只给数量(steering.length / followUp.length),不给排队消息的文本。agent_end → done,不带 messages、不带 willRetry。
最后那条 default → undefined 是白名单语义:上游新增的事件类型,在我们没有显式映射之前,默认不透传。这比「黑名单」安全得多。
测试里有一组断言专门盯这条边界,给 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.raw 的 close 事件是客户端断开的信号,接到它就调 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.json 里 apiKey 用 $COURSE_OPENAI_API_KEY 引用环境变量,不落 key 值;然后还是第 1 课学过的三道闸门——ModelRuntime.create → getModel → checkAuth,都过了才 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_delta、queue_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_update 的 delta(命令输出)是 payload,记了就泄漏了,所以只记类型。
now 可注入。 默认 Date.now,测试里用一个手动 tick(ms) 的确定性时钟,这样耗时可精确断言(turnDurationsMs === [60] 这种)。
插曲:smoke 为什么会打印 ~/.agents/skills
跑 smoke 时,控制台里除了课程自己的输出,还混进了几个跟本课无关的 skill——比如本机装的一个 ticket 工具和一个终端会话工具——它们来自 ~/.agents/skills。排查过程是一次很好的「读源码理解资源发现」的练习。
链路是这样的。smoke.ts 调 createAgentSession({...}),没有传 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.ts 和 homework.ts(作业)。
参考
代码固定使用 @earendil-works/[email protected],避免上游快速变化破坏可复现性。