mirror of
https://github.com/BerriAI/litellm.git
synced 2026-09-14 23:21:35 +00:00
fix(logging): keep tracebacks for callback-raised timeouts and flush the summary on stop
This commit is contained in:
parent
68509daeaf
commit
defead0eab
2 changed files with 124 additions and 7 deletions
|
|
@ -147,6 +147,10 @@ class LoggingWorker:
|
|||
self._sem = None
|
||||
self._worker_task = None
|
||||
self._running_tasks.clear()
|
||||
# The summary task is bound to the old loop; drop it and its pending burst so
|
||||
# timeouts on the new loop arm a fresh summary instead of a stale, stuck one.
|
||||
self._timeout_summary_task = None
|
||||
self._timeout_burst_count = 0
|
||||
self._queue = new_queue
|
||||
self._bound_loop = current_loop
|
||||
return
|
||||
|
|
@ -167,14 +171,17 @@ class LoggingWorker:
|
|||
"""Runs the logging task and handles cleanup. Releases semaphore when done."""
|
||||
try:
|
||||
if self._queue is not None:
|
||||
# Run the coroutine in its original context
|
||||
callback_task: Final = task["context"].run(asyncio.create_task, task["coroutine"])
|
||||
try:
|
||||
# Run the coroutine in its original context
|
||||
await asyncio.wait_for(
|
||||
task["context"].run(asyncio.create_task, task["coroutine"]),
|
||||
timeout=self.timeout,
|
||||
)
|
||||
except asyncio.TimeoutError:
|
||||
self._record_callback_timeout(task["coroutine"])
|
||||
await asyncio.wait_for(callback_task, timeout=self.timeout)
|
||||
except asyncio.TimeoutError as e:
|
||||
# wait_for cancels the callback when our own deadline expires; a callback that
|
||||
# raised TimeoutError itself did not hit our deadline, so keep its traceback.
|
||||
if callback_task.cancelled():
|
||||
self._record_callback_timeout(task["coroutine"])
|
||||
else:
|
||||
verbose_logger.exception("LoggingWorker error: %s", e)
|
||||
except Exception as e:
|
||||
verbose_logger.exception("LoggingWorker error: %s", e)
|
||||
finally:
|
||||
|
|
@ -197,6 +204,10 @@ class LoggingWorker:
|
|||
async def _flush_timeout_summary(self) -> None:
|
||||
"""After the burst settles, log one bounded summary covering every timeout in it."""
|
||||
await asyncio.sleep(self.timeout_summary_window)
|
||||
self._emit_timeout_summary()
|
||||
|
||||
def _emit_timeout_summary(self) -> None:
|
||||
"""Log one bounded summary for the current burst and reset the burst counter."""
|
||||
burst_count: Final = self._timeout_burst_count
|
||||
self._timeout_burst_count = 0
|
||||
if burst_count <= 0:
|
||||
|
|
@ -444,6 +455,13 @@ class LoggingWorker:
|
|||
|
||||
async def stop(self) -> None:
|
||||
"""Stop the logging worker and clean up resources."""
|
||||
if self._timeout_summary_task is not None:
|
||||
# Cancel the debounced sleeper and emit any pending burst now, so a summary
|
||||
# armed just before shutdown is not lost when the loop tears down.
|
||||
self._timeout_summary_task.cancel()
|
||||
self._timeout_summary_task = None
|
||||
self._emit_timeout_summary()
|
||||
|
||||
if self._worker_task is None and not self._running_tasks:
|
||||
# No worker launched and no in-flight tasks to drain.
|
||||
return
|
||||
|
|
|
|||
|
|
@ -597,3 +597,102 @@ class TestLoggingWorker:
|
|||
assert errors[0].exc_info is not None, "the traceback must be preserved"
|
||||
assert warnings == [], "a single real error is not a timeout summary"
|
||||
assert worker._timeout_summary_task is None
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_callback_raised_timeout_keeps_traceback(self):
|
||||
"""A callback that raises TimeoutError on its own did not hit the worker's deadline,
|
||||
so it is a real failure and must keep its traceback rather than being folded into the
|
||||
bounded burst summary.
|
||||
"""
|
||||
worker = LoggingWorker(timeout=5.0, max_queue_size=10, timeout_summary_window=1.0)
|
||||
|
||||
async def raises_own_timeout() -> None:
|
||||
raise asyncio.TimeoutError("callback's own downstream timeout")
|
||||
|
||||
logger = logging.getLogger("LiteLLM")
|
||||
collector = _RecordCollector()
|
||||
previous_level = logger.level
|
||||
logger.addHandler(collector)
|
||||
logger.setLevel(logging.DEBUG)
|
||||
try:
|
||||
worker.start()
|
||||
worker.enqueue(raises_own_timeout())
|
||||
|
||||
await worker.flush()
|
||||
await worker.stop()
|
||||
finally:
|
||||
logger.removeHandler(collector)
|
||||
logger.setLevel(previous_level)
|
||||
|
||||
errors = [r for r in collector.records if r.levelno >= logging.ERROR]
|
||||
warnings = [r for r in collector.records if r.levelno == logging.WARNING]
|
||||
assert len(errors) == 1, "a callback-raised TimeoutError must still be logged once"
|
||||
assert errors[0].exc_info is not None, "the traceback must be preserved"
|
||||
assert warnings == [], "a callback-raised TimeoutError is not a worker-deadline timeout"
|
||||
assert worker._timeout_summary_task is None
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_stop_flushes_pending_timeout_summary(self):
|
||||
"""A burst still inside its summary window when the worker stops must still emit its one
|
||||
summary, instead of losing it when the event loop tears down.
|
||||
"""
|
||||
worker = LoggingWorker(
|
||||
timeout=0.05,
|
||||
max_queue_size=10,
|
||||
concurrency=10,
|
||||
timeout_summary_window=30.0,
|
||||
)
|
||||
|
||||
async def slow_callback() -> None:
|
||||
await asyncio.sleep(10)
|
||||
|
||||
logger = logging.getLogger("LiteLLM")
|
||||
collector = _RecordCollector()
|
||||
previous_level = logger.level
|
||||
logger.addHandler(collector)
|
||||
logger.setLevel(logging.DEBUG)
|
||||
try:
|
||||
worker.start()
|
||||
for _ in range(3):
|
||||
worker.enqueue(slow_callback())
|
||||
|
||||
await worker.flush()
|
||||
assert worker._timeout_summary_task is not None
|
||||
assert not worker._timeout_summary_task.done(), "precondition: the summary window has not elapsed"
|
||||
await worker.stop()
|
||||
finally:
|
||||
logger.removeHandler(collector)
|
||||
logger.setLevel(previous_level)
|
||||
|
||||
warnings = [r for r in collector.records if r.levelno == logging.WARNING]
|
||||
assert len(warnings) == 1, "stop must flush the pending summary exactly once"
|
||||
assert "3 callback(s) timed out" in warnings[0].getMessage()
|
||||
assert worker._timeout_summary_task is None
|
||||
|
||||
def test_loop_change_resets_timeout_summary_state(self):
|
||||
"""On an event-loop change the summary task is bound to the dead loop; it and the pending
|
||||
burst count must reset so timeouts on the new loop arm a fresh summary rather than a stale,
|
||||
stuck one that silently drops later summaries.
|
||||
"""
|
||||
worker = LoggingWorker(timeout=1.0, max_queue_size=10, timeout_summary_window=30.0)
|
||||
|
||||
async def timed_out_callback() -> None:
|
||||
await asyncio.sleep(10)
|
||||
|
||||
async def arm_on_first_loop() -> None:
|
||||
worker._ensure_queue()
|
||||
coro = timed_out_callback()
|
||||
worker._record_callback_timeout(coro)
|
||||
coro.close()
|
||||
assert worker._timeout_summary_task is not None
|
||||
assert worker._timeout_burst_count == 1
|
||||
|
||||
asyncio.run(arm_on_first_loop())
|
||||
assert worker._timeout_summary_task is not None
|
||||
|
||||
async def rebind_on_second_loop() -> None:
|
||||
worker._ensure_queue()
|
||||
assert worker._timeout_summary_task is None, "stale summary task must drop on loop change"
|
||||
assert worker._timeout_burst_count == 0, "stale burst count must reset on loop change"
|
||||
|
||||
asyncio.run(rebind_on_second_loop())
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue