Coverage for /usr/local/lib/python3.12/site-packages/prefect/server/services/pause_expirations.py: 65%

29 statements  

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

1""" 

2The FailExpiredPauses service. Responsible for putting Paused flow runs in a Failed state if they are not resumed on time. 

3""" 

4 

5from datetime import timedelta 

6from typing import Annotated 

7from uuid import UUID 

8 

9import sqlalchemy as sa 

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

11 

12import prefect.server.models as models 

13from prefect.server.database import PrefectDBInterface, provide_database_interface 

14from prefect.server.schemas import states 

15from prefect.server.services.perpetual_services import perpetual_service 

16from prefect.settings.context import get_current_settings 

17from prefect.types._datetime import now 

18 

19 

20async def fail_expired_pause( 

21 flow_run_id: Annotated[UUID, Logged], 

22 pause_timeout: Annotated[str, Logged], 

23 *, 

24 db: PrefectDBInterface = Depends(provide_database_interface), 

25) -> None: 

26 """Mark a single expired paused flow run as failed (docket task).""" 

27 async with db.session_context(begin_transaction=True) as session: 

28 result = await session.execute( 

29 sa.select(db.FlowRun).where(db.FlowRun.id == flow_run_id) 

30 ) 

31 flow_run = result.scalar_one_or_none() 

32 

33 if not flow_run: 

34 return 

35 

36 if ( 

37 flow_run.state is not None 

38 and flow_run.state.state_details.pause_timeout is not None 

39 and flow_run.state.state_details.pause_timeout < now("UTC") 

40 ): 

41 await models.flow_runs.set_flow_run_state( 

42 session=session, 

43 flow_run_id=flow_run.id, 

44 state=states.Failed(message="The flow was paused and never resumed."), 

45 force=True, 

46 ) 

47 

48 

49@perpetual_service( 

50 enabled_getter=lambda: ( 

51 get_current_settings().server.services.pause_expirations.enabled 

52 ), 

53) 

54async def monitor_expired_pauses( 

55 docket: Docket = CurrentDocket(), 

56 db: PrefectDBInterface = Depends(provide_database_interface), 

57 perpetual: Perpetual = Perpetual( 

58 automatic=True, 

59 every=timedelta( 

60 seconds=get_current_settings().server.services.pause_expirations.loop_seconds 

61 ), 

62 ), 

63) -> None: 

64 """Monitor for expired paused flow runs and schedule failure tasks.""" 

65 batch_size = 200 

66 

67 async with db.session_context() as session: 

68 query = ( 

69 sa.select(db.FlowRun) 

70 .where(db.FlowRun.state_type == states.StateType.PAUSED) 

71 .limit(batch_size) 

72 ) 

73 

74 result = await session.execute(query) 

75 runs = result.scalars().all() 

76 

77 for run in runs: 77 ↛ 78line 77 didn't jump to line 78 because the loop on line 77 never started

78 if ( 78 ↛ 77line 78 didn't jump to line 77 because the condition on line 78 was always true

79 run.state is not None 

80 and run.state.state_details.pause_timeout is not None 

81 and run.state.state_details.pause_timeout < now("UTC") 

82 ): 

83 await docket.add(fail_expired_pause)( 

84 run.id, str(run.state.state_details.pause_timeout) 

85 )