mirror of
https://github.com/langbot-app/LangBot.git
synced 2026-10-02 06:16:41 +08:00
feat: bounded execution telemetry for all processing modes (#2617)
* feat: report bounded execution observations via existing telemetry * fix: tolerate optional telemetry context in routing
This commit is contained in:
@@ -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')
|
||||
|
||||
@@ -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(
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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):
|
||||
|
||||
@@ -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
|
||||
@@ -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)
|
||||
@@ -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:
|
||||
|
||||
Reference in New Issue
Block a user