Coverage for /usr/local/lib/python3.12/site-packages/prefect/server/services/late_runs.py: 67%
35 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
1"""
2The MarkLateRuns service. Responsible for putting flow runs in a Late state if they are not started on time.
3The threshold for a late run can be configured by changing `PREFECT_API_SERVICES_LATE_RUNS_AFTER_SECONDS`.
4"""
6from __future__ import annotations
8from datetime import timedelta
9from typing import Annotated
10from uuid import UUID
12import sqlalchemy as sa
13from docket import CurrentDocket, Depends, Docket, Logged, Perpetual
15import prefect.server.models as models
16from prefect.server.database import PrefectDBInterface, provide_database_interface
17from prefect.server.exceptions import ObjectNotFoundError
18from prefect.server.orchestration.core_policy import MarkLateRunsPolicy
19from prefect.server.schemas import states
20from prefect.server.services.perpetual_services import perpetual_service
21from prefect.settings.context import get_current_settings
22from prefect.types._datetime import now
25async def mark_flow_run_late(
26 flow_run_id: Annotated[UUID, Logged],
27 *,
28 db: PrefectDBInterface = Depends(provide_database_interface),
29) -> None:
30 """Mark a single flow run as late (docket task)."""
31 async with db.session_context(begin_transaction=True) as session:
32 result = await session.execute(
33 sa.select(db.FlowRun.id, db.FlowRun.next_scheduled_start_time).where(
34 db.FlowRun.id == flow_run_id
35 )
36 )
37 flow_run = result.one_or_none()
39 if not flow_run:
40 return
42 try:
43 await models.flow_runs.set_flow_run_state(
44 session=session,
45 flow_run_id=flow_run.id,
46 state=states.Late(scheduled_time=flow_run.next_scheduled_start_time),
47 flow_policy=MarkLateRunsPolicy, # type: ignore
48 )
49 except ObjectNotFoundError:
50 return
53@perpetual_service(
54 enabled_getter=lambda: get_current_settings().server.services.late_runs.enabled,
55)
56async def monitor_late_runs(
57 docket: Docket = CurrentDocket(),
58 db: PrefectDBInterface = Depends(provide_database_interface),
59 perpetual: Perpetual = Perpetual(
60 automatic=True,
61 every=timedelta(
62 seconds=get_current_settings().server.services.late_runs.loop_seconds
63 ),
64 ),
65) -> None:
66 """Monitor for late flow runs and schedule marking tasks."""
67 settings = get_current_settings().server.services.late_runs
68 batch_size = 400
69 scheduled_to_start_before = now("UTC") - settings.after_seconds
71 async with db.session_context() as session:
72 query = (
73 sa.select(db.FlowRun.id, db.FlowRun.next_scheduled_start_time)
74 .where(
75 (db.FlowRun.next_scheduled_start_time <= scheduled_to_start_before),
76 db.FlowRun.state_type == states.StateType.SCHEDULED,
77 db.FlowRun.state_name == "Scheduled",
78 )
79 .limit(batch_size)
80 )
81 result = await session.execute(query)
82 runs = result.all()
84 for run in runs: 84 ↛ 85line 84 didn't jump to line 85 because the loop on line 84 never started
85 await docket.add(mark_flow_run_late, key=f"mark-flow-run-late:{run.id}")(
86 run.id
87 )