recovery
Recovery for metadata-runtime flow runs whose process is gone.
The sibling of bitfount.flows.dag.recovery, separate from it not because of how
a run executes — both are deployment runs — but because of what replaying one
means. These runs are found by deployment_id, replayed verbatim from their own
parameters, and need their Prefect state and bookkeeping rows reconciled. A
background DAG run is found by its lineage tags — its deployment is per project
and is pruned when that project is no longer served — and is rebuilt by the pod
from a stored flow spec under an attempt bound. What the two share — the
candidate query, the ordering rule, the supersession verdict — lives in
bitfount.runtimes.discovery.
The background sweep does not cover these runs: it discovers work by the
run_type:background_dag tag and replays from a persisted flow spec, and these
runs need neither.
What they lacked was anyone to notice. The submitters are one-shot — the
orchestrator fires _trigger_file_metadata_runtime once per pod start, run_pod
once at startup — and each gates on is_flow_run_active. Until that gate became
heartbeat-aware it read the raw Prefect state, so a pod restarting inside the
zombie-reaper's ~90s window found its own orphaned run still sitting at RUNNING,
called it live, skipped, and never asked again. The datasource simply stopped
being indexed, silently, and scan_metadata never ran either because the
scan-chain automation fires on a Completed event that would now never arrive.
The gate fix alone would leave recovery dependent on a submitter running after the reaper had acted. This module removes that dependency: it runs once per process start, judges liveness itself, and resubmits.
Three things happen per lost run, and the order matters:
- Establish there is nothing live. Per
(task_hash, run_type), askis_flow_run_active. A live run means this pod is looking at its own healthy work and must not touch it. - Reconcile. End the dead run at
Crashedand fail itsrunningbookkeeping row. The Prefect transition is not needed for correctness — the gate already ignores the run — but a permanently-RUNNINGrow is read by people as "still working", and the reaper cannot be relied on to fix it: it cannot observe a crash that took the Prefect server down with it, and its pending window is swept fromautomation_bucketon server restart. The bookkeeping row matters more:_is_run_activeis the only thing that reaps stalerunningrows, it is called only from the refresh cron, and the refresh cron only visits datasources that already have rows — so a first run that crashed before banking anything leaves a row nothing will ever clean. - Resubmit. A
file_metadataresubmission normally supersedes ascan_metadataone, so at most one run pertask_hash: the scan-chain automation runs scan when file_metadata completes, and running scan concurrently with the indexer would parse a partial inventory, bank rows and finishCOMPLETED— under-coverage that looks like success and never trips any bound. That withholding is conditional on the chain existing, which this sweep ensures and then checks rather than assumes: with no chain, withholding is the drop, so both run types are replayed andscan_metadata_runtime's own deferral sequences them.
Recovery is not attempt-bounded, and that asymmetry with
bitfount.flows.dag.recovery is deliberate. That module polls every 60s and so
can burn a machine unattended; this one fires once per process start, and for any
datasource still in the pod config _trigger_file_metadata_runtime already
resubmits on every pod start regardless. A bound here would guard a loop that
the restart cadence already bounds. Resubmissions do carry an incrementing
attempt: tag, so a run that is genuinely looping is visible as such.
Module
Functions
collapse_scheduled_backlog
async def collapse_scheduled_backlog(client: PrefectClient) ‑> int:Cancel every superseded SCHEDULED run of the metadata deployments.
Runs once per process start, beside recover_lost_metadata_runs, which
discovers RUNNING and CRASHED only — so nothing has ever collected a
SCHEDULED run.
A backlog forms because a deployment's global_limit enqueues on collision:
when the slots are taken the server returns the run to
Scheduled(AwaitingConcurrencySlot), timed 30s out, and the runner logs
"Server returned a non-pending state 'SCHEDULED'" and leaves it. That much
is by design. What is not is that nothing drains it: ACTIVE_STATE_TYPES
counts SCHEDULED as live and is_flow_run_lost will not judge anything
pre-RUNNING, so a wedged run reads as live forever and permanently
suppresses resubmission for its task_hash — one more orphan per day under a
cron.
Keeps the newest of each group rather than ageing runs out. Every threshold
available here is wrong: a run waiting its turn is indistinguishable by age
from a wedged one, and liveness_silence_window() (90s) is orders of
magnitude shorter than the daily cron it would judge, so sizing from it would
cancel healthy queued work on every start.
Grouped per (deployment, task_hash, _work_key). Dropping any of the
three loses work: per deployment alone would cancel datasource B's queued
index because A had a newer one, and without the work key a scoped watcher
run would cancel a queued full walk.
Cancelling needs no bound, unlike the resubmission in
recover_lost_metadata_runs: it is terminal, so a collapsed group has one
run left and a second pass finds nothing. force is left off, matching
_end_run_as_crashed.
Arguments
client: An openPrefectClient.
Returns How many runs were cancelled.
end_stuck_refresh_runs
async def end_stuck_refresh_runs( client: PrefectClient, *, dead_before: datetime | None = None, own_pod_name: str | None = None,) ‑> int:End every stuck Submitting run of the daily refresh deployment.
This one's cron-scheduled and pod-wide rather than per-datasource, so it falls outside the usual recovery and slot-repair sweeps. Nothing needs resubmitting here, though — Prefect will create the next occurrence on its own schedule regardless, so all we need to do is end the stuck run and free up the slot it's holding.
Arguments
client: An openPrefectClient.dead_before: This process's own start time, so a stuck run left by our own recent restart is caught immediately. Only usable when own_pod_name isNone— see there for why.own_pod_name: Set by a caller that might be sharing this Prefect server with other pods. A refresh run has no per-pod identity to check it against (it refreshes every pod's stale datasources in one pass), so there's no way to tell "our" run from someone else's — the safe choice is to fall back to age alone rather than risk ending a run another pod's process still owns.
Returns How many runs were ended.
reconcile_metadata_runs_at_startup
async def reconcile_metadata_runs_at_startup( client: PrefectClient, *, dead_before: datetime | None = None, own_pod_name: str | None = None,) ‑> int:Run this module's whole start-of-process sweep, in order.
One entrypoint so that every process that recovers metadata runs recovers them the same way. The order is what makes each step useful:
collapse_scheduled_backlog— a queued run reads as live tois_flow_run_active, so a backlog suppresses the resubmission step 2 would otherwise make.recover_lost_metadata_runs— end and replay every per-datasource run whose process is gone. Ending a run is also what releases the concurrency slot it still held, so most of what step 4 exists for is gone by the time it runs.end_stuck_refresh_runs— the same idea, for the one deployment step 2 doesn't cover: the daily refresh run. Without this, a stuck refresh run would go unnoticed by every other step too.release_orphaned_deployment_slots— zero any slot counter that no surviving run can account for. Last, because a crash's runs sit atRUNNINGuntil step 2 ends them and aRUNNINGrun has to be read as holding its slot.
Putting step 4 last costs nothing even when a replacement run is submitted
while a counter is still stranded: such a run parks at
AwaitingConcurrencySlot with a 30-second scheduled retry rather than
failing, so it starts on its own once the counter is repaired.
Each step is independent of the others' success: one failing is logged and the rest still run.
Arguments
client: An openPrefectClient.dead_before: Passed torecover_lost_metadata_runsandend_stuck_refresh_runs.own_pod_name: Passed torecover_lost_metadata_runsandend_stuck_refresh_runs.
Returns How many replacement runs were submitted.
recover_lost_metadata_runs
async def recover_lost_metadata_runs( client: PrefectClient, *, dead_before: datetime | None = None, own_pod_name: str | None = None,) ‑> int:Reconcile and replay every metadata run whose process is gone.
Runs once per process start, from the same place the other server-side
guarantees are established (prefect_bootstrap's ensure_* calls). Callers
must treat a failure here as non-fatal — indexing recovery is worth
attempting, never worth blocking startup for.
The mutating half of this module carries a sleep hazard that is closed by
where it is called from rather than by a threshold, because no threshold
could close it: a laptop suspended for hours has its healthy run's heartbeat
thread suspended too, so that run looks exactly as dead as a crashed one.
Three things make it safe anyway. First, this runs only at process start, and
a suspended machine is not starting processes — reap and sleep are mutually
exclusive on one machine. Second, generate_prefect_task_hash is per-pod, so
this can only ever touch runs belonging to this pod's datasources; a
genuinely live run of the same task_hash would require a second process
serving the same pod name, which the orchestrator's single-pod lock prevents
and a container deployment has no desktop app beside it to create. Third,
dead_before only ever raises the cutoff, and a resumed laptop's run has
heartbeats from before this process started either way — so the sharper
signal cannot make the sleep case any worse than the window already does.
Arguments
client: An openPrefectClient.dead_before: This process's own start time. Without it the sweep cannot act on a fast restart: a run killed seconds before the restart has a heartbeat well inside the liveness window, so it reads as healthy, the sweep leaves it, the reaperCrashedes it moments later, and nothing re-indexes that datasource until the next start. A last heartbeat predating this process settles it — that run cannot be executing in a process that did not exist yet.own_pod_name: The pod this process serves, when it serves exactly one. dead_before is applied only to that pod's runs. Omitting it asserts this process owns every metadata run on this Prefect server — true of the orchestrator, whose Prefect server is its own child on a private port, and untrue of several pods sharing one sidecar server, where one pod starting would otherwise judge another's live runs dead.
Returns How many replacement runs were submitted.
release_orphaned_deployment_slots
async def release_orphaned_deployment_slots(client: PrefectClient) ‑> int:Zero each metadata deployment's slot counter when no run can hold one.
Runs last in the start-of-process sweep, and that is load-bearing. The
crash this repairs leaves its runs sitting at RUNNING in the database, and
a RUNNING run is indistinguishable here from one that is really executing —
so run before recover_lost_metadata_runs this would find a slot-holder for
every stranded slot and refuse to touch a single one. Ending those runs is
what recover_lost_metadata_runs does, and Prefect releases each slot as it
goes: ReleaseFlowConcurrencySlots fires on the way out of RUNNING, and
when the run's lease is gone it decrements the deployment's counter directly
instead. So by the time this runs, the ordinary case has already healed and
what is left is the residue.
A deployment concurrency slot is two things that can disagree: an
active_slots counter in the Prefect database, and a lease the running flow
renews. A counter can therefore outlive every run that could account for it —
a lease store lost with the server, a PENDING run that died before it ever
reached RUNNING and so never transitions again, rows aged out from under
the counter. The repossessor service revokes expired leases, so with no
lease left there is nothing to expire and nothing ever decrements, and the
deployment reads as permanently full: every run submitted after it parks in
Scheduled(AwaitingConcurrencySlot). Cancelling that queue does not help,
because a queued run holds no slot — it was refused one.
Safe because of what it checks, not because of any threshold: the counter is
only zeroed when Prefect reports no run in a slot-holding state at all, so
anything recover_lost_metadata_runs judged live and left alone protects its
own slot. A lease that outlives this and expires later decrements from zero,
which the server floors rather than taking negative.
Never raises: a failure here must not stop the sweep it belongs to.
Arguments
client: An openPrefectClient.
Returns How many deployments had their counter reset.