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

1# pyright: reportMissingTypeStubs=false # apscheduler ships no type information 

2 

3import asyncio 

4from collections.abc import Collection 

5from typing import Final, Protocol 

6 

7from apscheduler.executors.asyncio import AsyncIOExecutor 

8 

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) 

14 

15 

16class StoppableScheduler(Protocol): 

17 """The slice of ``AsyncIOScheduler`` shutdown uses, which ships no type information""" 

18 

19 @property 

20 def running(self) -> bool: ... 20 ↛ exitline 20 didn't return from function 'running' because

21 

22 def pause(self) -> None: ... 22 ↛ exitline 22 didn't return from function 'pause' because

23 

24 def shutdown(self, wait: bool = ...) -> None: ... 24 ↛ exitline 24 didn't return from function 'shutdown' because

25 

26 

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""" 

29 

30 _pending_futures: Collection["asyncio.Future[object]"] 

31 

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()) 

35 

36 

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() 

41 

42 

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. 

53 

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 )