Skip to content

Commit df64f83

Browse files
committed
fix(telemetry): schedule the trace payload send instead of dropping it
close_trace() called the coroutine TelemetryManager.start_send_task without awaiting or scheduling it, so every payload closed on that path was silently dropped and the runtime logged "coroutine ... was never awaited". The only other caller awaits it correctly, so both call styles stay supported: call the manager, schedule the returned coroutine on the running loop, and close it with a debug log when no loop is available. Regression case added in tests/unit_tests/telemetry/test_trace.py: it fails before this change (payload never delivered) and passes after.
1 parent 0efc1fb commit df64f83

2 files changed

Lines changed: 44 additions & 1 deletion

File tree

‎src/langbot/pkg/telemetry/execution.py‎

Lines changed: 19 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -167,10 +167,28 @@ def close_trace(self, state: TraceState, reason: str = 'event_done') -> None:
167167
return
168168
payload = self._build_trace_payload(state, reason)
169169
if payload is not None:
170-
self.manager.start_send_task(payload)
170+
self._dispatch_payload(payload)
171171
except Exception:
172172
return
173173

174+
def _dispatch_payload(self, payload: dict) -> None:
175+
"""Hand one built payload to the telemetry manager from this sync context.
176+
177+
TelemetryManager.start_send_task is a coroutine, so calling it without
178+
scheduling dropped every trace closed here. Stand-ins used by tests and
179+
manual tools schedule synchronously, hence the coroutine check.
180+
"""
181+
result = self.manager.start_send_task(payload)
182+
if not asyncio.iscoroutine(result):
183+
return
184+
try:
185+
asyncio.get_running_loop().create_task(result)
186+
except RuntimeError:
187+
result.close()
188+
logger = getattr(getattr(self.manager, 'ap', None), 'logger', None)
189+
if logger is not None:
190+
logger.debug('Execution trace payload dropped: no running event loop')
191+
174192
def _trace_emitted(self, state: TraceState) -> bool:
175193
mode = self.trace_mode()
176194
if mode == 'off':

‎tests/unit_tests/telemetry/test_trace.py‎

Lines changed: 25 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -344,3 +344,28 @@ def test_ingress_without_telemetry_is_a_no_op(self):
344344
_, execution = get_modules()
345345
with execution.ingress(types.SimpleNamespace(), 'event_done'):
346346
pass
347+
348+
349+
class AsyncSendManager(FakeManager):
350+
"""Stand-in whose start_send_task is a coroutine, like TelemetryManager."""
351+
352+
async def start_send_task(self, payload: dict) -> None:
353+
self.sent.append(payload)
354+
355+
356+
class TestTraceDispatch:
357+
async def test_close_trace_schedules_the_coroutine_send(self):
358+
trace, execution = get_modules()
359+
manager = AsyncSendManager(trace_config())
360+
counters = execution.ExecutionCounters(manager)
361+
binding = trace.bind()
362+
try:
363+
counters.record(CONTEXT, **STAGE)
364+
counters.close_trace(binding.state, 'event_done')
365+
# The payload is only delivered once the scheduled task runs.
366+
assert manager.sent == []
367+
await asyncio.sleep(0)
368+
finally:
369+
trace.unbind_root(binding)
370+
assert len(manager.sent) == 1
371+
assert manager.sent[0]['event_type'] == 'feature_execution'

0 commit comments

Comments
 (0)