task
Prefect task for the patient_eligibility reduce step.
Re-aggregates every patient in the partition: reads their stored per-scan verdicts
(scan_eligibility) and their patient-level verdict (patient_level_eligibility)
— each under its own partition, located via parent_task_hashes — picks the
determining scan, combines the two verdicts, and merges a patient_eligibility
row. The trials_published_data_pointer is flipped last, with one entry per
table this run wrote — an atomic cutover so readers never see a half-written
partition.
The reads are deliberately of the tables, not of the rows this run happened to produce: rows accumulate across runs within a partition, so reading back is what keeps a partial run from publishing partial rollups.
Module
Functions
patient_eligibility_task
def patient_eligibility_task( *, scan_eligibility: CacheAccessor | None = None, patient_level_eligibility: CacheAccessor | None, cache: CacheProtocol, task_hash: str, project_id: str, parent_task_hashes: dict[str, str] | None = None, config: PatientEligibilityConfig | None = None, run_id: str | None = None, filenames: list[str] | None = None, scan_coverage_probe: ScanCoverageProbeProtocol | None = None,) ‑> PatientEligibilityResult:Combine per-scan + patient-level verdicts into per-patient rows; flip pointer.
The rollup rows are always written. The pointer flip is withheld for an
empty run, and for one run when this run would publish zero eligible
patients on EHR coverage worse than the currently-published partition's —
see _withhold_publish, which also documents the case this does not
cover. A flip that happens records that coverage and the partition's EHR
staleness span in the pointer's tags.
A flow that wires no scan_eligibility input is EHR-only: every
patient-level row is served on its own verdict, with no scan rollup and no
scan-coverage tags. A flow that wires it is scan-driven, and a patient with
no stored scan is not served.
Arguments
scan_eligibility: The upstreamscan_eligibilityedge, orNonewhen the flow wires none (EHR-only). Its value is unused — this step reduces the table through direct store calls, which an accessor's grouped queries cannot express — but declaring it is this step's dependency contract, and the name is the key its partition appears under inparent_task_hashes.patient_level_eligibility: The upstreampatient_level_eligibilityedge, declared on the same basis.cache: The pod background cache to read verdicts from and write rollups to.task_hash: This step's own partition key — where the rollup rows are written, and what the published pointer names for this table.project_id: The trial (project) ID (a runtime param).parent_task_hashes: Partition key of each upstream step, keyed by this step's input parameter name (bitfount.flows.dag.hashing). Where the two source tables are read from — this step's owntask_hashwould find nothing in them.Nonebecause the executor injects it only for a stamped DAG; an unstamped one then reaches the explicit error below rather than aTypeErrorabout a missing argument.config: Optional step config (trial display name); nothing in it is persisted by this step.run_id: Optional provenance run ID.filenames: The datasource's indexed file paths as resolved for this run ($file_metadata.cache). Read only for its length, as the numerator of the published scan-coverage figure — this step reduces stored rows and never reads a file.scan_coverage_probe: Injected probe reporting how much of the datasource is indexed now, for the denominator of that figure.Noneleaves the coverage tags unmeasured rather than failing the publish.
Returns
PatientEligibilityResult with the number of patient_eligibility rows
written and a cache accessor scoped to this task_hash.
Raises
ValueError: Ifproject_idis not supplied, or parent_task_hashes is missing thepatient_level_eligibilitypartition — without which this would silently roll up zero patients and leave the pointer on the previous run's data.