models
Data models for the background DAG.
These are the core primitives that represent a parsed and validated v9 background DAG. They are pure data — no I/O, no registry access.
Input references
Two kinds of reference can appear in a step's inputs: block:
FromRef
Refers to a field on the result of a previously executed step.
Syntax in YAML: step_name.field e.g. ga_inference.cache
ContextRef
Refers to a value provided by a named ContextProvider (see
context.py). The $ prefix in YAML is stripped when parsing.
Syntax in YAML: $provider_name.field e.g. $file_metadata.cache
BackgroundRef
Used only in interactive DAGs. Refers to a field on the result of a
background step that ran in a separate execution and whose output lives
in the cache. Same YAML syntax as FromRef (step.field); the parser
emits a BackgroundRef instead when the referenced step is a background
step. Resolved at runtime via the BackgroundResultsContext provider.
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.
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
EmptySelection
class EmptySelection(**data: Any):Why a background run has no files to carry through its DAG.
Attributes
reason:"no_files_selected"when the run's selection was empty before any scope step ran,"scope_selected_none"when a scope step was given files and kept none of them.files_considered: How many files the selection held before the scope step;0for"no_files_selected".scope_step: The scope step that kept nothing, orNone.
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
files_considered : int
- static
model_config
- static
reason : Literal['no_files_selected', 'scope_selected_none']
- static
scope_step : str | None
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.