Coverage for /usr/local/lib/python3.12/site-packages/prefect/server/api/background_workers.py: 92%

34 statements  

« prev     ^ index     » next       coverage.py v7.15.2, created at 2026-10-07 02:04 +0000

1import asyncio 

2from contextlib import asynccontextmanager 

3from logging import Logger 

4from typing import Any, AsyncGenerator, Callable 

5 

6from docket import Docket, Worker 

7 

8from prefect.logging import get_logger 

9from prefect.server.api.flow_runs import delete_flow_run_logs 

10from prefect.server.api.task_runs import delete_task_run_logs 

11from prefect.server.events.services import triggers as _triggers_module # noqa: F401 

12from prefect.server.models.deployments import mark_deployments_ready 

13from prefect.server.models.work_queues import mark_work_queues_ready 

14from prefect.server.services.cancellation_cleanup import ( 

15 cancel_child_task_runs, 

16 cancel_subflow_run, 

17 handle_cancelling_timeout, 

18) 

19from prefect.server.services.db_vacuum import ( 

20 vacuum_events_with_retention_overrides, 

21 vacuum_old_events, 

22 vacuum_old_flow_runs, 

23 vacuum_orphaned_artifacts, 

24 vacuum_orphaned_logs, 

25 vacuum_stale_artifact_collections, 

26) 

27from prefect.server.services.late_runs import mark_flow_run_late 

28from prefect.server.services.pause_expirations import fail_expired_pause 

29from prefect.server.services.perpetual_services import ( 

30 register_and_schedule_perpetual_services, 

31) 

32from prefect.server.services.repossessor import revoke_expired_lease 

33 

34logger: Logger = get_logger(__name__) 

35 

36# Task functions to register with docket for background processing 

37task_functions: list[Callable[..., Any]] = [ 

38 # Simple background tasks (from Alex's PR #19377) 

39 mark_work_queues_ready, 

40 mark_deployments_ready, 

41 delete_task_run_logs, 

42 delete_flow_run_logs, 

43 # Find-and-flood pattern tasks used by perpetual services 

44 handle_cancelling_timeout, 

45 cancel_child_task_runs, 

46 cancel_subflow_run, 

47 fail_expired_pause, 

48 mark_flow_run_late, 

49 revoke_expired_lease, 

50 vacuum_orphaned_logs, 

51 vacuum_orphaned_artifacts, 

52 vacuum_stale_artifact_collections, 

53 vacuum_old_flow_runs, 

54 vacuum_events_with_retention_overrides, 

55 vacuum_old_events, 

56] 

57 

58 

59@asynccontextmanager 

60async def background_worker( 

61 docket: Docket, 

62 ephemeral: bool = False, 

63 webserver_only: bool = False, 

64) -> AsyncGenerator[None, None]: 

65 worker_task: asyncio.Task[None] | None = None 

66 async with Worker(docket) as worker: 

67 # Register background task functions 

68 docket.register_collection( 

69 "prefect.server.api.background_workers:task_functions" 

70 ) 

71 

72 # Register and schedule enabled perpetual services 

73 await register_and_schedule_perpetual_services( 

74 docket, ephemeral=ephemeral, webserver_only=webserver_only 

75 ) 

76 

77 try: 

78 worker_task = asyncio.create_task(worker.run_forever()) 

79 yield 

80 

81 finally: 

82 if worker_task: 82 ↛ exitline 82 didn't jump to the function exit

83 worker_task.cancel() 

84 try: 

85 await worker_task 

86 except asyncio.CancelledError: 

87 pass