diff --git a/src/langbot/pkg/telemetry/execution.py b/src/langbot/pkg/telemetry/execution.py index 99a6f9c48..ff5ef28ff 100644 --- a/src/langbot/pkg/telemetry/execution.py +++ b/src/langbot/pkg/telemetry/execution.py @@ -167,10 +167,28 @@ class ExecutionCounters: return payload = self._build_trace_payload(state, reason) if payload is not None: - self.manager.start_send_task(payload) + self._dispatch_payload(payload) except Exception: return + def _dispatch_payload(self, payload: dict) -> None: + """Hand one built payload to the telemetry manager from this sync context. + + TelemetryManager.start_send_task is a coroutine, so calling it without + scheduling dropped every trace closed here. Stand-ins used by tests and + manual tools schedule synchronously, hence the coroutine check. + """ + result = self.manager.start_send_task(payload) + if not asyncio.iscoroutine(result): + return + try: + asyncio.get_running_loop().create_task(result) + except RuntimeError: + result.close() + logger = getattr(getattr(self.manager, 'ap', None), 'logger', None) + if logger is not None: + logger.debug('Execution trace payload dropped: no running event loop') + def _trace_emitted(self, state: TraceState) -> bool: mode = self.trace_mode() if mode == 'off': diff --git a/tests/unit_tests/telemetry/test_trace.py b/tests/unit_tests/telemetry/test_trace.py index 0a567f152..f19ab2668 100644 --- a/tests/unit_tests/telemetry/test_trace.py +++ b/tests/unit_tests/telemetry/test_trace.py @@ -344,3 +344,28 @@ class TestIngress: _, execution = get_modules() with execution.ingress(types.SimpleNamespace(), 'event_done'): pass + + +class AsyncSendManager(FakeManager): + """Stand-in whose start_send_task is a coroutine, like TelemetryManager.""" + + async def start_send_task(self, payload: dict) -> None: + self.sent.append(payload) + + +class TestTraceDispatch: + async def test_close_trace_schedules_the_coroutine_send(self): + trace, execution = get_modules() + manager = AsyncSendManager(trace_config()) + counters = execution.ExecutionCounters(manager) + binding = trace.bind() + try: + counters.record(CONTEXT, **STAGE) + counters.close_trace(binding.state, 'event_done') + # The payload is only delivered once the scheduled task runs. + assert manager.sent == [] + await asyncio.sleep(0) + finally: + trace.unbind_root(binding) + assert len(manager.sent) == 1 + assert manager.sent[0]['event_type'] == 'feature_execution'