dag
Background DAG builder for v9 task templates.
This package converts a v9 FlowSpec (parsed from a task template YAML)
into an executable background DAG.
Typical usage
import desert
import yaml
from bitfount.flows.dag import DAGRunContext, build_background_dag, run_dag
from bitfount.flows.dag.context import FileMetadataContext
from bitfount.flows.dag.lifecycle import LifecycleNotifier
from bitfount.flows.schema import FlowSpec
# 1. Deserialise the v9 YAML into a FlowSpec.
with open("ga_trial_bronze_v9.yaml") as f:
raw = yaml.safe_load(f)
flow_spec = desert.schema(FlowSpec).load(raw)
# 2. Parse + validate → BackgroundDAG.
dag = build_background_dag(flow_spec)
# 3. Build the run context.
run_ctx = DAGRunContext(
context_providers={
"file_metadata": FileMetadataContext(cache=cache, task_hash=task_hash)
},
runtime_params={"cache": cache, "task_hash": task_hash, "datasource": datasource},
lifecycle_notifier=LifecycleNotifier(cache=cache, run_id=run_id),
)
# 4. Execute (fire-and-forget in a background thread / asyncio task).
await run_dag(dag, run_ctx)
Public API
build_background_dag(flow_spec)
Check the spec's declared version against this SDK's ceiling, then parse +
validate a FlowSpec → BackgroundDAG. Defined in builder.py.
execute_dag(dag, run_ctx)
Run all steps inside an already-active Prefect flow.
Defined in executor.py.
run_dag(dag, run_ctx)
Run all steps, creating a Prefect flow context if needed.
Defined in executor.py.
build_flow(dag, run_ctx)
Wrap execution in a named Prefect @flow for deployment scenarios.
Defined in executor.py.
stamp_step_task_hashes(dag, datasource_hash)
Assign each step its own cache partition key — a Merkle hash over its
semantic config and its upstream steps' hashes — plus
parent_task_hashes, the partition of each upstream step keyed by the
consuming step's input parameter name. Called by the run-context builders
in setup.py; defined in hashing.py.
Module
Submodules
- bitfount.flows.dag.builder - Builder: combines parsing and validation into a single callable.
- bitfount.flows.dag.child_auth - How a background DAG child process obtains its credentials and configuration.
- bitfount.flows.dag.child_datasource - Building the one datasource a background DAG child process needs.
- bitfount.flows.dag.context - Context providers for ambient DAG inputs.
- bitfount.flows.dag.coverage - Sampling how much of a datasource is indexed, for the published pointer.
- bitfount.flows.dag.dag_process - Run a background v9 DAG in a dedicated child process.
- bitfount.flows.dag.deployment - The served Prefect deployment that runs a background v9 DAG.
- bitfount.flows.dag.executor - Executor: runs a
RunnableDAGinside a Prefect flow. - bitfount.flows.dag.hashing - Per-step
task_hashderivation — a Merkle hash over a step's config chain. - bitfount.flows.dag.indexing - When a reader of the file inventory has to wait for one to exist.
- bitfount.flows.dag.interactive_handler - Interactive-DAG JOB_REQUEST handler.
- bitfount.flows.dag.lifecycle - Lifecycle signalling for DAG execution.
- bitfount.flows.dag.models - Data models for the background DAG.
- bitfount.flows.dag.parser - Parser: converts a v9
FlowSpecinto aBackgroundDAG. - bitfount.flows.dag.periodic_reruns - Rerunning healthy background lineages on a schedule.
- bitfount.flows.dag.progress - How much work a background run has banked, and when to stop replaying it.
- bitfount.flows.dag.recovery - Restarting background DAG runs whose process is gone.
- bitfount.flows.dag.scope_steps - Identifying a DAG's scope steps.
- bitfount.flows.dag.scoping - Running a DAG's scope step and rescoping the run to what it selects.
- bitfount.flows.dag.serving - Serving the background DAG deployment from inside the pod process.
- bitfount.flows.dag.setup - Shared assembly of a background-DAG
DAGRunContext. - bitfount.flows.dag.validator - Validator for a parsed
BackgroundDAG. - bitfount.flows.dag.version_gate - Version gate: refuse a v9 flow spec this SDK is too old to run.
- bitfount.flows.dag.waves - Planning a waved background run: which files go through the DAG, and when.
Functions
build_background_dag
def build_background_dag(flow_spec: FlowSpec) ‑> BackgroundDAG:Parse and validate a v9 FlowSpec into an executable BackgroundDAG.
This is the canonical entry point for building a background DAG. It
combines parse_background_dag and validate_background_dag into
a single call so callers do not need to orchestrate the two steps
themselves.
Arguments
flow_spec: A deserialisedFlowSpec(frombitfount.flows.schema). Must be a v9 flow — i.e. it must carryfederation,roles, andworkerfields.
Returns
A BackgroundDAG ready to be passed to execute_dag.
Raises
FlowSpecVersionError: If flow_spec declares a task YAML version newer than this SDK knows, or one that cannot be parsed.LegacyTaskYAMLError: If flow_spec is not a v9 flow.StepNotFoundError: If any step's(task, version)is absent from the step registry.ReferenceOrderError: If anyFromRefinput references a step declared later in the DAG (forward reference).ValueError: If a step config cannot be validated against its Pydantic class.
build_flow
def build_flow( dag: RunnableDAG, run_ctx: DAGRunContext,) ‑> collections.abc.Callable[[], collections.abc.Coroutine[typing.Any, typing.Any, dict[str, typing.Any]]]:Return a named Prefect @flow callable for dag.
Useful for deployment scenarios where a first-class flow object is
needed. For most runtime cases prefer run_dag.
Note that the flow's name is dag.name, which is shared by every project
and datasource using the same template — it does not identify a run. Flow
runs that need to be findable again are tagged at call time by
_run_dag_flow_and_flush; see bitfount.runtimes.dedup.background_dag_tags.
Arguments
dag: A validatedRunnableDAG(BackgroundDAGorInteractiveDAG).run_ctx: Runtime context passed through toexecute_dag.
Returns
An async Prefect @flow callable with no required arguments.
build_interactive_dag
def build_interactive_dag( flow_spec: FlowSpec,) ‑> InteractiveDAG:Parse and validate a v9 FlowSpec into an executable InteractiveDAG.
Canonical entry point for building the interactive (post-processing) DAG.
Combines parse_interactive_dag and validate_interactive_dag.
Arguments
flow_spec: A deserialised v9FlowSpec(must carryfederation,roles, andworker).
Returns
An InteractiveDAG ready to be executed (e.g. via
build_flow, which runs any DAG-shaped object).
Raises
FlowSpecVersionError: If flow_spec declares a task YAML version newer than this SDK knows, or one that cannot be parsed.LegacyTaskYAMLError: If flow_spec is not a v9 flow.StepNotFoundError: If any step's(task, version)is absent from the step registry.ReferenceOrderError: If anyFromRefreferences a later interactive step (forward reference).BackgroundStepNotFoundError: If anyBackgroundRefreferences a step not declared in the background phase.ValueError: If a step config cannot be validated against its Pydantic class.
execute_dag
async def execute_dag(dag: RunnableDAG, run_ctx: DAGRunContext) ‑> dict[str, typing.Any]:Run dag inside an already-active Prefect flow context.
Raises RuntimeError if no flow context is active — use
run_dag when you cannot guarantee one.
Arguments
dag: A validatedRunnableDAG(BackgroundDAGorInteractiveDAG).run_ctx: Runtime context: providers, runtime params, reporter.
Returns Mapping of step name → step result for every executed step.
Raises
RuntimeError: If called outside a Prefect flow context.Exception: Re-raises any step exception after callingrun_ctx.lifecycle_notifier.on_failure.
run_dag
async def run_dag(dag: RunnableDAG, run_ctx: DAGRunContext) ‑> dict[str, typing.Any]:Run dag, creating a Prefect flow context if one is not already active.
Safe to call from anywhere — no existing flow context required.
Arguments
dag: A validatedRunnableDAG(BackgroundDAGorInteractiveDAG).run_ctx: Runtime context: providers, runtime params, reporter.
Returns Mapping of step name → step result for every executed step.
stamp_step_task_hashes
def stamp_step_task_hashes( dag: BackgroundDAG | InteractiveDAG, datasource_hash: str, resource_identities: Mapping[str, str] | None = None, datastructure_identity: str | None = None, datasource_type_identity: str | None = None,) ‑> dict[str, str]:Compute and assign DAGStep.task_hash and .parent_task_hashes for dag.
Mutates the steps in place so each carries its own partition key: the executor
injects it as the step's task_hash kwarg, and BackgroundResultsContext
reads a background step's hash straight off the DAGStep it already holds —
which is what keeps the two phases in agreement with no extra plumbing.
Each step also gets parent_task_hashes (see _resolve_parents), the same
resolution the hashed payload used rather than a second one.
An interactive DAG's background_steps are hashed first and in the same pass,
so a BackgroundRef's parent hash is the value the background phase computed.
That holds only while both phases pass the same resource_identities and
datastructure_identity; both derive them from the same pod config and
FlowSpec, so they do.
Arguments
dag: ABackgroundDAGorInteractiveDAG. AnInteractiveDAG'sbackground_stepsare hashed and stamped too.resource_identities: Identity string per external resource kind, for steps declaringtask_hash_resources.datasource_hash: The datasource leaf hash —generate_prefect_task_hash(pod_name, datasource_name).datastructure_identity: Canonical identity of the task'sdata_structure, for steps declaring adatastructureparameter.datasource_type_identity: Identity of the datasource's type, for steps whose partition is keyed per file.
Returns Mapping of step name → assigned hash (interactive steps win on a name clash with a background step), for logging and assertions.
Classes
BackgroundDAG
class BackgroundDAG(**data: Any):A fully parsed and validated background DAG ready for execution.
Attributes
name: Human-readable name of the flow (fromFlowSpec.name).federation_strategy: The federation strategy declared in the YAML (e.g.worker_only).steps: Ordered list of steps. Execution order is declaration order; parallel steps are submitted as futures and flushed on demand.extra: Any additional top-level metadata from theFlowSpecthat callers may find useful (e.g. for logging or telemetry).
Create a new model by parsing and validating input data from keyword arguments.
Raises [ValidationError][pydantic_core.ValidationError] if the input data cannot be validated to form a valid model.
self is explicitly positional-only to allow self as a field name.
Variables
- static
extra : dict[str, typing.Any]
- static
federation_strategy : str
- static
model_config
- static
name : str
- static
steps : list[DAGStep]
BackgroundRef
class BackgroundRef(**data: Any):A cross-phase reference from an interactive step to a background step.
Background steps run in a separate execution; their outputs are read back
from the cache at interactive-DAG runtime via BackgroundResultsContext.
Example YAML (inside an interactive step): ga_inference.cache
Create a new model by parsing and validating input data from keyword arguments.
Raises [ValidationError][pydantic_core.ValidationError] if the input data cannot be
validated to form a valid model.
self is explicitly positional-only to allow self as a field name.
Variables
- static
model_config
- static
output_field : str- Attribute name on the background step's result object (e.g.cache).
- static
step : str- Name of the background step whose result is referenced.
BackgroundStepNotFoundError
class BackgroundStepNotFoundError(*args, **kwargs):Raised when an interactive BackgroundRef targets an unknown step.
Interactive steps may reference background step outputs (e.g.
ga_inference.cache). The referenced step must be declared in the
flow's worker.background phase; otherwise its results can never be
resolved from the cache.
ContextRef
class ContextRef(**data: Any):A reference to a value provided by an ambient ContextProvider.
Example YAML: $file_metadata.cache
Create a new model by parsing and validating input data from keyword arguments.
Raises [ValidationError][pydantic_core.ValidationError] if the input data cannot be
validated to form a valid model.
self is explicitly positional-only to allow self as a field name.
Variables
- static
field : str- Field name to resolve from the provider.
- static
model_config
- static
provider : str- Name of the context provider (e.g.file_metadata).
DAGRunContext
class DAGRunContext(**data: Any):Runtime context for a single background DAG execution.
Bundles everything the executor needs beyond the DAG itself so functions
take (dag, run_ctx) instead of four separate parameters.
Attributes
context_providers: Mapping of provider name →ContextProviderused to resolveContextRefinputs (e.g.file_metadata).runtime_params: Runtime objects injected into every step's kwargs as a base layer (e.g.datasource,cache,task_hash). YAML-resolved refs override any overlapping keys.lifecycle_notifier: Signals DAG lifecycle transitions (accept / success / failure) to whichever destinations are wired (cache Run row, initiator mailbox). Defaults to a no-op notifier.background_results: Provider used to resolveBackgroundRefinputs (cross-phase references from interactive steps to background step outputs).Nonefor background DAGs, which have no such refs.dataset_name: Name of the datasource this run is working through, for telemetry. A first-class field rather than aruntime_paramsentry, for the same reasonwave_planis one:runtime_paramsis copied into every step's kwargs, and this is not a step input.wave_plan: The waves this run should carry through the DAG, and the planner that banks them —Nonewhen the run is not waved, which is every interactive run, every run on a pod with waving disabled, and every run whose DAG has no file-scoped step or whose scan sweep has not finished. A first-class field rather than aruntime_paramsentry, becauseruntime_paramsis copied into every step's kwargs and this is not a step input.empty_selection: Set when the run has no files to carry, in which case the executor closes the run as complete without running any step. Running them instead is unsafe:model_inferencereads the datasource rather than the selection, and with no override walks every file on disk.awaiting_scan_sweep: Set when the run landed before the first scan sweep settled, in which case the executor closes the run as complete without running any step or banking the inventory cursor. The cold-start wave continuation runs the sweep once it can.
Create a new model by parsing and validating input data from keyword arguments.
Raises [ValidationError][pydantic_core.ValidationError] if the input data cannot be validated to form a valid model.
self is explicitly positional-only to allow self as a field name.
Variables
- static
awaiting_scan_sweep : bool
- static
background_results : BackgroundResultsContext | None
- static
context_providers : dict[str, ContextProvider]
- static
dataset_name : str | None
- static
empty_selection : EmptySelection | None
- static
lifecycle_notifier : LifecycleNotifier
- static
model_config
- static
runtime_params : dict[str, typing.Any]
- static
wave_plan : WavePlan | None
DAGStep
class DAGStep(**data: Any):A single step in the background DAG — one node to execute.
Attributes
name: Unique step name within the DAG (e.g.fovea_inference).task: Registry key for the step (e.g.model_inference).version: Step version to look up in the registry.config: Instantiated Pydantic config for this step, orNoneif the step has no config class or no config was provided.inputs: Mapping of parameter name → resolvedInputRef.parallel: IfTruethe step is submitted as a Prefect future and its result is not awaited until a downstream step needs it.save_to_cache: Declarative hint from the YAML. Parsed and stored but not currently acted upon — steps handle their own caching.task_hash: This step's cache partition key — a Merkle hash over its own semantic config and its upstream steps' hashes, assigned byhashing.stamp_step_task_hasheswhen the run context is assembled.Noneuntil then (parsing and validation do not need it, and neither does a caller that only inspects the DAG); the executor then falls back to the run's datasource-levelruntime_params["task_hash"].parent_task_hashes: Thetask_hashof each upstream step this step declares an input from, keyed by this step's input parameter name — its routing table for reading an upstream partition through a direct store call, which aCacheAccessorcannot express. Assigned alongsidetask_hash;Noneuntil then. Keyed by parameter rather than upstream step name because the parameter is part of this step's own contract, while the upstream step's name is a label the flow author may change. Not part of the hashed payload (seehashing), so its contents never affect any partition key.
Create a new model by parsing and validating input data from keyword arguments.
Raises [ValidationError][pydantic_core.ValidationError] if the input data cannot be validated to form a valid model.
self is explicitly positional-only to allow self as a field name.
Variables
- static
config : pydantic.main.BaseModel | None
- static
inputs : dict[str, FromRef | ContextRef | BackgroundRef]
- static
model_config
- static
name : str
- static
parallel : bool
- static
parent_task_hashes : dict[str, str] | None
- static
save_to_cache : bool
- static
task : str
- static
task_hash : str | None
- static
version : int
DuplicateStepNameError
class DuplicateStepNameError(*args, **kwargs):Raised when two or more steps share the same name.
Node names are used as keys in execution result mappings and for
FromRef(step=...) lookups. Duplicate names are therefore ambiguous
and must be rejected at validation time.
FlowSpecVersionError
class FlowSpecVersionError(*args, **kwargs):Raised when a flow spec's declared version cannot be run by this SDK.
Covers three cases: the version is newer than this SDK's ceiling, the version cannot be read as a semantic version, or this SDK has no readable known versions to check against.
FromRef
class FromRef(**data: Any):A reference to a field on the result of a prior step.
Example YAML: ga_inference.cache
Create a new model by parsing and validating input data from keyword arguments.
Raises [ValidationError][pydantic_core.ValidationError] if the input data cannot be
validated to form a valid model.
self is explicitly positional-only to allow self as a field name.
Variables
- static
model_config
- static
output_field : str- Attribute name on the step's result object.
- static
step : str- Name of the step whose result is referenced.
InteractiveDAG
class InteractiveDAG(**data: Any):A fully parsed and validated interactive DAG ready for execution.
The interactive DAG is the post-processing phase of a v9 flow (GA
calculation, criteria matching, report generation). It runs as a separate
execution from the background DAG and may reference background step outputs
via BackgroundRef inputs, which are resolved from the cache at runtime.
Attributes
name: Human-readable name of the flow (fromFlowSpec.name).federation_strategy: The federation strategy declared in the YAML (e.g.worker_only).steps: Ordered list of interactive steps. Execution order is declaration order; parallel steps are submitted as futures and flushed on demand.background_steps: The parsed backgroundDAGSteplist (carrying each step's instantiated config). Used to build theBackgroundResultsContextand to validateBackgroundReftargets; not executed by the interactive run.extra: Any additional top-level metadata from theFlowSpec.
Create a new model by parsing and validating input data from keyword arguments.
Raises [ValidationError][pydantic_core.ValidationError] if the input data cannot be validated to form a valid model.
self is explicitly positional-only to allow self as a field name.
Variables
- static
background_steps : list[DAGStep]
- static
extra : dict[str, typing.Any]
- static
federation_strategy : str
- static
model_config
- static
name : str
- static
steps : list[DAGStep]
background_step_names : frozenset[str]- Names of the background steps interactive steps may reference.
LegacyTaskYAMLError
class LegacyTaskYAMLError(*args, **kwargs):Raised when a non-v9 FlowSpec is passed to the background DAG parser.
v9 flows are identified by the presence of federation, roles,
and worker fields. Older task YAML formats (v8 and below) use a
flat ModellerConfig structure and are not supported by this parser.
ReferenceOrderError
class ReferenceOrderError(*args, **kwargs):Raised when a FromRef references a step declared later in the DAG.
The background DAG is executed in declaration order. A step cannot consume the output of a step that has not yet run.
StepNotFoundError
class StepNotFoundError(*args, **kwargs):Raised when a DAG step references a task that is not in the registry.
Includes the task name and version in the message so the user can identify which YAML entry needs correcting.