How Flow Run Suspension Works in Prefect: Implementation and Enforcement
Flow run suspension in Prefect is a coordinated state transition that pauses execution by moving a run from Running to Suspended, enforced through server-side policy validation and client-side polling in the flow engine.
Prefect allows long-running workflows to pause mid-execution and wait for external input or manual approval. This flow run suspension mechanism involves a sophisticated handshake between the Prefect SDK, the orchestration server, and the local flow engine. Understanding this architecture is essential for building resilient data pipelines that can wait indefinitely without consuming worker resources.
Initiating Suspension from the SDK
A flow initiates suspension by calling the public helper suspend_flow_run() (or its async counterpart asuspend_flow_run()). These functions are defined in src/prefect/flow_runs.py at lines 599–616.
The helper generates a unique pause key to distinguish multiple suspension points within the same flow run. It then issues a request to the orchestration server requesting a transition to the Suspended state. If the server rejects the transition—such as when the flow run has already been resumed—the SDK raises a RuntimeError (handled in lines 574–585 of the same file).
from prefect import flow, suspend_flow_run
@flow
def approval_workflow():
# Business logic here...
suspend_flow_run(timeout=3600) # Pause for up to 1 hour
# Execution resumes here after manual intervention
Server-Side State Transition Validation
The orchestration server validates every suspension request through its core policy logic. The relevant implementation resides in src/prefect/server/orchestration/core_policy.py at lines 1361–1444, where the system processes the Suspend verb for the Suspended state target.
This policy layer ensures that:
- Only valid state transitions are permitted
- Race conditions between multiple suspension requests are handled
- The flow run meets preconditions for entering a paused state
Once the policy accepts the transition, the server persists the flow run with the Suspended state, making it visible to the UI and API.
Engine Enforcement Through State Polling
The flow engine enforces suspension by checking the flow run's state before scheduling each task. In src/prefect/flow_engine.py at lines 1066–1068, the engine evaluates:
if (state := self.flow_run.state) and is_suspended_flow_run_state(state):
# Pause execution until the flow is resumed
The helper function is_suspended_flow_run_state is implemented in src/prefect/_flow_run_suspension.py and simply validates that state.name == "Suspended".
When a suspended state is detected, the engine enters a wait loop, periodically re-querying the server until the state transitions back to Running or a terminal state. This prevents task execution while maintaining the process connection for eventual resumption.
Background Coordination and Observers
To maintain accurate state synchronization without blocking the main thread, Prefect uses an internal observer pattern. The FlowRunSuspensionObserver class in src/prefect/_internal/observers.py (lines 256–342) maintains a registry of suspended flow run IDs and triggers callbacks when server-side state changes occur.
A background task named _check_for_suspended_flow_runs (defined at lines 415–438 in the same file) periodically queries the Prefect API for flow runs that have entered the Suspended state. This polling mechanism updates the observer's bookkeeping, ensuring the engine's local view remains consistent with the server state even across process restarts or network interruptions.
Resuming a Suspended Flow Run
Resumption reverses the suspension process. Calling resume_flow_run() (or aresume_flow_run()) sends a request to the server to transition the state away from Suspended. The server validates this transition through the same core policy logic, then updates the flow run record.
On the next polling cycle, the flow engine detects the state change from Suspended to Running and continues normal task scheduling. The flow resumes execution at the point immediately following the suspend_flow_run() call.
Summary
- Flow run suspension requires explicit SDK calls (
suspend_flow_run()) that negotiate with the orchestration server. - State transitions are validated by core policy logic in
src/prefect/server/orchestration/core_policy.pybefore the server commits theSuspendedstate. - The flow engine enforces suspension by checking
is_suspended_flow_run_state()before every task scheduling decision insrc/prefect/flow_engine.py. - Background observers in
src/prefect/_internal/observers.pymaintain synchronization between the client and server through periodic polling. - Resumption requires a separate API call to transition the state back to
Running, at which point the engine continues execution.
Frequently Asked Questions
What happens if I try to suspend a flow run that is already suspended?
According to the source code in src/prefect/flow_runs.py, attempting to suspend an already suspended flow run will result in a RuntimeError. The SDK validates the server's response and raises this exception if the state transition is rejected, preventing duplicate pause keys and conflicting suspension states.
Does suspension consume worker resources while waiting?
No. When the flow engine detects a suspended state in src/prefect/flow_engine.py, it enters a passive wait loop rather than executing tasks. The process remains alive to maintain the flow run context, but it does not consume CPU cycles for task execution or hold database connections open for active work.
Can a suspended flow run resume automatically without manual intervention?
Yes. While the examples show manual resumption via resume_flow_run(), the suspension mechanism only checks for state transitions on the server. Any process with appropriate permissions—including automated services, webhooks, or scheduled jobs—can call the resume API endpoint to transition the state from Suspended to Running, triggering the engine to continue execution.
Have a question about this repo?
These articles cover the highlights, but your codebase questions are specific. Give your agent direct access to the source. Share this with your agent to get started:
curl -s "https://instagit.com/install.md" Maintain an open-source project? Get it listed too →