diff --git a/src/langbot/pkg/agent/runner/orchestrator.py b/src/langbot/pkg/agent/runner/orchestrator.py index 55aba94e7..b9d34e460 100644 --- a/src/langbot/pkg/agent/runner/orchestrator.py +++ b/src/langbot/pkg/agent/runner/orchestrator.py @@ -3,6 +3,7 @@ from __future__ import annotations import time +import asyncio import contextlib import typing @@ -36,6 +37,7 @@ 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.execution import record as record_execution ACTIVATED_SKILL_NAMES_STATE_KEY = 'host.activated_skills' @@ -208,6 +210,7 @@ class AgentRunOrchestrator: terminal_status: str | None = None terminal_reason: str | None = None terminal_usage: dict[str, typing.Any] | None = None + execution_outcome = 'unknown' try: await self.journal.create_run( @@ -389,7 +392,14 @@ class AgentRunOrchestrator: status_reason=terminal_reason, usage=terminal_usage, ) + execution_outcome = {'completed': 'success', 'failed': 'failed', 'cancelled': 'cancelled'}.get( + terminal_status or '', 'unknown' + ) + except asyncio.CancelledError: + execution_outcome = 'cancelled' + raise except Exception as exc: + execution_outcome = 'timeout' if self._is_deadline_exhausted(context) else 'failed' failed_usage = terminal_usage await self.journal.finalize_run( run_id=run_id, @@ -399,6 +409,16 @@ class AgentRunOrchestrator: ) raise finally: + record_execution( + self.ap, + execution_context, + family='runner', + operation='execute', + mode=binding.processor_type, + runner=descriptor.id, + outcome=execution_outcome, + synthetic=event.source == 'webui', + ) 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/agent/runner/platform_tools.py b/src/langbot/pkg/agent/runner/platform_tools.py index c3ba18458..5599c4953 100644 --- a/src/langbot/pkg/agent/runner/platform_tools.py +++ b/src/langbot/pkg/agent/runner/platform_tools.py @@ -689,6 +689,17 @@ async def execute_platform_tool( result = _execute_mock_platform_tool(definition, context, normalized) if message_chain is not None: result['parameters']['message'] = message_chain.model_dump(mode='json') + from ...telemetry.execution import record + + record( + ap, + execution_context, + family='platform_api', + operation=definition.api, + mode=authorization.get('processor_type', 'none'), + synthetic=True, + outcome='success', + ) return result bot_id = authorization.get('bot_id') if not bot_id: @@ -709,7 +720,13 @@ async def execute_platform_tool( if message_chain is not None else platform_message.MessageChain([platform_message.Plain(text=_require_string(normalized, 'text'))]), } - return await api_func(**normalized) + from ...telemetry.platform import processing_mode + + token = processing_mode.set(authorization.get('processor_type', 'none')) + try: + return await api_func(**normalized) + finally: + processing_mode.reset(token) def _execute_mock_platform_tool( diff --git a/src/langbot/pkg/pipeline/pipelinemgr.py b/src/langbot/pkg/pipeline/pipelinemgr.py index 2f429bc60..a32ee7145 100644 --- a/src/langbot/pkg/pipeline/pipelinemgr.py +++ b/src/langbot/pkg/pipeline/pipelinemgr.py @@ -367,6 +367,15 @@ class RuntimePipeline: i += 1 async def process_query(self, query: pipeline_query.Query): + from ..telemetry.platform import processing_mode + + token = processing_mode.set('pipeline') + try: + return await self._process_query(query) + finally: + processing_mode.reset(token) + + async def _process_query(self, query: pipeline_query.Query): """处理请求""" await self._assert_execution_active(query) # Get monitoring metadata diff --git a/src/langbot/pkg/platform/botmgr.py b/src/langbot/pkg/platform/botmgr.py index 401abc2db..8f068e37c 100644 --- a/src/langbot/pkg/platform/botmgr.py +++ b/src/langbot/pkg/platform/botmgr.py @@ -382,6 +382,18 @@ 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.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'), + ) return metadata def get_pipeline_target_for_event_type(self, event_type: str = 'message.received') -> str | None: @@ -841,6 +853,16 @@ class RuntimeBot: adapter: abstract_platform_adapter.AbstractMessagePlatformAdapter, ) -> None: event.bot_uuid = self.bot_entity.uuid + from ..telemetry.execution import record + + record( + getattr(self, 'ap', None), + getattr(self, 'execution_context', None), + family='platform_event', + operation=event.type, + adapter=adapter.__class__.__name__, + outcome='success', + ) await self._record_adapter_event(event, adapter) primary = ( @@ -1468,6 +1490,10 @@ class RuntimeBot: ) async def initialize(self): + from ..telemetry.platform import observe_adapter + + observe_adapter(self.ap, self.execution_context, self.adapter) + def tenant_scoped_listener(listener): @functools.wraps(listener) async def wrapped(*args, **kwargs): diff --git a/src/langbot/pkg/telemetry/execution.py b/src/langbot/pkg/telemetry/execution.py new file mode 100644 index 000000000..201037368 --- /dev/null +++ b/src/langbot/pkg/telemetry/execution.py @@ -0,0 +1,143 @@ +"""Bounded, content-free execution counters for the existing telemetry sender. + +This module reports observations only. Coverage catalogs and acceptance rules +belong to Space. Aggregation keys include both immutable execution identities. +""" + +from __future__ import annotations + +import asyncio +from datetime import datetime, timezone +from uuid import uuid4 + +from .identity import workspace_identity + +MAX_KEYS = 512 +MAX_BATCH = 32 +FLUSH_SECONDS = 60 +MODES = frozenset({'pipeline', 'agent', 'event_processor', 'none'}) +OUTCOMES = frozenset({'success', 'failed', 'cancelled', 'timeout', 'skipped', 'unknown'}) + + +class ExecutionCounters: + def __init__(self, manager): + self.manager = manager + self.pending: dict[tuple, dict] = {} + self.task: asyncio.Task | None = None + self.dropped = 0 + + def record( + self, + context, + *, + family: str, + operation: str, + mode: str = 'none', + adapter: str = '', + runner: str = '', + outcome: str = 'unknown', + synthetic: bool = False, + ): + try: + cfg = self.manager.telemetry_config + if not cfg or cfg.get('disable_telemetry', False) or not cfg.get('url'): + return + if family not in {'platform_event', 'event_route', 'runner', 'platform_api'}: + return + if mode not in MODES or outcome not in OUTCOMES: + return + # Only code-defined identifiers are accepted; never pass user values. + if any(not isinstance(v, str) or len(v) > 160 for v in (operation, adapter, runner)): + return + identity = workspace_identity(context) + key = ( + identity['instance_id'], + identity['workspace_uuid'], + family, + operation, + mode, + adapter, + runner, + outcome, + bool(synthetic), + ) + row = self.pending.get(key) + if row is None: + if len(self.pending) >= MAX_KEYS: + self.dropped += 1 + return + row = {'count': 0, 'first_seen': datetime.now(timezone.utc).isoformat()} + 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()) + except Exception: + # Observability must never change execution behavior. + return + + async def _loop(self): + while self.pending: + await asyncio.sleep(FLUSH_SECONDS) + await self.flush() + + async def flush(self): + from ..utils import constants + + # Remove only one bounded batch; subsequent windows drain the remainder. + # Drain one tenant per minute, at most one request, round-robin by insertion. + if not self.pending: + return + tenant = next(iter(self.pending))[:2] + keys = [key for key in self.pending if key[:2] == tenant][:MAX_BATCH] + groups: dict[tuple[str, str], list[dict]] = {} + for key in keys: + row = self.pending.pop(key) + instance, workspace, family, operation, mode, adapter, runner, outcome, synthetic = key + groups.setdefault((instance, workspace), []).append( + { + **row, + 'family': family, + 'operation': operation, + 'mode': mode, + 'adapter': adapter, + 'runner': runner, + 'outcome': outcome, + 'synthetic': synthetic, + } + ) + for (instance, workspace), observations in groups.items(): + payload = { + 'event_type': 'feature_execution', + 'query_id': str(uuid4()), + 'instance_id': instance, + 'workspace_uuid': workspace, + 'version': constants.semantic_version, + 'edition': constants.edition, + 'timestamp': datetime.now(timezone.utc).isoformat(), + 'features': {'schema': 1, 'observations': observations}, + } + if not await self.manager.send(payload): + await asyncio.sleep(1) + await self.manager.send(payload) + + async def shutdown(self): + if self.task is not None: + self.task.cancel() + await asyncio.gather(self.task, return_exceptions=True) + self.task = None + try: + await asyncio.wait_for(self.flush(), timeout=2) + except (Exception, asyncio.CancelledError): + pass + self.pending.clear() + + +def record(ap, context, **observation): + """Best-effort bridge usable with optional telemetry and test doubles.""" + try: + counters = getattr(getattr(ap, 'telemetry', None), 'execution', None) + if isinstance(counters, ExecutionCounters): + counters.record(context, **observation) + except Exception: + pass diff --git a/src/langbot/pkg/telemetry/platform.py b/src/langbot/pkg/telemetry/platform.py new file mode 100644 index 000000000..a7e636bcf --- /dev/null +++ b/src/langbot/pkg/telemetry/platform.py @@ -0,0 +1,99 @@ +"""Observe adapter calls without collecting their arguments or response content.""" + +from __future__ import annotations + +import asyncio +import functools +from contextvars import ContextVar + +from .execution import record + +processing_mode: ContextVar[str] = ContextVar('telemetry_processing_mode', default='none') + + +def result_outcome(result): + # Empty returns do not prove a remote operation succeeded. + if result is None: + return 'unknown' + if result is False: + return 'failed' + raw = getattr(result, 'raw', result) + if isinstance(raw, dict): + if raw.get('ok') is False or raw.get('success') is False or raw.get('status') == 'failed': + return 'failed' + for key in ('errcode', 'retcode'): + if key in raw and raw[key] not in (0, '0', None): + return 'failed' + if raw.get('error'): + return 'failed' + if isinstance(raw.get('results'), list): + outcomes = [result_outcome(item) for item in raw['results']] + if 'failed' in outcomes: + return 'failed' + if not outcomes or 'unknown' in outcomes: + return 'unknown' + return 'success' + # Arbitrary error envelopes or empty mappings are not affirmative evidence. + if isinstance(raw, dict): + if raw.get('ok') is True or raw.get('success') is True or raw.get('status') == 'ok': + return 'success' + if any(raw.get(key) not in (None, '') for key in ('message_id', 'id')): + return 'success' + if any(key in raw and raw[key] in (0, '0') for key in ('errcode', 'retcode')): + return 'success' + return 'unknown' + return 'success' if result is True or getattr(result, 'message_id', None) is not None else 'unknown' + + +def observe_adapter(ap, context, adapter): + """Install once on a concrete bot; identity never comes from task-local tenants.""" + if getattr(adapter, '_execution_observed', False): + return + try: + declared = frozenset(adapter.get_supported_apis()) + except Exception: + return + for name in declared: + if '.' in name: + continue + original = getattr(adapter, name, None) + if not asyncio.iscoroutinefunction(original): + continue + + def make_wrapper(method, operation): + @functools.wraps(method) + async def wrapped(*args, **kwargs): + observed_operation = operation + if operation == 'call_platform_api': + action = args[0] if args else kwargs.get('action') + if action in declared: + observed_operation = action + outcome = 'unknown' + try: + result = await method(*args, **kwargs) + outcome = result_outcome(result) + return result + except asyncio.CancelledError: + outcome = 'cancelled' + raise + except TimeoutError: + outcome = 'timeout' + raise + except Exception: + outcome = 'failed' + raise + finally: + record( + ap, + context, + family='platform_api', + operation=observed_operation, + adapter=adapter.__class__.__name__, + mode=processing_mode.get(), + outcome=outcome, + ) + + return wrapped + + setattr(adapter, name, make_wrapper(original, name)) + setattr(adapter, '_execution_observed', True) diff --git a/src/langbot/pkg/telemetry/telemetry.py b/src/langbot/pkg/telemetry/telemetry.py index 3d3e29898..e4aaf6fd8 100644 --- a/src/langbot/pkg/telemetry/telemetry.py +++ b/src/langbot/pkg/telemetry/telemetry.py @@ -9,6 +9,7 @@ import httpx from ..core import app as core_app from ..utils import httpclient +from .execution import ExecutionCounters _MAX_INFLIGHT_TELEMETRY_TASKS = 8 @@ -28,6 +29,7 @@ class TelemetryManager: self.telemetry_config: dict[str, typing.Any] = {} self.send_tasks: list[asyncio.Task] = [] self._client: httpx.AsyncClient | None = None + self.execution = ExecutionCounters(self) async def initialize(self): self.telemetry_config = self.ap.instance_config.data.get('space', {}) @@ -48,6 +50,7 @@ class TelemetryManager: pass async def shutdown(self) -> None: + await self.execution.shutdown() tasks = list(self.send_tasks) for task in tasks: task.cancel() @@ -194,6 +197,8 @@ class TelemetryManager: self.ap.logger.debug( f'Telemetry posted to {url}, status {resp.status_code} - response: {body}' ) + if not app_err: + return True except asyncio.TimeoutError: self.ap.logger.warning(f'Telemetry post to {url} timed out') except Exception as e: