feat(telemetry): trace one inbound event end to end

Execution telemetry already reported window counters. This adds one bounded,
content-free chain per inbound platform event, so a single event can be
followed from its source platform through routing and processing to the
platform API calls it caused.

- telemetry/trace.py: ContextVar trace identity with route and run scopes, and
  a 32-stage bound per chain.
- telemetry/execution.py: keeps at most 64 chains for 120s, decides sampling
  from the observed outcome, and sends one payload per chain keyed by
  query_id = trace_id, so Space fetches a chain by primary identity. Window
  counters are unchanged and remain the source of coverage statistics.
- botmgr / pipelinemgr / orchestrator: bind the trace at the ingress boundary,
  scope routing identity per dispatch, and attach the run identity, so nested
  lanes (event -> route -> pipeline/runner -> platform API) reuse one chain. A
  run reached without an ingress (WebUI debug, service API) owns its own chain.
- space.execution_trace selects off | failures | sampled | all (default
  sampled: failures, WebUI debug runs and every N-th success).

Chains carry only code-defined identifiers (adapter and runner types, route
identity as type:uuid, run id); never message content, tool arguments or
platform user identifiers.
This commit is contained in:
Hyu
2026-10-01 02:12:44 +08:00
parent ae516d2290
commit 0efc1fb1cc
9 changed files with 1082 additions and 38 deletions
@@ -37,6 +37,8 @@ from .run_journal import AgentRunJournal
from .session_registry import AgentRunSessionRegistry, get_session_registry
from .state_scope import build_state_context
from ...provider.tools.loaders import skill as skill_loader
from ...telemetry import trace as trace_mod
from ...telemetry.execution import close_trace as close_execution_trace
from ...telemetry.execution import record as record_execution
@@ -211,8 +213,14 @@ class AgentRunOrchestrator:
terminal_reason: str | None = None
terminal_usage: dict[str, typing.Any] | None = None
execution_outcome = 'unknown'
# A run reached without a platform ingress (WebUI debug, service API)
# owns its own chain; a run inside an ingress reuses that chain.
trace_binding: trace_mod.TraceBinding | None = None
run_token: typing.Any = None
try:
trace_binding = trace_mod.bind()
run_token = trace_mod.set_run(run_id)
await self.journal.create_run(
event=event,
binding=binding,
@@ -419,6 +427,9 @@ class AgentRunOrchestrator:
outcome=execution_outcome,
synthetic=event.source == 'webui',
)
trace_mod.reset_run(run_token)
if trace_binding is not None and trace_mod.unbind_root(trace_binding):
close_execution_trace(self.ap, trace_binding.state, 'runner_done')
binding_box = getattr(execution_query, '_box_binding', None)
if binding_box is not None and binding_box.run_id == run_id:
object.__delattr__(execution_query, '_box_binding')
+5 -1
View File
@@ -367,11 +367,15 @@ class RuntimePipeline:
i += 1
async def process_query(self, query: pipeline_query.Query):
from ..telemetry.execution import ingress
from ..telemetry.platform import processing_mode
token = processing_mode.set('pipeline')
try:
return await self._process_query(query)
# Callers without a platform event (Webchat, HTTP API) still get one
# trace for the whole Pipeline lane; nested calls reuse the trace.
with ingress(self.ap, 'pipeline_done'):
return await self._process_query(query)
finally:
processing_mode.reset(token)
+71 -32
View File
@@ -352,6 +352,20 @@ class RuntimeBot:
"""Return the selected event binding plus per-binding diagnostic steps."""
return self._evaluate_eba_event_bindings(self._get_event_bindings(), event, event_type)
@staticmethod
def _route_ref(
binding: dict | None,
target_type: str | None = None,
target_uuid: str | None = None,
) -> str:
"""Code-defined route identity for telemetry; never a user-facing name."""
binding = binding or {}
kind = str(target_type or binding.get('target_type') or '').strip()
target = str(target_uuid or binding.get('target_uuid') or '').strip()
if not kind or not target:
return ''
return f'{kind}:{target}'[:160]
async def _record_event_route_trace(
self,
*,
@@ -383,17 +397,19 @@ class RuntimeBot:
log_method = getattr(self.logger, level, self.logger.info)
await log_method(text, metadata=metadata)
if status in {'delivered', 'failed', 'discarded', 'not_matched'}:
from ..telemetry import trace as trace_mod
from ..telemetry.execution import record
record(
getattr(self, 'ap', None),
getattr(self, 'execution_context', None),
family='event_route',
operation=event_type,
adapter=type(getattr(self, 'adapter', None)).__name__,
mode=target_type if target_type in {'pipeline', 'agent', 'event_processor'} else 'none',
outcome={'delivered': 'success', 'failed': 'failed'}.get(status, 'skipped'),
)
with trace_mod.scope(route_ref=self._route_ref(binding, target_type, target_uuid)):
record(
getattr(self, 'ap', None),
getattr(self, 'execution_context', None),
family='event_route',
operation=event_type,
adapter=type(getattr(self, 'adapter', None)).__name__,
mode=target_type if target_type in {'pipeline', 'agent', 'event_processor'} else 'none',
outcome={'delivered': 'success', 'failed': 'failed'}.get(status, 'skipped'),
)
return metadata
def get_pipeline_target_for_event_type(self, event_type: str = 'message.received') -> str | None:
@@ -852,7 +868,19 @@ class RuntimeBot:
event: platform_events.EBAEvent,
adapter: abstract_platform_adapter.AbstractMessagePlatformAdapter,
) -> None:
# One inbound event owns one execution trace; every stage recorded while
# it is handled (routing, runner, platform API calls) joins that trace.
from ..telemetry.execution import ingress
event.bot_uuid = self.bot_entity.uuid
with ingress(getattr(self, 'ap', None), 'event_done'):
await self._handle_platform_event_body(event, adapter)
async def _handle_platform_event_body(
self,
event: platform_events.EBAEvent,
adapter: abstract_platform_adapter.AbstractMessagePlatformAdapter,
) -> None:
from ..telemetry.execution import record
record(
@@ -951,12 +979,15 @@ class RuntimeBot:
)
if target_type == 'discard':
if isinstance(event, platform_events.MessageReceivedEvent):
await self._dispatch_eba_message_to_pipeline(
event,
adapter,
pipeline_uuid=self.PIPELINE_DISCARD,
routed_by_event_binding=True,
)
from ..telemetry import trace as trace_mod
with trace_mod.scope(route_ref=self._route_ref(event_binding)):
await self._dispatch_eba_message_to_pipeline(
event,
adapter,
pipeline_uuid=self.PIPELINE_DISCARD,
routed_by_event_binding=True,
)
return await self._record_event_route_trace(
event_type=event_type,
status='discarded',
@@ -984,12 +1015,17 @@ class RuntimeBot:
reason='Pipeline targets only support message events',
text=f'Event {event_type} ignored Pipeline target for non-message event',
)
await self._dispatch_eba_message_to_pipeline(
event,
adapter,
pipeline_uuid=event_binding.get('target_uuid'),
routed_by_event_binding=True,
)
from ..telemetry import trace as trace_mod
with trace_mod.scope(
route_ref=self._route_ref(event_binding, target_type, event_binding.get('target_uuid'))
):
await self._dispatch_eba_message_to_pipeline(
event,
adapter,
pipeline_uuid=event_binding.get('target_uuid'),
routed_by_event_binding=True,
)
return await self._record_event_route_trace(
event_type=event_type,
status='delivered',
@@ -1068,18 +1104,21 @@ class RuntimeBot:
envelope = self._eba_event_to_agent_envelope(event, adapter)
if target_type == 'event_processor':
envelope.data = event.model_dump(mode='json', exclude={'source_platform_object', 'legacy_event'})
from ..telemetry import trace as trace_mod
try:
async for _ in self.ap.agent_run_orchestrator.run(
envelope,
binding,
adapter_context={
'_delivery_adapter': adapter,
'_platform_event': event,
'_execution_context': self.execution_context,
},
):
# Results are journaled by the orchestrator; platform sends require explicit actions.
pass
with trace_mod.scope(route_ref=self._route_ref(event_binding, target_type, target_uuid)):
async for _ in self.ap.agent_run_orchestrator.run(
envelope,
binding,
adapter_context={
'_delivery_adapter': adapter,
'_platform_event': event,
'_execution_context': self.execution_context,
},
):
# Results are journaled by the orchestrator; platform sends require explicit actions.
pass
except Exception:
return await self._record_event_route_trace(
event_type=event_type,
+195 -5
View File
@@ -2,15 +2,29 @@
This module reports observations only. Coverage catalogs and acceptance rules
belong to Space. Aggregation keys include both immutable execution identities.
Two shapes share one payload type:
* window counters, aggregated by ``(family, operation, mode, adapter, runner,
outcome, synthetic)`` — unchanged coverage reporting;
* per-event traces, one bounded payload per traced event, keyed by
``query_id = trace_id`` so Space can fetch a whole chain by primary identity.
Traces are additive: they never alter the counters, and both are bounded in
memory, batch size and send concurrency.
"""
from __future__ import annotations
import asyncio
import contextlib
import time
from datetime import datetime, timezone
from uuid import uuid4
from . import trace as trace_mod
from .identity import workspace_identity
from .trace import TraceState
MAX_KEYS = 512
MAX_BATCH = 32
@@ -18,6 +32,13 @@ FLUSH_SECONDS = 60
MODES = frozenset({'pipeline', 'agent', 'event_processor', 'none'})
OUTCOMES = frozenset({'success', 'failed', 'cancelled', 'timeout', 'skipped', 'unknown'})
# Trace bounds: stages per trace, traces buffered per process, trace lifetime.
MAX_TRACES = 64
TRACE_TTL_SECONDS = 120
TRACE_MODES = frozenset({'off', 'failures', 'sampled', 'all'})
DEFAULT_TRACE_MODE = 'sampled'
DEFAULT_TRACE_SAMPLE = 20
class ExecutionCounters:
def __init__(self, manager):
@@ -25,6 +46,24 @@ class ExecutionCounters:
self.pending: dict[tuple, dict] = {}
self.task: asyncio.Task | None = None
self.dropped = 0
self.traces: dict[str, TraceState] = {}
self.trace_deadlines: dict[str, float] = {}
self.dropped_traces = 0
# ------------------------------------------------------------------ config
def trace_mode(self) -> str:
mode = str(self.manager.telemetry_config.get('execution_trace', DEFAULT_TRACE_MODE) or '').strip().lower()
return mode if mode in TRACE_MODES else DEFAULT_TRACE_MODE
def trace_sample(self) -> int:
try:
sample = int(self.manager.telemetry_config.get('execution_trace_sample', DEFAULT_TRACE_SAMPLE))
except (TypeError, ValueError):
sample = DEFAULT_TRACE_SAMPLE
return sample if 1 <= sample <= 100000 else DEFAULT_TRACE_SAMPLE
# --------------------------------------------------------------- recording
def record(
self,
@@ -50,6 +89,19 @@ class ExecutionCounters:
if any(not isinstance(v, str) or len(v) > 160 for v in (operation, adapter, runner)):
return
identity = workspace_identity(context)
state = trace_mod.current()
if state is not None:
self._record_trace_stage(
state,
identity=identity,
family=family,
operation=operation,
mode=mode,
adapter=adapter,
runner=runner,
outcome=outcome,
synthetic=synthetic,
)
key = (
identity['instance_id'],
identity['workspace_uuid'],
@@ -70,16 +122,124 @@ class ExecutionCounters:
self.pending[key] = row
row['count'] = min(row['count'] + 1, 2147483647)
row['last_seen'] = datetime.now(timezone.utc).isoformat()
if self.task is None or self.task.done():
self.task = asyncio.create_task(self._loop())
self._ensure_loop()
except Exception:
# Observability must never change execution behavior.
return
def _record_trace_stage(
self, state: TraceState, *, identity, family, operation, mode, adapter, runner, outcome, synthetic
) -> None:
if state.abandoned:
return
if state.trace_id not in self.traces:
# Never buffer traces this configuration would discard anyway.
if self.trace_mode() == 'off' or len(self.traces) >= MAX_TRACES:
self.dropped += 1
state.abandoned = True
return
self.traces[state.trace_id] = state
self.trace_deadlines[state.trace_id] = time.monotonic() + TRACE_TTL_SECONDS
state.identity = identity
state.append(
family=family,
operation=operation,
mode=mode,
adapter=adapter,
runner=runner,
outcome=outcome,
synthetic=bool(synthetic),
)
def _ensure_loop(self) -> None:
if self.task is None or self.task.done():
self.task = asyncio.create_task(self._loop())
# ------------------------------------------------------------------ traces
def close_trace(self, state: TraceState, reason: str = 'event_done') -> None:
"""Emit one trace payload when it matches the configured sample."""
try:
if self.traces.pop(state.trace_id, None) is None:
return
self.trace_deadlines.pop(state.trace_id, None)
if not state.stages or not self._trace_emitted(state):
return
payload = self._build_trace_payload(state, reason)
if payload is not None:
self.manager.start_send_task(payload)
except Exception:
return
def _trace_emitted(self, state: TraceState) -> bool:
mode = self.trace_mode()
if mode == 'off':
return False
if mode == 'all':
return True
# A failed, cancelled or timed-out stage is always worth keeping.
if any(stage['outcome'] not in ('success', 'skipped') for stage in state.stages):
return True
# WebUI debug runs are deliberately traced; they stay flagged synthetic.
if state.synthetic:
return True
if mode == 'failures':
return False
try:
return int(state.trace_id[:8], 16) % self.trace_sample() == 0
except ValueError:
return False
def _build_trace_payload(self, state: TraceState, reason: str) -> dict | None:
from ..utils import constants
identity = state.identity
if not identity.get('instance_id') or not identity.get('workspace_uuid'):
return None
observations = [{**stage, 'trace_id': state.trace_id, 'count': 1} for stage in state.stages]
return {
'event_type': 'feature_execution',
# One payload per trace: lookup uses the indexed query_id column.
'query_id': state.trace_id,
'instance_id': identity['instance_id'],
'workspace_uuid': identity['workspace_uuid'],
'version': constants.semantic_version,
'edition': constants.edition,
'timestamp': datetime.now(timezone.utc).isoformat(),
'features': {
'schema': 1,
'observations': observations,
'trace': {
'closed_by': reason,
'started_at': state.started_at,
'ended_at': datetime.now(timezone.utc).isoformat(),
'route_ref': next((s['route_ref'] for s in state.stages if s['route_ref']), ''),
'run_id': next((s['run_id'] for s in state.stages if s['run_id']), ''),
'outcome': state.outcome(),
'dropped_stages': state.dropped_stages,
},
},
}
def _sweep_traces(self, now: float) -> None:
for trace_id in [key for key, deadline in self.trace_deadlines.items() if deadline <= now]:
state = self.traces.get(trace_id)
if state is not None:
self.close_trace(state, 'ttl')
# ----------------------------------------------------------------- flushing
async def _loop(self):
while self.pending:
await asyncio.sleep(FLUSH_SECONDS)
await self.flush()
last_flush = time.monotonic()
while self.pending or self.traces:
# Traces expire on a short TTL; counters keep their minute cadence.
await asyncio.sleep(1 if self.traces else FLUSH_SECONDS)
now = time.monotonic()
self._sweep_traces(now)
if self.pending and now - last_flush >= FLUSH_SECONDS:
last_flush = now
await self.flush()
self.task = None
async def flush(self):
from ..utils import constants
@@ -131,6 +291,11 @@ class ExecutionCounters:
except (Exception, asyncio.CancelledError):
pass
self.pending.clear()
# Traces of in-flight events cannot be completed during shutdown.
self.dropped += len(self.traces)
self.dropped_traces += len(self.traces)
self.traces.clear()
self.trace_deadlines.clear()
def record(ap, context, **observation):
@@ -141,3 +306,28 @@ def record(ap, context, **observation):
counters.record(context, **observation)
except Exception:
pass
def close_trace(ap, state, reason: str = 'event_done') -> None:
"""Best-effort trace close usable with optional telemetry and test doubles."""
try:
counters = getattr(getattr(ap, 'telemetry', None), 'execution', None)
if isinstance(counters, ExecutionCounters):
counters.close_trace(state, reason)
except Exception:
pass
@contextlib.contextmanager
def ingress(ap, reason: str = 'event_done'):
"""Trace one ingress boundary; only the owner of the trace closes it.
Nested boundaries (an event routed into a Pipeline, a Runner invoked from a
dispatch) reuse the in-flight trace instead of starting a second one.
"""
binding = trace_mod.bind()
try:
yield binding
finally:
if trace_mod.unbind_root(binding):
close_trace(ap, binding.state, reason)
+166
View File
@@ -0,0 +1,166 @@
"""Content-free per-event execution traces.
One trace identity is minted per inbound platform event and survives until the
ingress handler returns. Stage records appended under it describe how that one
event was routed and processed, using only code-defined identifiers. Nothing
here is aggregated: a trace is a bounded, ordered sequence for one event.
Three bounds keep memory independent of traffic volume:
* ``MAX_STAGES`` stages per trace (further stages are counted, not kept);
* ``MAX_TRACES`` traces in flight, enforced by the sender that owns the buffer;
* a wall-clock TTL, enforced by the sender's sweep.
The identity itself lives in a ContextVar so that asynchronous work spawned
while handling one event inherits it, mirroring ``telemetry.platform``.
"""
from __future__ import annotations
import contextlib
import contextvars
import typing
from datetime import datetime, timezone
from uuid import uuid4
MAX_STAGES = 32
_current: contextvars.ContextVar['TraceState | None'] = contextvars.ContextVar('telemetry_trace', default=None)
_route: contextvars.ContextVar[str] = contextvars.ContextVar('telemetry_trace_route', default='')
_run: contextvars.ContextVar[str] = contextvars.ContextVar('telemetry_trace_run', default='')
def _now() -> str:
return datetime.now(timezone.utc).isoformat()
class TraceState:
"""Bounded stage buffer for exactly one event."""
__slots__ = (
'trace_id',
'started_at',
'stages',
'dropped_stages',
'synthetic',
'sequence',
'identity',
'abandoned',
)
def __init__(self) -> None:
self.trace_id = str(uuid4())
self.started_at = _now()
self.stages: list[dict[str, typing.Any]] = []
self.dropped_stages = 0
self.synthetic = False
self.sequence = 0
self.identity: dict[str, str] = {}
# Set when the sender refused to buffer this trace: stop appending.
self.abandoned = False
def append(
self,
*,
family: str,
operation: str,
mode: str,
adapter: str,
runner: str,
outcome: str,
synthetic: bool,
) -> None:
if len(self.stages) >= MAX_STAGES:
self.dropped_stages += 1
return
seen = _now()
self.stages.append(
{
'family': family,
'operation': operation,
'mode': mode,
'adapter': adapter,
'runner': runner,
'outcome': outcome,
'synthetic': bool(synthetic),
'seq': self.sequence,
'route_ref': _route.get(),
'run_id': _run.get(),
'first_seen': seen,
'last_seen': seen,
}
)
self.sequence += 1
if synthetic:
self.synthetic = True
def outcome(self) -> str:
"""Terminal outcome of the last recorded stage, for space-side display."""
return self.stages[-1]['outcome'] if self.stages else 'unknown'
class TraceBinding(typing.NamedTuple):
state: TraceState
created: bool
token: typing.Any
def bind() -> TraceBinding:
"""Start a trace unless one is already in flight in this context."""
existing = _current.get()
if existing is not None:
return TraceBinding(existing, False, None)
state = TraceState()
return TraceBinding(state, True, _current.set(state))
def unbind(binding: TraceBinding) -> bool:
"""Detach this binding. Returns True when this caller owns the trace."""
if binding.token is not None:
try:
_current.reset(binding.token)
except (ValueError, RuntimeError):
# A token from another context must never break execution.
pass
return binding.created
def unbind_root(binding: TraceBinding) -> bool:
"""Detach only when this caller started the trace, else leave it in place."""
if not binding.created:
return False
return unbind(binding)
def current() -> TraceState | None:
return _current.get()
def set_run(run_id: str) -> typing.Any:
return _run.set(run_id)
def reset_run(token: typing.Any) -> None:
if token is None:
return
try:
_run.reset(token)
except (ValueError, RuntimeError):
# A token from another context must never break execution.
pass
@contextlib.contextmanager
def scope(*, route_ref: str | None = None) -> typing.Iterator[None]:
"""Pin routing identity for stages recorded inside this block."""
if not route_ref:
yield
return
token = _route.set(route_ref)
try:
yield
finally:
try:
_route.reset(token)
except (ValueError, RuntimeError):
pass
+5
View File
@@ -429,3 +429,8 @@ space:
# Optional telemetry is enabled by default. Set true to disable usage,
# heartbeat and execution reporting; restart after changing this setting.
disable_telemetry: false
# Per-event execution traces: off | failures | sampled | all.
# 'sampled' keeps every failed run plus every N-th successful one
# (N is execution_trace_sample) and always keeps WebUI debug runs.
execution_trace: sampled
execution_trace_sample: 20
@@ -0,0 +1,97 @@
"""Dump one real execution-trace payload for the Space wire-format contract test.
Run from the repository root with the project interpreter, e.g.:
PYTHONPATH=src python tests/manual/dump_execution_trace_payload.py
`langbot-space` keeps the output as
`internal/service/testdata/execution_trace_payload.json` and asserts that its
ingest validator accepts it. Regenerate that file whenever the payload shape
changes; the Space test documents the same command.
"""
from __future__ import annotations
import asyncio
import json
import sys
import types
from importlib import import_module
class CaptureManager:
"""Captures what TelemetryManager would have posted."""
def __init__(self):
self.telemetry_config = {'url': 'https://space.langbot.test', 'execution_trace': 'all'}
self.sent: list[dict] = []
async def send(self, payload: dict) -> bool:
self.sent.append(payload)
return True
def start_send_task(self, payload: dict) -> None:
self.sent.append(payload)
async def main() -> int:
execution = import_module('langbot.pkg.telemetry.execution')
trace = import_module('langbot.pkg.telemetry.trace')
manager = CaptureManager()
counters = execution.ExecutionCounters(manager)
context = types.SimpleNamespace(instance_uuid='instance-1', workspace_uuid='workspace-1')
route_ref = 'agent:11111111-1111-4111-8111-111111111111'
run_id = '22222222-2222-4222-8222-222222222222'
with execution.ingress(types.SimpleNamespace(telemetry=types.SimpleNamespace(execution=counters))):
# Inbound platform event.
counters.record(
context,
family='platform_event',
operation='message.received',
adapter='AiocqhttpAdapter',
outcome='success',
)
# Route decision for that event.
with trace.scope(route_ref=route_ref):
counters.record(
context,
family='event_route',
operation='message.received',
mode='agent',
adapter='AiocqhttpAdapter',
outcome='success',
)
# Runner execution of the routed processor.
run_token = trace.set_run(run_id)
try:
counters.record(
context,
family='runner',
operation='execute',
mode='agent',
runner='langbot/runner-demo',
outcome='success',
)
# Outbound platform API call made by the runner.
counters.record(
context,
family='platform_api',
operation='send_message',
mode='agent',
adapter='AiocqhttpAdapter',
outcome='success',
)
finally:
trace.reset_run(run_token)
if len(manager.sent) != 1:
print(f'expected exactly one payload, captured {len(manager.sent)}', file=sys.stderr)
return 1
print(json.dumps(manager.sent[0], indent=2, ensure_ascii=False, sort_keys=True))
return 0
if __name__ == '__main__':
raise SystemExit(asyncio.run(main()))
@@ -0,0 +1,186 @@
"""RuntimeBot ingress tracing: one inbound event owns one execution trace."""
from __future__ import annotations
from types import SimpleNamespace
from unittest.mock import AsyncMock, Mock
import pytest
from langbot.pkg.api.http.context import ExecutionContext
TEST_CONTEXT = ExecutionContext(
instance_uuid='instance-test',
workspace_uuid='workspace-test',
placement_generation=1,
bot_uuid='bot-1',
)
class FakeTelemetryManager:
def __init__(self, config=None):
self.telemetry_config = (
{'url': 'https://space.example.test', 'execution_trace': 'all'} if config is None else config
)
self.sent: list[dict] = []
async def send(self, payload: dict) -> bool:
self.sent.append(payload)
return True
def start_send_task(self, payload: dict) -> None:
self.sent.append(payload)
def make_bot(event_bindings: list[dict], config=None, agent=None):
from langbot.pkg.platform.botmgr import RuntimeBot
from langbot.pkg.telemetry.execution import ExecutionCounters
manager = FakeTelemetryManager(config)
counters = ExecutionCounters(manager)
bot = object.__new__(RuntimeBot)
bot.bot_entity = SimpleNamespace(
uuid='bot-1',
workspace_uuid=TEST_CONTEXT.workspace_uuid,
event_bindings=event_bindings,
plugin_processors=[],
)
bot.execution_context = TEST_CONTEXT
bot.workspace_uuid = TEST_CONTEXT.workspace_uuid
bot.placement_generation = TEST_CONTEXT.placement_generation
bot.logger = SimpleNamespace(info=AsyncMock(), warning=AsyncMock(), error=AsyncMock(), debug=AsyncMock())
bot.adapter = FakeAdapter()
bot.ap = SimpleNamespace(
telemetry=SimpleNamespace(execution=counters),
agent_service=SimpleNamespace(get_agent=AsyncMock(return_value=agent)),
pipeline_service=SimpleNamespace(get_pipeline=AsyncMock(return_value=None)),
)
return bot, manager, counters
def message_received_event():
from langbot_plugin.api.entities.builtin.platform import entities, events, message
return events.MessageReceivedEvent(
message_id='message-1',
message_chain=message.MessageChain([message.Plain(text='hello')]),
sender=entities.User(id='user-1', nickname='QA User'),
chat_type=entities.ChatType.PRIVATE,
chat_id='user-1',
)
class FakeAdapter:
async def get_supported_apis(self):
return []
@pytest.mark.asyncio
async def test_route_miss_emits_one_closed_trace():
bot, manager, counters = make_bot([])
await bot._handle_platform_event(message_received_event(), FakeAdapter())
assert len(manager.sent) == 1
payload = manager.sent[0]
assert payload['event_type'] == 'feature_execution'
assert len(payload['query_id']) == 36
assert payload['instance_id'] == 'instance-test'
assert payload['workspace_uuid'] == 'workspace-test'
features = payload['features']
assert features['schema'] == 1
assert [(row['family'], row['operation'], row['outcome']) for row in features['observations']] == [
('platform_event', 'message.received', 'success'),
('event_route', 'message.received', 'skipped'),
]
assert [row['seq'] for row in features['observations']] == [0, 1]
assert all(row['trace_id'] == payload['query_id'] for row in features['observations'])
assert all(row['adapter'] == 'FakeAdapter' for row in features['observations'])
assert features['trace']['closed_by'] == 'event_done'
# Window counters are still aggregated for coverage.
assert counters.pending
@pytest.mark.asyncio
async def test_unavailable_route_target_carries_route_identity():
bot, manager, _ = make_bot(
[
{
'id': 'agent-binding',
'enabled': True,
'event_pattern': 'message.received',
'target_type': 'agent',
'target_uuid': 'agent-1',
}
]
)
await bot._handle_platform_event(message_received_event(), FakeAdapter())
stages = manager.sent[0]['features']['observations']
assert [(stage['family'], stage['outcome'], stage['mode']) for stage in stages] == [
('platform_event', 'success', 'none'),
('event_route', 'failed', 'agent'),
]
assert stages[1]['route_ref'] == 'agent:agent-1'
assert stages[0]['route_ref'] == ''
@pytest.mark.asyncio
async def test_discarded_route_is_visible_without_user_values():
bot, manager, _ = make_bot(
[
{
'id': 'discard-binding',
'enabled': True,
'event_pattern': '*',
'target_type': 'discard',
}
]
)
await bot._handle_platform_event(SimpleNamespace(type='platform.member.joined'), FakeAdapter())
stages = manager.sent[0]['features']['observations']
assert [stage['outcome'] for stage in stages] == ['success', 'skipped']
assert stages[1]['operation'] == 'platform.member.joined'
assert stages[1]['mode'] == 'none'
@pytest.mark.asyncio
async def test_telemetry_opt_out_emits_nothing():
bot, manager, counters = make_bot([], config={'url': 'https://space.example.test', 'disable_telemetry': True})
await bot._handle_platform_event(message_received_event(), FakeAdapter())
assert manager.sent == []
assert counters.pending == {}
assert counters.traces == {}
@pytest.mark.asyncio
async def test_route_trace_scope_is_restored_after_dispatch():
from langbot.pkg.telemetry import trace as trace_mod
bot, _, _ = make_bot([])
assert trace_mod.current() is None
await bot._handle_platform_event(message_received_event(), FakeAdapter())
assert trace_mod.current() is None
# A later record outside the ingress must not join the closed trace.
from langbot.pkg.telemetry.execution import record
record(bot.ap, TEST_CONTEXT, family='platform_api', operation='send_message', outcome='success')
assert list(bot.ap.telemetry.execution.pending.values())[-1]['count'] == 1
@pytest.mark.asyncio
async def test_adapter_call_without_ingress_still_counts():
"""Mock adapter objects used by other tests must not break the ingress path."""
bot, manager, _ = make_bot([])
adapter = Mock()
await bot._handle_platform_event(message_received_event(), adapter)
assert manager.sent
+346
View File
@@ -0,0 +1,346 @@
"""Unit tests for per-event execution traces (pkg/telemetry/trace.py + execution.py)."""
from __future__ import annotations
import asyncio
import time
import types
from importlib import import_module
def get_modules():
return (
import_module('langbot.pkg.telemetry.trace'),
import_module('langbot.pkg.telemetry.execution'),
)
class FakeManager:
"""Stand-in for TelemetryManager: records payloads instead of posting them."""
def __init__(self, config=None):
self.telemetry_config = (
{'url': 'https://space.example.test', 'disable_telemetry': False} if config is None else config
)
self.sent: list[dict] = []
async def send(self, payload: dict) -> bool:
self.sent.append(payload)
return True
def start_send_task(self, payload: dict) -> None:
self.sent.append(payload)
def make_counters(config=None):
_, execution = get_modules()
manager = FakeManager(config)
return manager, execution.ExecutionCounters(manager)
def make_ap(counters):
return types.SimpleNamespace(telemetry=types.SimpleNamespace(execution=counters))
def trace_config(**overrides):
config = {'url': 'https://space.example.test', 'execution_trace': 'all'}
config.update(overrides)
return config
CONTEXT = types.SimpleNamespace(instance_uuid='instance-1', workspace_uuid='workspace-1')
STAGE = {
'family': 'platform_event',
'operation': 'message.received',
'mode': 'none',
'adapter': 'AiocqhttpAdapter',
'outcome': 'success',
}
class TestTraceIdentity:
def test_nested_bind_reuses_the_in_flight_trace(self):
trace, _ = get_modules()
outer = trace.bind()
inner = trace.bind()
try:
assert outer.created is True
assert inner.created is False
assert inner.state is outer.state
assert trace.current() is outer.state
finally:
trace.unbind_root(inner)
trace.unbind_root(outer)
assert trace.current() is None
def test_unbind_root_leaves_a_nested_trace_attached(self):
trace, _ = get_modules()
outer = trace.bind()
inner = trace.bind()
assert trace.unbind_root(inner) is False
assert trace.current() is outer.state
assert trace.unbind_root(outer) is True
assert trace.current() is None
# A repeated unbind of the same binding must never raise.
assert trace.unbind_root(outer) is True
def test_route_and_run_are_only_pinned_inside_the_block(self):
trace, _ = get_modules()
binding = trace.bind()
token = trace.set_run('run-1')
try:
with trace.scope(route_ref='agent:agent-1'):
inside = trace.current()
inside.append(
family='runner',
operation='execute',
mode='agent',
adapter='',
runner='langbot/runner',
outcome='success',
synthetic=False,
)
outside = trace.current()
outside.append(
family='platform_api',
operation='send_message',
mode='agent',
adapter='AiocqhttpAdapter',
runner='',
outcome='success',
synthetic=False,
)
finally:
trace.reset_run(token)
trace.unbind_root(binding)
assert [(s['route_ref'], s['run_id']) for s in binding.state.stages] == [
('agent:agent-1', 'run-1'),
('', 'run-1'),
]
class TestTraceRecording:
def test_stage_carries_sequence_identity_and_timestamps(self):
trace, _ = get_modules()
_, counters = make_counters()
binding = trace.bind()
try:
counters.record(CONTEXT, **STAGE)
counters.record(CONTEXT, **{**STAGE, 'family': 'runner', 'operation': 'execute', 'mode': 'pipeline'})
finally:
trace.unbind_root(binding)
stages = binding.state.stages
assert [stage['seq'] for stage in stages] == [0, 1]
assert stages[0]['first_seen'] == stages[0]['last_seen']
assert binding.state.identity == {'instance_id': 'instance-1', 'workspace_uuid': 'workspace-1'}
assert binding.state.outcome() == 'success'
def test_stage_buffer_is_bounded_and_counts_overflow(self):
trace, _ = get_modules()
_, counters = make_counters()
binding = trace.bind()
try:
for _ in range(trace.MAX_STAGES + 5):
counters.record(CONTEXT, **STAGE)
finally:
trace.unbind_root(binding)
assert len(binding.state.stages) == trace.MAX_STAGES
assert binding.state.dropped_stages == 5
def test_oversized_identifier_never_becomes_a_stage(self):
trace, _ = get_modules()
_, counters = make_counters()
binding = trace.bind()
try:
counters.record(CONTEXT, **{**STAGE, 'operation': 'x' * 200})
counters.record(CONTEXT, **{**STAGE, 'family': 'not-a-family'})
finally:
trace.unbind_root(binding)
assert binding.state.stages == []
def test_trace_limit_abandons_further_traces(self):
trace, execution = get_modules()
_, counters = make_counters()
bindings = []
for _ in range(execution.MAX_TRACES + 1):
binding = trace.bind()
bindings.append(binding)
counters.record(CONTEXT, **STAGE)
trace.unbind_root(binding)
assert trace.current() is None
assert len(counters.traces) == execution.MAX_TRACES
assert counters.dropped == 1
assert bindings[-1].state.abandoned is True
assert bindings[-1].state.stages == []
def test_counters_keep_aggregating_while_tracing(self):
trace, _ = get_modules()
_, counters = make_counters()
binding = trace.bind()
try:
for _ in range(3):
counters.record(CONTEXT, **STAGE)
finally:
trace.unbind_root(binding)
assert list(counters.pending.values())[0]['count'] == 3
class TestTraceDelivery:
def test_close_emits_one_payload_per_trace(self):
trace, _ = get_modules()
manager, counters = make_counters(trace_config())
binding = trace.bind()
try:
counters.record(CONTEXT, **STAGE)
counters.record(
CONTEXT, **{**STAGE, 'family': 'runner', 'operation': 'execute', 'mode': 'agent', 'runner': 'runner-1'}
)
counters.close_trace(binding.state, 'event_done')
counters.close_trace(binding.state, 'event_done')
finally:
trace.unbind_root(binding)
assert len(manager.sent) == 1
payload = manager.sent[0]
assert payload['event_type'] == 'feature_execution'
assert payload['query_id'] == binding.state.trace_id
assert len(payload['query_id']) == 36
assert payload['instance_id'] == 'instance-1'
assert payload['workspace_uuid'] == 'workspace-1'
features = payload['features']
assert features['schema'] == 1
assert [row['count'] for row in features['observations']] == [1, 1]
assert [row['seq'] for row in features['observations']] == [0, 1]
assert [row['trace_id'] for row in features['observations']] == [binding.state.trace_id] * 2
assert features['trace'] == {
'closed_by': 'event_done',
'started_at': binding.state.started_at,
'ended_at': features['trace']['ended_at'],
'route_ref': '',
'run_id': '',
'outcome': 'success',
'dropped_stages': 0,
}
assert counters.traces == {}
def test_sampling_decision_is_deterministic_and_keeps_failures(self):
trace, _ = get_modules()
_, counters = make_counters(trace_config(execution_trace='sampled', execution_trace_sample=2))
binding = trace.bind()
try:
binding.state.trace_id = '00000001-0000-4000-8000-000000000000'
counters.record(CONTEXT, **STAGE)
assert counters._trace_emitted(binding.state) is False
binding.state.trace_id = '00000000-0000-4000-8000-000000000000'
assert counters._trace_emitted(binding.state) is True
binding.state.trace_id = '00000001-0000-4000-8000-000000000001'
counters.record(CONTEXT, **{**STAGE, 'outcome': 'timeout'})
assert counters._trace_emitted(binding.state) is True
finally:
trace.unbind_root(binding)
def test_failures_mode_keeps_failed_and_debug_traces_only(self):
trace, _ = get_modules()
manager, counters = make_counters(trace_config(execution_trace='failures', execution_trace_sample=1000000))
failed = trace.bind()
counters.record(CONTEXT, **{**STAGE, 'outcome': 'failed'})
counters.close_trace(failed.state, 'event_done')
trace.unbind_root(failed)
debug = trace.bind()
counters.record(CONTEXT, **{**STAGE, 'synthetic': True})
counters.close_trace(debug.state, 'event_done')
trace.unbind_root(debug)
sampled_out = trace.bind()
counters.record(CONTEXT, **STAGE)
counters.close_trace(sampled_out.state, 'event_done')
trace.unbind_root(sampled_out)
assert [payload['features']['observations'][0]['outcome'] for payload in manager.sent] == ['failed', 'success']
assert manager.sent[1]['features']['observations'][0]['synthetic'] is True
def test_off_mode_and_opt_out_emit_nothing(self):
trace, _ = get_modules()
manager, counters = make_counters(trace_config(execution_trace='off'))
binding = trace.bind()
try:
counters.record(CONTEXT, **{**STAGE, 'outcome': 'failed'})
counters.close_trace(binding.state, 'event_done')
finally:
trace.unbind_root(binding)
assert manager.sent == []
disabled_manager, disabled = make_counters(
{'url': 'https://space.example.test', 'disable_telemetry': True, 'execution_trace': 'all'}
)
binding = trace.bind()
try:
disabled.record(CONTEXT, **STAGE)
disabled.close_trace(binding.state, 'event_done')
finally:
trace.unbind_root(binding)
assert disabled_manager.sent == []
assert disabled.pending == {}
assert disabled.traces == {}
def test_ttl_sweep_closes_stale_trace(self):
trace, _ = get_modules()
manager, counters = make_counters(trace_config())
binding = trace.bind()
try:
counters.record(CONTEXT, **{**STAGE, 'outcome': 'failed'})
counters.trace_deadlines[binding.state.trace_id] = 0
counters._sweep_traces(time.monotonic())
finally:
trace.unbind_root(binding)
assert manager.sent[0]['features']['trace']['closed_by'] == 'ttl'
assert counters.traces == {}
assert counters.trace_deadlines == {}
async def test_shutdown_drops_open_traces(self):
trace, _ = get_modules()
manager, counters = make_counters(trace_config())
binding = trace.bind()
try:
counters.record(CONTEXT, **{**STAGE, 'outcome': 'failed'})
await counters.shutdown()
finally:
trace.unbind_root(binding)
# Shutdown may still flush window counters; it must not claim a trace.
assert [payload for payload in manager.sent if 'trace' in payload['features']] == []
assert counters.traces == {}
assert counters.dropped_traces == 1
async def test_loop_sweeps_and_sends_expired_trace(self):
trace, _ = get_modules()
manager, counters = make_counters(trace_config())
binding = trace.bind()
try:
counters.record(CONTEXT, **{**STAGE, 'outcome': 'failed'})
counters.trace_deadlines[binding.state.trace_id] = 0
await asyncio.sleep(1.2)
assert [payload['features']['trace']['closed_by'] for payload in manager.sent] == ['ttl']
assert counters.traces == {}
finally:
await counters.shutdown()
trace.unbind_root(binding)
class TestIngress:
def test_ingress_closes_only_the_trace_it_started(self):
_, execution = get_modules()
manager, counters = make_counters(trace_config())
ap = make_ap(counters)
with execution.ingress(ap, 'event_done'):
with execution.ingress(ap, 'pipeline_done'):
counters.record(CONTEXT, **{**STAGE, 'outcome': 'failed'})
assert len(manager.sent) == 1
assert manager.sent[0]['features']['trace']['closed_by'] == 'event_done'
def test_ingress_without_telemetry_is_a_no_op(self):
_, execution = get_modules()
with execution.ingress(types.SimpleNamespace(), 'event_done'):
pass