Skip to content

Commit 0efc1fb

Browse files
committed
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.
1 parent ae516d2 commit 0efc1fb

9 files changed

Lines changed: 1082 additions & 38 deletions

File tree

‎src/langbot/pkg/agent/runner/orchestrator.py‎

Lines changed: 11 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -37,6 +37,8 @@
3737
from .session_registry import AgentRunSessionRegistry, get_session_registry
3838
from .state_scope import build_state_context
3939
from ...provider.tools.loaders import skill as skill_loader
40+
from ...telemetry import trace as trace_mod
41+
from ...telemetry.execution import close_trace as close_execution_trace
4042
from ...telemetry.execution import record as record_execution
4143

4244

@@ -211,8 +213,14 @@ async def run(
211213
terminal_reason: str | None = None
212214
terminal_usage: dict[str, typing.Any] | None = None
213215
execution_outcome = 'unknown'
216+
# A run reached without a platform ingress (WebUI debug, service API)
217+
# owns its own chain; a run inside an ingress reuses that chain.
218+
trace_binding: trace_mod.TraceBinding | None = None
219+
run_token: typing.Any = None
214220

215221
try:
222+
trace_binding = trace_mod.bind()
223+
run_token = trace_mod.set_run(run_id)
216224
await self.journal.create_run(
217225
event=event,
218226
binding=binding,
@@ -419,6 +427,9 @@ async def run(
419427
outcome=execution_outcome,
420428
synthetic=event.source == 'webui',
421429
)
430+
trace_mod.reset_run(run_token)
431+
if trace_binding is not None and trace_mod.unbind_root(trace_binding):
432+
close_execution_trace(self.ap, trace_binding.state, 'runner_done')
422433
binding_box = getattr(execution_query, '_box_binding', None)
423434
if binding_box is not None and binding_box.run_id == run_id:
424435
object.__delattr__(execution_query, '_box_binding')

‎src/langbot/pkg/pipeline/pipelinemgr.py‎

Lines changed: 5 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -367,11 +367,15 @@ async def _execute_from_stage(
367367
i += 1
368368

369369
async def process_query(self, query: pipeline_query.Query):
370+
from ..telemetry.execution import ingress
370371
from ..telemetry.platform import processing_mode
371372

372373
token = processing_mode.set('pipeline')
373374
try:
374-
return await self._process_query(query)
375+
# Callers without a platform event (Webchat, HTTP API) still get one
376+
# trace for the whole Pipeline lane; nested calls reuse the trace.
377+
with ingress(self.ap, 'pipeline_done'):
378+
return await self._process_query(query)
375379
finally:
376380
processing_mode.reset(token)
377381

‎src/langbot/pkg/platform/botmgr.py‎

Lines changed: 71 additions & 32 deletions
Original file line numberDiff line numberDiff line change
@@ -352,6 +352,20 @@ def diagnose_eba_event_binding(
352352
"""Return the selected event binding plus per-binding diagnostic steps."""
353353
return self._evaluate_eba_event_bindings(self._get_event_bindings(), event, event_type)
354354

355+
@staticmethod
356+
def _route_ref(
357+
binding: dict | None,
358+
target_type: str | None = None,
359+
target_uuid: str | None = None,
360+
) -> str:
361+
"""Code-defined route identity for telemetry; never a user-facing name."""
362+
binding = binding or {}
363+
kind = str(target_type or binding.get('target_type') or '').strip()
364+
target = str(target_uuid or binding.get('target_uuid') or '').strip()
365+
if not kind or not target:
366+
return ''
367+
return f'{kind}:{target}'[:160]
368+
355369
async def _record_event_route_trace(
356370
self,
357371
*,
@@ -383,17 +397,19 @@ async def _record_event_route_trace(
383397
log_method = getattr(self.logger, level, self.logger.info)
384398
await log_method(text, metadata=metadata)
385399
if status in {'delivered', 'failed', 'discarded', 'not_matched'}:
400+
from ..telemetry import trace as trace_mod
386401
from ..telemetry.execution import record
387402

388-
record(
389-
getattr(self, 'ap', None),
390-
getattr(self, 'execution_context', None),
391-
family='event_route',
392-
operation=event_type,
393-
adapter=type(getattr(self, 'adapter', None)).__name__,
394-
mode=target_type if target_type in {'pipeline', 'agent', 'event_processor'} else 'none',
395-
outcome={'delivered': 'success', 'failed': 'failed'}.get(status, 'skipped'),
396-
)
403+
with trace_mod.scope(route_ref=self._route_ref(binding, target_type, target_uuid)):
404+
record(
405+
getattr(self, 'ap', None),
406+
getattr(self, 'execution_context', None),
407+
family='event_route',
408+
operation=event_type,
409+
adapter=type(getattr(self, 'adapter', None)).__name__,
410+
mode=target_type if target_type in {'pipeline', 'agent', 'event_processor'} else 'none',
411+
outcome={'delivered': 'success', 'failed': 'failed'}.get(status, 'skipped'),
412+
)
397413
return metadata
398414

399415
def get_pipeline_target_for_event_type(self, event_type: str = 'message.received') -> str | None:
@@ -852,7 +868,19 @@ async def _handle_platform_event(
852868
event: platform_events.EBAEvent,
853869
adapter: abstract_platform_adapter.AbstractMessagePlatformAdapter,
854870
) -> None:
871+
# One inbound event owns one execution trace; every stage recorded while
872+
# it is handled (routing, runner, platform API calls) joins that trace.
873+
from ..telemetry.execution import ingress
874+
855875
event.bot_uuid = self.bot_entity.uuid
876+
with ingress(getattr(self, 'ap', None), 'event_done'):
877+
await self._handle_platform_event_body(event, adapter)
878+
879+
async def _handle_platform_event_body(
880+
self,
881+
event: platform_events.EBAEvent,
882+
adapter: abstract_platform_adapter.AbstractMessagePlatformAdapter,
883+
) -> None:
856884
from ..telemetry.execution import record
857885

858886
record(
@@ -951,12 +979,15 @@ async def _dispatch_eba_event_to_processor(
951979
)
952980
if target_type == 'discard':
953981
if isinstance(event, platform_events.MessageReceivedEvent):
954-
await self._dispatch_eba_message_to_pipeline(
955-
event,
956-
adapter,
957-
pipeline_uuid=self.PIPELINE_DISCARD,
958-
routed_by_event_binding=True,
959-
)
982+
from ..telemetry import trace as trace_mod
983+
984+
with trace_mod.scope(route_ref=self._route_ref(event_binding)):
985+
await self._dispatch_eba_message_to_pipeline(
986+
event,
987+
adapter,
988+
pipeline_uuid=self.PIPELINE_DISCARD,
989+
routed_by_event_binding=True,
990+
)
960991
return await self._record_event_route_trace(
961992
event_type=event_type,
962993
status='discarded',
@@ -984,12 +1015,17 @@ async def _dispatch_eba_event_to_processor(
9841015
reason='Pipeline targets only support message events',
9851016
text=f'Event {event_type} ignored Pipeline target for non-message event',
9861017
)
987-
await self._dispatch_eba_message_to_pipeline(
988-
event,
989-
adapter,
990-
pipeline_uuid=event_binding.get('target_uuid'),
991-
routed_by_event_binding=True,
992-
)
1018+
from ..telemetry import trace as trace_mod
1019+
1020+
with trace_mod.scope(
1021+
route_ref=self._route_ref(event_binding, target_type, event_binding.get('target_uuid'))
1022+
):
1023+
await self._dispatch_eba_message_to_pipeline(
1024+
event,
1025+
adapter,
1026+
pipeline_uuid=event_binding.get('target_uuid'),
1027+
routed_by_event_binding=True,
1028+
)
9931029
return await self._record_event_route_trace(
9941030
event_type=event_type,
9951031
status='delivered',
@@ -1068,18 +1104,21 @@ async def _dispatch_eba_event_to_processor(
10681104
envelope = self._eba_event_to_agent_envelope(event, adapter)
10691105
if target_type == 'event_processor':
10701106
envelope.data = event.model_dump(mode='json', exclude={'source_platform_object', 'legacy_event'})
1107+
from ..telemetry import trace as trace_mod
1108+
10711109
try:
1072-
async for _ in self.ap.agent_run_orchestrator.run(
1073-
envelope,
1074-
binding,
1075-
adapter_context={
1076-
'_delivery_adapter': adapter,
1077-
'_platform_event': event,
1078-
'_execution_context': self.execution_context,
1079-
},
1080-
):
1081-
# Results are journaled by the orchestrator; platform sends require explicit actions.
1082-
pass
1110+
with trace_mod.scope(route_ref=self._route_ref(event_binding, target_type, target_uuid)):
1111+
async for _ in self.ap.agent_run_orchestrator.run(
1112+
envelope,
1113+
binding,
1114+
adapter_context={
1115+
'_delivery_adapter': adapter,
1116+
'_platform_event': event,
1117+
'_execution_context': self.execution_context,
1118+
},
1119+
):
1120+
# Results are journaled by the orchestrator; platform sends require explicit actions.
1121+
pass
10831122
except Exception:
10841123
return await self._record_event_route_trace(
10851124
event_type=event_type,

0 commit comments

Comments
 (0)