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

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

5 

6from __future__ import annotations 

7 

8from datetime import timedelta 

9from typing import Annotated 

10from uuid import UUID 

11 

12import sqlalchemy as sa 

13from docket import CurrentDocket, Depends, Docket, Logged, Perpetual 

14 

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 

23 

24 

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

38 

39 if not flow_run: 

40 return 

41 

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 

51 

52 

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 

70 

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

83 

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 )