Telemetry cn
Jianmu 的可观测性体系由两层独立的事件基础设施构成:TelemetryHub 负责跨模块的遥测事件采集、采样与路由;RuntimeEventBus 负责引擎运行时的结构化事件发布。两者通过独立的 subscriber/sink 模型解耦,但共享 contextvars 驱动的 TraceContext 实现端到端的 Span 追踪级联。本章系统阐述两者的架构设计、事件生命周期、Sink 扩展机制以及与配置体系的联动。
双层事件架构总览¶
Jianmu 的事件系统遵循"遥测层 + 运行时层"的双总线设计。遥测层面向跨模块的可观测信号(LLM 调用、工具执行、BT 技能状态),由 TelemetryHub 集中管理采样和路由;运行时层面向引擎内部的生命周期事件(Agent 创建/暂停/恢复、消息路由、调度 tick),由 RuntimeEventBus 承载。两层之间的桥梁是 TraceContext——一个基于 contextvars 的协程安全上下文,携带 trace_id 和 span_stack,使得遥测事件能够自动继承当前 Span 的身份标识,形成完整的调用链。
graph TD
subgraph "应用层"
A[Agent.run]
B[ToolNode.update_async]
C[ModelClient.invoke]
D[BTSkillTool.execute]
end
subgraph "遥测层 (jianmu.telemetry)"
H[TelemetryHub]
S1[DebugLogSink]
S2[JsonlFileSink]
S3[UsageAggregationSink]
S4[ConsoleStreamSink]
L[Listener Callbacks]
H --> S1
H --> S2
H --> S3
H --> S4
H --> L
end
subgraph "运行时层 (jianmu.engine.events)"
EB[RuntimeEventBus]
RS[Runtime Subscribers]
EB --> RS
end
subgraph "上下文层"
TC[TraceContext<br/>contextvars]
SP[Span 栈]
TC --- SP
end
A -->|emit| H
B -->|trace_emit| H
C -->|trace_emit| H
D -->|trace_emit| H
A -->|emit_runtime_event| EB
H -.->|注入 trace_id / span_id| TC
EB -.->|独立上下文| EB
TelemetryHub:核心事件中枢¶
TelemetryHub 是整个遥测体系的中心调度器,提供线程安全的事件发射、监听器管理、Sink 注册以及多层采样控制。每个进程维护一个全局默认 Hub 实例,通过 get_default_hub() 获取或 set_default_hub() 替换。长生命周期集成可使用 subscribe_default_hub_changes() 跟随替换,并安全地从旧 Hub 解绑。
采样与限流策略¶
Hub 实现了三级过滤机制,按优先级从高到低依次为:
| 层级 | 配置项 | 作用 | 默认值 |
|---|---|---|---|
| 白名单 | always_events |
该集合内的事件无条件通过,跳过采样和限流 | trace.span.start, trace.span.end, llm.error, tool.error, tool.denied, budget.exceeded, bt.skill.* |
| 按事件采样 | event_sample_rates |
特定事件名的独立采样率,优先级高于全局采样率 | {}(空,回退到全局) |
| 全局采样 | sample_rate |
未匹配到按事件采样时的兜底采样率 | 1.0(100%) |
| 速率限制 | max_per_sec |
每个 trace/agent 维度的每秒最大事件数 | 200 |
过滤逻辑由 _should_emit() 方法统一执行:先检查白名单,再检查按事件采样率(使用随机数判断),最后检查速率限制窗口。速率限制的 scope key 解析优先级为:当前 TraceContext 的 trace_id → 上下文 metadata 中的 agent_id → payload 中的 trace_id/agent_id → 全局 "global"。
事件发射流程¶
emit(event, payload) 是发射事件的核心入口。其完整生命周期如下:
sequenceDiagram
participant C as 调用方
participant H as TelemetryHub
participant CTX as TraceContext
participant SL as Sink/Listener
C->>H: emit("tool.call", {"tool": "calc"})
H->>H: _should_emit() 检查采样/限流
alt 被过滤
H-->>C: return (静默丢弃)
end
H->>CTX: current_context() 获取 trace_id / span_id
CTX-->>H: trace_id="t1", span_id="s2"
H->>H: 合并上下文 metadata 到 payload
H->>H: safe_serialize() 安全序列化(safe_mode)
H->>H: 构造 TelemetryEvent(name, ts, trace_id, span_id, payload)
H->>SL: 遍历 sinks → handle(event)
H->>SL: 遍历 listeners → callback(event)
关键细节:当 safe_mode=True(默认开启),所有 payload 会经过 safe_serialize() 递归转换为 JSON 安全结构——处理 dataclass、Pydantic model(model_dump())、循环引用(标记为 "<recursion>"),最大深度 4 层。
默认白名单事件¶
DEFAULT_ALWAYS_EVENTS 定义了始终不被过滤的关键事件集,覆盖 Span 生命周期(trace.span.start、trace.span.end)、错误信号(llm.error、tool.error、tool.denied、budget.exceeded)以及 BT 技能执行边界(bt.skill.started、bt.skill.completed、bt.skill.node.status、bt.skill.node.io)。这些事件对调试和审计至关重要,因此绕过采样。
Span 追踪:TraceContext 与 span 上下文管理器¶
Span 追踪是实现分布式调用链可观测性的核心机制。Jianmu 使用 Python 的 contextvars 模块实现协程安全的上下文传递,无需显式参数穿透。
TraceContext 数据结构¶
TraceContext 是一个轻量 dataclass,包含三个字段:
- trace_id:一次完整运行的唯一标识(UUID)
- span_stack:不可变元组,记录从根到当前 Span 的 ID 栈
- metadata:附加键值对(如
agent_id),在 Span 嵌套时会自动继承合并
Span 对象则记录单个 Span 的完整信息:id、trace_id、parent_id、name、start_time、end_time、status("running" / "success" / "error")以及 duration_ms 属性。
span 上下文管理器的工作原理¶
span(name, **metadata) 同时支持作为上下文管理器和装饰器使用。其内部流程在两个关键阶段展开:
进入(__enter__):获取当前 TraceContext(若无则生成新 UUID 作为 trace_id),以当前栈顶 span_id 为 parent_id,生成新的 span_id 并压入 span_stack,通过 set_context() 将更新后的上下文安装到当前协程,最后发射 trace.span.start 事件。
退出(__exit__):记录结束时间和状态(有异常则为 "error"),发射 trace.span.end 事件(携带 duration_ms 和异常信息),调用 reset_context() 恢复之前的上下文——弹出当前 Span,将栈恢复到进入前的状态。
# 典型用法示例
from jianmu.telemetry import span, set_context, reset_context
token = set_context(trace_id="trace-1", agent_id="agent-1")
try:
with span("tool_execution", tool="calculator"):
# 此时 trace_id="trace-1", span_stack=("span-1",)
with span("sub_operation", step="validate"):
# 此时 trace_id="trace-1", span_stack=("span-1", "span-2")
pass
# 退出 sub_operation 后 span_stack 恢复为 ("span-1",)
finally:
reset_context(token)
这种栈式设计确保 Span 的父子关系严格对应于代码的嵌套结构,即使存在异步并发也不会发生上下文泄漏。
内置 Sink 体系¶
Jianmu 提供四种开箱即用的 Sink 实现,均遵循 TelemetrySink Protocol(仅需实现 handle(event: TelemetryEvent) -> None)。
| Sink | 用途 | 输出目标 | 关键行为 |
|---|---|---|---|
DebugLogSink |
调试时镜像事件到日志 | Loguru debug 输出 | 每条事件写入一行 📡 [Telemetry] {name} {payload} |
JsonlFileSink |
持久化事件到文件 | JSONL 文件(线程安全追加) | 每条事件序列化为一行 JSON,包含 event、ts、trace_id、span_id、payload |
UsageAggregationSink |
按 trace 聚合 Token 用量 | 内存字典 | 仅处理 llm.usage 事件,提供 snapshot() 和 reset() 方法 |
ConsoleStreamSink |
实时打印 LLM 流式输出 | stdout | 仅处理 llm.stream.delta 事件,支持 enabled/flush 控制 |
UsageAggregationSink 是成本追踪的关键组件。它监听 llm.usage 事件,使用 normalize_usage_dict() 将不同 Provider 返回的异构用量格式(prompt_tokens/input_tokens、completion_tokens/output_tokens)统一为 {"prompt_tokens": int, "completion_tokens": int, "total_tokens": int} 格式,按 trace_id 维度累加。
RuntimeEventBus:运行时事件总线¶
与 TelemetryHub 的遥测定位不同,RuntimeEventBus 服务于引擎运行时的工作流事件。RuntimeEvent 拥有更丰富的结构化字段:event_id(自动生成的 UUID)、run_id、agent_id、node_id、node_name、tool_call_id、thread_id——这些字段使得下游消费者(如 TUI Chat、Tree Studio)能够精确关联事件到具体的 Agent、节点或工具调用。
RuntimeEventBus 内置有界缓冲区(默认 1000 条),支持 subscribe()(返回 unsubscribe 回调)、get_since(timestamp) 和 get_all() 方法。emit_runtime_event() 是便捷的发射辅助函数,在 bus 为 None 时静默跳过,避免了空值检查的样板代码。
在 Swarm 运行时中,RuntimeEventBus 承载了 Agent 生命周期事件(agent_created、agent_started、agent_stopped、agent_paused、agent_resumed、agent_killed)、消息路由事件(message_sent、message_received)以及调度事件(agent_wake、tick_completed),为可视化调试面板提供了完整的事件流。
生产级事件目录¶
以下是 Jianmu 框架内部发射的全部遥测事件,按来源模块分类:
模型调用事件¶
| 事件名 | 发射位置 | Payload 关键字段 |
|---|---|---|
llm.stream.delta |
jianmu/model/stream.py |
content: 流式文本块 |
llm.usage |
jianmu/model/invoke.py |
node: 节点名, usage: 标准化用量字典 |
流式 delta 事件的采样由 token_sample_rate 独立控制,默认值为 0.0(完全不采样),以避免高频 token 事件淹没遥测通道。如需启用,在配置中设置 trace.token_sample_rate 为大于 0 的值。
工具执行事件¶
| 事件名 | 发射位置 | Payload 关键字段 |
|---|---|---|
tool.call |
jianmu/node/builtin/tool.py |
node, tool, mode ("workflow"/"agent"), args, tool_call_id |
tool.result |
jianmu/node/builtin/tool.py |
node, tool, ok, result/error, tool_call_id |
工具事件区分两种执行模式:mode="workflow" 表示通过行为树端口绑定直接调用的工具节点,mode="agent" 表示 LLM 驱动通过 tool_call_id 关联的 Agent 工具调用。两种模式共享相同的 tool.call / tool.result 事件名,通过 mode 字段区分。
行为树技能事件¶
| 事件名 | 发射位置 | Payload 关键字段 |
|---|---|---|
bt.skill.started |
jianmu/tool/builtin/bt_skill.py |
run_id, skill_name, namespace |
bt.skill.completed |
jianmu/tool/builtin/bt_skill.py |
run_id, status, duration_ms, tick_count |
bt.skill.node.status |
jianmu/tool/builtin/bt_skill.py |
run_id, tick, changes: 状态变更列表 |
bt.skill.node.io |
jianmu/node/builtin/tool.py |
run_id, stage, tool_name, ok, input_preview, output_preview |
BT 技能的事件体系通过 _BTSkillTraceVisitor(py_trees 的 VisitorBase 子类)在每个 tick 后遍历行为树,对比前后状态差异,仅发射发生变化的节点。这种差异发射策略显著减少了状态稳定期的事件量。bt.skill.node.io 则由 ToolNode._emit_bt_skill_node_io() 在存在 _bt_trace_context 时发射,携带截断预览(字典递归限制 300 字符,列表限制 20 项)。
Span 生命周期事件¶
| 事件名 | 发射位置 | Payload 关键字段 |
|---|---|---|
trace.span.start |
jianmu/telemetry/context.py |
span_id, parent_id, trace_id, name, **metadata |
trace.span.end |
jianmu/telemetry/context.py |
span_id, trace_id, name, status, duration_ms, error (如有) |
这两个事件始终在白名单中,不会被采样过滤,确保 Span 的边界完整性。
配置体系与引导启动¶
遥测系统的配置入口位于 jianmu.yaml 的 trace 段,对应 Pydantic 模型 TraceConfig:
| 配置键 | 类型 | 默认值 | 说明 |
|---|---|---|---|
log_events |
bool |
false |
是否同时将事件镜像到 Loguru debug 日志 |
safe_mode |
bool |
true |
是否对 payload 执行安全序列化 |
sample_rate |
float |
1.0 |
全局采样率(0.0 ~ 1.0) |
token_sample_rate |
float |
0.0 |
LLM 流式 token 事件的独立采样率 |
max_per_sec |
int |
200 |
每 trace/agent 每秒最大事件数 |
event_sample_rates |
dict[str, float] |
{} |
按事件名的采样率覆盖 |
always_events |
list[str] |
[] |
白名单事件名列表 |
router |
str \| None |
null |
Sink 路由 URI(如 "file:.outputs/trace.jsonl") |
bootstrap_telemetry(config) 是标准化的初始化入口,执行顺序为:配置 Loguru 日志级别 → 构造 TelemetryHub(应用采样参数)→ 根据配置注册 DebugLogSink(若 log_events=true)→ 解析 router 注册 JsonlFileSink(若以 "file:" 开头)→ 调用 set_default_hub() 安装为全局默认。配置的 always_events 会扩展而不是替换强制生命周期白名单,因此成对的 span start/end 事件始终绕过采样和速率限制。
扩展自定义 Sink¶
实现自定义 Sink 仅需满足 TelemetrySink Protocol——即实现 handle(self, event: TelemetryEvent) -> None 方法。以下是一个将事件写入 InfluxDB 的示意:
from jianmu.telemetry.types import TelemetryEvent
from jianmu.telemetry import register_sink
class InfluxDBSink:
def __init__(self, url: str, token: str, bucket: str):
self._client = influxdb_client.InfluxDBClient(url=url, token=token)
self._bucket = bucket
def handle(self, event: TelemetryEvent) -> None:
point = (
influxdb_client.Point("jianmu_events")
.tag("event", event.name)
.tag("trace_id", event.trace_id or "unknown")
.field("payload", str(event.payload))
.time(int(event.ts * 1e9))
)
self._client.write_api().write(bucket=self._bucket, record=point)
register_sink(InfluxDBSink(url="...", token="...", bucket="jianmu"))
对于运行时事件总线,类似地可通过 RuntimeEventBus.subscribe() 注册回调,返回的 unsubscribe 函数可用于生命周期管理。
阅读建议¶
遥测与可观测性是 Jianmu 运行时诊断的最后一块拼图。理解本章后,您可以结合以下页面深化对特定模块的认知:
- ModelClient 外观:深入理解
llm.usage和llm.stream.delta事件的发射时机与 Provider 集成 - Guard 体系:了解
tool.denied和budget.exceeded事件的触发条件 - Swarm 运行时:查看 RuntimeEventBus 在多 Agent 协作中的事件路由机制
- TUI Chat 和 Tree Studio:了解上层应用如何消费遥测与运行时事件进行可视化