outbox
Transactional outbox over std.sql: write your business row and the message you owe in the same transaction, then relay it with at-least-once delivery, retry and backoff.
ecko get github.com/ecko-lang/outbox
import outbox
Pure computation: it declares no capabilities, so it cannot touch the network, the filesystem or the environment.
Version 0.5.0 - source - MIT.
init(db)
init(db) -> null. Creates the outbox table (and its relay index) if they do not already exist. Safe to call every time your program starts.
add(db, topic, payload, opts = empty_map())
add(db, topic, payload, opts?) -> the new message's id.
Call this inside the same sql.transaction as your business write - the row you are recording and the message you owe both commit together, or neither does. payload is any JSON-serializable value.
opts: id (a caller-supplied id, default a fresh uuid.v7()), max_attempts (default 5 - how many delivery attempts relay makes before dead-lettering this message), now (Unix ms, default time.now() - inject a fixed value in tests so behavior does not depend on wall time).
relay(db, deliver_fn, opts = empty_map())
relay(db, deliver_fn, opts?) -> { delivered, failed, dead } counts.
Reads the pending messages whose next_attempt_at has arrived, oldest first, and calls deliver_fn(msg) for each - msg has id, topic, payload, attempts, max_attempts, created_at, next_attempt_at, delivered_at and last_error. deliver_fn's return value is ignored; it signals failure by throwing.
A message that returns normally is marked delivered. A message that throws has its attempt recorded and is rescheduled with exponential backoff (base_backoff_ms * 2^(attempts - 1), capped at max_backoff_ms) until it reaches the max_attempts set on add, at which point it is marked dead instead of rescheduled. Each deliver_fn call is individually caught, so one bad message is recorded and skipped rather than blocking the batch or stopping the relay.
Delivery is at-least-once: a message already marked delivered is never picked up again, but a crash between a successful deliver_fn call and that update means the same message is retried on the next relay. Make deliver_fn idempotent - dedupe on msg.id.
opts: batch_size (default 100), base_backoff_ms (default 1000), max_backoff_ms (default 60000), now (Unix ms, default time.now()).
pending(db)
pending(db) -> messages awaiting delivery, oldest first.
Includes messages that already failed once and are backing off, not only ones never attempted - compare next_attempt_at if you need to tell them apart.
dead(db)
dead(db) -> messages that reached max_attempts and will not be retried automatically, oldest first. Use retry_dead to give one another chance.
retry_dead(db, id, opts = empty_map())
retry_dead(db, id, opts?) -> null. Moves a dead message back to pending with its attempt count reset, so the next relay picks it up immediately.
Throws { kind: "not_found", message: ... } if id names no dead message.
opts: now (Unix ms, default time.now()).
prune(db, older_than)
prune(db, older_than) -> the number of rows removed.
Deletes delivered and dead messages last touched before older_than (Unix ms). Pending messages are never pruned, however old - they are still waiting on a real outcome.