checkpoint

Resumable long-running jobs over std.sql: batch a source through a handler with the cursor committed in the same transaction as the writes, and memoise named pipeline steps. A crashed job resumes where it stopped and applies nothing twice.

ecko get github.com/ecko-lang/checkpoint
import checkpoint

Pure computation: it declares no capabilities, so it cannot touch the network, the filesystem or the environment.

Version 0.5.0 - source - MIT.


each_batch(db, name, source, limit, handler)

each_batch(db, name, source, limit, handler) -> { batches, items, done, cursor }. Runs the job name to completion, one batch per transaction, resuming from the last committed cursor. Each round calls source(cursor, limit), which must return { items: [...], next: cursor } (the cursor is null on the first call), then handler(items, tx), then stores next - the handler's writes and the new cursor commit together or not at all. If the handler throws, that batch rolls back, the error re-raises, and the next call resumes at the same batch. The run ends when the source returns no items, or returns items with next: null (a last page); the job is then marked done and later calls return at once without calling the source. batches and items count this call's work only; status has the totals. Write through tx (the same handle as db) and do not open another transaction on it inside the handler.

step(db, name, f)

step(db, name, f) -> the value f() returned. Runs one named step of a pipeline once. The first successful run stores the result and later calls return the stored value without calling f. f runs inside a transaction on db together with the write that records it, so writes f makes through db commit only if the step is recorded; if f throws, nothing is stored, its writes roll back and the error re-raises. The result must survive a JSON round trip unchanged (null, bool, int, float, string, and lists and maps of those); anything else is refused. Both calls return the decoded form, so the first run and a resumed run see the same value. f must not open its own transaction on db.

status(db, name)

status(db, name) -> map or null. The stored state of a job or step, or null if name has never run (or was reset). Fields: name, kind ("batch" or "step"), done, cursor (the last committed cursor of a batch job), value (a step's stored result), batches and items (totals across every run), and started_at, updated_at, finished_at in milliseconds since the Unix epoch.

reset(db, name)

reset(db, name) -> Bool. Forgets a job or step, so the next each_batch starts from the beginning and the next step runs its function again. Returns whether there was anything to forget. It does not undo the writes the job made; clear those yourself first if a rerun would duplicate them.

sql_source(db, table, key, where = null, columns = null)

sql_source(db, table, key, where = null, columns = null) -> source function. A keyset-paginated source over table for each_batch: each batch is the next limit rows with key > cursor, ordered by key, and the cursor is the last row's key. It never uses OFFSET, so each batch costs the same however far in the job is, and gaps in the key are harmless. key must be unique on its own - the table's single-column primary key or the only column of a non-partial unique index - and this is checked when the source is built, because keyset over a non-unique key silently skips rows. Rows whose key is NULL are never visited. where narrows the rows with a sql { ... } block (its holes bind as parameters); columns is a list of column names, and the key is always included. Items are row maps.