Coverage for .venv/lib/python3.13/site-packages/litellm/proxy/shutdown/scheduled_jobs.py: 69%
32 statements
« prev ^ index » next coverage.py v7.15.2, created at 2026-10-10 12:01 +0000
« prev ^ index » next coverage.py v7.15.2, created at 2026-10-10 12:01 +0000
1# pyright: reportMissingTypeStubs=false # apscheduler ships no type information
3import asyncio
4from collections.abc import Collection
5from typing import Final, Protocol
7from apscheduler.executors.asyncio import AsyncIOExecutor
9from litellm._logging import verbose_proxy_logger
10from litellm.constants import (
11 SCHEDULED_JOB_SHUTDOWN_CANCEL_TIMEOUT_SECONDS,
12 SCHEDULED_JOB_SHUTDOWN_FINISH_TIMEOUT_SECONDS,
13)
16class StoppableScheduler(Protocol):
17 """The slice of ``AsyncIOScheduler`` shutdown uses, which ships no type information"""
19 @property
20 def running(self) -> bool: ... 20 ↛ exitline 20 didn't return from function 'running' because
22 def pause(self) -> None: ... 22 ↛ exitline 22 didn't return from function 'pause' because
24 def shutdown(self, wait: bool = ...) -> None: ... 24 ↛ exitline 24 didn't return from function 'shutdown' because
27class AwaitableAsyncIOExecutor(AsyncIOExecutor): # pyright: ignore[reportUntypedBaseClass] # apscheduler ships no type information and is absent from the type-check env
28 """``AsyncIOExecutor`` whose in-flight job tasks can be awaited after ``shutdown`` cancels them"""
30 _pending_futures: Collection["asyncio.Future[object]"]
32 def in_flight_jobs(self) -> tuple["asyncio.Future[object]", ...]:
33 """The job tasks that are running right now, as a snapshot"""
34 return tuple(future for future in self._pending_futures if not future.done())
37def pause_scheduled_jobs(scheduler: StoppableScheduler) -> None:
38 """Stop the scheduler from starting jobs that shutdown would only cancel; running jobs continue"""
39 if scheduler.running: 39 ↛ exitline 39 didn't return from function 'pause_scheduled_jobs' because the condition on line 39 was always true
40 scheduler.pause()
43async def stop_in_flight_scheduler_jobs(
44 scheduler: StoppableScheduler,
45 executor: AwaitableAsyncIOExecutor,
46 *,
47 finish_timeout_seconds: float = SCHEDULED_JOB_SHUTDOWN_FINISH_TIMEOUT_SECONDS,
48 cancel_timeout_seconds: float = SCHEDULED_JOB_SHUTDOWN_CANCEL_TIMEOUT_SECONDS,
49) -> None:
50 """
51 Let in-flight jobs finish for up to finish_timeout_seconds, then stop the scheduler and wait, bounded by
52 cancel_timeout_seconds, for the jobs it cancels.
54 Must run before the database is disconnected: a write job that finishes needs its connection,
55 and a job's cancellation handler is what records the run's outcome.
56 """
57 if not scheduler.running: 57 ↛ 58line 57 didn't jump to line 58 because the condition on line 57 was never true
58 return
59 in_flight: Final = executor.in_flight_jobs()
60 if in_flight: 60 ↛ 61line 60 didn't jump to line 61 because the condition on line 60 was never true
61 verbose_proxy_logger.info(
62 "Waiting up to %ss for %d in-flight scheduled job(s) to finish",
63 finish_timeout_seconds,
64 len(in_flight),
65 )
66 still_running: Final = (
67 (await asyncio.wait(in_flight, timeout=finish_timeout_seconds))[1] if in_flight else frozenset()
68 )
69 scheduler.shutdown(wait=False)
70 if not still_running: 70 ↛ 72line 70 didn't jump to line 72 because the condition on line 70 was always true
71 return
72 verbose_proxy_logger.info("Cancelling %d in-flight scheduled job(s) for shutdown", len(still_running))
73 _done, pending = await asyncio.wait(still_running, timeout=cancel_timeout_seconds)
74 if pending:
75 verbose_proxy_logger.warning(
76 "%d scheduled job(s) did not finish within %ss of cancellation; giving up on them",
77 len(pending),
78 cancel_timeout_seconds,
79 )