Skip to main content

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:

  1. The scan sweep must be complete. scan_metadata chunk-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.
  2. 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, or None.

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:

  1. Waving is disabled (background_dag_wave_size is zero).
  2. The DAG has no file-scoped step to hand a wave of (is_waveable).
  3. The scan_metadata sweep has not finished, so patient identity is not final and a wave cut would split a patient.
  4. 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, or None.
  • 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, or None. 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. None when 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, or None.
  • refresh_before: The current freshness horizon, or None to 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, or None.

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:

  1. 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.
  2. Patients whose every file is ledgered but whose oldest completion predates refresh_before, oldest refresh first. Empty without a horizon.
  3. 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, None where 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. None for 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.cache and .records to every step that reads them; and
  • datasource.selected_file_names_override, because model_inference does not read filenames at all — it iterates datasource.selected_file_names_iter(), and selection.select_for_datasource never 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, or None when there is none. Typed loosely because this probes for the override attribute rather than requiring a BaseSource: 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, and selected_file_names_iter reads 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 as missing_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 what WavePlan.left_work_outstanding turns 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 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, as file_metadata spells it.
  • bitfount_patient_id: The patient it belongs to, or None when the scan sweep could not identify one — an orphan.
  • scan_datetime: The acquisition time, or None.
  • 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

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. A DATASET_PROJECT_LINKED run is scheduled immediately and is gated on neither metadata runtime, so it can arrive before the scan sweep has settled, when awaiting_first_sweep stands 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.

Variables​

  • static COLD_START
  • static REFRESH

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.

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 cumulative covered list the published coverage figure reads.
  • datasource: The datasource the steps share, the other channel. Typed loosely because DAGRunContext is a pydantic model carrying this, and a BaseSource annotation 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, where banked_waves is 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, or None for no bound. What stops a refresh holding the machine for hours; see config.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. None when 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 datasource : Any
  • static deferred_files : int
  • static last_error : Exception | None
  • static retry_limit : int
  • 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, or None.
  • wave_size: Patients per wave.
  • run_id: The flow run, banked on the ledger rows.
  • rerun_horizon: The freshness horizon, or None when this run is not a refresh. Passed to plan_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_plan has, 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.

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.