Skip to main content

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; None reads it off trigger. Recorded separately from the trigger because the two come apart: a recovery replacement for a dead rerun is tagged trigger:recovery and is still a rerun. Defaulting to the trigger rather than to False is 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 open PrefectClient.
  • 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 open PrefectClient.
  • 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 to is_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 open PrefectClient.
  • 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 to is_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 open PrefectClient.
  • 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 to is_flow_run_lost.
  • dead_before: Passed through to is_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 RUNNING but 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 stays RUNNING forever. A pre-RUNNING run is reported not lost, and callers must not ask about one: there is no heartbeat to judge and its age says nothing useful. Callers query RUNNING/CRASHED only, so such a run never reaches here — is_background_flow_run_active widens its own query to SCHEDULED and 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 open PrefectClient.
  • 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 to max(window cutoff, dead_before); None leaves 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 open PrefectClient.
  • 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 actually Submitting first — 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 open PrefectClient.
  • 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 open PrefectClient.
  • 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 open PrefectClient.
  • 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.