TT Lab
Get started
Learn Learning paths Courses

The snack machine died before ACK

The Weekend Display and the Missing Snack Log

Continue in TT Lab

Goal

Recover the snack stock display board that was off over the weekend. When the log's retention range has ended, install a snapshot from the same point in time and follow from the next event.

Why it matters

Even when the connection recovers, events that were already deleted do not come back. If the stock and the cursor point to different points in time, you display a wrong number even when the replay succeeds. This time you implement the retention boundary, the epoch, same-point-in-time reads, and atomic install yourself and verify them with the on-disk state after a real termination.

It is estimated at 100 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 the Python exceptions, SQLite transactions, cursors, and outbox recovery from the earlier modules.

Data contract

The deliverable is a single file, /root/snapshot/worker.py. The source and the replica are separate disposable local SQLite files with the same schema below. The checker creates and cleans them up, so do not hardcode a DB path. All of them are trusted local files with a normal schema, and the con lent to a function has no open transaction at the start. Functions other than open_store and sync_files do not close the borrowed connection and leave no transaction after success or failure.

CREATE TABLE checkpoint (id INTEGER PRIMARY KEY CHECK(id=1),
    epoch TEXT NOT NULL, last INTEGER NOT NULL, floor INTEGER NOT NULL);
CREATE TABLE stock (id INTEGER PRIMARY KEY CHECK(id=1), total INTEGER NOT NULL);
CREATE TABLE events (seq INTEGER PRIMARY KEY, delta INTEGER NOT NULL);

The initial checkpoint row is (1, the epoch you pass in, -1, 0) and the initial stock row is (1,0). Within the same epoch, numbers are not reused, and the retained events are consecutive from floor to last. If everything was truncated, floor=last+1. total is the business state that adds each delta starting from an initial 0. expected_epoch comes from a trusted connection setting, and you do not extract it from the snapshot itself to use as the expected value.

The lab's writes run one at a time, sequentially. You check point-in-time isolation by having another connection commit a write between read transactions. It is not a test of concurrent checkpoints in separate threads or of multi-writer performance. Set wal_autocheckpoint=0 on every connection, and close the connections in order after writing is done. Do not use this as is as a WAL management policy for long-term operation.

Steps

  1. Store the epoch and the recovery state — In worker.py, define the Exception subclasses Conflict, ResyncRequired, Gap, and CoveredBySnapshot and implement validate_snapshot(snapshot) and open_store(path,epoch). A snapshot is a dict with only version, epoch, last, and total. version is exactly the int 1, epoch is 1–64 characters of ASCII letters, digits, underscores, and hyphens, last is an int from -1 to 2147483647 excluding bool, and total is an int excluding bool whose absolute value is at most 1000*(last+1). If valid, return a new dict copy; otherwise it is a ValueError. After validating the epoch, open_store atomically creates the three tables below and the initial rows only if they do not exist. Open with isolation_level=None, timeout=1, WAL, synchronous=FULL, and wal_autocheckpoint=0 and return a sqlite3.Connection. It does not overwrite existing content, and on an initialization failure it closes the connection.
  2. Finalize the stock and the next number together — append(con,delta,fault=None) adds a change of an int from -1000 to 1000 excluding bool. Inside BEGIN IMMEDIATE it checks the next seq=last+1 and proceeds in the order events insert, then fault('after-event'), then adding delta to stock.total and updating checkpoint.last, then fault('after-state'), then COMMIT, then fault('after-commit'). seq is 0 to 2147483647, and anything outside the range is a ValueError. It calls fault when given and returns seq. An error before the commit rolls back everything, an error after the commit preserves the finalized state, and the original error is propagated.
  3. Export a same-point-in-time snapshot — export_snapshot(con,between=None) returns a version=1 snapshot in a BEGIN read transaction in the order query epoch and last, then between(), then query total, then COMMIT. It calls the hook only once and only when given, and another connection's write must be able to commit there. On failure it cleans up the read transaction and propagates the original error. Do not use a COMMIT or BEGIN IMMEDIATE between the two SELECTs.
  4. Compare the retention range with the replay cursor — trim(con,through) takes an int from -1 to 2147483647 excluding bool and checks in one write transaction whether floor-1 <= through <= last. If not, it is a ValueError; if so, it deletes seq<=through, updates only floor=through+1, and returns None. replay(con,epoch,last,limit=16) validates the epoch format, last from -1 to 2147483647 excluding bool, and limit from 1 to 16. In the same read transaction, a different epoch is ResyncRequired, a cursor ahead of the server is a ValueError, and lastlast in number order, at most limit of them. Distinguish an empty result from a missing error.
  5. Install the snapshot in one go — install_snapshot(con,snapshot,expected_epoch,fault=None) validates the format of snapshot and expected_epoch. If snapshot.epoch differs from the expected epoch, it is ResyncRequired. In one write transaction, a smaller last of the same epoch, or the same last with a different total, is a Conflict, and the same last and total is False with no change. Otherwise it returns True after deleting all events, then fault('after-clear'), then replacing total and updating the checkpoint's epoch, last, and floor=last+1, then fault('after-install'), then COMMIT, then fault('after-commit'). An error before the commit rolls back everything, and an error after the commit preserves the finalized state. A different allowed epoch can be replaced even with a smaller number.
  6. Apply verifiable duplicates and a new batch — apply_batch(con,epoch,events,fault=None) validates epoch and events first. events is a list of at most 16, and an empty list is allowed. Each element is a tuple/list of two, seq and delta, where seq is an int from 0 to 2147483647 excluding bool and delta is an int from -1000 to 1000 excluding bool. A format error is a ValueError, and if the numbers inside the batch do not increase by 1 it is a Gap. In one write transaction, an epoch mismatch is ResyncRequired. For seq<=last, the same retained delta has no effect, a different delta is a Conflict, and if the individual record is gone it is CoveredBySnapshot. A new seq allows only last+1, otherwise it is a Gap. Each time it updates a new event, the stock, and the cursor, it calls fault('after-one'), and after the whole COMMIT it calls fault('after-commit'). It returns the number of new effects this time and does not modify the input. An error before the commit rolls back the whole batch, and later errors preserve the finalized state.
  7. End a single recovery attempt in a finite way — sync_once(source,replica,expected_epoch,after_snapshot=None) validates the expected epoch and raises ResyncRequired if it differs from the source's actual epoch. It requests one replay with the replica's epoch and last. Only on ResyncRequired does it recover with export_snapshot(source), then install_snapshot(replica,...,expected_epoch), then after_snapshot(), then a replay after the snapshot's last. It calls the hook only when given. A batch is at most 16, and after apply_batch it returns export_snapshot(replica). If the second replay is truncated again, it propagates the error as is and keeps the valid installed state. It propagates other errors too, and does no internal infinite retry.
  8. Recover into a separate file after a real termination — sync_files(source_path,replica_path,expected_epoch) validates the expected epoch and raises FileNotFoundError if the local source file does not exist. If the realpath is the same or the two existing paths are samefile, it is a ValueError before opening. It opens the separate source and replica with open_store, returns the result of running sync_once once, and closes the connections it owns on every path. Even if it cannot open the replica, it closes the source. The checker ends a real child process with exit code 73 before and after the commit of append, install, and batch and checks the on-disk state with a new connection.

Notes

Store the epoch and the recovery state

In worker.py, define the Exception subclasses Conflict, ResyncRequired, Gap, and CoveredBySnapshot and implement validate_snapshot(snapshot) and open_store(path,epoch). A snapshot is a dict with only version, epoch, last, and total. version is exactly the int 1, epoch is 1–64 characters of ASCII letters, digits, underscores, and hyphens, last is an int from -1 to 2147483647 excluding bool, and total is an int excluding bool whose absolute value is at most 1000*(last+1). If valid, return a new dict copy; otherwise it is a ValueError. After validating the epoch, open_store atomically creates the three tables below and the initial rows only if they do not exist. Open with isolation_level=None, timeout=1, WAL, synchronous=FULL, and wal_autocheckpoint=0 and return a sqlite3.Connection. It does not overwrite existing content, and on an initialization failure it closes the connection.

If last=-1, total can only be 0. Split CREATE IF NOT EXISTS from the initial row insert, and exclude bool from every int contract.

Finalize the stock and the next number together

append(con,delta,fault=None) adds a change of an int from -1000 to 1000 excluding bool. Inside BEGIN IMMEDIATE it checks the next seq=last+1 and proceeds in the order events insert, then fault('after-event'), then adding delta to stock.total and updating checkpoint.last, then fault('after-state'), then COMMIT, then fault('after-commit'). seq is 0 to 2147483647, and anything outside the range is a ValueError. It calls fault when given and returns seq. An error before the commit rolls back everything, an error after the commit preserves the finalized state, and the original error is propagated.

Even after truncating the log, the numbers continue from last. Do not use MAX(events.seq) as the basis for the next number.

Export a same-point-in-time snapshot

export_snapshot(con,between=None) returns a version=1 snapshot in a BEGIN read transaction in the order query epoch and last, then between(), then query total, then COMMIT. It calls the hook only once and only when given, and another connection's write must be able to commit there. On failure it cleans up the read transaction and propagates the original error. Do not use a COMMIT or BEGIN IMMEDIATE between the two SELECTs.

The combination of values is what gets graded, not a single value. The new stock that arose inside the hook must show up in the next export.

Compare the retention range with the replay cursor

trim(con,through) takes an int from -1 to 2147483647 excluding bool and checks in one write transaction whether floor-1 <= through <= last. If not, it is a ValueError; if so, it deletes seq<=through, updates only floor=through+1, and returns None. replay(con,epoch,last,limit=16) validates the epoch format, last from -1 to 2147483647 excluding bool, and limit from 1 to 16. In the same read transaction, a different epoch is ResyncRequired, a cursor ahead of the server is a ValueError, and lastlast in number order, at most limit of them. Distinguish an empty result from a missing error.

Even after deleting everything, last and floor must remain. A single equals sign in an inequality can include an already processed event again.

Install the snapshot in one go

install_snapshot(con,snapshot,expected_epoch,fault=None) validates the format of snapshot and expected_epoch. If snapshot.epoch differs from the expected epoch, it is ResyncRequired. In one write transaction, a smaller last of the same epoch, or the same last with a different total, is a Conflict, and the same last and total is False with no change. Otherwise it returns True after deleting all events, then fault('after-clear'), then replacing total and updating the checkpoint's epoch, last, and floor=last+1, then fault('after-install'), then COMMIT, then fault('after-commit'). An error before the commit rolls back everything, and an error after the commit preserves the finalized state. A different allowed epoch can be replaced even with a smaller number.

Do not ask the snapshot itself whether its epoch is allowed. expected_epoch is a baseline that comes from outside the thing being validated.

Apply verifiable duplicates and a new batch

apply_batch(con,epoch,events,fault=None) validates epoch and events first. events is a list of at most 16, and an empty list is allowed. Each element is a tuple/list of two, seq and delta, where seq is an int from 0 to 2147483647 excluding bool and delta is an int from -1000 to 1000 excluding bool. A format error is a ValueError, and if the numbers inside the batch do not increase by 1 it is a Gap. In one write transaction, an epoch mismatch is ResyncRequired. For seq<=last, the same retained delta has no effect, a different delta is a Conflict, and if the individual record is gone it is CoveredBySnapshot. A new seq allows only last+1, otherwise it is a Gap. Each time it updates a new event, the stock, and the cursor, it calls fault('after-one'), and after the whole COMMIT it calls fault('after-commit'). It returns the number of new effects this time and does not modify the input. An error before the commit rolls back the whole batch, and later errors preserve the finalized state.

After validating the whole input, process one small batch atomically. You cannot work out an individual past delta from a snapshot's total.

End a single recovery attempt in a finite way

sync_once(source,replica,expected_epoch,after_snapshot=None) validates the expected epoch and raises ResyncRequired if it differs from the source's actual epoch. It requests one replay with the replica's epoch and last. Only on ResyncRequired does it recover with export_snapshot(source), then install_snapshot(replica,...,expected_epoch), then after_snapshot(), then a replay after the snapshot's last. It calls the hook only when given. A batch is at most 16, and after apply_batch it returns export_snapshot(replica). If the second replay is truncated again, it propagates the error as is and keeps the valid installed state. It propagates other errors too, and does no internal infinite retry.

The first miss is a signal to choose the recovery path, and a second miss during recovery is a signal that this attempt has not finished.

Recover into a separate file after a real termination

sync_files(source_path,replica_path,expected_epoch) validates the expected epoch and raises FileNotFoundError if the local source file does not exist. If the realpath is the same or the two existing paths are samefile, it is a ValueError before opening. It opens the separate source and replica with open_store, returns the result of running sync_once once, and closes the connections it owns on every path. Even if it cannot open the replica, it closes the source. The checker ends a real child process with exit code 73 before and after the commit of append, install, and batch and checks the on-disk state with a new connection.

With sqlite3.Connection, transaction management and closing the connection are different. Be clear about who owns the open connection.