Skip to main content

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:

  1. Its process is provably gone — CRASHED, or RUNNING with a stale heartbeat (is_flow_run_lost). Heartbeat staleness is what survives an outage that took the Prefect server with it.
  2. 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.
  3. 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_lineage drops the older crashes discovery returned, and _is_superseded drops a crash whose replacement has since failed or completed — invisible to a CRASHED/RUNNING query, and the case that otherwise replays one crash on every tick forever.
  4. Indexing is not in flight for it — a file_metadata re-index moves the input under the DAG's feet, so recovery waits for the next tick.
  5. The lineage has banked progress recently — should_recover bounds 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 open PrefectClient.
  • 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 open PrefectClient.
  • 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 to is_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 open PrefectClient.
  • 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 open PrefectClient.
  • served_project_ids: The projects whose deployments to keep.

Returns The names of the deployments that were deleted.

recover_lost_runs​

async def recover_lost_runs(    *,    client: PrefectClient,    cache: CacheProtocol,    in_flight_task_hashes: Container[str],    replay: ReplayFn,    progress_for: ProgressFn,    attempt_limit: int,    silence_window: timedelta | None = None,) ‑> list[LostRun]:

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 open PrefectClient.
  • 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 to is_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 is attempt + 1.
  • trigger: Why the dead run was started, read off its trigger: tag. None for 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 its rerun: tag. None for a run created before the tag existed, leaving was_rerun to fall back to the trigger.

Variables​

  • static attempt : int
  • static datasource_name : str
  • static is_rerun : bool | None
  • static project_id : str
  • static task_hash : str
  • 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 keeps trigger: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.