context
Context providers for ambient DAG inputs.
A ContextProvider is an object that can resolve a named field into a
concrete Python value at DAG execution time. They are the runtime
counterpart to ContextRef in models.py.
The only provider shipped today is FileMetadataContext, which backs
the $file_metadata reference used in v9 task templates:
inputs: filenames: $file_metadata.cache
Any step that references $file_metadata.cache will receive a
list[str] of the datasource's indexed file paths.
Extending context providers
Implement the ContextProvider protocol and register an instance in
DAGRunContext.context_providers passed to execute_dag. The
dict key must match the provider name used in the YAML (without the $).
Classes
BackgroundResultsContext
class BackgroundResultsContext( cache: CacheProtocol, task_hash: str, background_steps: list[DAGStep],):Resolves cross-phase references to background step outputs from cache.
Interactive steps reference background step results with step.field
inputs (e.g. ga_inference.cache or ehr_query.cache). Those background
steps ran in a separate execution; their results live in a cache table
(e.g. model_inferences, ehr_data), not in the interactive DAG's
in-memory results mapping.
Only cache-backed steps (those declaring an orm_model) can be referenced
cross-phase; referencing any other step raises KeyError.
Resolved results are memoised so repeated references to the same background
step (e.g. ga_inference.cache from both ga_calculation and
pdf_report) reuse a single accessor.
Each step is read from the partition its own DAGStep.task_hash names (the
per-step Merkle hash, see flows/dag/hashing.py); the task_hash passed here
is the datasource-level fallback for steps that were never stamped.
Arguments
cache: An open cache instance satisfyingCacheProtocol.task_hash: The task hash identifying the datasource partition.background_steps: The parsed backgroundDAGSteplist (with instantiated configs) from the sameFlowSpec.
Methods
resolve_step
def resolve_step(self, step_name: str) ‑> CacheBackedResult:Return a result-like object for background step_name.
The returned object exposes the same cache (a CacheAccessor)
attribute that the live background step produced, so the executor can
resolve <step_name>.cache against it unchanged.
The step's (table, filters) partition and result type are sourced from
the step registry (its declared orm_model and Result).
Arguments
step_name: Name of the background step being referenced.
Returns
The step's registered result type (e.g. ModelInferenceResult,
EHRQueryResult) whose cache accessor reads the background step's
rows from its cache table.
Raises
KeyError: If step_name is not a known background step, or the step is not cache-backed (declares noorm_model).BackgroundResultsNotFoundError: If the cache holds no rows for the referenced background step (background phase has not run) — unless that step declaresallows_empty_partition = True(bitfount.steps.registry.allows_empty_partition), in which case an empty partition is a legitimate outcome and a result is returned as normal.
BackgroundResultsNotFoundError
class BackgroundResultsNotFoundError(*args, **kwargs):Raised when an interactive step references a background step with no results.
Interactive DAGs run as a separate execution from the background DAG and
read background step outputs straight from the cache. If the cache holds no
rows for a referenced background step, the background phase has not run (or
not yet completed) for this task_hash — the interactive run cannot
proceed and fails fast with this error rather than producing empty reports.
ContextProvider
class ContextProvider(*args, **kwargs):Resolves a named field into a runtime value.
Implementations are registered by name in the context_providers
DAGRunContext.context_providers. The name matches the
provider key in the YAML (the part after $ and before .).
Ancestors
Methods
resolve
def resolve(self, field: str) ‑> Any:Return the value for field.
Arguments
field: The field name to resolve (the part after.in the YAML reference, e.g."cache"from$file_metadata.cache).
Returns The resolved value for the given field.
Raises
KeyError: If field is not supported by this provider.
FileMetadataContext
class FileMetadataContext( cache: CacheProtocol, task_hash: str, datasource: BaseSource | None = None,):Resolves $file_metadata.<field> references.
Currently supports a single field:
cache
Returns a sorted list[str] of the datasource's indexed file paths.
Results are loaded eagerly on first access and cached for the lifetime
of this object so repeated resolutions do not re-query the database.
records
Returns a CollectedFileMetadataQueryResult with the indexed metadata
rows for callers that need more than file paths.
covered
Returns a sorted list[str] of every file the sweep has carried
through the whole DAG — the waves already banked plus the one in
progress. A step publishing a figure about the cohort wants this;
a step deciding what to compute wants cache.
One instance serves every step in a run, so scoping it is how a file scope reaches all of them at once. There are two scoping layers, and they intersect:
restrict_to
The run-level layer. Set once by a scope step before waves are planned,
and in force for the whole run.
narrow_to
The per-wave layer, set and cleared by wave_scope around each wave.
Both scope cache and records together and neither invalidates the
memoised selection — the inventory pass and the task filters are paid once
per run, not once per wave. They are separate slots because wave_scope
clears its own with narrow_to(None), which would otherwise lift a run-level
restriction at the end of the first wave.
Selection is computed from the datasource — its walk root and accepted
extensions — not from the inventory row's task_hash column. file_metadata
is a pod-wide inventory keyed on file_path, so task_hash on a
row is provenance: the datasource that indexed it last. Scoping by it drops
a file shared between two datasources from the earlier one's listing, and a
file missing here is never evaluated at all. See
bitfount.runtimes.file_metadata.selection.
Arguments
cache: An open cache instance satisfyingCacheProtocol.task_hash: The datasource-level task hash (fallback scope).datasource: The datasource whose files are being listed.
Methods
mark_covered
def mark_covered(self, file_paths: Collection[str]) ‑> None:Record file_paths as carried through the whole DAG by an earlier wave.
Called once a wave has been banked in the ledger — and seeded, at the
start of a sweep, from what previous runs banked — so covered counts
only work that actually landed. A wave that failed and was abandoned
must not inflate the published coverage figure.
Arguments
file_paths: The files a completed wave carried.
narrow_to
def narrow_to(self, file_paths: Collection[str] | None) ‑> None:Scope cache and records to file_paths until narrowed again.
How a wave's file scope reaches every step that reads
$file_metadata.cache: one provider instance serves the whole run, so
narrowing it narrows every consumer at once. Paths outside the
datasource's own selection are ignored rather than added, so a scope
can only ever shrink what a step sees.
Does not touch covered, and does not invalidate the memoised
selection — re-narrowing is free.
Arguments
file_paths: The wave's files, orNoneto restore the full selection.
resolve
def resolve(self, field: str) ‑> Any:Return the value for field.
Lazily loads file metadata on first call and caches the result.
Arguments
field: The field to resolve."cache","records"and"covered"are supported.
Returns
For field="cache", a list[str] of indexed file paths — the
current wave's when this provider has been narrowed.
For field="records", a CollectedFileMetadataQueryResult over
the same files.
For field="covered", a list[str] of every file the sweep has
carried through the DAG — the waves already banked plus the one in
progress. A step publishing a figure about the cohort wants this
rather than cache, which is only the current wave.
Raises
KeyError: If field is not in_SUPPORTED_FIELDS.
restrict_to
def restrict_to(self, file_paths: Collection[str] | None) ‑> None:Scope cache and records to file_paths for the rest of the run.
The run-level layer, set once by a scope step before waves are planned,
and deliberately not the same slot as narrow_to. wave_scope
exits every wave with narrow_to(None): sharing a slot would mean the
first wave to finish silently restored the whole datasource, and
model_inference — which is scoped by file name, not by filenames —
would then run over every file for the rest of the sweep.
The two layers intersect, so a wave can only narrow further within the restriction, never widen past it. Paths outside the datasource's own selection are ignored rather than added.
Arguments
file_paths: The reduced selection, orNoneto lift the restriction — what astrategy: allreduction leaves in place.