deployments
Deployment identity and readiness for the background metadata runtimes.
Both metadata flows register under the same Prefect flow name (background)
with distinct deployment names, and are addressed by the
"<flow-name>/<deployment-name>" strings run_deployment accepts. Those
strings are load-bearing in four places — the flows' own serve() calls, the
orchestrator's triggers and chain automation, run_pod's triggers, and the
refresh flow — so they live here rather than being retyped at each site: a
mismatched literal does not fail loudly, it silently produces a deployment
nothing ever triggers.
wait_for_deployment lives here for the same reason. serve() registers a
deployment asynchronously from a daemon thread, so any code that triggers a run
may get there first; every serving process therefore needs the same
poll-until-registered probe.
Also home to the run_type values every flow run this system creates is tagged
with, and the map from those to the deployment each is submitted against. They
sit beside the handles because they are the same kind of thing — identity every
submitter, every dedup check and both recovery sweeps must agree on — and
because what has to hold of them is mutual: they must all be distinct, and in
particular none may collide with BACKGROUND_DAG_RUN_TYPE, since
flows.dag.recovery finds its work by matching that tag and would otherwise
start replaying metadata runs it has no spec for. That invariant is checkable at
a glance in one list and unverifiable spread across the runtime packages.
Background DAG runs are deployment runs too, but their handle is per project
rather than one constant, so they are resolved through
background_dag_deployment and deployment_for_run instead of being named in
DEPLOYMENT_BY_RUN_TYPE. Keyed on the project id and not its name, because a
rename must not orphan a live deployment.
Also home to ConfigError, the misconfiguration signal both flows raise (see
the class docstring for why the type matters).
Module
Functions
background_dag_deployment
def background_dag_deployment(project_id: str) ‑> str:Return the full "<flow-name>/<deployment-name>" handle for project_id.
This is the string run_deployment and read_deployment_by_name accept.
Arguments
project_id: The project the deployment serves.
Returns The fully-qualified deployment handle.
background_dag_deployment_name
def background_dag_deployment_name(project_id: str) ‑> str:Return the deployment name serving project_id's background DAG runs.
Keyed on the project id and never on the project name. A project renamed
in the Hub would otherwise orphan its live deployment and silently gain a
second one, leaving every outstanding run pointing at a name nothing serves.
The readable name rides along as the deployment's description and its
project_name: tag, so the UI still reads well without identity moving.
Arguments
project_id: The project the deployment serves.
Returns The deployment name, without the flow-name prefix.
deployment_for_run
def deployment_for_run(run_type: str, project_id: str | None) ‑> str | None:Return the deployment handle a run of run_type is submitted against.
DEPLOYMENT_BY_RUN_TYPE cannot answer this alone once background DAG runs
are deployment runs: their handle is derived per project rather than being
one constant, so the recovery sweep needs the project to resolve it. Kept
here beside the mapping it generalises, so a caller never has to know which
run types are per-project.
Arguments
run_type: Therun_type:tag value carried by the run.project_id: The run's project, where it has one.
Returns
The handle, or None when run_type has no deployment or a
background DAG run carries no project to resolve one from.
project_id_from_deployment_name
def project_id_from_deployment_name(deployment_name: str) ‑> str | None:Return the project a background DAG deployment serves.
The inverse of background_dag_deployment_name, used by the start-up
reconciliation to decide which registered deployments still have a linked
project behind them.
Arguments
deployment_name: A deployment name, with or without the flow-name prefix.
Returns
The project id, or None when deployment_name is not a background
DAG deployment.
wait_for_deployment
async def wait_for_deployment( deployment_name: str, timeout: float = 30.0, poll_interval: float = 0.2,) ‑> bool:Poll the Prefect API until deployment_name is registered.
Intended to be awaited by a process that has just started serve() in a
background thread, before it submits any run. serve() registers the
deployment asynchronously, so a triggerer that does not wait can call
run_deployment against a deployment that does not exist yet.
Coroutine rather than a blocking function so the caller owns the event loop:
the orchestrator drives it through its _run_async bridge (Flask-SocketIO
runs handlers on plain threads) and run_pod through asyncio.run. Probing
inside one loop also avoids standing up and tearing down an event loop on
every poll for the whole timeout.
Arguments
deployment_name: The full"<flow-name>/<deployment-name>"string, e.g.FILE_METADATA_DEPLOYMENT.timeout: Maximum seconds to wait before giving up.poll_interval: Seconds between consecutive API probes.
Returns
True when the deployment was found within timeout seconds, False
if the timeout was reached first. Never raises: a caller that cannot
confirm readiness should degrade (skip indexing) rather than abort
startup, so the failure is reported by return value.
Classes
BackgroundTrigger
class BackgroundTrigger(*args, **kwds):Why a background DAG run was started.
Lives here for the same reason the run_type values do: it is run identity
that submitters, the Prefect UI and support all have to agree on, and the
values must stay distinct.
is_rerun is the load-bearing distinction rather than the individual
values. A rerun re-executes a lineage that already ran to completion, so it
is discretionary work that must queue behind a resource limit; a link or a
recovery is work that has not happened yet and must not. A future manual or
data-arrival trigger becomes a rerun by being added to _RERUN_TRIGGERS,
with no other code change.
Attributes
LINK: TheDATASET_PROJECT_LINKEDmessage that first armed the lineage.RECOVERY: A replacement for a run whose process is provably gone.PERIODIC_RERUN: The cron-scheduled rerun of a still-valid lineage.WAVE_CONTINUATION: A waved sweep that could not start when its lineage was linked, picked up once it can. A link-triggered run is scheduled immediately and is not gated on either metadata runtime, so on a cold start it arrives before the scan sweep has settled and stands down without running a step. This trigger is what starts the sweep once that changes. Deliberately not arerun: it is the lineage's first pass over work that has not happened yet, so it must not queue behind the rerun limit.
Ancestors
Variables
- static
LINK
- static
PERIODIC_RERUN
- static
RECOVERY
- static
WAVE_CONTINUATION
is_rerun : bool- Whether this trigger re-runs a lineage that already ran.
-
refreshes_stale : bool- Whether this trigger re-carries work the ledger already holds.Separate from
is_rerunbecause the two answer different questions and will not always agree.is_rerundecides scheduling — whether the run queues behind the rerun limit and stamps the lineage. This decides planning: whetherbitfount.flows.dag.wavesis given a freshness horizon, and so whether a fully-ledgered patient is outstanding again.A link must never inherit refresh work. Its whole purpose is to make the newest patients recruitable in minutes, and giving it the datasource's stale backlog as well would bury that behind a re-sweep.
ConfigError
class ConfigError(*args, **kwargs):Raised when a metadata runtime's datasource configuration is invalid.
A dedicated type (rather than a bare ValueError) is what lets each flow's
except Exception handler call mark_run_failed, store the message, and
re-raise unconditionally — so Prefect records the run as failed rather
than successful. Operators watching Prefect flow state then get an accurate
signal without inspecting the cache DB or a callback payload.
Subclasses ValueError rather than BitfountError: the misconfiguration is
always a bad argument to the flow (an unsupported datasource type, a missing
path or connection string), and callers that already catch ValueError
around flow invocation should keep catching it.