Skip to content

latent.prefect.resumable

Crash-resumable batch work: a keyed work list with per-item durability.

@checkpoint is an argument-hash memoizer and is deliberately a no-op outside development mode — a disk cache keyed on arguments has no business serving production results. That leaves the case its own example describes uncovered: a production run that walks thousands of rows through a live agent and must not restart from zero when it dies 40 minutes in.

:func:resumable covers that case with a different model. Work is a list of items with caller-declared keys; each completed item is appended to a JSONL file the moment it finishes; a rerun subtracts the recorded keys from the work list and hands back only what is left. The checkpoint is crash state, not a cache: it is deleted as soon as the batch completes, so a finished run never serves its results to the next one.

.map(concurrency=N) owns the batching half of this problem. This owns the durability half only — it does not run anything, so the caller keeps whatever execution model they already have (a plain loop, asyncio.gather, .map over pending).

Example: >>> from latent.prefect import resumable, task >>> >>> @task("score_rows", input="eval_set", output="scores") ... async def score_rows(eval_set: list[dict]) -> list[dict]: ... with resumable("score_rows", eval_set, key=lambda r: r["id"]) as run: ... for row in run.pending: ... run.record(row, await agent.invoke(row["prompt"])) ... return run.results()

One writer per checkpoint name, enforced with an exclusive advisory lock for the life of the block — a second holder is refused rather than allowed to interleave. Threads and coroutines inside that process are safe. The lock is flock, so the guarantee is only as strong as the filesystem underneath it: a mount that emulates flock per client (NFS with local_lock, some SMB configurations) lets two hosts both acquire, so keep the workspace — and with it .latent/checkpoints — on local disk.

Classes

ResumableBatch

ResumableBatch(name: str, items: Sequence[T], key: Callable[[T], str], path: Path, recorded: dict[str, Any], handle: TextIO)

The work list of one :func:resumable block.

Built by :func:resumable; not constructed directly. Results are typed Any on purpose — a resumed result comes back through JSON, so the batch cannot promise the type the producer returned.

Functions

resumable

resumable(name: str, items: Sequence[T], key: Callable[[T], str]) -> Iterator[ResumableBatch[T]]

Run a batch that survives a crash and resumes where it stopped.

Unlike @checkpoint this is active in every environment — resuming a long production run is the whole point. It stays safe there because the checkpoint is deleted the moment the batch completes, so it can only ever replay work from a run that died.

Args: name: Checkpoint identity, unique per flow. Names the file under .latent/checkpoints/. items: The full work list, including items already done. Editing it between runs is safe: recorded results are subtracted per key, so an appended item is computed and a removed one is reported and ignored. key: Item identity. Two runs must agree on it, and anything that invalidates a recorded result belongs in it — a key that ignores a changed field will resume onto stale results.

Yields: The batch: iterate pending, record() each result, then read results().

Example: >>> with resumable("score", rows, key=lambda r: r["id"]) as run: ... for row in run.pending: ... run.record(row, score(row)) >>> scored = run.results()

Methods

ResumableBatch.pending

Items with no recorded result, in input order.

ResumableBatch.record

record(item: T, result: Any) -> None

Durably record result for item before moving to the next one.

The line is flushed on the way out, so everything recorded before a crash or a kill survives it. Machine-level power loss can still lose the tail: that costs a rerun of those items, never a wrong result.

ResumableBatch.recorded_count

How many of this batch's items already have a result.

ResumableBatch.results

results() -> list[Any]

Every item's result, in input order.

Raises: RuntimeError: If any item was never recorded.