waves
Planning a waved background run: which files go through the DAG, and when.
A background run is normally handed the datasource's entire file list, and each step runs to exhaustion over it before the next one starts. Nothing downstream exists until the slowest step has swept everything, which on a large datasource is days. Waving cuts the sweep into batches that each go through the whole DAG, newest first, so the most recently scanned patients are recruitable in minutes.
The unit is a patient, not a file. A patient's eligibility verdict is reduced from every scan stored for them, so splitting their files across two waves would publish a verdict computed from part of their evidence and correct it silently several waves later. Cutting on patient boundaries removes that class of error entirely: within a sweep, no patient is ever judged on a subset of what the sweep found.
Two preconditions follow from that, and both are checked by WavePlanner
rather than assumed here:
- The scan sweep must be complete.
scan_metadatachunk-commits as it parses, so before it finishes a patient can be partly identified — three of their five scans have rows and two do not — and a wave cut then would split them after all. A patient's file set is only final once the sweep is. - The DAG must have something file-scoped to narrow. A background phase
made entirely of patient-grain steps (
ehr_trial_v9: list, query, enrich) has no file list to hand a wave of, and waving it would re-issue a full external patient listing per wave.
Recency is the patient's most recent acquisition time, from
scan_metadata.scan_datetime, falling back to the file's mtime. The fallback
matters less than the primary: the realistic onboarding case is a customer bulk
-copying a PACS archive, and cp without -p stamps every file with the copy
time, so an mtime-ordered sweep would deliver its waves in essentially arbitrary
order and the whole feature would be pointless on the datasource that needs it
most.
Which files are already finished comes from the dag_wave_ledger, because no
step's own partition answers it: the per-file calculation steps write a row for
every candidate including the ones they failed, and the eligibility steps are
patient-grain and hold no file id at all.
"Finished" is relative to a freshness horizon, which is what puts periodic reruns on this mechanism rather than beside it. A rerun exists to re-pull EHR data for patients whose files have not changed, so on a fully ledgered lineage a planner that only looked for unledgered files would plan nothing and the rerun would do nothing. Given a horizon, a ledger row older than it counts as outstanding again, and the ordinary patient chunking cuts the re-sweep into waves like any other — rather than the single unbounded pass over the whole datasource that a rerun used to get.
Re-swept patients are ordered by staleness, oldest refresh first, where new patients are ordered by recency. The horizon comes from the cron calendar and so moves at every occurrence; under recency ordering a datasource too large to re-sweep within one occurrence would redo its freshest patients every night and never reach its tail. Staleness ordering makes the sweep rotate instead, which is also what lets the horizon stay derived rather than persisted: there is no epoch to carry between runs.
Module
Functions
awaiting_first_sweep
def awaiting_first_sweep( cache: CacheProtocol, dag: BackgroundDAG, task_hash: str, project_id: str | None,) ‑> bool:Return whether a run must stand down until the first scan sweep settles.
True for a waveable DAG, with waving enabled, whose scan sweep has not read
a settled inventory and whose ledger holds nothing for this scope. That is
a run landing during a first index: build_wave_plan would decline to wave
it, and the single pass it fell back to would compute and publish over a
partial inventory with no scan identities.
The exact complement of the COLD_START case of due_for_wave_continuation
on the sweep and the ledger, which is what guarantees a stood-down run is
picked up once the sweep lands. A non-empty ledger is deliberately excluded:
no cold start would follow it.
Never raises: a cache it cannot read answers False, which keeps the single
pass rather than stranding the run.
Arguments
cache: The pod's background cache.dag: The lineage's background DAG, with its step hashes stamped.task_hash: The datasource-level task hash.project_id: The lineage's project, orNone.
Returns Whether the run should end without running any step.
build_wave_plan
def build_wave_plan( *, dag: BackgroundDAG, cache: CacheProtocol, task_hash: str, project_id: str | None, run_id: str | None, context: FileMetadataContext, datasource: Any = None, rerun_horizon: datetime | None = None, should_yield: Callable[[], Awaitable[bool]] | None = None, now: datetime | None = None,) ‑> WavePlan | None:Decide whether this run is waved, and plan its waves if so.
Returns None — the executor's untouched path — in four cases, each logged
once with its reason, because "why is this pod not waving?" is otherwise
unanswerable from the outside:
- Waving is disabled (
background_dag_wave_sizeis zero). - The DAG has no file-scoped step to hand a wave of (
is_waveable). - The
scan_metadatasweep has not finished, so patient identity is not final and a wave cut would split a patient. - Planning itself failed. A run that cannot plan should sweep the whole datasource as it always did, not fall over.
A plan with no waves is not None: it means the run is waved and has
nothing outstanding, which the executor ends cleanly rather than sweeping.
Arguments
dag: The validated background DAG, with its step hashes stamped.cache: The pod's background cache.task_hash: The datasource-level task hash.project_id: The lineage's project, orNone.run_id: The flow run, banked on the ledger rows.context: The run's file-metadata provider.datasource: The datasource the steps share.rerun_horizon: The freshness horizon for a refresh run, orNone. A ledger row older than it is outstanding again, which is what makes a periodic rerun re-carry patients whose files have not changed. Also selects the wave size, since a refresh and a first sweep want different ones, and whether the run is budgeted.should_yield: Awaited between waves; truthy when non-rerun work is live elsewhere on this pod, so a refresh should stand aside.Nonewhen the caller has nothing to yield to.now: The current time, for the future clamp. Defaults to the wall clock; a parameter so tests need not freeze it.
Returns
The plan, or None when the run is not waved.
due_for_wave_continuation
def due_for_wave_continuation( cache: CacheProtocol, dag: BackgroundDAG, task_hash: str, project_id: str | None, refresh_before: datetime | None = None,) ‑> WaveContinuationReason | None:Return why dag's lineage needs a run scheduled, or None.
COLD_START when the scan sweep has finished and this DAG's ledger is
empty — no wave has ever been banked for this configuration, so the
sweep has not started. A finished sweep leaves a non-empty ledger, so this
cannot re-fire.
REFRESH when the ledger is not empty but holds rows older than
refresh_before. This is what resumes a budgeted refresh across runs, and
it terminates: every carried wave moves its files' completed_at past the
horizon, so the stale count strictly decreases. A lineage whose waves keep
failing is bounded elsewhere — its work stays outstanding, on_success
never fires, and the growth trigger's give-up streak stands it down.
Cold start wins when both could apply: a lineage with an empty ledger has no stale rows to find anyway, and the distinction decides how the run is scheduled.
The last precondition is what keeps this out of recovery's way. A sweep
that crashes after banking a wave is excluded from COLD_START by the
ledger, but one that crashes before its first wave lands leaves the ledger
empty and looks exactly like a cold start. Firing on it starts a fresh
attempt=1 lineage, and recovery decides supersession on start time across
the whole task hash (runtimes.recovery._is_superseded) — so the new run
masks the crashed one, the replay never happens, and the attempt chain
never advances.
Never raises: a cache it cannot read answers "not due", which costs a tick rather than a poller.
Arguments
cache: The pod's background cache.dag: The lineage's background DAG, with its step hashes stamped.task_hash: The datasource-level task hash.project_id: The lineage's project, orNone.refresh_before: The current freshness horizon, orNoneto ask about the cold start only.
Returns
The reason a run should be scheduled now, or None.
is_waveable
def is_waveable(dag: BackgroundDAG) ‑> bool:Return whether dag has any file-scoped work to cut into waves.
A background phase made only of patient-grain steps has no file list to hand a wave of, and the scope step does not count.
Arguments
dag: The parsed background phase.
Returns Whether any step computes over the file inventory.
partition_key_for
def partition_key_for(dag: BackgroundDAG, project_id: str | None) ‑> str:Return the ledger scope for dag under project_id.
A digest over every runnable step's stamped Merkle task_hash, so a config
change anywhere in the DAG moves the key and the datasource re-waves. A
step with no stamped hash contributes its name; the scope step contributes
neither.
Arguments
dag: The background DAG, ideally with its step hashes stamped.project_id: The lineage's project, orNone.
Returns A 32-character hex digest.
plan_waves
def plan_waves( candidates: Sequence[WaveCandidate], ledgered: Mapping[str, datetime | None], wave_size: int, now: datetime, refresh_before: datetime | None = None,) ‑> list[Wave]:Cut candidates into the waves a run should carry through the DAG.
Patients come first, then orphans — files the scan sweep could not identify,
which are waved rather than skipped: they still settle as missing_data:
rows as they do today, and they are in the published coverage denominator,
so a sweep that never banked them would report itself permanently
incomplete. Within each, work the ledger has never carried comes before
work it carried too long ago.
Four groups, in this order:
- Patients with any unledgered file, newest scan first. This is what makes a newly linked dataset recruitable in minutes, and it is the only group a link, a recovery or a cold-start continuation ever has.
- Patients whose every file is ledgered but whose oldest completion predates refresh_before, oldest refresh first. Empty without a horizon.
- Unledgered orphans, then 4. stale orphans, on the same rule.
A patient with any outstanding file is re-waved whole. The boundary is the point — a verdict must be computed with the patient's entire scan set in front of it — and the files already done cost almost nothing, since every step skips its own cached work.
Group 2 is how a rerun is expressed. It replaces a trailing wave over every candidate, which was unbounded by construction and carried the run's new patients a second time, having already carried them in group 1.
Arguments
candidates: Every file in the run's selection.ledgered: What the ledger holds for this scope: file path to completion time,Nonewhere the stored time could not be read.wave_size: Patients per wave, and files per orphan wave.now: The current time, for the future clamp.refresh_before: The freshness horizon. A file carried earlier is outstanding again, which is what makes a periodic rerun re-carry patients whose files have not changed.Nonefor every trigger that is not a refresh, leaving groups 2 and 4 empty and this function's behaviour exactly what it was before horizons existed.
Returns The waves, indexed from 1. Empty when nothing is outstanding.
Raises
ValueError: If wave_size is not positive. Zero disables waving upstream and must never reach here; waves of no patients would never terminate.
scan_sweep_complete
def scan_sweep_complete(cache: CacheProtocol, task_hash: str) ‑> bool:Return whether scan_metadata has swept a settled inventory.
Waving may not start before it has, for two reasons that are really one: a wave is cut on patient boundaries, and a patient's file set is only final once everything that produces it has finished.
scan_metadata chunk-commits as it parses, so mid-sweep a patient can be
partly identified — three of their five scans have rows and two do not —
and a wave cut then would split them, publishing a verdict computed from
part of their evidence and correcting it silently several waves later.
The inventory underneath moves the same way. file_metadata persists its
walk in chunks, so rows exist long before all of them do, and "the
inventory has rows" is true within seconds of a first index starting. Only
the run finishing says the file set is settled.
Hence the ordering rather than two existence checks: the sweep counts only if it finished after the index it was meant to read. That is one question, and it rejects every way the pair can be out of step — an index still running or never run, a sweep that stood down for a half-built inventory (which completes its bookkeeping row having written nothing, and so predates the index it stood down for), and a re-index that has landed new files no sweep has identified yet.
A sweep that ran and wrote nothing still passes here — the datasource may
simply not be ophthalmology. That is a question about this run's own files
against a pod-wide table, so WavePlanner.candidates asks it instead and
raises NoScanCoverageError.
Two completions that share a timestamp are not ordered, so they answer
"not swept" as well. Windows stamps datetime.now from a ~15.6ms-grained
clock, and a sweep that stood down writes its bookkeeping row within that
of the index it stood down for — the tie that must not read as a sweep. A
sweep that truly read the inventory parses it first, so it lands ticks
later and keeps its order.
Never raises. This is a gate on an optimisation, and a cache it cannot read answers "not ready" — which falls back to the single pass the run would have done anyway, rather than failing a run over a scheduling question.
scan_metadata chunk-commits as it parses, so mid-sweep a patient can be
partly identified — three of their five scans have rows and two do not —
and a wave cut then would split them, publishing a verdict computed from
part of their evidence and correcting it silently several waves later.
The inventory underneath moves the same way. file_metadata persists its
walk in chunks, so rows exist long before all of them do, and "the
inventory has rows" is true within seconds of a first index starting. Only
the run finishing says the file set is settled.
Hence the ordering rather than two existence checks: the sweep counts only if it finished after the index it was meant to read. That is one question, and it rejects every way the pair can be out of step — an index still running or never run, a sweep that stood down for a half-built inventory (which completes its bookkeeping row having written nothing, and so predates the index it stood down for), and a re-index that has landed new files no sweep has identified yet.
A sweep that ran and wrote nothing still passes here — the datasource may
simply not be ophthalmology. That is a question about this run's own files
against a pod-wide table, so WavePlanner.candidates asks it instead and
raises NoScanCoverageError.
Arguments
cache: The pod's background cache.task_hash: The datasource-level task hash.
Returns
Whether a scan_metadata run completed strictly after the last
file_metadata run did.
wave_scope
def wave_scope( context: WaveScopeProvider, datasource: Any, wave: Wave,) ‑> collections.abc.Iterator[None]:Scope a run to wave for the duration of the block.
Sets both channels a file scope travels down, from one value:
FileMetadataContext, which serves$file_metadata.cacheand.recordsto every step that reads them; anddatasource.selected_file_names_override, becausemodel_inferencedoes not readfilenamesat all — it iteratesdatasource.selected_file_names_iter(), andselection.select_for_datasourcenever consults the override. Scoping only the context would leave the most expensive step in the DAG running full-sweep on every wave, which is worse than not waving at all.
Both are restored on the way out, on the raising path too. The override in
particular leaks for the lifetime of the process otherwise:
clear_task_specific_configs does not clear it, which is why
bitfount.federated.worker already resets it in a finally.
Arguments
context: The run's file-metadata provider.datasource: The datasource the steps share, orNonewhen there is none. Typed loosely because this probes for the override attribute rather than requiring aBaseSource: a source that cannot be selected by file name simply does not carry it.wave: The wave to scope to.
Raises
ValueError: If wave holds no files. An empty override is falsy, andselected_file_names_iterreads that as "no override" and walks the whole datasource — so an empty wave would silently process everything rather than nothing.
Classes
NoScanCoverageError
class NoScanCoverageError(*args, **kwargs):Raised when the scan sweep holds nothing for the files this run selected.
The scan table is pod-wide but a sweep is per-datasource, and the sweep only runs for ophthalmology sources. A run whose own files have no scan rows has no patient identities to wave by, so it is not waved at all — it sweeps the whole selection in one pass, as it did before waving existed.
Wave
class Wave( index: int, candidates: tuple[WaveCandidate, ...], is_orphan_wave: bool = False, is_refresh_wave: bool = False,):One batch of files to carry through the whole DAG.
Attributes
index: The wave's position in its sweep, from 1. Reported in the logs and banked in the ledger.candidates: The files, in a stable order.is_orphan_wave: Whether this wave holds files with no patient identity. They still go through every step and still settle asmissing_data:rows, but a step that logs loudly when nothing loads is expected to do so here, so the loop can downgrade that one line.is_refresh_wave: Whether this wave re-carries files the ledger already holds, rather than carrying them for the first time. What makes the two recoverable by different triggers, and so whatWavePlan.left_work_outstandingturns on: a refresh wave left undone is visible to the scheduler as a stale ledger row, while a first-carry wave left undone is visible only as inventory the growth cursor has not reached.
Variables
- static
candidates : tuple[WaveCandidate, ...]
- static
index : int
- static
is_orphan_wave : bool
- static
is_refresh_wave : bool
file_paths : tuple[str, ...]- The wave's files, in the order they were planned.
Methods
ledger_records
def ledger_records( self, completed_at: datetime, run_id: str | None,) ‑> list[WaveLedgerRecord]:Build the ledger rows banking this wave as finished.
Arguments
completed_at: When the wave finished.run_id: The flow run that carried it.
Returns One record per file in the wave.
WaveCandidate
class WaveCandidate( file_path: str, bitfount_patient_id: str | None = None, scan_datetime: datetime | None = None, modified_at: datetime | None = None,):One file the sweep may carry through the DAG.
Attributes
file_path: The file, asfile_metadataspells it.bitfount_patient_id: The patient it belongs to, orNonewhen the scan sweep could not identify one — an orphan.scan_datetime: The acquisition time, orNone.modified_at: The file's mtime, used only when there is no acquisition time to order by.
Variables
- static
bitfount_patient_id : str | None
- static
file_path : str
- static
modified_at : datetime.datetime | None
- static
scan_datetime : datetime.datetime | None
WaveContinuationReason
class WaveContinuationReason(*args, **kwds):Why a lineage needs a background run scheduled outside its triggers.
Two different situations, deliberately not collapsed into one boolean, because they are scheduled differently. A cold start is a lineage's first pass over work nobody has done, so it must not queue behind the rerun limit; a refresh is discretionary re-work and must.
Attributes
COLD_START: The sweep has never begun. ADATASET_PROJECT_LINKEDrun is scheduled immediately and is gated on neither metadata runtime, so it can arrive before the scan sweep has settled, whenawaiting_first_sweepstands it down without running a step. Without this the lineage would sit idle until the next cron occurrence.REFRESH: The sweep has begun but the lineage still holds ledger rows older than the current freshness horizon. The normal state of a budgeted refresh part way through: the run that carried the last slice ended on purpose, and this is what carries the next.
Ancestors
WaveLedgerWriter
class WaveLedgerWriter(*args, **kwargs):The part of the planner a running sweep drives.
The executor's only dependency on planning is "this wave finished, record it". Naming that one method keeps the loop testable without a cache, and keeps the loop from reaching into the planner for anything else.
Ancestors
Methods
mark_complete
def mark_complete(self, wave: Wave, completed_at: datetime) ‑> None:Bank wave as carried through the whole DAG.
WavePlan
class WavePlan( waves: list[Wave], planner: WaveLedgerWriter, context: WaveScopeProvider, datasource: Any = None, retry_limit: int = 2, last_error: Exception | None = None, banked_waves: int = 0, waves_started: int = 0, deferred_files: int = 0, wave_budget: int | None = None, should_yield: Callable[[], Awaitable[bool]] | None = None, unstarted_waves: list[Wave] = [],):A run's planned waves, plus what it needs to run and bank them.
Carried on DAGRunContext so the executor can loop without knowing how a
wave was chosen or where it is recorded.
Attributes
waves: The waves, in order.planner: Banks a finished wave in the ledger.context: The run's file-metadata provider — one of the two channels a wave's scope travels down, and the holder of the cumulativecoveredlist the published coverage figure reads.datasource: The datasource the steps share, the other channel. Typed loosely becauseDAGRunContextis a pydantic model carrying this, and aBaseSourceannotation would force it to resolve a name this module only imports for type checking.retry_limit: How many times a failing wave is retried before the run gives up on it and moves to the next.last_error: The most recent wave failure, kept so a run in which every wave failed can re-raise something specific rather than a summary.banked_waves: How many waves this run actually finished and recorded. Fewer than it planned means a wave exhausted its retries and its files were not carried.waves_started: How many waves the loop has entered, banked or not. The sweep's position, wherebanked_wavesis its success count: a run that abandoned a wave is still past it, and the coverage figure reporting what is left has to count down from here or it reports the last wave of the sweep as having one still to come.deferred_files: How many selected files the planner left out because the scan sweep had not reached them. They are outstanding work like an abandoned wave's files, and for the same reason: nothing has carried them yet.wave_budget: How many waves this run may carry before it ends and leaves the rest to a later one, orNonefor no bound. What stops a refresh holding the machine for hours; seeconfig.background_dag_rerun_wave_budget.should_yield: Awaited after each banked wave, and truthy when work that is not a rerun is live elsewhere on the pod. The rerun then ends early rather than competing with it. Awaitable because the DAG runs in its own process, so the only thing that can answer it is a Prefect query.Nonewhen nothing can be yielded to.unstarted_waves: Waves this run deliberately did not begin, having spent its budget or stood aside. Counted apart from failures: they are a scheduling decision, not a fault.
Variables
- static
banked_waves : int
- static
context : WaveScopeProvider
- static
datasource : Any
- static
deferred_files : int
- static
last_error : Exception | None
- static
planner : WaveLedgerWriter
- static
retry_limit : int
- static
should_yield : collections.abc.Callable[[], collections.abc.Awaitable[bool]] | None
- static
unstarted_waves : list[Wave]
- static
wave_budget : int | None
- static
waves : list[Wave]
- static
waves_started : int
-
abandoned_waves : int- How many waves this run started and could not carry.A wave that exhausted its retries, as distinct from one the run never began: deferring work is a scheduling decision, abandoning it is a fault. The failure streak counts on that distinction — a lineage that abandons a wave has not finished what it started, however much else it got through.
Returns: The count, zero when every started wave was banked.
-
left_work_outstanding : bool- Whether this run finished with first-carry work it did not carry.What the inventory cursor turns on: it records how much of the datasource a completed run saw, so a run that failed a wave, or left files out for want of a scan row, must not advance it. Banking the whole inventory there would make the files it never reached stop looking new, and the growth trigger — which fires on the inventory growing past the cursor — would never bring them back.
Refresh waves are the exception, and the distinction is load-bearing rather than a nicety. The two kinds of unfinished work are recovered by different triggers: a refresh wave left undone is a stale ledger row, which the scheduler can see and reschedule directly, while a first-carry wave left undone has no ledger row at all and is invisible to that check — its only safety net is the growth trigger, which the cursor disarms. So a budgeted run that carried every first-carry wave has covered the datasource and may bank the cursor, however much refreshing is left; one that stopped short of a first-carry wave has not.
A plan with no waves at all left nothing behind: every file was already carried and fresh, so the run covered the datasource by having nothing to do.
Returns: Whether work this run was meant to carry is still outstanding.
WavePlanner
class WavePlanner( cache: CacheProtocol, partition_key: str, project_id: str | None, wave_size: int, run_id: str | None = None, rerun_horizon: datetime | None = None,):Plans a run's waves and scopes the run to one at a time.
Owns the cache reads planning needs — the scan identities, the ledger — and
the two channels a wave's scope has to reach. Built once per run in
bitfount.flows.dag.setup, and None when waving is disabled or the DAG is
not waveable, which is how the executor takes its untouched path.
Arguments
cache: The pod's background cache.partition_key: The digest over the DAG's stamped step hashes, scoping the ledger so a config change re-waves the datasource.project_id: The lineage's project, orNone.wave_size: Patients per wave.run_id: The flow run, banked on the ledger rows.rerun_horizon: The freshness horizon, orNonewhen this run is not a refresh. Passed toplan_waves, where it decides whether a fully-ledgered patient is outstanding again.
Initialise the planner (see the class docstring).
Methods
candidates
def candidates( self, records: Iterable[FileMetadataRecord],) ‑> list[WaveCandidate]:Build the run's candidates from its inventory rows.
Joins each file to what the scan sweep found for it: the patient it belongs to and when it was acquired. A file whose rows carry no usable name and date of birth becomes an orphan — ordered last and waved by file count rather than by patient.
A file the sweep has produced no row for is left out of the run entirely. It is not an orphan, it is unfinished: the sweep runs on its own schedule and the inventory moves under it, so a file indexed since the last sweep has nothing to identify it yet. Waving it would bank it in the ledger under no patient, and the ledger is what stops a file being waved twice — so a returning patient's new scans would be carried without them and never judged alongside the rest of their evidence. Left out, they wave with the patient on the run after the sweep catches up.
Arguments
records: The run's selected inventory rows, from$file_metadata.records.
Returns One candidate per inventory row the scan sweep has reached.
Raises
NoScanCoverageError: If the sweep holds no row for any of records. The scan table is pod-wide, so another datasource's rows are no evidence that this one was swept: a pod running one ophthalmology datasource alongside one that is not would otherwise wave the second on the first's rows, defer every file for want of an identity, and process nothing at all.
identities
def identities(self) ‑> ScanIdentitySource:Return per-file identity and recency, from one pass over the scans.
Returns The source the candidates are built from.
ledgered
def ledgered(self) ‑> dict[str, datetime.datetime | None]:Return what the ledger holds for this scope, with its timestamps.
One read answering both of the caller's questions: which files have ever been carried — the cumulative coverage figure's starting point — and when, which is what a refresh measures against its horizon.
Returns
File path to completion time, None where it could not be read.
mark_complete
def mark_complete(self, wave: Wave, completed_at: datetime) ‑> None:Bank wave in the ledger as carried through the whole DAG.
Called only once the wave's steps have all returned, so a wave that failed and was abandoned leaves no trace and is retried by the next run.
Arguments
wave: The finished wave.completed_at: When it finished.
waves
def waves( self, records: Iterable[FileMetadataRecord], now: datetime, ledgered: Mapping[str, datetime | None] | None = None,) ‑> list[Wave]:Plan this run's waves.
Arguments
records: The run's selected inventory rows.now: The current time, for the future clamp.ledgered: What the ledger holds, when the caller has already read it.build_wave_planhas, because it also seeds the run's coverage figure from it, and the ledger is the largest read planning does.
Returns The waves, indexed from 1.
WaveScopeProvider
class WaveScopeProvider(*args, **kwargs):The part of the file-metadata provider a wave scope drives.
Runtime-checkable because WavePlan carries one and rides on
DAGRunContext, a pydantic model, which validates an arbitrary type with
isinstance — and isinstance refuses a Protocol that is not.
Narrower than FileMetadataContext on purpose: a scope sets and clears a
wave's file list and records what a finished wave covered, and nothing
else. Naming that surface lets a caller be tested against it without
standing up a provider over a real cache.
Ancestors
Methods
mark_covered
def mark_covered(self, file_paths: Collection[str]) ‑> None:Record file_paths as carried through the whole DAG.
narrow_to
def narrow_to(self, file_paths: Collection[str] | None) ‑> None:Scope the provider to file_paths, or clear the scope with None.