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 openPrefectClient.scope: What identifies the caller's runs.description: What is being read, for the truncation warning.crashed_horizon: How far back aCRASHEDrun is worth looking at.include_submitting: Whether runs wedged inSubmittingare 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: Theexpected_start_timebeing tested.than: Theexpected_start_timeto 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 openPrefectClient.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'sexpected_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 openPrefectClient.flow_run_filter: The filter to page through.description: What is being read, for the truncation warning.
Returns The runs, newest first.