Skip to main content

registry

Task registry for the YAML DAG loader.

Maps string task names (as used in YAML task: fields) and task versions (as used in YAML version: fields) to the corresponding Prefect @task-decorated callables and their associated Config/Result classes.

Discovery is automatic: each step sub-package exposes a versioned sub-package (e.g. v1/) whose __init__.py exports Config, Result, and task_fn.

Module​

Functions​

allows_empty_partition​

def allows_empty_partition(name: str, version: int = 1) ‑> bool:

Whether step name at version treats an empty partition as legitimate.

False for the vast majority of steps: a zero-row cross-phase read of a background step's partition normally means the background phase has not run, and BackgroundResultsContext fails fast on that. A step opts out of that guard by declaring a module-level allows_empty_partition = True because its background run can legitimately persist nothing (e.g. an EHR query that degraded rather than aborted) — see ehr_query.v3.

Returns False for unknown steps, mirroring the leniency of the other optional-attribute accessors: the validators reject unregistered tasks before hashing, so this only ever means "keep the fail-fast guard".

declares_runtime_param​

def declares_runtime_param(name: str, version: int, param: str) ‑> bool:

Return whether step name vversion declares the parameter param.

Lets the hashing layer ask "does this step actually consume this runtime input?" without a per-step declaration to keep in sync — see StepEntry.runtime_params.

Arguments

  • name: Registered step name.
  • version: Step version.
  • param: The parameter name to look for.

Returns True when the step declares param (or accepts **kwargs). False for an unregistered step, matching the other optional-attribute accessors: the validators reject unknown tasks before hashing.

emits_batch_progress​

def emits_batch_progress(name: str, version: int = 1) ‑> bool:

Whether step name at version emits its own NUMBER_OF_BATCHES signal.

Returns False for unknown steps (the executor's fallback then applies), mirroring the leniency of the other optional-attribute accessors.

get_cache_accessor_fn​

def get_cache_accessor_fn(    name: str, version: int = 1,) ‑> collections.abc.Callable[..., typing.Any] | None:

Get the cache-accessor builder for name at version, or None.

Returns the step's (cache, task_hash, config) -> CacheAccessor function, or None if the step is not cache-backed (i.e. declares no record binding).

get_cache_binding​

def get_cache_binding(name: str, version: int = 1) ‑> tuple[str | None, int]:

Return the (cache_table, cache_table_version) step name binds to.

The cache record (and version) the step persists to, declared via the step's cache_table/cache_table_version module vars. Returns (None, 1) for steps with no cache view (or an unknown step/version).

Lets DAG setup enumerate the record versions a task pins — to migrate each record's physical table up to the needed version — without importing every step module.

get_config_cls​

def get_config_cls(name: str, version: int = 1) ‑> type[pydantic.main.BaseModel] | None:

Get the Config class for name at version, or None.

get_result_cls​

def get_result_cls(name: str, version: int = 1) ‑> type[pydantic.main.BaseModel] | None:

Get the Result class for name at version, or None.

get_task​

def get_task(name: str, version: int = 1) ‑> collections.abc.Callable[..., typing.Any]:

Get the Prefect task callable for name at version.

get_task_hash_resources​

def get_task_hash_resources(name: str, version: int = 1) ‑> tuple[str, ...]:

Return the external resources step name's task_hash must reflect.

Empty for the vast majority of steps: only a step reading an external system that its config does not name needs one. Unknown steps get () for the same reason as the other optional-attribute accessors — the validators reject unregistered tasks before hashing. See bitfount.flows.dag.hashing.

partition_is_shareable​

def partition_is_shareable(name: str, version: int = 1) ‑> bool:

Return whether step name's partition may be shared across datasources.

See _partition_is_shareable for a step that writes a table. A step with no cache table is shareable: it writes no partition, so it has none to mis-scope, and keeping it on the per-datasource leaf would re-partition every cache-backed step downstream of it — which is most of the eligibility chain, since criteria_matching and the tabulation steps sit between the calculations and scan_eligibility. A descendant that must stay datasource-scoped is not affected: it mixes the leaf hash in on its own account.

Arguments

  • name: The registered step name.
  • version: The step version to look up.

Returns True when the step's partition may be shared; False for an unregistered step, the conservative answer since nothing is known about what it writes.

produces_file_scope​

def produces_file_scope(name: str, version: int = 1) ‑> bool:

Whether step name at version narrows the run's file selection.

Returns False for unknown steps, mirroring the leniency of the other optional-attribute accessors: an unrecognised step is an ordinary one, and is run in the wave loop like any other.

Arguments

  • name: The step's task name.
  • version: The step version.

Returns Whether the step declares produces_file_scope.

requires_ehr_data_resource​

def requires_ehr_data_resource(name: str, version: int = 1) ‑> bool:

Whether step name at version needs a live ehr_data_resource.

False for the vast majority of steps, which never touch the EHR at all. A step opts in by declaring a module-level requires_ehr_data_resource = True because it raises rather than degrades when the resource is None — ehr_query v1/v2 and ehr_patient_lister v1 each raise ValueError on entry. ehr_query v3 does not declare it, being the version that serves stored rows instead.

Returns False for unknown steps, mirroring the leniency of the other optional-attribute accessors: an unrecognised step is not one this can claim will abort.

Classes​

StepEntry​

class StepEntry(    task_fn: Callable[..., Any],    config_cls: type[BaseModel] | None = None,    result_cls: type[BaseModel] | None = None,    cache_accessor_fn: Callable[..., Any] | None = None,    cache_table: str | None = None,    cache_table_version: int = 1,    task_hash_resources: tuple[str, ...] = (),    emits_batch_progress: bool = False,    produces_file_scope: bool = False,    allows_empty_partition: bool = False,    requires_ehr_data_resource: bool = False,    partition_is_shareable: bool = False,    runtime_params: frozenset[str] = frozenset(),):

A single registered step version.

Variables​

  • static allows_empty_partition : bool
  • static cache_table : str | None
  • static cache_table_version : int
  • static config_cls : type[pydantic.main.BaseModel] | None
  • static emits_batch_progress : bool
  • static partition_is_shareable : bool
  • static produces_file_scope : bool
  • static requires_ehr_data_resource : bool
  • static result_cls : type[pydantic.main.BaseModel] | None
  • static runtime_params : frozenset[str]
  • static task_hash_resources : tuple[str, ...]