periodic_reruns
Rerunning healthy background lineages on a schedule.
The third way a background DAG run starts, after a dataset-project link and
bitfount.flows.dag.recovery's crash replay — and the only one that fires for a
lineage where nothing went wrong. A registered background task otherwise runs
exactly once ever: files that arrive afterwards are never processed, and EHR data
that changes for an already-processed patient is never re-fetched.
Why the calendar is read backwards from now
Due-ness is "the most recent cron occurrence is recent enough, and we have not run since it", not "one interval has passed since the last run". The two differ for exactly the machine this runs on: a customer laptop that is off at midnight most nights. An interval-since-last-run rule fires the moment the lid opens, whatever the hour and however long the gap; an occurrence-based rule fires only if the missed occurrence is still inside the grace window, and fires once however many ticks observe it. That also means no calendar cursor has to be persisted — the wall clock and one timestamp are the whole state.
The cron/timezone/grace/cooldown values are passed in rather than read
from config here, so every function in this module is pure and testable against
a fixed clock. The pod supplies them.
Module
Functions
due_for_growth_rerun
def due_for_growth_rerun( records: Sequence[FlowSpecRecord], inventory_count: Callable[[FlowSpecRecord], int | None], *, now: datetime, cooldown: timedelta, failure_limit: int,) ‑> GrowthReruns:Return the records whose datasource has been indexed further since last run.
The complement to due_for_periodic_rerun, and the reason the two coexist:
the cron exists partly to refresh EHR data that changes with no new files at
all, so growth cannot replace it. Growth covers what a daily clock is too
coarse for — a first index of a large share, where the cohort should grow as
the walk proceeds rather than once a night.
A lineage with no banked count is due. That is the "no baseline yet" case, and failing toward doing the work is deliberate: reading it as "nothing has grown" would wedge a lineage whose every run so far has failed, leaving it to the cron for ever. The first run to complete banks a count and the lineage settles.
Two bounds, because growth has no occurrence to fire once per. Unlike the cron, it re-fires on every tick of the poller for as long as its condition holds, and both of its due conditions can hold indefinitely:
- A lineage whose runs fail terminally never banks a count, because the
cursor is written on success only — so "no baseline yet" stays true for
ever. The same lineage loops on the other branch too once a single run has
banked something: a growing inventory it can never process keeps
currentabovebankedjust as permanently. - A first index of a large share grows the count on every tick, so each rerun succeeds, banks, and is immediately due again — the intended direction at an unintended cadence, a full reduce and inference pass per poll for the length of the walk.
cooldown bounds the rate, measured from last_rerun_at falling back to
linked_at — the same clock due_for_periodic_rerun reads. That stamp is
written when a rerun takes its concurrency slot rather than when it
succeeds, which is what makes it the right one here: it advances even for
the failing lineage above, where the inventory cursor by design does not.
failure_limit bounds the total, and is why growth needs a give-up of its
own rather than recovery's. Recovery's bound is on rows banked and is
explicitly exempt for rerun-triggered attempts — a rerun over unchanged
data banks nothing by construction, and counting those readings once
abandoned healthy lineages outright. It also never sees this failure at all:
it acts on runs whose process is lost, and a run that failed terminally is
not one. So the signal here is reruns_since_completion instead — runs that
started and never finished — and what it stands down is growth alone. The
cron keeps its schedule, since it refreshes EHR data this trigger knows
nothing about and may well succeed where growth's reruns did not, and
gave_up_at is left alone, so crash recovery stays armed. Any completed run
resets the count, as does re-linking the dataset.
Both bounds are checked before the count is taken, so a lineage held back by either costs no inventory query at all.
Increases only. The inventory also shrinks — prune_missing_under_root
deletes rows for files gone from disk after a successful walk — but a rerun
does not actually remove a deleted scan from the served cohort:
scan_eligibility only ever upserts, with no delete or prune, so the old
row survives and keeps contributing to its patient's rollup. Spending a full
reduce pass on something that cannot fix it is not worth it, so a shrink is
logged and left to the cron. Making deletions propagate is its own piece of
work, and this comment is where to start from.
Arguments
records: Candidate records, typicallylist_flow_specs(replayable_only=True).inventory_count: How many files are indexed for a record's datasource now, orNonewhen that cannot be measured — an unconfigured datasource, or one with no filesystem root. Unmeasurable is neverdue: there is nothing to compare.now: The current time.cooldown: How long after a lineage's last rerun growth may fire again.failure_limit: How many reruns may start without any of them completing before growth stands down for the lineage. Zero or less disables the trigger outright.
Returns The due records and the stood-down ones, both in the order given.
due_for_periodic_rerun
def due_for_periodic_rerun( records: Sequence[FlowSpecRecord], *, now: datetime, cron: str, timezone: tzinfo, grace: timedelta,) ‑> list[FlowSpecRecord]:Return the records whose next scheduled rerun is due and not stale.
A record is due when the most recent occurrence T satisfies both:
now - T <= grace— the occurrence has not gone stale, which is what stops a pod that was off for a week firing at an arbitrary hour; and- the lineage has not run since
T— measured fromlast_rerun_at, orlinked_atwhen it has never been rerun, because the link itself ran the DAG.
last_rerun_at is written by any rerun trigger, not only this one, so a
manual rerun hours before a scheduled occurrence satisfies it rather than
having the DAG redone.
Arguments
records: Candidate records, typicallylist_flow_specs(replayable_only=True).now: The current time.cron: The shared schedule.timezone: Zone the schedule is evaluated in.grace: How stale an occurrence may be and still fire.
Returns The due records, least-recently-rerun first.
Raises
ValueError: If cron is not a valid expression.
latest_occurrence
def latest_occurrence( *, cron: str, timezone: tzinfo, now: datetime,) ‑> datetime.datetime | None:Return the most recent occurrence of cron at or before now.
Evaluated in timezone so that "midnight" means midnight where the pod is,
which is what a customer means by it. CronSim in reverse yields
occurrences strictly before the datetime it is seeded with, so a tick
landing exactly on an occurrence reads back to the previous one and the
current one is picked up by the next tick a minute later.
Arguments
cron: A five-field cron expression.timezone: Zone the expression is evaluated in.now: The current time; may be in any zone, and is converted.
Returns
The occurrence, in timezone, or None if the expression has no past
occurrence (a fixed date already gone by, say).
Raises
ValueError: If cron is not a valid expression.
rerun_horizon
def rerun_horizon( *, cron: str, timezone: tzinfo, now: datetime,) ‑> datetime.datetime | None:Return the freshness horizon a refresh run should re-sweep against.
The current cron occurrence, in UTC. bitfount.flows.dag.waves treats a
ledger row older than this as outstanding again, which is what makes a
rerun re-carry patients whose files have not changed — the one thing the
nightly cron exists for.
Deliberately the occurrence, and not the run's own start time nor
last_rerun_at. Growth-triggered reruns carry the same trigger and fire as
often as the growth cooldown, so a per-run epoch would re-sweep the whole
datasource every hour during a first index — worse than the unbounded
single wave this replaced. Against an occurrence-derived horizon a growth
rerun plans only the files with no ledger row at all, which is exactly the
new ones. It also makes every run inside one occurrence plan against the
same horizon, so a sweep resumed by a later run continues rather than
restarts, with no cursor persisted anywhere.
Normalised to UTC because the ledger stores UTC. latest_occurrence
answers in the evaluation zone, and an hour's offset silently reads stale
rows as fresh.
Arguments
cron: The shared schedule.timezone: Zone the schedule is evaluated in.now: The current time.
Returns
The occurrence, or None when the expression has no past occurrence —
which the caller reads as "no horizon", planning outstanding work only.
Raises
ValueError: If cron is not a valid expression.
resolve_timezone
def resolve_timezone(configured: str | None) ‑> datetime.tzinfo:Return the zone the rerun calendar is evaluated in.
configured is an IANA name, already validated at settings load; None
means "the pod's local zone", which has to be a real zone rather than the
offset that zone happens to be at right now.
datetime.now().astimezone().tzinfo gives the latter — a fixed
timezone(timedelta(...), 'BST') sampled at call time — and a fixed offset
double-fires across a DST fall-back: with TZ=Europe/London on 2026-10-25,
23:30 UTC computes midnight as 00:00+01:00 and fires, then 01:30 UTC
computes it as 00:00+00:00, an hour earlier, and fires again. Spring
forward has the mirror symptom, a catch-up falling outside a grace window
that should have admitted it.
zoneinfo has no local-zone lookup, so tzlocal does it. Its failure is
survivable and must not take the poller down — a pod with an unreadable zone
database should still rerun on the offset it can see — so the fixed offset
stays as a warned fallback.
Arguments
configured:background_periodic_rerun_timezone, orNonefor local.
Returns The zone to evaluate cron occurrences in.
Classes
GrowthReruns
class GrowthReruns(due: list[FlowSpecRecord], stood_down: list[FlowSpecRecord]):What one growth check decided.
Two lists rather than one, because standing a lineage down is a fault worth
an operator's attention while being due is routine — and this module cannot
report the first itself. It is called once per poll, so a logger.warning
here would repeat for the pod's whole life; the pod owns the once-per-
occurrence suppression that makes such a report readable, and reports these
through it.
Attributes
due: Lineages to rerun, in the order given.stood_down: Lineages growth has given up on, having startedfailure_limitreruns since any of them last completed. Not due, and not reported anywhere else.