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
« 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
6from docket import Docket, Worker
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
34logger: Logger = get_logger(__name__)
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]
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 )
72 # Register and schedule enabled perpetual services
73 await register_and_schedule_perpetual_services(
74 docket, ephemeral=ephemeral, webserver_only=webserver_only
75 )
77 try:
78 worker_task = asyncio.create_task(worker.run_forever())
79 yield
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