TT Lab
Get started
Learn Learning paths Courses

The snack machine died before ACK

Two Snack Couriers and Never-Ending Retries

Continue in TT Lab

Goal

Even when two festival snack dispatchers compete and one of them goes down right after sending, recover the remaining work with a finite budget. Keep a stale worker from overwriting the new worker's completion record.

Why it matters

A network timeout does not mean the work was not processed. Even after the lease ends, the earlier request can keep running. This time you distinguish the queue's claim token from the receiving inbox, and you design the send outside the lock, the failure budget, and the quarantine record together. You extend the duplicate-effect prevention of the earlier outbox lab to a multi-dispatcher situation.

It is estimated at 110 minutes. It is longer than the default session, so press +time to extend it before it expires. Files disappear after the session ends. Keep any code you need separately. You start knowing Python exception handling, SQLite transactions, and the outbox and receiver deduplication from the earlier modules.

Data contract

The deliverable is /root/lease/worker.py. The checker creates and cleans up the queue file and the HTTP receiving server temporarily, so do not hardcode a path or a port. Every queue is a trusted local file with a normal schema. Use the schema below.

CREATE TABLE jobs (seq INTEGER PRIMARY KEY AUTOINCREMENT,
  id TEXT NOT NULL UNIQUE, qty INTEGER NOT NULL,
  status TEXT NOT NULL CHECK(status IN ('pending','leased','done','dead')),
  attempts INTEGER NOT NULL, max_attempts INTEGER NOT NULL, token INTEGER NOT NULL,
  available_at INTEGER NOT NULL, deadline INTEGER NOT NULL,
  owner TEXT, lease_until INTEGER, reason TEXT);
CREATE TABLE redrives (seq INTEGER PRIMARY KEY AUTOINCREMENT,
  id TEXT NOT NULL, token INTEGER NOT NULL, at_ms INTEGER NOT NULL, note TEXT NOT NULL);

An identifier is an exact str of 1–64 characters of ASCII letters, digits, underscores, and hyphens. now_ms and the clock return value are an exact int from 0 to (10**15-60000) on a common time baseline, and every integer contract rejects bool. Each DB function rejects a format or range error as a ValueError before changing anything and does not modify the input. If the second clock of run_once is invalid or goes backward, the claim already committed is preserved and the error is propagated. It is not a consensus algorithm that synchronizes clocks between hosts.

The borrowed con has no transaction at the start of the call, and the function does not close it. After success or failure, leave no open transaction. Tie the write selection and the change together with BEGIN IMMEDIATE, roll back everything on a fault error before the commit and propagate the original error. An error after the commit preserves the state that was already finalized. If fault=None, it is not called, and the hook is called only on the path that performed that change. There is no hook on no-change return paths such as None or False. claim commits the quarantine update it already performed even when there is no target.

Each work item has independent meaning. This is not a strictly ordered queue in which the failure of earlier work blocks later work. token is a per-work monotonically increasing claim number, within 10**15 in the checked range. The queue's marker check does not replace an external server's permission check or fencing of an external resource. The receiving server is provided so that it removes duplicate effects for the same ID and the same quantity. redrive assumes it is called by a privileged operator, and no login or role management feature is implemented.

Steps

  1. Calculate the retry delay cap and jitter — In worker.py, define the Exception subclasses Conflict, Retryable, Permanent, and BadAck and implement retry_delay(attempt,base_ms,cap_ms,jitter). attempt is an exact int from 1–16, base_ms is 1–60000, and cap_ms is base_ms–60000. jitter is an exact int/float excluding bool and is a finite 0–1. Return floor(min(cap_ms,base_ms*2**(attempt-1))*jitter). Invalid input is a ValueError. Calculate with the given ratio, without internal random numbers or sleep.
  2. Open a queue that survives a restart — open_queue(path) creates the jobs and redrives tables below in one transaction only if they do not exist and returns a sqlite3.Connection. Use isolation_level=None, timeout=1 second, journal_mode=DELETE, and synchronous=FULL. It preserves existing work and audit rows, and on an initialization failure it closes the connection. It takes a disposable local file, not an in-memory DB.
  3. Do not reset the budget on a duplicate intake — enqueue(con,event,now_ms,ttl_ms,max_attempts=3) takes an exact dict with only id and qty. id is the identifier below, qty is an exact int 1–1000, ttl_ms is 1–60000, and max_attempts is 1–16. It inserts new work as pending, attempts=0, token=0, available_at=now_ms, deadline=now_ms+ttl_ms, and owner, lease_until, and reason=NULL, and returns True. The same ID and the same quantity preserves the entire existing row even with a different input budget and returns False, and a different quantity is a Conflict. Validate all inputs first and handle the selection and insert in a single write transaction.
  4. Keep two dispatchers from taking the same work at the same time — Implement claim(con,owner,now_ms,lease_ms=1000,fault=None). owner is an identifier and lease_ms is an exact int 1–60000. Inside BEGIN IMMEDIATE, among the pending or expired leased rows, move the rows with deadline<=now_ms or attempts>=max_attempts to dead. The reason is deadline for a missed deadline and otherwise exhausted, and it sets owner and lease_until to NULL. Among the remaining rows, select one piece of work in ascending seq that is pending with available_at<=now_ms or leased with lease_until<=now_ms. If there is none, it is None. If there is one, raise attempts and token by 1 each and save leased, owner, lease_until=min(now_ms+lease_ms,deadline), and reason=NULL. Return a dict with only id, qty, owner, token, attempt, and lease_until. attempt is the updated attempts. After the update call fault('after-claim'), and after COMMIT call fault('after-commit').
  5. Reject the completion of a stale dispatcher — Implement finish(con,ticket,outcome,now_ms,delay_ms=0,fault=None). ticket is a dict with only the 6 keys of claim, where id and owner are identifiers and qty=1–1000, token=1–1015, attempt=1–16, and lease_until=1–1015 are exact ints. outcome is a str among ok, retry, and permanent, and delay_ms is an exact int 0–60000. In one write transaction, if the work does not exist, is not leased, the marker's quantity, owner, token, attempt, or expiry differs, or now_ms>=lease_until or the deadline, it returns False with no change. A valid ok is done with a NULL reason, and permanent is dead/permanent. For retry, if it reached the maximum attempts it is dead/exhausted, otherwise if now_ms+delay_ms>=deadline it is dead/deadline, and the rest is pending/retry. Only when pending is available_at=now_ms+delay_ms, otherwise now_ms, and it sets owner and lease_until to NULL. It preserves the other fields and returns the new status string. After the update call fault('after-finish'), and after COMMIT call fault('after-commit').
  6. Leave the quarantine release and the audit record together — redrive(con,job_id,now_ms,ttl_ms,note,fault=None) redrives only existing dead work and returns True. job_id and note are identifiers and ttl_ms is an exact int 1–60000. If the target does not exist or is not dead, it is a ValueError. In one write transaction, the order is: insert id, the current token, at_ms=now_ms, and note into redrives, then fault('after-audit'), then change the work to pending, attempts=0, available_at=now_ms, deadline=now_ms+ttl_ms, and owner, lease_until, and reason=NULL, then fault('after-redrive'), then COMMIT, then fault('after-commit'). It preserves qty, max_attempts, and token.
  7. Perform the network send outside the write lock — run_once(con,owner,clock,send,jitter,lease_ms=1000) validates jitter first and claims with the first time from clock(). If there is no work, it is None and does not call send. If there is, it passes a new dict with only id and qty to send once, outside a transaction. Only an exact True is ok, and any other return value is a BadAck. Only Retryable is classified as retry and uses retry_delay(ticket.attempt,100,1000,jitter), and Permanent is classified as permanent. Other errors and BadAck are propagated as they are and preserve the leased state. After a normal result or a classified error, call clock() again, and if the second value is smaller than the first, it is a ValueError. After passing finish the second time, the outcome, and the delay, return an id, token, status dict. If finish is False, status is stale. Even if the send callback changes the dict it received or puts in work over a separate connection, the current marker must be kept.
  8. Have another process recover the work that died right after the send — run_file(path,owner,clock,send,jitter,lease_ms=1000) rejects with FileNotFoundError if the existing ordinary queue file does not exist and does not create a new empty queue. It owns the connection with open_queue, returns the run_once result, and closes the connection on both success and failure. The final check verifies that only one of the simultaneous claims of two real child processes succeeds and that the state after a forced termination before and after the commit is atomic. Right after a separate HTTP receiving server commits the quantity, it terminates the first dispatcher and resends from the second dispatcher's process. The HTTP request arrives twice, but the inbox and the receiving effect must be once. The checker provides the receiving server and the temporary DB.

Notes

Calculate the retry delay cap and jitter

In worker.py, define the Exception subclasses Conflict, Retryable, Permanent, and BadAck and implement retry_delay(attempt,base_ms,cap_ms,jitter). attempt is an exact int from 1–16, base_ms is 1–60000, and cap_ms is base_ms–60000. jitter is an exact int/float excluding bool and is a finite 0–1. Return floor(min(cap_ms,base_ms*2**(attempt-1))*jitter). Invalid input is a ValueError. Calculate with the given ratio, without internal random numbers or sleep.

In Python, bool is a subtype of int. Apply the cap first, multiply by the injected ratio, and then round down.

Open a queue that survives a restart

open_queue(path) creates the jobs and redrives tables below in one transaction only if they do not exist and returns a sqlite3.Connection. Use isolation_level=None, timeout=1 second, journal_mode=DELETE, and synchronous=FULL. It preserves existing work and audit rows, and on an initialization failure it closes the connection. It takes a disposable local file, not an in-memory DB.

Coordinate the workers only with short write transactions. Do not confuse creating the tables with resetting the existing contents.

Do not reset the budget on a duplicate intake

enqueue(con,event,now_ms,ttl_ms,max_attempts=3) takes an exact dict with only id and qty. id is the identifier below, qty is an exact int 1–1000, ttl_ms is 1–60000, and max_attempts is 1–16. It inserts new work as pending, attempts=0, token=0, available_at=now_ms, deadline=now_ms+ttl_ms, and owner, lease_until, and reason=NULL, and returns True. The same ID and the same quantity preserves the entire existing row even with a different input budget and returns False, and a different quantity is a Conflict. Validate all inputs first and handle the selection and insert in a single write transaction.

A retried intake is not new work. Compare the UNIQUE ID and the existing quantity, and do not change the input dict either.

Keep two dispatchers from taking the same work at the same time

Implement claim(con,owner,now_ms,lease_ms=1000,fault=None). owner is an identifier and lease_ms is an exact int 1–60000. Inside BEGIN IMMEDIATE, among the pending or expired leased rows, move the rows with deadline<=now_ms or attempts>=max_attempts to dead. The reason is deadline for a missed deadline and otherwise exhausted, and it sets owner and lease_until to NULL. Among the remaining rows, select one piece of work in ascending seq that is pending with available_at<=now_ms or leased with lease_until<=now_ms. If there is none, it is None. If there is one, raise attempts and token by 1 each and save leased, owner, lease_until=min(now_ms+lease_ms,deadline), and reason=NULL. Return a dict with only id, qty, owner, token, attempt, and lease_until. attempt is the updated attempts. After the update call fault('after-claim'), and after COMMIT call fault('after-commit').

Do not release the lock between the selection and the update. Distinguish the equals sign at the expiry boundary from the quarantine change that must be finalized even when there is nothing to select.

Reject the completion of a stale dispatcher

Implement finish(con,ticket,outcome,now_ms,delay_ms=0,fault=None). ticket is a dict with only the 6 keys of claim, where id and owner are identifiers and qty=1–1000, token=1–1015, attempt=1–16, and lease_until=1–1015 are exact ints. outcome is a str among ok, retry, and permanent, and delay_ms is an exact int 0–60000. In one write transaction, if the work does not exist, is not leased, the marker's quantity, owner, token, attempt, or expiry differs, or now_ms>=lease_until or the deadline, it returns False with no change. A valid ok is done with a NULL reason, and permanent is dead/permanent. For retry, if it reached the maximum attempts it is dead/exhausted, otherwise if now_ms+delay_ms>=deadline it is dead/deadline, and the rest is pending/retry. Only when pending is available_at=now_ms+delay_ms, otherwise now_ms, and it sets owner and lease_until to NULL. It preserves the other fields and returns the new status string. After the update call fault('after-finish'), and after COMMIT call fault('after-commit').

Even a matching owner alone can be a stale run. Check the whole marker and the time boundary, and do not delete the current work just because the result was late.

Leave the quarantine release and the audit record together

redrive(con,job_id,now_ms,ttl_ms,note,fault=None) redrives only existing dead work and returns True. job_id and note are identifiers and ttl_ms is an exact int 1–60000. If the target does not exist or is not dead, it is a ValueError. In one write transaction, the order is: insert id, the current token, at_ms=now_ms, and note into redrives, then fault('after-audit'), then change the work to pending, attempts=0, available_at=now_ms, deadline=now_ms+ttl_ms, and owner, lease_until, and reason=NULL, then fault('after-redrive'), then COMMIT, then fault('after-commit'). It preserves qty, max_attempts, and token.

Give a fresh attempt budget, but do not reuse the claim generation. Prevent a half redrive in which only the audit row remains or only the state is released.

Perform the network send outside the write lock

run_once(con,owner,clock,send,jitter,lease_ms=1000) validates jitter first and claims with the first time from clock(). If there is no work, it is None and does not call send. If there is, it passes a new dict with only id and qty to send once, outside a transaction. Only an exact True is ok, and any other return value is a BadAck. Only Retryable is classified as retry and uses retry_delay(ticket.attempt,100,1000,jitter), and Permanent is classified as permanent. Other errors and BadAck are propagated as they are and preserve the leased state. After a normal result or a classified error, call clock() again, and if the second value is smaller than the first, it is a ValueError. After passing finish the second time, the outcome, and the delay, return an id, token, status dict. If finish is False, status is stale. Even if the send callback changes the dict it received or puts in work over a separate connection, the current marker must be kept.

Do not hold a DB transaction while waiting for an external response. Judge validity again with the second time, and do not hide an unknown failure as a success.

Have another process recover the work that died right after the send

run_file(path,owner,clock,send,jitter,lease_ms=1000) rejects with FileNotFoundError if the existing ordinary queue file does not exist and does not create a new empty queue. It owns the connection with open_queue, returns the run_once result, and closes the connection on both success and failure. The final check verifies that only one of the simultaneous claims of two real child processes succeeds and that the state after a forced termination before and after the commit is atomic. Right after a separate HTTP receiving server commits the quantity, it terminates the first dispatcher and resends from the second dispatcher's process. The HTTP request arrives twice, but the inbox and the receiving effect must be once. The checker provides the receiving server and the temporary DB.

The next process reopens the same file. Do not claim to have closed the gap between the external receiving effect and the queue completion; recover from the duplicate send.