Skip to content

Commit 923ebe1

Browse files
committed
feat(telemetry): attach called platform APIs to their step node on execution traces
Every host step now opens a stage scope, so observations recorded inside it (including cross-task plugin/RPC calls resolved through the execution registry) carry the owning step as parent; the pipeline lane records its own pipeline/run node so acks and replies issued outside the runner still hang off a node.
1 parent df64f83 commit 923ebe1

12 files changed

Lines changed: 1803 additions & 494 deletions

File tree

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

Lines changed: 132 additions & 121 deletions
Large diffs are not rendered by default.

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

Lines changed: 62 additions & 40 deletions
Original file line numberDiff line numberDiff line change
@@ -684,49 +684,71 @@ async def execute_platform_tool(
684684
normalized = _normalize_platform_params(definition, parameters)
685685
if definition.scope == 'event':
686686
normalized = _event_params(definition, context, normalized)
687-
# This flag is frozen by the Host from the synthetic debug envelope, not tool arguments.
688-
if delivery.get('surface') == 'webui' and (delivery.get('platform_capabilities') or {}).get('debug_mock') is True:
689-
result = _execute_mock_platform_tool(definition, context, normalized)
690-
if message_chain is not None:
691-
result['parameters']['message'] = message_chain.model_dump(mode='json')
692-
from ...telemetry.execution import record
693-
694-
record(
695-
ap,
696-
execution_context,
697-
family='platform_api',
698-
operation=definition.api,
699-
mode=authorization.get('processor_type', 'none'),
700-
synthetic=True,
701-
outcome='success',
702-
)
703-
return result
704-
bot_id = authorization.get('bot_id')
705-
if not bot_id:
706-
raise ValueError('This run is not associated with a platform bot')
707-
bot = await ap.platform_mgr.get_bot_by_uuid(execution_context, bot_id)
708-
if bot is None:
709-
raise ValueError(f'Bot {bot_id} is not running')
710-
if definition.api not in set(bot.adapter.get_supported_apis() or []):
711-
raise ValueError(f'Platform API {definition.api} is no longer supported by bot {bot_id}')
712-
api_func = getattr(bot.adapter, definition.api, None)
713-
if not callable(api_func):
714-
raise ValueError(f'Platform API {definition.api} is declared but not implemented')
715-
if definition.api == 'send_message':
716-
normalized = {
717-
'target_type': _require_string(normalized, 'target_type'),
718-
'target_id': _require_string(normalized, 'target_id'),
719-
'message': message_chain
720-
if message_chain is not None
721-
else platform_message.MessageChain([platform_message.Plain(text=_require_string(normalized, 'text'))]),
722-
}
723-
from ...telemetry.platform import processing_mode
687+
# Plugin/RPC actions run outside the ingress context, so bind the owning
688+
# execution id here: nested adapter observations then join its trace.
689+
from ...telemetry.execution import record, reset_execution_id, set_execution_id
724690

725-
token = processing_mode.set(authorization.get('processor_type', 'none'))
691+
execution_id = str(session.get('run_id') or '').strip()
692+
execution_token = set_execution_id(execution_id)
726693
try:
727-
return await api_func(**normalized)
694+
# This flag is frozen by the Host from the synthetic debug envelope, not tool arguments.
695+
mock = (
696+
delivery.get('surface') == 'webui'
697+
and (delivery.get('platform_capabilities') or {}).get('debug_mock') is True
698+
)
699+
if mock:
700+
outcome = 'unknown'
701+
error_detail = ''
702+
try:
703+
result = _execute_mock_platform_tool(definition, context, normalized)
704+
if message_chain is not None:
705+
result['parameters']['message'] = message_chain.model_dump(mode='json')
706+
outcome = 'success'
707+
except Exception as exc:
708+
outcome = 'failed'
709+
error_detail = str(exc)
710+
raise
711+
finally:
712+
record(
713+
ap,
714+
execution_context,
715+
family='platform_api',
716+
operation=definition.api,
717+
mode=authorization.get('processor_type', 'none'),
718+
synthetic=True,
719+
outcome=outcome,
720+
error=error_detail,
721+
execution_id=execution_id or None,
722+
)
723+
return result
724+
bot_id = authorization.get('bot_id')
725+
if not bot_id:
726+
raise ValueError('This run is not associated with a platform bot')
727+
bot = await ap.platform_mgr.get_bot_by_uuid(execution_context, bot_id)
728+
if bot is None:
729+
raise ValueError(f'Bot {bot_id} is not running')
730+
if definition.api not in set(bot.adapter.get_supported_apis() or []):
731+
raise ValueError(f'Platform API {definition.api} is no longer supported by bot {bot_id}')
732+
api_func = getattr(bot.adapter, definition.api, None)
733+
if not callable(api_func):
734+
raise ValueError(f'Platform API {definition.api} is declared but not implemented')
735+
if definition.api == 'send_message':
736+
normalized = {
737+
'target_type': _require_string(normalized, 'target_type'),
738+
'target_id': _require_string(normalized, 'target_id'),
739+
'message': message_chain
740+
if message_chain is not None
741+
else platform_message.MessageChain([platform_message.Plain(text=_require_string(normalized, 'text'))]),
742+
}
743+
from ...telemetry.platform import processing_mode
744+
745+
token = processing_mode.set(authorization.get('processor_type', 'none'))
746+
try:
747+
return await api_func(**normalized)
748+
finally:
749+
processing_mode.reset(token)
728750
finally:
729-
processing_mode.reset(token)
751+
reset_execution_id(execution_token)
730752

731753

732754
def _execute_mock_platform_tool(

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

Lines changed: 94 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -3,6 +3,7 @@
33
import dataclasses
44
import typing
55
import traceback
6+
import asyncio
67

78
import sqlalchemy
89

@@ -257,6 +258,9 @@ async def _check_output(self, query: pipeline_query.Query, result: pipeline_enti
257258
self.ap.logger.error(result.error_notice)
258259
# Mark query as having error
259260
query.variables['_monitoring_has_error'] = True
261+
# The lane reports failures as a value instead of raising, so record the
262+
# reason here: without it the uploaded trace would not explain the break.
263+
self._record_lane_failure(query, str(result.error_notice))
260264
# Record error to monitoring system
261265
try:
262266
await self._assert_execution_active(query)
@@ -368,14 +372,59 @@ async def _execute_from_stage(
368372

369373
async def process_query(self, query: pipeline_query.Query):
370374
from ..telemetry.execution import ingress
375+
from ..telemetry.execution import record as record_execution
371376
from ..telemetry.platform import processing_mode
377+
from ..telemetry.trace import stage_scope
372378

373379
token = processing_mode.set('pipeline')
380+
# The Workspace-scoped opaque query uuid is the execution identity Space
381+
# shows for this lane; it is a stable UUID for every pooled query.
382+
execution_id = str(getattr(query, 'query_uuid', '') or '').strip() or str(query.query_id)
374383
try:
375384
# Callers without a platform event (Webchat, HTTP API) still get one
376385
# 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)
386+
with ingress(self.ap, 'pipeline_done', getattr(self, 'execution_context', None), execution_id=execution_id):
387+
# The lane is its own workflow step: acknowledgements and replies it
388+
# sends directly, and the runner it drives, hang under this node.
389+
# When the lane runs under a platform route the node nests there.
390+
with stage_scope() as node:
391+
lane_outcome = 'success'
392+
lane_error = ''
393+
try:
394+
return await self._process_query(query)
395+
except asyncio.CancelledError:
396+
lane_outcome = 'cancelled'
397+
lane_error = 'cancelled'
398+
raise
399+
except BaseException as exc:
400+
lane_outcome = 'failed'
401+
lane_error = str(exc) or type(exc).__name__
402+
raise
403+
finally:
404+
try:
405+
variables = getattr(query, 'variables', None) or {}
406+
lane_has_error = bool(variables.get('_monitoring_has_error'))
407+
except Exception:
408+
lane_has_error = False
409+
if lane_outcome == 'success' and lane_has_error:
410+
# The lane reported the failure as a value, not an exception.
411+
lane_outcome = 'failed'
412+
config = getattr(query, 'pipeline_config', None)
413+
try:
414+
runner_id = (RunnerConfigResolver.resolve_runner_id(config) or '') if config else ''
415+
except Exception:
416+
runner_id = ''
417+
record_execution(
418+
self.ap,
419+
getattr(self, 'execution_context', None),
420+
family='pipeline',
421+
operation='run',
422+
mode='pipeline',
423+
runner=runner_id,
424+
outcome=lane_outcome,
425+
error=lane_error,
426+
node=node,
427+
)
379428
finally:
380429
processing_mode.reset(token)
381430

@@ -490,6 +539,9 @@ async def _process_query(self, query: pipeline_query.Query):
490539
inst_name = query.current_stage_name if query.current_stage_name else 'unknown'
491540
self.ap.logger.error(f'Error processing query {query.query_id} stage={inst_name} : {e}')
492541
self.ap.logger.error(f'Traceback: {traceback.format_exc()}')
542+
# The lane itself broke: land the reason on the execution trace so the
543+
# chain is uploaded even though the error never became a StageProcessResult.
544+
self._record_lane_failure(query, str(e) or type(e).__name__)
493545

494546
# Record query error
495547
try:
@@ -513,6 +565,46 @@ async def _process_query(self, query: pipeline_query.Query):
513565
self.ap.logger.debug(f'Query {query.query_id} processed')
514566
await self.ap.query_pool.remove_query(query)
515567

568+
def _record_lane_failure(self, query: pipeline_query.Query, reason: str) -> None:
569+
"""Land one failed ``runner/execute`` stage on the owning execution trace."""
570+
try:
571+
from ..telemetry.execution import record as record_execution
572+
from ..telemetry.trace import current as current_trace
573+
from ..telemetry.trace import stage_scope
574+
575+
detail = reason or 'pipeline_lane_failed'
576+
state = current_trace()
577+
# The runner orchestrator already lands its own failed ``runner/execute``
578+
# node when the break happens inside a runner; in that case only carry the
579+
# reason over instead of repeating the node on the trace.
580+
already_recorded = bool(
581+
state is not None
582+
and any(
583+
stage.get('family') == 'runner' and stage.get('outcome') not in ('success', 'skipped')
584+
for stage in state.stages
585+
)
586+
)
587+
if not already_recorded:
588+
# The lane failure is its own workflow step, hanging under whatever
589+
# step routed into it (or a root when nothing did).
590+
with stage_scope() as node:
591+
record_execution(
592+
self.ap,
593+
getattr(self, 'execution_context', None),
594+
family='runner',
595+
operation='execute',
596+
mode='pipeline',
597+
runner=(RunnerConfigResolver.resolve_runner_id(query.pipeline_config) or ''),
598+
outcome='failed',
599+
error=detail,
600+
node=node,
601+
)
602+
state = current_trace()
603+
if state is not None:
604+
state.mark_failure(detail)
605+
except Exception:
606+
pass
607+
516608

517609
class PipelineManager:
518610
"""流水线管理器"""

‎src/langbot/pkg/pipeline/process/handlers/chat.py‎

Lines changed: 34 additions & 62 deletions
Original file line numberDiff line numberDiff line change
@@ -3,8 +3,6 @@
33
import uuid
44
import typing
55
import traceback
6-
import time
7-
from datetime import datetime
86

97

108
from .. import handler
@@ -14,9 +12,7 @@
1412
import langbot_plugin.api.entities.events as events
1513
from ....agent.runner.config_resolver import RunnerConfigResolver
1614
from ....agent.runner import config_schema
17-
from ....utils import constants, runner as runner_utils
18-
from ....telemetry import features as telemetry_features
19-
from ....telemetry.identity import workspace_identity
15+
from ....utils import runner as runner_utils
2016
import langbot_plugin.api.entities.builtin.provider.session as provider_session
2117
import langbot_plugin.api.entities.builtin.pipeline.query as pipeline_query
2218
import langbot_plugin.api.entities.builtin.provider.message as provider_message
@@ -115,9 +111,6 @@ async def handle(
115111
text_length = 0
116112
runner = None
117113
try:
118-
# Mark start time for telemetry
119-
start_ts = time.time()
120-
121114
try_claim_steering = getattr(
122115
self.ap.agent_run_orchestrator,
123116
'try_claim_steering_from_query',
@@ -258,69 +251,48 @@ async def handle(
258251
debug_notice=traceback.format_exc(),
259252
)
260253
finally:
261-
# Telemetry reporting
254+
# Telemetry: the per-execution record is built when the owning
255+
# trace closes, so only attach the execution-scoped fields this
256+
# lane knows here.
262257
try:
263-
end_ts = time.time()
264-
duration_ms = None
265-
if 'start_ts' in locals():
266-
duration_ms = int((end_ts - start_ts) * 1000)
258+
from ....telemetry.trace import current as current_trace
259+
260+
state = current_trace()
261+
if state is not None:
262+
adapter_name = query.adapter.__class__.__name__ if hasattr(query, 'adapter') else ''
263+
runner_name = self.ap.agent_run_orchestrator.resolve_runner_id_for_telemetry(query)
264+
state.adapter = state.adapter or adapter_name
265+
state.runner = state.runner or (runner_name or '')
266+
state.runner_category = state.runner_category or runner_utils.get_runner_category_from_runner(
267+
runner_name, None, query.pipeline_config
268+
)
269+
model_name = ''
270+
try:
271+
if getattr(query, 'use_llm_model_uuid', None):
272+
m = await self.ap.model_mgr.get_model_by_uuid(
273+
get_query_execution_context(query),
274+
query.use_llm_model_uuid,
275+
)
276+
if m and getattr(m, 'model_entity', None):
277+
model_name = getattr(m.model_entity, 'name', '') or ''
278+
except Exception:
279+
model_name = ''
280+
state.model_name = state.model_name or model_name
281+
if state.pipeline_plugins is None:
282+
state.pipeline_plugins = query.variables.get('_pipeline_bound_plugins', None)
283+
except Exception as ex:
284+
self.ap.logger.warning(f'Failed to attach execution telemetry fields: {ex}')
267285

286+
# Trigger survey events on successful non-WebSocket responses
287+
try:
268288
adapter_name = query.adapter.__class__.__name__ if hasattr(query, 'adapter') else None
269-
270-
# Use orchestrator to resolve runner ID for telemetry
271-
runner_name = self.ap.agent_run_orchestrator.resolve_runner_id_for_telemetry(query)
272-
273-
# Model name if available
274-
model_name = None
275-
try:
276-
if getattr(query, 'use_llm_model_uuid', None):
277-
m = await self.ap.model_mgr.get_model_by_uuid(
278-
get_query_execution_context(query),
279-
query.use_llm_model_uuid,
280-
)
281-
if m and getattr(m, 'model_entity', None):
282-
model_name = getattr(m.model_entity, 'name', None)
283-
except Exception:
284-
model_name = None
285-
286-
pipeline_plugins = query.variables.get('_pipeline_bound_plugins', None)
287-
288-
runner_category = runner_utils.get_runner_category_from_runner(
289-
runner_name, None, query.pipeline_config
290-
)
291-
292-
# Feature usage collected during query processing (tool calls,
293-
# knowledge base usage, sandbox executions, activated skills, ...)
294-
features = telemetry_features.collect_features(query)
295-
296-
payload = {
297-
'event_type': 'query',
298-
'query_id': query.query_id,
299-
'adapter': adapter_name,
300-
'runner': runner_name,
301-
'runner_category': runner_category,
302-
'duration_ms': duration_ms,
303-
'model_name': model_name,
304-
'version': constants.semantic_version,
305-
**workspace_identity(get_query_execution_context(query)),
306-
'runtime_instance_id': constants.instance_id,
307-
'edition': constants.edition,
308-
'pipeline_plugins': pipeline_plugins,
309-
'features': features,
310-
'error': locals().get('error_info', None),
311-
'timestamp': datetime.utcnow().isoformat(),
312-
}
313-
314-
await self.ap.telemetry.start_send_task(payload)
315-
316-
# Trigger survey events on successful non-WebSocket responses
317289
if not locals().get('error_info') and adapter_name and 'WebSocket' not in adapter_name:
318290
if self.ap.survey:
319291
await self.ap.survey.trigger_event('first_bot_response_success')
320292
# Counts toward the bot_response_success_100 milestone event
321293
await self.ap.survey.record_bot_response_success()
322294
except Exception as ex:
323-
self.ap.logger.warning(f'Failed to send telemetry: {ex}')
295+
self.ap.logger.warning(f'Failed to trigger survey event: {ex}')
324296

325297
async def _ensure_conversation_for_history(
326298
self,

0 commit comments

Comments
 (0)