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
« 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"""
5from datetime import timedelta
6from typing import Annotated
7from uuid import UUID
9import sqlalchemy as sa
10from docket import CurrentDocket, Depends, Docket, Logged, Perpetual
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
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()
33 if not flow_run:
34 return
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 )
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
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 )
74 result = await session.execute(query)
75 runs = result.scalars().all()
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 )