scan_coverage
Scan-coverage probe protocol for step tasks.
A step that publishes results for a subset of a datasource has no way, from
its own inputs, to say how large that subset was relative to the datasource.
Its filenames list is resolved once per run and memoised
(bitfount.flows.dag.context.FileMetadataContext), so every step in the run
sees the same list and a step comparing its own inputs against it can only ever
conclude that it processed all of them.
This protocol is the seam for asking that question of something outside the run:
how many files the inventory holds now, and whether it is still being built.
The implementation lives in bitfount.flows.dag.coverage because answering it
needs the runtimes layer, which no step imports.
Classes
ScanCoverage
class ScanCoverage( files_indexed: int | None, indexing_in_flight: bool | None, waves_remaining: int | None = None, wave_in_progress: bool | None = None,):A point-in-time reading of how much of a datasource is indexed.
Attributes
files_indexed: Inventory rows selected for this datasource, counted at the moment of sampling rather than when the run started, orNonewhen the inventory could not be counted. A run that has been executing for hours is compared against what exists now, not against the list it was handed — which is the whole point of sampling, since the two are equal by construction otherwise.indexing_in_flight: Whether afile_metadatarun for this datasource was executing when the sample was taken, orNonewhen that could not be determined. Without it a coverage ratio reads as complete whenever a run happens to catch up with the inventory, which includes a walk that has merely paused between chunks — or stalled.waves_remaining: How many waves of this run's sweep are still to come after the one publishing, orNonewhen that could not be determined. Zero on a run that is not waved, whose single implicit pass is the one publishing.wave_in_progress: Whether more waves follow the one publishing — i.e. whether this cohort is a sweep still in progress rather than a finished one.Nonewhen it could not be determined.
This answers a different question from indexing_in_flight, and both are needed. That one says the denominator is still moving; this one says the numerator is. A waved run over a fully-indexed datasource has indexing_in_flight false throughout while publishing a cohort that grows with every wave, so without this a wave-1 cohort would read as complete.
Variables
- static
files_indexed : int | None
- static
indexing_in_flight : bool | None
- static
wave_in_progress : bool | None
- static
waves_remaining : int | None
ScanCoverageProbeProtocol
class ScanCoverageProbeProtocol(*args, **kwargs):An object that reports how much of a datasource is currently indexed.
Ancestors
Methods
sample
def sample(self) ‑> ScanCoverage:Return the inventory size and indexing state as of now.
Returns
A ScanCoverage reading taken at call time.