Skip to main content

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 current above banked just 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, typically list_flow_specs(replayable_only=True).
  • inventory_count: How many files are indexed for a record's datasource now, or None when that cannot be measured — an unconfigured datasource, or one with no filesystem root. Unmeasurable is never
  • due: 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 from last_rerun_at, or linked_at when 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, typically list_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, or None for 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 started failure_limit reruns since any of them last completed. Not due, and not reported anywhere else.

Variables​