07 — 分布式与可观测性 (chord, server, client, telemetry)

问题

当 Agent 需要以 client-server 模式运行(例如远程会话、多 worker 并行)时,需要:

  • 一个服务组合运行时来管理多个服务及其依赖
  • 一个 client-server 协议实现来传输消息
  • 一个可观测性系统来追踪跨进程的请求

Pi 将这些放在 chordserverclienttelemetry 四个包中。这些目前是实验性的(受 PI_EXPERIMENTAL=1 环境变量控制),但设计成熟。

第一部分:chord — 应用组合运行时

packages/chord/src/

核心概念

概念 说明
Facet 一个模块,声明它提供和消费的服务
Service 类型化的服务契约(singleton 或 keyed 模式)
ReplicatedState 可变复制状态,跨进程同步
FacetHost 管理一组 facet 的激活、热重载

Facet 模型

packages/chord/src/facets/host.ts(~32K)

一个 Facet 通过 env 对象声明依赖:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
defineFacet({
name: "agent-controller",
activate(env) {
// 提供服务
env.provide("AgentController", agentControllerImpl);

// 消费服务
const models = env.use("Models");
const transcript = env.use("Transcript");

// 观察服务(不强制激活)
env.observe("SessionManagement", (sm) => { ... });
},
});

FacetKernel 解析依赖图:provider 在 consumer 之前激活。如果 A 消费 B 提供的服务,B 会先于 A 激活。

FacetHost.reload() 支持热重载——替换 facet 实现而不重启整个进程。

Service 模型

packages/chord/src/services/

服务是类型化的契约:

  • Singleton 模式:一个 host 中一个实例
  • Keyed 模式:按 key 多实例(如按 session ID)

远程服务通过 RemoteServiceProvider 发布,通过 RemoteServiceBindingImpl 消费。状态编码/解码用于线路传输。

ReplicatedState

packages/chord/src/services/state.ts

1
2
3
const state = replicatedState(initialValue);
state.set(newValue); // 本地修改
state.subscribe(snapshot => { ... }); // 接收远程更新

MutableReplicatedStateImpl 跟踪本地修改,发布不可变快照。消费者接收 hydration(初始全量)+ updates(后续增量)。

Delta 编码

packages/chord/src/delta/index.ts(~45K):结构化 delta 编码,用于高效的状态同步。只传输变化的部分,而非整个状态。

第二部分:server — 服务端会话路由

packages/server/src/

Server 类

server.ts(~20K)— Server<TMetadata>

1
2
3
4
5
6
7
8
9
10
11
12
13
客户端连接


Server
├─ 握手 (版本协商)
│ ClientHello(version) → ServerHello(serverId)
│ 或 ServerHelloError

└─ SessionRouter
├─ 路由 RequestEnvelope 到正确的 session
├─ 管理 client attachments (哪个 client 挂在哪个 session)
├─ 转发 ServiceEventEnvelope (订阅更新)
└─ 转发 AttachmentEnvelope (附件路由更新)

SessionRouter

session-router.ts(~12K):

  • 跟踪哪个 client 附加到哪个 session
  • 将 RPC 调用路由到正确的 session
  • 管理 session 生命周期

传输抽象

connection.tsByteConnection 接口抽象传输层。transports/ 目录提供具体实现(WebSocket、Unix socket 等)。

第三部分:client — 客户端连接

packages/client/src/

Client 类

client.ts(~14K):

1
2
3
4
5
6
7
8
9
10
11
12
const client = new Client({ transport });
await client.connect(); // 握手

// RPC 调用
const result = await client.call(target, method, params);

// 订阅服务更新
const subscription = client.subscribe(target, method, params);
subscription.on("update", (data) => { ... });

// 附件状态
client.on("attachment", (attachment) => { ... });

Connection

connection.ts(~8K):管理字节级连接,处理帧的发送和接收(使用 pi-protocol 的帧编解码器)。

状态管理

Client 使用 chord 的服务状态解码器从线路更新重建复制状态。跟踪:

  • 连接状态变化(ConnectionState
  • 附件变化(当前附加到哪个 session)

第四部分:实验性多进程架构

packages/coding-agent/src/experimental/

这是 Pi 的下一代架构(受 PI_EXPERIMENTAL=1 控制):

1
2
3
4
5
6
7
┌──────────────┐     ┌──────────────────┐     ┌──────────────────┐
│ Client │ │ Coordinator │ │ Session Worker │
│ (TUI) │────▶│ (Unix socket │────▶│ (Agent 进程) │
│ │ │ 路由) │ │ │
│ pi-tui 渲染 │ │ 转发消息 │ │ AgentSession │
│ 事件订阅 │ │ 管理 workers │ │ Agent loop │
└──────────────┘ └──────────────────┘ └──────────────────┘

Coordinator

coordinator.ts(~20K):CoordinatorConnection — 基于 Unix socket 的对等消息路由。管理多个 session worker 进程。

Session Worker

session-worker.ts(~30K)/ session-worker-manager.ts(~30K):

Worker 进程拥有 agent session。每个 session 可以在独立 worker 中运行,实现隔离和并行。

Chord 服务

services/(16 个文件)— 定义为 chord 服务:

服务 说明
AgentController Agent 操作控制
Models 模型管理
Transcript 对话记录
SessionManagement 会话管理
SlashCommands Slash 命令
PresentationUI 演示 UI

这些服务通过 chord 的 facet 系统组合,Client TUI 远程消费这些服务。

插件

plugins/:插件包加载——bundled.ts(内置插件)和 package.ts(npm/git 包插件)。

第五部分:telemetry — 遥测契约

packages/telemetry/src/

设计理念

厂商中立的遥测抽象——类似 OpenTelemetry 的 span,但不耦合任何具体后端。零运行时依赖。

核心契约

index.ts

1
2
3
4
5
6
7
8
9
10
11
// 遥测上下文 — 基本能力
interface TelemetryContext {
startSpan<T>(options: SpanOptions, callback: (span: TelemetrySpan) => T | Promise<T>): Promise<T>;
}

// Span — 一个操作的时间段
interface TelemetrySpan extends TelemetryContext {
addEvent(name: string, attributes?: SpanAttributes): void;
setAttributes(attributes: SpanAttributes): void;
setStatus(status: SpanStatus): void;
}

关键设计startSpan 是回调式的——span 在回调 resolve 时自动结算(ok),在 reject 时自动标记 error。不需要手动 endSpan()

1
2
3
4
5
6
7
8
await telemetry.startSpan({ name: "llm.request", attributes: { model: "claude-4" } },
async (span) => {
const result = await callLLM();
span.addEvent("response_received", { tokens: result.usage.output });
return result;
}
);
// span 在回调返回后自动结算

TelemetrySpan extends TelemetryContext——span 本身也是一个 context,可以从 span 启动子 span,形成调用树。

Schema 类型推断

packages/telemetry/src/index.ts:76-322

一套条件类型系统从 schema 定义推断类型安全的属性:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
const schema = defineTelemetrySchema({
version: 1,
spans: {
"llm.request": {
description: "LLM API request",
parents: { kind: "any" },
startAttributes: {
model: { type: "string", required: true, description: "Model ID" },
},
endAttributes: {
tokens: { type: "number", required: false },
},
events: {
"response_received": {
description: "Response received",
attributes: {
tokens: { type: "number", required: true },
},
},
},
status: { default: "ok", errorWhen: "request_failed" },
},
},
});

// 创建类型安全的 span starter
const startSpan = createTypedSpanStarter(telemetryContext, schema);

// 编译时检查属性名和类型
await startSpan("llm.request",
{ model: "claude-4" }, // ✓ 正确
async (span) => {
span.addEvent("response_received", { tokens: 100 }); // ✓ 正确
// span.addEvent("unknown_event", {}); // ✗ 编译错误
}
);

关键设计:schema 值仅用于类型推断——没有运行时 schema 验证index.ts:347 注释明确说明)。这避免了运行时开销。

UniqueTelemetrySchemas<Schemas>(line 293)是编译时检查——如果多个 schema 共享 span 名,会产生 "duplicate telemetry span names" 类型错误。

实现

实现 文件 说明
NOOP_TELEMETRY_CONTEXT noop.ts 空操作上下文,不记录任何内容
InMemoryTelemetryContext memory.ts 内存记录,用于测试

InMemoryTelemetryContextmemory.ts:192):

  • 创建 MutableRecordedTelemetrySpan(含 id、parentId、name、attributes、events、status、endSequence)
  • 回调 resolve 时自动结算(ok),reject 时自动标记 error
  • 子 span 在已结算的 parent 上启动时降级为 NOOP_TELEMETRY_CONTEXT
  • getSpans() 返回按启动顺序的快照,含 endSequence 用于按完成顺序排序

测试

testing/TelemetryAdapterFixture + TelemetryAdapterConformanceCase——runner 无关的一致性测试,可与任何测试框架配合。

与其他包的连接

1
2
3
4
5
6
7
8
9
10
telemetry (契约)

├── ai (ProviderRequestOptions.telemetryContext)
│ └─ 每次 provider 请求可以追踪

├── agent (harness/telemetry.ts)
│ └─ 定义 agent 级 span schema

└── 具体后端 (用户实现)
└─ OpenTelemetry adapter / 自定义后端

整体架构图

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
┌─ Client 进程 ──────────┐     ┌─ Server 进程 ──────────────────────┐
│ │ │ │
│ pi-tui (渲染) │ │ Coordinator │
│ │ │ │ │ │
│ ▼ │ │ ├── Session Worker 1 │
│ Client (pi-client) │◄───►│ │ ├── AgentSession │
│ │ │ │ │ ├── Agent (agent-core) │
│ ▼ │ │ │ └── Models (pi-ai) │
│ chord RemoteBinding │ │ │ │
│ │ │ │ ├── Session Worker 2 │
│ ▼ │ │ │ └── ... │
│ ReplicatedState │ │ │ │
│ (Transcript, UI, ...) │ │ └── Session Worker N │
│ │ │ │
└────────────────────────┘ │ chord FacetHost │
│ └── Services (AgentController, │
│ Models, Transcript, ...) │
│ │
│ pi-protocol (消息编解码) │
│ session-backends (SQLite 存储) │
│ telemetry (span 追踪) │
└─────────────────────────────────────┘

学习要点

  • Facet 模式实现服务组合:通过 provide/use/observe 声明依赖,FacetKernel 自动解析激活顺序。比手动管理依赖更可靠,且支持热重载
  • ReplicatedState + Delta 编码实现跨进程状态同步:全量 hydration + 增量 update,delta 编码减少传输量
  • 传输中立协议 + 可插拔传输实现pi-protocol 定义消息格式,server/client 提供传输实现,同一套协议可用于本地和远程
  • Telemetry 的回调式 span 设计startSpan(callback) 自动结算,避免忘记 endSpan 导致的泄漏。类型安全的 schema 系统在编译时检查属性,无运行时开销
  • 实验性多进程架构的价值:每个 session 在独立 worker 中运行,实现进程级隔离和并行。coordinator 负责路由,client 只负责 UI 渲染
  • 零依赖契约包telemetryprotocol 都是零运行时依赖的纯契约包——只定义接口和类型,不包含实现。这让上层可以自由选择后端

上一篇:06-终端 UI-pi-tui
返回:00-导读与学习路线