From 0efc1fb1ccb88c642ff881e77b86bc67d5f6b243 Mon Sep 17 00:00:00 2001 From: Hyu Date: Thu, 1 Oct 2026 02:12:44 +0800 Subject: [PATCH] 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. --- src/langbot/pkg/agent/runner/orchestrator.py | 11 + src/langbot/pkg/pipeline/pipelinemgr.py | 6 +- src/langbot/pkg/platform/botmgr.py | 103 ++++-- src/langbot/pkg/telemetry/execution.py | 200 +++++++++- src/langbot/pkg/telemetry/trace.py | 166 +++++++++ src/langbot/templates/config.yaml | 5 + tests/manual/dump_execution_trace_payload.py | 97 +++++ .../platform/test_botmgr_execution_trace.py | 186 ++++++++++ tests/unit_tests/telemetry/test_trace.py | 346 ++++++++++++++++++ 9 files changed, 1082 insertions(+), 38 deletions(-) create mode 100644 src/langbot/pkg/telemetry/trace.py create mode 100644 tests/manual/dump_execution_trace_payload.py create mode 100644 tests/unit_tests/platform/test_botmgr_execution_trace.py create mode 100644 tests/unit_tests/telemetry/test_trace.py diff --git a/src/langbot/pkg/agent/runner/orchestrator.py b/src/langbot/pkg/agent/runner/orchestrator.py index b9d34e460..e207e743a 100644 --- a/src/langbot/pkg/agent/runner/orchestrator.py +++ b/src/langbot/pkg/agent/runner/orchestrator.py @@ -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') diff --git a/src/langbot/pkg/pipeline/pipelinemgr.py b/src/langbot/pkg/pipeline/pipelinemgr.py index a32ee7145..a737594d4 100644 --- a/src/langbot/pkg/pipeline/pipelinemgr.py +++ b/src/langbot/pkg/pipeline/pipelinemgr.py @@ -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) diff --git a/src/langbot/pkg/platform/botmgr.py b/src/langbot/pkg/platform/botmgr.py index 8f068e37c..d4430f55d 100644 --- a/src/langbot/pkg/platform/botmgr.py +++ b/src/langbot/pkg/platform/botmgr.py @@ -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, diff --git a/src/langbot/pkg/telemetry/execution.py b/src/langbot/pkg/telemetry/execution.py index 201037368..99a6f9c48 100644 --- a/src/langbot/pkg/telemetry/execution.py +++ b/src/langbot/pkg/telemetry/execution.py @@ -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) diff --git a/src/langbot/pkg/telemetry/trace.py b/src/langbot/pkg/telemetry/trace.py new file mode 100644 index 000000000..a239c707e --- /dev/null +++ b/src/langbot/pkg/telemetry/trace.py @@ -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 diff --git a/src/langbot/templates/config.yaml b/src/langbot/templates/config.yaml index edd64149f..19c46ae35 100644 --- a/src/langbot/templates/config.yaml +++ b/src/langbot/templates/config.yaml @@ -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 diff --git a/tests/manual/dump_execution_trace_payload.py b/tests/manual/dump_execution_trace_payload.py new file mode 100644 index 000000000..0cd852f09 --- /dev/null +++ b/tests/manual/dump_execution_trace_payload.py @@ -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())) diff --git a/tests/unit_tests/platform/test_botmgr_execution_trace.py b/tests/unit_tests/platform/test_botmgr_execution_trace.py new file mode 100644 index 000000000..12f55fde6 --- /dev/null +++ b/tests/unit_tests/platform/test_botmgr_execution_trace.py @@ -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 diff --git a/tests/unit_tests/telemetry/test_trace.py b/tests/unit_tests/telemetry/test_trace.py new file mode 100644 index 000000000..0a567f152 --- /dev/null +++ b/tests/unit_tests/telemetry/test_trace.py @@ -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