Skip to main content

discovery

Reading recoverable flow runs off a Prefect server.

Shared by the two recovery sweeps — bitfount.runtimes.recovery for the metadata runtimes and bitfount.flows.dag.recovery for the background DAG. Both are deployment runs, so what separates them is not how a run executes but what replaying one means: a metadata run is resubmitted verbatim from its own parameters; a background DAG run is rebuilt by the pod from a stored flow spec, under an attempt bound, with a fresh run id and token grant. Neither sweep can hold that policy for the other.

What they share is the query: the paging loop, the two constants, the pair of state queries, the rule about which of two runs is later, and the supersession verdict built on it.

Two things here are easy to get wrong in a way nothing notices:

Paging. read_flow_runs with no limit does not mean "all of them": the server applies PREFECT_SERVER_API_DEFAULT_LIMIT, 200 by default. Both sweeps query the entire crash history of something, and nothing prunes terminal runs, so past 200 rows they silently stopped seeing whole lineages — and the supersession checks then read a stale row as the latest and declined to replay the newer crash they had never seen.

Ordering. "Which of these two runs is later" has to be answered by expected_start_time and nothing else. A run's attempt: tag looks like an ordering and is not: it is a label, and a replay that derived its number from an older crash repeats one. Where two guards ranked the same runs by different keys they each deferred to the other's loser and dropped the lineage between them.

Module​

Functions​

discover_recovery_candidates​

async def discover_recovery_candidates(    client: PrefectClient,    *,    scope: FlowRunFilter,    description: str,    crashed_horizon: timedelta = datetime.timedelta(days=7),    include_submitting: bool = False,) ‑> list[FlowRun]:

Return every recovery candidate within scope, before judging any.

Separate queries rather than one any_=[CRASHED, RUNNING], because the halves want different bounds. CRASHED is an ever-growing archive that has to be bounded to stay queryable, and can be — a crash older than crashed_horizon has been superseded many times over. RUNNING must not be bounded: a run wedged there since long before the horizon is what both sweeps exist to find.

Submitting is asked for only when include_submitting says so. It is a PENDING run by state type, and the two sweeps read that state differently: a wedged metadata run holds a concurrency slot nothing else will release, while a pre-RUNNING background DAG run is the litter the Windows 409 retry leaves behind, which dedup deliberately ignores.

scope says which runs are the caller's — a deployment id for the metadata runtimes, the run_type: tag for the background DAG. Any state or expected_start_time it carries is replaced.

Arguments

  • client: An open PrefectClient.
  • scope: What identifies the caller's runs.
  • description: What is being read, for the truncation warning.
  • crashed_horizon: How far back a CRASHED run is worth looking at.
  • include_submitting: Whether runs wedged in Submitting are candidates too. Judged on age rather than heartbeat by the caller, since a run that never started has no heartbeat to read.

Returns The running runs, then the crashed ones, then any Submitting ones, each newest first.

is_later​

def is_later(candidate: datetime | None, than: datetime | None) ‑> bool:

Whether candidate is the later of two runs' expected_start_times.

Unknown either side answers False — "not later", so not a successor. That asymmetry is deliberate: a replay that was not needed is a duplicate the submission-time dedup absorbs, while a replay wrongly withheld is a lineage that stops for good.

Arguments

  • candidate: The expected_start_time being tested.
  • than: The expected_start_time to beat.

Returns Whether candidate is strictly later.

newer_run_for_lineage​

async def newer_run_for_lineage(    client: PrefectClient,    *,    lineage: FlowRunFilter,    flow_run_id: uuid.UUID,    expected_start_time: datetime | None,) ‑> FlowRun | None:

Return the run that has already taken over from this one, if any.

Prefect keeps terminal runs indefinitely and nothing prunes them, so a CRASHED run stays discoverable long after its replacement has come and gone and "is this run lost" is permanently true of it. Unchecked, the same crash is replayed every tick forever, each replay minting the same attempt number.

No matching runs is not superseded: for the metadata sweep that is a run predating deployment-level tagging, which is exactly the one to recover. One recovery later the replacement is tagged, so it settles.

The verdict compares timestamps, not identity. "The newest run is not me" is wrong whenever the query cannot see the run being judged — metadata discovery is by deployment_id, so it also returns runs with no task_hash: tag, and by identity any older tagged run makes those look replaced. By time, an older run cannot supersede a newer one whether or not it sees it.

expected_start_time and nothing else: every run has one from its first state transition, unlike start_time, which a run that never started has not got. Two guards ranking by different keys each defer to the other's loser and drop the lineage between them.

Arguments

  • client: An open PrefectClient.
  • lineage: What identifies the runs of this lineage.
  • flow_run_id: The run being judged, so it cannot supersede itself.
  • expected_start_time: That run's expected_start_time.

Returns The successor, or None when this run is still the lineage's latest.

read_all_pages​

async def read_all_pages(    client: PrefectClient, flow_run_filter: FlowRunFilter, description: str,) ‑> list[FlowRun]:

Return every run matching flow_run_filter, walking the pages.

Sorted newest-first, so a fixed offset means the same row between requests: the set being paged is dominated by terminal runs, which do not move.

Arguments

  • client: An open PrefectClient.
  • flow_run_filter: The filter to page through.
  • description: What is being read, for the truncation warning.

Returns The runs, newest first.