dedup
Prefect-native liveness and dedup for background work.
Two kinds of run rely on this module, and they differ in how they are addressable.
Deployment runs (file_metadata, scan_metadata) are triggered as
fire-and-forget Prefect deployment runs, tagged at submission with
task_hash:<task_hash> (see task_hash_tag) and
run_type:<file_metadata|scan_metadata> (see run_type_tag). Before submitting
a new run, callers query Prefect for a live flow-run carrying the same task_hash
tag on the same deployment and skip submission if one is found.
Background DAG runs are deployment runs too, but their deployment is per
project and is deleted once that project is no longer served — so a run outlives
the deployment_id it was created under. They are addressed purely by their
lineage tags (see background_dag_tags), which are written at creation and
outlive everything; the flow name is no use as a key either, being shared by
every project and datasource using the same template.
This makes Prefect's own flow-run state the source of truth, replacing the
Bitfount runs table's stale-running gate. But state alone is never enough:
the zombie-reaper automation
(bitfount.runtimes.prefect_bootstrap.ensure_zombie_automation) cannot observe a
crash that took the Prefect server down with it, and its pending detection window
does not survive a server restart, so a machine-wide outage leaves a dead run
sitting at RUNNING forever. Both liveness queries are therefore state
plus heartbeat recency over RUNNING runs — see is_flow_run_lost,
is_flow_run_active and is_background_flow_run_active. There is exactly one
definition of "dead" in the codebase, sized from
prefect_bootstrap.liveness_silence_window(), so poller, reaper and dedup cannot
drift into disagreeing about it.
Judging liveness is all this module does. Acting on a dead run — ending it in
Prefect, failing its bookkeeping row, replaying it — is
bitfount.runtimes.recovery for the metadata runtimes and
bitfount.flows.dag.recovery for background DAG runs.
Tagging
Tags are the only identity a flow run carries, so every run this system creates
gets run_type: and, where one exists, task_hash:. There are two ways they
arrive, and the rule is:
Tag at creation if you control creation; the in-flow self-tag
(ensure_flow_run_tags) exists only because Prefect's RunDeployment automation
action cannot.
Creation-time tags are what let the liveness queries above see a run before its
body starts — a run tagged only from inside itself is invisible to dedup for as
long as it sits SCHEDULED/PENDING, which under a concurrency limit can be a
while. Prefect unions deployment-level tags with per-run tags server-side, so
the two routes compose without duplicates.
Module
Functions
attempt_from_tags
def attempt_from_tags(tags: list[str]) ‑> int | None:Return the attempt number carried in tags, or None if absent.
Each recovery attempt is a new flow run carrying its own tags rather than a mutation of the previous one's — Prefect exposes no tag-mutation action — so reading the lineage back means parsing it off the run.
Returns None for an untagged or malformed run rather than raising: a run
whose lineage cannot be read is one this code did not create, and the caller
treats that as "not part of a lineage" rather than an error.
attempt_tag
def attempt_tag(attempt: int) ‑> str:Return the flow-run tag encoding this run's position in its lineage.
background_dag_tags
def background_dag_tags( task_hash: str, project_id: str, datasource_name: str, *, trigger: BackgroundTrigger, attempt: int = 1, is_rerun: bool | None = None,) ‑> list[str]:Return the full tag set identifying one background DAG flow run.
The trigger: tag is deliberately not part of any dedup filter.
is_background_flow_run_active matches on task_hash, project_id and
run_type only, so a nightly rerun and a dataset-project link still dedup
against each other — which is what stops the two double-running one
lineage. The tag exists to answer "why did this run start?" in the Prefect
UI, where the answer was previously unrecoverable.
Arguments
task_hash: The(pod, datasource)task hash.project_id: Project the datasource is linked to.datasource_name: Datasource the run is processing.trigger: Why this run was started.attempt: Position in the recovery lineage; 1 for the run a dataset-project link or a rerun started.is_rerun: Whether the run carries rerun semantics;Nonereads it off trigger. Recorded separately from the trigger because the two come apart: a recovery replacement for a dead rerun is taggedtrigger:recoveryand is still a rerun. Defaulting to the trigger rather than toFalseis what stops a caller that says nothing from writing a tag that contradicts the trigger it did pass.
Returns The tags to attach to the Prefect flow run.
cancel_unconfigured_metadata_runs
async def cancel_unconfigured_metadata_runs( client: PrefectClient, run_type: str, pod_name: str, configured_datasources: Collection[str],) ‑> dict[str, int]:Cancel pod_name's runs of run_type for datasources it no longer has.
Prefect is the only place a run of this kind is visible for its whole life.
A run is tagged and given its parameters at submission, so this sees it
while it is still SCHEDULED or PENDING — queued behind a concurrency
limit, having written no bookkeeping row and touched no file yet — as well
as once it is walking. Asking the pod's own runs table instead would
catch only the latter and let a queued run start on a share the config has
already dropped.
Scoped to one pod by the pod_name parameter every metadata run carries;
the task_hash tag is a digest and cannot be attributed to a pod. A run
whose parameters do not name both a pod and a datasource is left alone:
unattributable is not the same as unconfigured, and cancelling on a guess
would stop another pod's work.
Arguments
client: An openPrefectClient.run_type: The run type to scope to, e.g.FILE_METADATA_RUN_TYPE.pod_name: The pod whose runs to consider.configured_datasources: Every datasource currently in that pod's config. A run for anything else is cancelled, so this must be the configured set and not the servable one — a datasource whose hub registration failed has not been removed.
Returns How many runs were cancelled, keyed by datasource name. Empty when there was nothing to cancel.
datasource_tag
def datasource_tag(datasource_name: str) ‑> str:Return the flow-run tag encoding datasource_name.
ensure_flow_run_tags
async def ensure_flow_run_tags(*tags: str) ‑> None:Add tags to the flow run this coroutine is executing inside.
The in-flow half of the tagging rule in this module's docstring: a run whose
creator could not tag it labels itself once its body starts. The only
creator in that position is Prefect's RunDeployment automation action, which
has no tags field at all — so the scan_metadata runs spawned by the
scan-chain automation arrive with an empty tag set and no way for the
automation to fix it.
Best-effort by construction. Tagging is an observability affordance; a run that indexed a datasource successfully but failed to label itself has done the job, and raising here would turn a cosmetic failure into a lost run.
The client is opened here, inside the same try, rather than taken as an
argument — because a caller opening one for us puts the riskiest part of the
operation outside the guarantee above. A flow that did that had this tagging
call at the top of its body, so an unreachable Prefect server failed the run
before it did any work, three times over (retries=3), and the pod's
datasources got no schemas at all. Nothing about labelling a run should be
able to do that.
Arguments
*tags: Tags to ensure are present. Already-present tags are a no-op.
ensure_flow_run_tags_sync
def ensure_flow_run_tags_sync(*tags: str) ‑> None:Synchronous ensure_flow_run_tags, for use from a sync flow body.
The metadata runtimes' flows are sync (def file_metadata_runtime), so they
cannot await. Prefect's own sync client is used rather than driving the async
one, so this never touches the ambient event loop the flow is running on.
Arguments
*tags: Tags to ensure are present. Already-present tags are a no-op.
is_background_flow_run_active
async def is_background_flow_run_active( client: PrefectClient, task_hash: str, project_id: str, silence_window: timedelta | None = None, exclude_flow_run_id: uuid.UUID | None = None,) ‑> bool:Return True if a background DAG run for this lineage is already live.
The tags-only counterpart to is_flow_run_active. Background DAG runs are
served from a per-project deployment, so they do have a deployment_id
now — but filtering on it would scope the answer to one project's
deployment, and the lineage is what must not double-run. The tags stay the
whole key.
Counts submitted-but-not-started runs, per
BACKGROUND_SUBMITTED_STATE_TYPES. A run sits SCHEDULED until the Runner
polls for it, and longer when queued behind the project's concurrency limit;
ignoring those would have the wave-continuation and periodic-rerun pollers
submit a fresh duplicate every tick for as long as the first stayed queued.
Stale-aware for RUNNING runs, and that is load-bearing. A bare state
check would treat a permanently-RUNNING zombie as live and reject every
subsequent dataset-project link as a duplicate, wedging that
(task_hash, project_id) indefinitely. Cache-level dedup avoided this with
a staleness threshold; a Prefect-level guard has to earn it back explicitly.
No such check is possible for a pre-RUNNING run — it has no heartbeat — so
an orphaned one wedges the lineage until the start-up sweep in
bitfount.runtimes.recovery cancels it.
Arguments
client: An openPrefectClient.task_hash: The(pod, datasource)task hash to check for.project_id: Project the run belongs to. Two projects on one datasource share a task_hash and must not dedup against each other.silence_window: Passed through tois_flow_run_lost.exclude_flow_run_id: A run to leave out of the answer. Required when the caller is one of the runs being counted: a flow asking this from inside its own body carries the same lineage tags and a heartbeat seconds old, so without the exclusion it would find itself, call itself a duplicate, and no background run could ever execute.
Returns Whether a live background DAG run for this pair was found.
is_flow_run_active
async def is_flow_run_active( client: PrefectClient, deployment_id: uuid.UUID, task_hash: str, silence_window: timedelta | None = None, dead_before: datetime | None = None,) ‑> bool:Return True if a genuinely live flow-run for task_hash already exists.
"Live" covers a RUNNING run with a recent heartbeat, a run stuck
submitting for less than a couple of minutes, or a run genuinely queued
and waiting its turn — a real queue can take as long as it needs, so it's
never judged by age.
This answers the broad question of "would submitting again duplicate
work", so a queued run still counts as live here. A transient-failure
follow-up that is not yet due does not (is_pending_transient_retry). A
narrower check exists elsewhere for deciding whether something is
actively executing right now, which matters when deciding whether it's
safe to clean up a seemingly-dead run.
Arguments
client: An openPrefectClient.deployment_id: The id of the deployment to scope the check to.task_hash: The(pod, datasource)task hash to check for.silence_window: Passed through tois_flow_run_lost.dead_before: This process's own start time, used to treat a run as dead sooner than the usual age cutoff would allow — needed so a check right after our own restart doesn't wait it out.
Returns Whether a live flow-run for this task_hash was found.
is_flow_run_actively_running
async def is_flow_run_actively_running( client: PrefectClient, deployment_id: uuid.UUID, task_hash: str, silence_window: timedelta | None = None, dead_before: datetime | None = None,) ‑> bool:Return True only if something is genuinely executing right now.
Unlike the broader liveness check, a run that's merely queued and waiting for its turn does not count here. That distinction matters when deciding whether it's safe to clean up a run that looks dead: a queue entry isn't proof anything is actually working, and if a dead run is hogging the slot a queued replacement needs, cleaning up the dead one is exactly what lets the replacement start.
Arguments
client: An openPrefectClient.deployment_id: The id of the deployment to scope the check to.task_hash: The(pod, datasource)task hash to check for.silence_window: Passed through tois_flow_run_lost.dead_before: Passed through tois_flow_run_lost.
Returns Whether a genuinely running flow-run for this task_hash was found.
is_flow_run_lost
async def is_flow_run_lost( client: PrefectClient, flow_run: FlowRun, silence_window: timedelta | None = None, dead_before: datetime | None = None,) ‑> bool:Return True when flow_run's process is provably gone.
Two ways a lost process presents, and both count:
- Prefect already says
CRASHED— the reaper saw the silence and acted. This is what happens when the application stayed up and only the run died. - Prefect says
RUNNINGbut the last heartbeat is older than silence_window. This is what happens after a reboot or power cut, where the Prefect server died in the same instant as the run, nothing was alive to notice the silence, and the reaper's pending window was swept on restart. Left to Prefect alone such a run staysRUNNINGforever. A pre-RUNNINGrun is reported not lost, and callers must not ask about one: there is no heartbeat to judge and its age says nothing useful. Callers queryRUNNING/CRASHEDonly, so such a run never reaches here —is_background_flow_run_activewidens its own query toSCHEDULEDand therefore short-circuits those candidates rather than asking here.
A run with no heartbeat events at all falls back to its start_time age.
Because the heartbeat cadence is pinned and the first beat is emitted
immediately, a genuinely live run with no events is necessarily seconds old,
so the same window is safe. The fallback also covers a run whose heartbeats
aged out of retention — a start_time that old is unambiguously stale.
This says nothing about ownership: a suspended laptop's own healthy run also looks lost, because its heartbeat thread suspended with it. Callers must additionally exclude runs they are themselves executing
Arguments
client: An openPrefectClient.flow_run: The run to judge.silence_window: How long silence may last before the run is presumed dead. Defaults to the same window the reaper automation uses, so the two cannot disagree about what dead means.dead_before: An instant before which any last signal proves the run is gone, regardless of the window. Raises the cutoff tomax(window cutoff, dead_before);Noneleaves the window alone.
For a caller that owns every run it is asking about, its own process start time is exactly such an instant: a run whose last heartbeat predates this process cannot be running in this process, so it is dead however recent that heartbeat looks. That is what lets a start-up sweep recover a run orphaned seconds ago, which the window alone cannot — a run killed 5s before a restart has a heartbeat well inside 90s and is indistinguishable from a healthy one by age.
It is not a substitute for the window and must not be used as one. The window is what identifies a dead run this process never owned — a sibling pod's, say — and it is also the only thing standing between a start-up sweep and a run this boot has just created: such a run's heartbeats postdate the process, so the raised cutoff spares it while a blanket "any RUNNING at boot is dead" rule would not.
Returns Whether this run's process should be presumed gone.
is_indexing_in_flight
async def is_indexing_in_flight(client: PrefectClient, prefect_task_hash: str) ‑> bool:Return whether file_metadata indexing is running for prefect_task_hash.
Anything that reads the file inventory has to ask this first, because a live
indexer is building that inventory: a reader that starts anyway sees a
partial one, banks whatever it found and finishes COMPLETED, which is
under-coverage wearing the costume of success. Two callers, for the same
reason: bitfount.flows.dag.recovery defers replaying a DAG (the next tick
finds the run again), and scan_metadata_runtime skips its own body (the
indexer's Completed event chains a fresh scan). Uses the same
deployment-scoped check file_metadata_refresh uses to avoid double-indexing.
A missing deployment answers False: nothing can be indexing if nothing is
registered to index, and the pod's own indexing path would have registered it.
Arguments
client: An openPrefectClient.prefect_task_hash: The datasource-level task hash to check.
Returns Whether indexing is in flight for this datasource.
is_pending_transient_retry
def is_pending_transient_retry(run: FlowRun, now: datetime) ‑> bool:Whether run is a transient-failure follow-up that is not yet due.
Such a run is a deferred retry, not queued work: counting it as live would stop every other submitter for hours, so dedup ignores it. A newer run that does the same work leaves it to fire later and find nothing to do.
Arguments
run: A flow run.now: The current time.
is_submitting_run_stale
def is_submitting_run_stale( flow_run: FlowRun, dead_before: datetime | None = None,) ‑> bool:Return True when a Submitting run has been stuck too long to be alive.
Submitting has no heartbeat to check, but it doesn't need one: a run only sits there for a split second, so age alone is enough to tell a genuine stall from a live one.
Arguments
flow_run: The run to judge. Caller must confirm it's actuallySubmittingfirst — every other state counts as alive here.dead_before: This process's own start time. Lets us treat a run as dead immediately after our own restart, rather than waiting out the full age cutoff — useful since this check normally only gets one chance, right when a process starts up.
Returns Whether this run is presumed dead.
last_heartbeat_at
async def last_heartbeat_at( client: PrefectClient, flow_run_id: uuid.UUID,) ‑> datetime | None:Return when flow_run_id last emitted a heartbeat, or None if never.
Prefect 3.x records no heartbeat timestamp on the flow-run row —
_emit_flow_run_heartbeat only publishes an event — so flow_run.updated
tracks the last state transition and is hours stale on a perfectly healthy
long run. The event stream is the only real per-run liveness signal, and this
query is backed by a purpose-built (event, resource_id, occurred) index.
None has two causes and the caller must not conflate them with "alive": the
run died before its first heartbeat, or its heartbeats aged out of the events
retention window (7 days by default). is_flow_run_lost handles both by
falling back to the run's own start time.
Arguments
client: An openPrefectClient.flow_run_id: The flow run to look up.
Returns
The occurred timestamp of the most recent heartbeat, or None.
non_rerun_background_run_active
async def non_rerun_background_run_active( client: PrefectClient, exclude_flow_run_id: uuid.UUID | None = None,) ‑> bool:Return whether a background run that is not a rerun is live anywhere.
Pod-wide rather than lineage-scoped, which is the opposite of
is_background_flow_run_active and deliberate. That one answers "is this
lineage already running", to stop a duplicate. This one answers "is
someone waiting on the machine", so a discretionary refresh can stand aside
at a wave boundary.
Nothing makes a link wait for a rerun — the duplicate guard is scoped to one lineage and a link takes no concurrency slot — so the two run at once and compete for the GPU, the memory and the model. A refresh's work is already banked wave by wave, so ending early costs nothing and hands the machine to whoever is trying to surface a newly linked dataset.
Asked over Prefect rather than from the pod's own state because the
background DAG runs in the subprocess its deployment was launched in
(bitfount.flows.dag.deployment): a closure over the pod's state could
neither be sent across that boundary nor see anything useful once there.
The tags it filters on are already written by background_dag_tags.
Stale-aware, like its sibling: a permanently-RUNNING zombie would
otherwise make every refresh on the pod stand aside for ever.
Arguments
client: An openPrefectClient.exclude_flow_run_id: The asking run, left out of the answer. A refresh continuation is scheduled as a rerun and so is not matched here anyway, but a caller that is itself a link must not find itself.
Returns Whether non-rerun background work is live on this machine.
pod_tag
def pod_tag(pod_name: str) ‑> str:Return the flow-run tag encoding pod_name.
project_id_tag
def project_id_tag(project_id: str) ‑> str:Return the flow-run tag encoding project_id.
project_name_tag
def project_name_tag(project_name: str) ‑> str:Return the flow-run tag carrying a project's readable name.
rerun_from_tags
def rerun_from_tags(tags: Sequence[str]) ‑> Optional[bool]:Return whether tags mark a run as carrying rerun semantics.
The read counterpart to rerun_tag, and the only thing that carries
rerun-ness across the process that scheduled it. trigger: cannot answer
this on its own: a recovery replacement for a dead rerun is tagged
trigger:recovery so the unproductive-attempt bound still applies to it,
and is a rerun all the same. Reading the trigger alone therefore demoted
that replacement's own replacement to an ordinary recovery — holding no
rerun slot, and writing no stamp.
None means the tag is absent, which is a run created before it existed;
the caller falls back to the trigger. An unrecognised value reads the same
way, for the reason trigger_from_tags gives.
Arguments
tags: The flow run's tags.
Returns
Whether the run carries rerun semantics, or None when the tag cannot
be read.
rerun_tag
def rerun_tag(is_rerun: bool) ‑> str:Return the rerun: tag recording whether a run has rerun semantics.
run_type_tag
def run_type_tag(run_type: str) ‑> str:Return the flow-run tag encoding run_type (e.g. "file_metadata").
submitted_background_runs
async def submitted_background_runs( client: PrefectClient, task_hashes: Collection[str],) ‑> list[FlowRun]:Return the background runs of task_hashes a restarted pod adopts.
Those in BACKGROUND_ADOPTABLE_STATE_TYPES: submitted and not finished,
including a run the runner is moving through PENDING.
What a pod asks for at start-up to find the runs it submitted before it restarted — the ones its runner is about to pick up again, and which it must therefore claim before recovery reads them as unclaimed work.
Scoped by task hash rather than by a pod tag because (pod, datasource) is
what a task hash is: another pod sharing this Prefect server hashes its own
datasources differently.
Paged, because a SCHEDULED backlog alone can exceed the 200 rows
read_flow_runs returns per request, and a run past the first page is a
lineage left unclaimed.
Arguments
client: An openPrefectClient.task_hashes: The(pod, datasource)hashes the caller owns.
Returns The matching runs, which may be empty.
tag_value
def tag_value(tags: Sequence[str], prefix: str) ‑> str | None:Return the value carried by the prefix tag in tags, or None.
Tags are the only lineage a background flow run carries — a flow name is
dag.name, shared across templates — so reading a run's identity back means
parsing them off it. None for a run that has no such tag, which the caller
should read as "this run is not one of ours" rather than as an error.
Arguments
tags: The flow run's tags.prefix: The tag prefix to look for, e.g.TASK_HASH_TAG_PREFIX.
Returns
The remainder of the first matching tag, or None.
task_hash_tag
def task_hash_tag(task_hash: str) ‑> str:Return the flow-run tag encoding task_hash.
top_level_runs_only
def top_level_runs_only() ‑> prefect.client.schemas.filters.FlowRunFilterParentTaskRunId:Return a filter excluding subflow runs, such as a background run's waves.
A subflow run is part of the run that owns it, never a run of its own, so a lineage query must not find one as a duplicate or recover one as lost.
transient_retry_from_tags
def transient_retry_from_tags(tags: Sequence[str]) ‑> int:Return the transient-retry attempt carried in tags, 0 when absent.
A malformed tag reads as 0: the chain then restarts rather than raising.
transient_retry_tag
def transient_retry_tag(attempt: int) ‑> str:Return the tag marking a run as follow-up attempt after a transient failure.
trigger_from_tags
def trigger_from_tags( tags: Sequence[str],) ‑> BackgroundTrigger | None:Return the trigger carried in tags, or None if absent or unknown.
The read counterpart to trigger_tag. None covers two cases a caller
must treat the same way — a run created before the tag existed, and a value
this SDK version does not know because a newer pod wrote it — because
guessing at either would decide a lineage's fate on a string it cannot
interpret.
Arguments
tags: The flow run's tags.
Returns
The trigger, or None when it cannot be read.
trigger_tag
def trigger_tag(trigger: BackgroundTrigger) ‑> str:Return the trigger: tag naming why a run was started.