Progress

progress.WorkSet

A resumable set of work units journaled as JSONL; errors are retried on resume.

Usage

Source

progress.WorkSet(
    path,
    _done,
)

Each unit is journaled as it lands: a done record marks it complete, an error record leaves it pending so it is retried on the next open.

Parameter Attributes

path: Path
_done: set[str]

Example

>>> work = WorkSet.open(Path("~/.athome/run.jsonl"))
>>> for unit in work.pending(["a", "b", "c"]):
...     await work.done(unit)

progress.RunSink

Crash-safe incremental JSONL sink with a failure budget.

Usage

Source

progress.RunSink(path, failure_budget, _failures)

Every record is one appended JSON line; fail tags its record "failed": true and raises FailureBudgetExceeded once the recorded failures exceed the budget. The budget survives a crash: open recounts failed records from the journal.

Parameter Attributes

path: Path
failure_budget: int
_failures: int

Example

>>> sink = RunSink.open(Path("~/.athome/out.jsonl"), failure_budget=3)
>>> await sink.append({"item": "a", "ok": True})

progress.Phases

Ordered phase markers gating multi-stage runs.

Usage

Source

progress.Phases(path, _marked)

mark journals a phase as reached; require raises PhaseMissing when a prerequisite has not been marked. Markers survive a crash via the journal.

Parameter Attributes

path: Path
_marked: set[str]

Example

>>> phases = Phases.open(Path("~/.athome/phases.jsonl"))
>>> phases.require("extract")

progress.FailureBudgetExceeded

Raised when a RunSink’s recorded failures exceed its budget.

Usage

Source

progress.FailureBudgetExceeded()

progress.PhaseMissing

Raised by Phases.require when a required phase has not been marked.

Usage

Source

progress.PhaseMissing()