跳转至

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:了解上层应用如何消费遥测与运行时事件进行可视化