recovery
Restarting background DAG runs whose process is gone.
The sibling of bitfount.runtimes.recovery, separate from it not because of how
a run executes — both are deployment runs — but because of what replaying one
means. A metadata run is resubmitted verbatim from its own parameters, once per
process start. A background DAG run is rebuilt by the pod from a stored flow
spec, with a fresh run id, token grant and attempt number, under a bound this
module needs because it polls. That is pod policy, and runtimes.recovery
holds no pod dependency. What the two share — the candidate query, the ordering
rule, the supersession verdict — lives in bitfount.runtimes.discovery.
A worker.background DAG runs on a customer's laptop. A Windows Update reboot,
a power cut or an OOM kill takes it down mid-run. Prefect's scheduler will not
re-submit it — a lost run is a run that already started — and the pod that would
have closed out its bookkeeping died with it.
This module is what notices, the next time a pod runs.
Why a poller and not an automation
Prefect's own zombie reaper is a proactive automation, and proactive trigger state
lives in automation_bucket, swept on server restart. The Prefect server here is
a child of the same application as the pod, so a machine-wide outage kills both
and the missed-heartbeat window is gone by the time anything is alive to act on
it — the crash this feature exists for is exactly the one an automation cannot
see. Detection therefore has to compare stored state against wall-clock time
after the fact, which is a poll wherever it is hosted; and it is hosted in the
pod because replaying a run needs the pod's cache, its datasources and its
credentials, which no automation action can reach.
The five guards
A lost run is replayed only when all of these hold, and each covers a case the others cannot:
- Its process is provably gone —
CRASHED, orRUNNINGwith a stale heartbeat (is_flow_run_lost). Heartbeat staleness is what survives an outage that took the Prefect server with it. - This pod is not executing it — a suspended laptop's own healthy run looks stale, because its heartbeat thread suspended too. Ownership is the only thing that distinguishes "asleep" from "dead", so a run in our in-flight set is never recovered no matter how silent it is.
- Nothing has replaced it yet — nothing here ends a dead run and Prefect
keeps terminal runs indefinitely, so every crash a lineage ever had stays
discoverable and "lost" forever. Two checks, because they see different
things:
_latest_per_lineagedrops the older crashes discovery returned, and_is_supersededdrops a crash whose replacement has since failed or completed — invisible to aCRASHED/RUNNINGquery, and the case that otherwise replays one crash on every tick forever. - Indexing is not in flight for it — a
file_metadatare-index moves the input under the DAG's feet, so recovery waits for the next tick. - The lineage has banked progress recently —
should_recoverbounds consecutive attempts that achieved nothing, so a poisoned spec cannot loop unattended forever.
Guards 1 and 2 are complementary in the same way heartbeats and ownership are throughout: staleness catches a dead process whoever owned it, including a sibling pod's; ownership catches our own live process being suspended.
Module
Functions
cancel_orphaned_scheduled_runs
async def cancel_orphaned_scheduled_runs( client: PrefectClient, *, served_project_ids: Container[str],) ‑> int:Cancel the scheduled background runs nothing will pick up.
Best-effort per run: one that cannot be cancelled is logged and the sweep moves on, because this is tidying that runs at start-up and must not be able to stop a pod coming up.
Arguments
client: An openPrefectClient.served_project_ids: The projects whose deployments this pod serves.
Returns How many runs were cancelled.
find_lost_background_runs
async def find_lost_background_runs( client: PrefectClient, *, in_flight_task_hashes: Container[str], silence_window: timedelta | None = None,) ‑> list[LostRun]:Return the background runs this pod should consider replaying.
Applies guards 1 to 3 (see the module docstring): the run's process must be
provably gone, it must not be one this pod is currently executing, and it must
be the run its lineage currently ends at — _latest_per_lineage for the runs
discovery can see, then _is_superseded for the ones it cannot.
Arguments
client: An openPrefectClient.in_flight_task_hashes: The task hashes this pod is running right now. A suspended laptop's own run looks exactly like a dead one, so this is the only thing that keeps it from being double-started.silence_window: Passed through tois_flow_run_lost.
Returns
At most one lost run per lineage: the RUNNING ones first, then the
CRASHED ones, each group newest first.
find_orphaned_scheduled_runs
async def find_orphaned_scheduled_runs( client: PrefectClient, *, served_project_ids: Container[str],) ‑> list[FlowRun]:Return scheduled background runs that nothing will ever pick up.
A run submitted against a project's deployment and not started before the
pod died sits SCHEDULED for ever unless something either serves that
deployment again or ends the run. A pod re-registers the deployments of
every project it still has work for, so what is left over is a run for a
project this pod no longer serves.
Arguments
client: An openPrefectClient.served_project_ids: The projects whose deployments this pod has just re-registered. Their scheduled runs are left alone: the runner picks them up within one poll interval.
Returns
The runs to cancel. Never includes one queued behind the per-project
concurrency limit, whose state type is also SCHEDULED.
prune_background_dag_deployments
async def prune_background_dag_deployments( client: PrefectClient, *, served_project_ids: Container[str],) ‑> list[str]:Delete background DAG deployments for projects this pod no longer serves.
A deployment outlives the pod that created it, so without this every project the pod has ever been linked to keeps a page in the Prefect UI and an entry the scheduler still considers. Deleting one is safe precisely because nothing is outstanding for it: the sweep above has already ended its scheduled runs.
Best-effort per deployment, for the reason the sweep is.
Arguments
client: An openPrefectClient.served_project_ids: The projects whose deployments to keep.
Returns The names of the deployments that were deleted.
recover_lost_runs
Find lost background runs and replay the ones that qualify.
One tick of the poller. Never raises for a single run's sake: a lineage that cannot be judged is skipped with its reason logged, so one bad spec cannot stop the others being recovered.
Arguments
client: An openPrefectClient.cache: The pod's background cache, holding both the specs and the rows progress is measured from.in_flight_task_hashes: Task hashes this pod is executing (guard 2).replay: Schedules a fresh attempt from a stored spec.progress_for: Measures what a lineage has banked so far.attempt_limit: How many consecutive unproductive attempts to tolerate.silence_window: Passed through tois_flow_run_lost.
Returns The runs a replacement attempt was started for.
Classes
LostRun
class LostRun( flow_run_id: uuid.UUID, task_hash: str, project_id: str, datasource_name: str, attempt: int, trigger: BackgroundTrigger | None = None, is_rerun: bool | None = None, expected_start_time: datetime | None = None,):A background flow run whose process is gone, identified by its tags.
Arguments
flow_run_id: The dead run, retained for logging and for excluding it from the liveness query when its replacement asks.task_hash: The(pod, datasource)task hash from the run's tags.project_id: Project from the run's tags.datasource_name: Datasource from the run's tags.attempt: The dead run's position in its lineage; its replacement isattempt + 1.trigger: Why the dead run was started, read off itstrigger:tag.Nonefor a run created before the tag existed, or one carrying a value this SDK version does not know.is_rerun: Whether the dead run carried rerun semantics, read off itsrerun:tag.Nonefor a run created before the tag existed, leavingwas_rerunto fall back to the trigger.
Variables
- static
attempt : int
- static
datasource_name : str
- static
expected_start_time : datetime.datetime | None
- static
flow_run_id : uuid.UUID
- static
is_rerun : bool | None
- static
project_id : str
- static
task_hash : str
- static
trigger : BackgroundTrigger | None
-
was_rerun : bool- Whether this run carried rerun semantics.Prefers the
rerun:tag, which is written explicitly, and falls back to the trigger for a run created before that tag existed. The fallback is not sufficient on its own: a replacement for a dead rerun keepstrigger:recovery, so reading the trigger alone silently demoted that replacement's own replacement to an ordinary recovery.
ReplayFn
class ReplayFn(*args, **kwargs):Rebuilds and schedules a fresh attempt from a stored spec.
Injected rather than imported: replaying needs the pod's datasources, Hub session and cache, none of which this module should know about — and keeping it a callback is what lets the pod route a recovered attempt through the same construction path a dataset-project link uses.