跳转至

jianmu.telemetry

适用对象:平台维护者 / 可观测性集成者 / 调试者 是否必读:按需 相关模块:jianmu.engine, jianmu.swarm

1. 模块职责

jianmu.telemetry 汇总日志、流式输出、trace 事件和 usage 统计的公开入口。

它主要服务于调试、观测和运行分析,而不是业务流程编排本身。

2. 适合查什么

  • 日志配置:configure_logging()
  • 流式输出:ConsoleStreamSink
  • trace:emit()、span()、subscribe()
  • usage:UsageAggregationSink

3. 使用建议

  • 想快速观察本地运行过程时,先看 ConsoleStreamSink
  • 想接自定义 trace 路由或埋点,再使用 trace 相关函数

4. 注意事项

  • 这里更偏运行观测层,而不是业务编排层
  • 如果只是想让 agent 正常执行,通常不需要先接 telemetry
  • trace 订阅和路由更适合平台集成或调试场景

5. 最小示例

from jianmu.telemetry import ConsoleStreamSink, emit, register_sink

printer = ConsoleStreamSink()
register_sink(printer)
emit("custom.event", {"phase": "start"})

6. 常见入口

  • 想打印或消费流式输出:看 ConsoleStreamSink
  • 想打点:看 emit()
  • 想包裹 trace span:看 span()
  • 想订阅 trace 事件:看 subscribe()

7. API 参考

Hub 与事件流

TelemetryHub

TelemetryHub(
    *,
    sample_rate: float = 1.0,
    max_per_sec: int = 200,
    event_sample_rates: dict[str, float] | None = None,
    always_events: set[str] | None = None,
    safe_mode: bool = True,
    log_enabled: bool = False,
)

Thread-safe event hub for typed telemetry events.

属性:

名称 类型 描述
_listeners list[Callable[[TelemetryEvent], None]]

In-process listener callbacks invoked for each event.

_sinks list[TelemetrySink]

Registered telemetry sinks consuming emitted events.

_lock

Lock protecting listeners, sinks, and config updates.

_rate_lock

Lock protecting rate-limiting windows.

_rate_windows dict[str, tuple[float, int]]

Per-scope rate-limiting counters.

_sample_rate

Default sample rate for events.

_max_per_sec

Per-scope maximum emitted events per second.

_event_sample_rates

Event-specific sampling overrides.

_always_events

Events that bypass sampling and rate limiting.

_safe_mode

Whether payloads are safe-serialized before emission.

_log_enabled

Whether emitted events are mirrored to debug logs.

源代码位于: jianmu/telemetry/hub.py
def __init__(
    self,
    *,
    sample_rate: float = 1.0,
    max_per_sec: int = 200,
    event_sample_rates: dict[str, float] | None = None,
    always_events: set[str] | None = None,
    safe_mode: bool = True,
    log_enabled: bool = False,
) -> None:
    self._listeners: list[Callable[[TelemetryEvent], None]] = []
    self._sinks: list[TelemetrySink] = []
    self._lock = threading.Lock()
    self._rate_lock = threading.Lock()
    self._rate_windows: dict[str, tuple[float, int]] = {}
    self._sample_rate = min(1.0, max(0.0, float(sample_rate)))
    self._max_per_sec = max(0, int(max_per_sec))
    self._event_sample_rates = {
        str(name): min(1.0, max(0.0, float(rate)))
        for name, rate in dict(event_sample_rates or {}).items()
    }
    self._always_events = set(always_events or DEFAULT_ALWAYS_EVENTS)
    self._safe_mode = bool(safe_mode)
    self._log_enabled = bool(log_enabled)

subscribe

subscribe(
    callback: Callable[[TelemetryEvent], None],
) -> None

Register a listener callback invoked for each emitted event.

参数:

名称 类型 描述 默认
callback Callable[[TelemetryEvent], None]

Callback invoked by the operation.

必需
源代码位于: jianmu/telemetry/hub.py
def subscribe(self, callback: Callable[[TelemetryEvent], None]) -> None:
    """Register a listener callback invoked for each emitted event.

    Args:
        callback: Callback invoked by the operation.
    """
    with self._lock:
        self._listeners.append(callback)

unsubscribe

unsubscribe(
    callback: Callable[[TelemetryEvent], None],
) -> None

Remove a previously registered listener callback if present.

参数:

名称 类型 描述 默认
callback Callable[[TelemetryEvent], None]

Callback invoked by the operation.

必需
源代码位于: jianmu/telemetry/hub.py
def unsubscribe(self, callback: Callable[[TelemetryEvent], None]) -> None:
    """Remove a previously registered listener callback if present.

    Args:
        callback: Callback invoked by the operation.
    """
    with self._lock:
        try:
            self._listeners.remove(callback)
        except ValueError:
            pass

register_sink

register_sink(sink: TelemetrySink) -> None

Register a sink that consumes emitted telemetry events.

参数:

名称 类型 描述 默认
sink TelemetrySink

Telemetry sink to register or remove.

必需
源代码位于: jianmu/telemetry/hub.py
def register_sink(self, sink: TelemetrySink) -> None:
    """Register a sink that consumes emitted telemetry events.

    Args:
        sink: Telemetry sink to register or remove.
    """
    with self._lock:
        self._sinks.append(sink)

unregister_sink

unregister_sink(sink: TelemetrySink) -> None

Remove a previously registered telemetry sink if present.

参数:

名称 类型 描述 默认
sink TelemetrySink

Telemetry sink to register or remove.

必需
源代码位于: jianmu/telemetry/hub.py
def unregister_sink(self, sink: TelemetrySink) -> None:
    """Remove a previously registered telemetry sink if present.

    Args:
        sink: Telemetry sink to register or remove.
    """
    with self._lock:
        try:
            self._sinks.remove(sink)
        except ValueError:
            pass

reset

reset() -> None

Clear listeners, sinks, and rate-limiting windows.

源代码位于: jianmu/telemetry/hub.py
def reset(self) -> None:
    """Clear listeners, sinks, and rate-limiting windows."""
    with self._lock:
        self._listeners = []
        self._sinks = []
    with self._rate_lock:
        self._rate_windows = {}

configure

configure(
    *,
    sample_rate: float | None = None,
    max_per_sec: int | None = None,
    event_sample_rates: dict[str, float] | None = None,
    always_events: set[str] | None = None,
    safe_mode: bool | None = None,
    log_enabled: bool | None = None,
) -> None

Update hub sampling, throttling, and serialization configuration.

参数:

名称 类型 描述 默认
sample_rate float | None

The sample_rate value.

None
max_per_sec int | None

The max_per_sec value.

None
event_sample_rates dict[str, float] | None

Collection of event sample rate values.

None
always_events set[str] | None

Collection of always event values.

None
safe_mode bool | None

The safe_mode value.

None
log_enabled bool | None

The log_enabled value.

None
源代码位于: jianmu/telemetry/hub.py
def configure(
    self,
    *,
    sample_rate: float | None = None,
    max_per_sec: int | None = None,
    event_sample_rates: dict[str, float] | None = None,
    always_events: set[str] | None = None,
    safe_mode: bool | None = None,
    log_enabled: bool | None = None,
) -> None:
    """Update hub sampling, throttling, and serialization configuration.

    Args:
        sample_rate: The `sample_rate` value.
        max_per_sec: The `max_per_sec` value.
        event_sample_rates: Collection of event sample rate values.
        always_events: Collection of always event values.
        safe_mode: The `safe_mode` value.
        log_enabled: The `log_enabled` value.
    """
    with self._lock:
        if sample_rate is not None:
            self._sample_rate = min(1.0, max(0.0, float(sample_rate)))
        if max_per_sec is not None:
            self._max_per_sec = max(0, int(max_per_sec))
        if event_sample_rates is not None:
            self._event_sample_rates = {
                str(name): min(1.0, max(0.0, float(rate)))
                for name, rate in dict(event_sample_rates).items()
            }
        if always_events is not None:
            self._always_events = set(always_events)
        if safe_mode is not None:
            self._safe_mode = bool(safe_mode)
        if log_enabled is not None:
            self._log_enabled = bool(log_enabled)
    with self._rate_lock:
        self._rate_windows = {}

emit

emit(
    event: str, payload: dict[str, Any] | None = None
) -> None

Emit an event to registered sinks and listeners if sampling allows it.

参数:

名称 类型 描述 默认
event str

Event name or payload to process.

必需
payload dict[str, Any] | None

Payload data for the operation.

None
源代码位于: jianmu/telemetry/hub.py
def emit(self, event: str, payload: dict[str, Any] | None = None) -> None:
    """Emit an event to registered sinks and listeners if sampling allows it.

    Args:
        event: Event name or payload to process.
        payload: Payload data for the operation.
    """
    payload_dict = dict(payload or {})
    if not self._should_emit(event, payload_dict):
        return

    ctx = current_context()
    trace_id = payload_dict.get("trace_id")
    if trace_id is None and ctx and ctx.trace_id:
        trace_id = ctx.trace_id

    span_id = payload_dict.get("span_id")
    if span_id is None and ctx and ctx.span_stack:
        span_id = ctx.span_stack[-1]

    if ctx and ctx.metadata:
        for key, value in ctx.metadata.items():
            payload_dict.setdefault(key, value)

    ts = time.time()
    if self._safe_mode:
        serialized_payload = safe_serialize(payload_dict)
        if isinstance(serialized_payload, dict):
            payload_dict = serialized_payload
        else:
            payload_dict = {"value": serialized_payload}

    telemetry_event = TelemetryEvent(
        name=event,
        ts=ts,
        trace_id=trace_id,
        span_id=span_id,
        payload=payload_dict,
    )

    with self._lock:
        listeners = list(self._listeners)
        sinks = list(self._sinks)
        log_enabled = self._log_enabled

    if log_enabled:
        from jianmu.telemetry.logging import logger

        logger.debug("📡 [Trace] {} {}", event, payload_dict)

    for sink in sinks:
        try:
            sink.handle(telemetry_event)
        except Exception as exc:
            from jianmu.telemetry.logging import logger

            logger.warning("⚠️ [Telemetry] Sink failed: {}", exc)

    for callback in listeners:
        try:
            callback(telemetry_event)
        except Exception as exc:
            from jianmu.telemetry.logging import logger

            logger.warning("⚠️ [Telemetry] Listener failed: {}", exc)

emit

emit(
    event: str, payload: dict[str, Any] | None = None
) -> None

Emit an event through the default hub.

参数:

名称 类型 描述 默认
event str

Event name or payload to process.

必需
payload dict[str, Any] | None

Payload data for the operation.

None
源代码位于: jianmu/telemetry/hub.py
def emit(event: str, payload: dict[str, Any] | None = None) -> None:
    """Emit an event through the default hub.

    Args:
        event: Event name or payload to process.
        payload: Payload data for the operation.
    """
    _default_hub.emit(event, payload)

subscribe

subscribe(
    callback: Callable[[TelemetryEvent], None],
) -> None

Register one listener on the default hub.

参数:

名称 类型 描述 默认
callback Callable[[TelemetryEvent], None]

Callback invoked by the operation.

必需
源代码位于: jianmu/telemetry/hub.py
def subscribe(callback: Callable[[TelemetryEvent], None]) -> None:
    """Register one listener on the default hub.

    Args:
        callback: Callback invoked by the operation.
    """
    _default_hub.subscribe(callback)

unsubscribe

unsubscribe(
    callback: Callable[[TelemetryEvent], None],
) -> None

Remove one listener from the default hub.

参数:

名称 类型 描述 默认
callback Callable[[TelemetryEvent], None]

Callback invoked by the operation.

必需
源代码位于: jianmu/telemetry/hub.py
def unsubscribe(callback: Callable[[TelemetryEvent], None]) -> None:
    """Remove one listener from the default hub.

    Args:
        callback: Callback invoked by the operation.
    """
    _default_hub.unsubscribe(callback)

register_sink

register_sink(sink: TelemetrySink) -> None

Register one sink on the default hub.

参数:

名称 类型 描述 默认
sink TelemetrySink

Telemetry sink to register or remove.

必需
源代码位于: jianmu/telemetry/hub.py
def register_sink(sink: TelemetrySink) -> None:
    """Register one sink on the default hub.

    Args:
        sink: Telemetry sink to register or remove.
    """
    _default_hub.register_sink(sink)

unregister_sink

unregister_sink(sink: TelemetrySink) -> None

Remove one sink from the default hub.

参数:

名称 类型 描述 默认
sink TelemetrySink

Telemetry sink to register or remove.

必需
源代码位于: jianmu/telemetry/hub.py
def unregister_sink(sink: TelemetrySink) -> None:
    """Remove one sink from the default hub.

    Args:
        sink: Telemetry sink to register or remove.
    """
    _default_hub.unregister_sink(sink)

Trace 上下文

current_context

current_context() -> TraceContext | None

Return the active trace context.

返回:

类型 描述
TraceContext | None

The resolved value, or None when no value is available.

源代码位于: jianmu/telemetry/context.py
def current_context() -> TraceContext | None:
    """Return the active trace context.

    Returns:
        The resolved value, or `None` when no value is available.
    """
    return _context.get()

set_context

set_context(
    trace_id: str | None = None, **metadata: Any
) -> contextvars.Token

Install a trace context and return the reset token.

参数:

名称 类型 描述 默认
trace_id str | None

Identifier for trace.

None
**metadata Any

The metadata value.

{}

返回:

类型 描述
Token

The resulting contextvars.Token value.

源代码位于: jianmu/telemetry/context.py
def set_context(trace_id: str | None = None, **metadata: Any) -> contextvars.Token:
    """Install a trace context and return the reset token.

    Args:
        trace_id: Identifier for trace.
        **metadata: The `metadata` value.

    Returns:
        The resulting `contextvars.Token` value.
    """
    ctx = _context.get()
    if trace_id is None:
        trace_id = ctx.trace_id if ctx else str(uuid.uuid4())
    merged_metadata: dict[str, Any] = dict(ctx.metadata) if ctx and ctx.metadata else {}
    merged_metadata.update(metadata)
    span_stack = ctx.span_stack if ctx else ()
    return _context.set(TraceContext(trace_id=trace_id, span_stack=span_stack, metadata=merged_metadata))

reset_context

reset_context(token: Token | None) -> None

Restore the previous trace context.

参数:

名称 类型 描述 默认
token Token | None

The token value.

必需
源代码位于: jianmu/telemetry/context.py
def reset_context(token: contextvars.Token | None) -> None:
    """Restore the previous trace context.

    Args:
        token: The `token` value.
    """
    if token is None:
        return
    _context.reset(token)

span

span(name: str, **kwargs: Any)

Context manager and decorator for tracing execution spans.

属性:

名称 类型 描述
name

Span name emitted to the telemetry hub.

metadata

Structured metadata attached to the emitted span.

span_obj Span | None

Active span object created on entry.

token Token | None

Context token used to restore the prior trace context.

源代码位于: jianmu/telemetry/context.py
def __init__(self, name: str, **kwargs: Any):
    self.name = name
    self.metadata = kwargs
    self.span_obj: Span | None = None
    self.token: contextvars.Token | None = None

Sink

ConsoleStreamSink

ConsoleStreamSink(
    *, enabled: bool = True, flush: bool = True
)

Print stream delta events to stdout.

属性:

名称 类型 描述
enabled

Whether console streaming output is active.

flush

Whether stdout is flushed after each printed delta.

源代码位于: jianmu/telemetry/sinks.py
def __init__(self, *, enabled: bool = True, flush: bool = True) -> None:
    self.enabled = enabled
    self.flush = flush

handle

handle(event: TelemetryEvent) -> None

Print streamed LLM delta content to stdout when enabled.

参数:

名称 类型 描述 默认
event TelemetryEvent

Event name or payload to process.

必需
源代码位于: jianmu/telemetry/sinks.py
def handle(self, event: TelemetryEvent) -> None:
    """Print streamed LLM delta content to stdout when enabled.

    Args:
        event: Event name or payload to process.
    """
    if not self.enabled or event.name != "llm.stream.delta":
        return
    delta = event.payload.get("content")
    if not delta:
        return
    print(str(delta), end="", flush=self.flush)

set_enabled

set_enabled(value: bool) -> None

Enable or disable console streaming output.

参数:

名称 类型 描述 默认
value bool

Value to apply.

必需
源代码位于: jianmu/telemetry/sinks.py
def set_enabled(self, value: bool) -> None:
    """Enable or disable console streaming output.

    Args:
        value: Value to apply.
    """
    self.enabled = value

reset

reset() -> None

Reset sink state.

源代码位于: jianmu/telemetry/sinks.py
def reset(self) -> None:
    """Reset sink state."""
    return None

DebugLogSink

Mirror telemetry events to loguru debug output.

handle

handle(event: TelemetryEvent) -> None

Write the event name and payload to the debug logger.

参数:

名称 类型 描述 默认
event TelemetryEvent

Event name or payload to process.

必需
源代码位于: jianmu/telemetry/sinks.py
def handle(self, event: TelemetryEvent) -> None:
    """Write the event name and payload to the debug logger.

    Args:
        event: Event name or payload to process.
    """
    logger.debug("📡 [Telemetry] {} {}", event.name, event.payload)

JsonlFileSink

JsonlFileSink(path: str)

Append telemetry events to a JSONL file.

属性:

名称 类型 描述
path

Target JSONL file path.

_lock

Lock protecting concurrent file appends.

源代码位于: jianmu/telemetry/sinks.py
def __init__(self, path: str):
    self.path = path
    self._lock = threading.Lock()

handle

handle(event: TelemetryEvent) -> None

Append the event payload as one JSON line to the configured file.

参数:

名称 类型 描述 默认
event TelemetryEvent

Event name or payload to process.

必需
源代码位于: jianmu/telemetry/sinks.py
def handle(self, event: TelemetryEvent) -> None:
    """Append the event payload as one JSON line to the configured file.

    Args:
        event: Event name or payload to process.
    """
    payload = {
        "event": event.name,
        "ts": event.ts,
        "trace_id": event.trace_id,
        "span_id": event.span_id,
        "payload": dict(event.payload),
    }
    line = json.dumps(payload, ensure_ascii=False)
    with self._lock:
        with open(self.path, "a", encoding="utf-8") as fh:
            fh.write(line + "\n")

UsageAggregationSink dataclass

UsageAggregationSink(
    totals: dict[str, dict[str, int]] = dict(),
)

Aggregate normalized usage by trace id.

属性:

名称 类型 描述
totals dict[str, dict[str, int]]

Aggregated token counters keyed by trace id.

handle

handle(event: TelemetryEvent) -> None

Accumulate usage counters from llm.usage telemetry events.

参数:

名称 类型 描述 默认
event TelemetryEvent

Event name or payload to process.

必需
源代码位于: jianmu/telemetry/sinks.py
def handle(self, event: TelemetryEvent) -> None:
    """Accumulate usage counters from `llm.usage` telemetry events.

    Args:
        event: Event name or payload to process.
    """
    if event.name != "llm.usage":
        return

    usage = _normalize_usage(event.payload.get("usage"))
    if not usage:
        return

    trace_id = event.trace_id or event.payload.get("trace_id") or "global"
    current = self.totals.setdefault(
        trace_id,
        {"prompt_tokens": 0, "completion_tokens": 0, "total_tokens": 0},
    )
    for key, value in usage.items():
        current[key] = current.get(key, 0) + int(value)

snapshot

snapshot(trace_id: str | None = None) -> dict[str, int]

Return aggregated token counters for one trace or the global bucket.

参数:

名称 类型 描述 默认
trace_id str | None

Identifier for trace.

None

返回:

类型 描述
dict[str, int]

The resulting mapping value.

源代码位于: jianmu/telemetry/sinks.py
def snapshot(self, trace_id: str | None = None) -> dict[str, int]:
    """Return aggregated token counters for one trace or the global bucket.

    Args:
        trace_id: Identifier for trace.

    Returns:
        The resulting mapping value.
    """
    key = trace_id or "global"
    return dict(self.totals.get(key, {}))

reset

reset(trace_id: str | None = None) -> None

Clear aggregated usage for one trace or for all traces.

参数:

名称 类型 描述 默认
trace_id str | None

Identifier for trace.

None
源代码位于: jianmu/telemetry/sinks.py
def reset(self, trace_id: str | None = None) -> None:
    """Clear aggregated usage for one trace or for all traces.

    Args:
        trace_id: Identifier for trace.
    """
    if trace_id is None:
        self.totals.clear()
    else:
        self.totals.pop(trace_id, None)

类型与日志

TelemetryEvent dataclass

TelemetryEvent(
    name: str,
    ts: float,
    trace_id: str | None = None,
    span_id: str | None = None,
    payload: dict[str, Any] = dict(),
)

One emitted telemetry event.

属性:

名称 类型 描述
name str

Event name.

ts float

Event timestamp in seconds since epoch.

trace_id str | None

Optional trace identifier associated with the event.

span_id str | None

Optional active span identifier associated with the event.

payload dict[str, Any]

Structured event payload.

TelemetrySink

Bases: Protocol

Protocol implemented by telemetry sinks.

handle

handle(event: TelemetryEvent) -> None

Consume one telemetry event.

参数:

名称 类型 描述 默认
event TelemetryEvent

Event name or payload to process.

必需
源代码位于: jianmu/telemetry/types.py
def handle(self, event: TelemetryEvent) -> None:
    """Consume one telemetry event.

    Args:
        event: Event name or payload to process.
    """
    ...

configure_logging

configure_logging(
    *,
    level: str = "INFO",
    colorize: bool = True,
    force: bool = False,
) -> None

Configure the shared Loguru logger.

参数:

名称 类型 描述 默认
level str

Logging level as a string. Defaults to "INFO".

'INFO'
colorize bool

Whether log output should be colorized. Defaults to True.

True
force bool

Force configuration update if already configured. Defaults to False.

False
源代码位于: jianmu/telemetry/logging.py
def configure_logging(*, level: str = "INFO", colorize: bool = True, force: bool = False) -> None:
    """Configure the shared Loguru logger.

    Args:
        level: Logging level as a string. Defaults to "INFO".
        colorize: Whether log output should be colorized. Defaults to True.
        force: Force configuration update if already configured. Defaults to False.
    """
    global _CONFIGURED
    if _CONFIGURED and not force:
        return
    logger.remove()
    logger.add(
        sys.stderr,
        format="<level>{level: <8}</level> | <cyan>{name}</cyan>:<cyan>{function}</cyan>:<cyan>{line}</cyan> - <level>{message}</level>",
        level=str(level or "INFO").upper(),
        colorize=bool(colorize),
    )
    _CONFIGURED = True