The snack machine died before ACK
Recover the snack machine that died before ACK
Goal
Implement a persistent cursor, atomic processing, and a replay range, and prevent duplicate processing with a real TCP reconnect.
Why it matters
Connection recovery and business recovery are different. You inspect a receiver that dies between saving and acknowledging to learn the boundary between the two. Python functions, exceptions, async/await, SQL basics, and the earlier realtime subscriber course are recommended prerequisites. No installation or internet is needed. This is a 75-minute lab, so extend the default 60-minute session with the +time button. All files disappear when the session ends, so keep any code you need separately.
Steps
- Do not accept a torn last line as a command — Implement parse_line(line). It accepts only bytes and returns an ASCII EVENT seq delta LF of at most 64 bytes as a (seq, delta) tuple. seq is 0–2147483647 and delta is -1000–1000. Leading zeros other than the number 0, a + sign, -0, CRLF, a missing newline, multiple lines, and non-ASCII are a ValueError.
- Open a store that remembers across restarts — open_store(path) returns a sqlite3 connection. Create the tables checkpoint(id INTEGER PRIMARY KEY CHECK(id=1), last INTEGER NOT NULL, total INTEGER NOT NULL) and ledger(seq INTEGER PRIMARY KEY, delta INTEGER NOT NULL) only if they do not exist, and insert the initial checkpoint row (1,-1,0) only if it does not exist. state(con) is the (last,total) tuple. Connect with isolation_level=None and a busy timeout of 1 second so that you can use explicit SQL transactions. On an initialization failure, close the connection and propagate the error.
- Save the total and the cursor together or cancel them together — Declare the Exception subclasses Gap and Conflict and implement apply_event(con,seq,delta,fault=None). Allow only an int in the step 1 range excluding bool; anything else is a ValueError. For seq=last+1, commit in one transaction the ledger insert, then the total update, then fault() if fault is given, then the last update, and return True. For seq<=last, it is False if the ledger has the same delta and a Conflict if the content differs. A larger sequence number is a Gap. On failure, roll back all changes, propagate the original error, and leave no transaction behind.
- Do not hide a truncated record as a successful recovery — Implement the Exception subclass ResyncRequired and replay(events,last,limit). events is a list of 1–128 (seq,delta) tuples whose sequence numbers are all consecutive and ascending, and whose values are ints in the step 1 range. last is an int from -1–2147483647 and limit is an int from 1–16. A bool, a format or range error, an empty log, or last>the final sequence number is a ValueError. lastlast as a new list. Validate the whole input first and do not modify it.
- Produce the business ACK only after the commit — consume_line(con,line,fault=None) calls parse_line and apply_event and, on success, returns b"ACK seq\n". Pass fault to apply_event. It ACKs a duplicate with the same content too, but for parsing, gap, conflict, or business failures it propagates the error and does not ACK.
- Close the connection on errors and cancellation too — Implement async consume_session(reader,writer,con,timeout,before_ack=None). timeout is a positive finite int/float excluding bool, and anything else is a ValueError. Limit each readline wait with timeout, and on EOF return state. After committing one line with consume_line, if before_ack is given, call before_ack(seq), then perform writer.write(ack) and a drain limited by timeout. Propagate errors and cancellation, and on normal, failure, and cancellation alike reclaim with close then wait_closed limited by timeout. On a TimeoutError or ConnectionError while waiting for the close, call writer.transport.abort and do not hide the existing error. Do not close the DB connection. For an invalid timeout, the close wait uses 1 second.
- Kill the receiver and pick up from the same DB — Implement async run_client(host,port,path,timeout,before_ack=None). Validate timeout first, then open in the order open_store(path), then asyncio.open_connection(host,port,limit=64) limited by timeout. Write the DB's last as b"RESUME last\n", and after a drain limited by timeout, pass it to consume_session and return its return value. Close the DB connection on every path, and for a failure before passing to consume_session, reclaim the socket yourself. Propagate errors and cancellation. The checker creates a temporary DB and a loopback server, terminates the first process between commit 1 and ACK 1, and then reconnects with a new process. It must not apply the resent 1 twice.
Notes
Implement it in /root/resume/client.py after mkdir -p /root/resume. The checker also runs the earlier steps, and it prepares and reclaims the temporary DB, the loopback port, and the receiving process. You do not need to keep a server or DB file running. Do not modify the submission file. Each step has a 15-second run limit, and this does not limit your study time. The last check terminates a real process, so do not assume that the cleanup finally runs. This lab does not implement production TLS, authentication, automatic backoff, snapshot recovery, or distributed exactly-once.
Do not accept a torn last line as a command
Implement parse_line(line). It accepts only bytes and returns an ASCII EVENT seq delta LF of at most 64 bytes as a (seq, delta) tuple. seq is 0–2147483647 and delta is -1000–1000. Leading zeros other than the number 0, a + sign, -0, CRLF, a missing newline, multiple lines, and non-ASCII are a ValueError.
Before you remove the last newline with strip, check that it is one complete frame.
Open a store that remembers across restarts
open_store(path) returns a sqlite3 connection. Create the tables checkpoint(id INTEGER PRIMARY KEY CHECK(id=1), last INTEGER NOT NULL, total INTEGER NOT NULL) and ledger(seq INTEGER PRIMARY KEY, delta INTEGER NOT NULL) only if they do not exist, and insert the initial checkpoint row (1,-1,0) only if it does not exist. state(con) is the (last,total) tuple. Connect with isolation_level=None and a busy timeout of 1 second so that you can use explicit SQL transactions. On an initialization failure, close the connection and propagate the error.
If you reset a cursor you already saved, a reconnect becomes duplicate processing from the beginning.
Save the total and the cursor together or cancel them together
Declare the Exception subclasses Gap and Conflict and implement apply_event(con,seq,delta,fault=None). Allow only an int in the step 1 range excluding bool; anything else is a ValueError. For seq=last+1, commit in one transaction the ledger insert, then the total update, then fault() if fault is given, then the last update, and return True. For seq<=last, it is False if the ledger has the same delta and a Conflict if the content differs. A larger sequence number is a Gap. On failure, roll back all changes, propagate the original error, and leave no transaction behind.
Whether you commit the cursor first or commit the total separately first, the two get out of step at a failure.
Do not hide a truncated record as a successful recovery
Implement the Exception subclass ResyncRequired and replay(events,last,limit). events is a list of 1–128 (seq,delta) tuples whose sequence numbers are all consecutive and ascending, and whose values are ints in the step 1 range. last is an int from -1–2147483647 and limit is an int from 1–16. A bool, a format or range error, an empty log, or last>the final sequence number is a ValueError. lastlast as a new list. Validate the whole input first and do not modify it.
A retention start of 10 and a cursor of 9 connect, but for a cursor of 8, the 9 is missing.
Produce the business ACK only after the commit
consume_line(con,line,fault=None) calls parse_line and apply_event and, on success, returns b"ACK seq\n". Pass fault to apply_event. It ACKs a duplicate with the same content too, but for parsing, gap, conflict, or business failures it propagates the error and does not ACK.
Right after the ACK, the processing result must be visible from another DB connection too.
Close the connection on errors and cancellation too
Implement async consume_session(reader,writer,con,timeout,before_ack=None). timeout is a positive finite int/float excluding bool, and anything else is a ValueError. Limit each readline wait with timeout, and on EOF return state. After committing one line with consume_line, if before_ack is given, call before_ack(seq), then perform writer.write(ack) and a drain limited by timeout. Propagate errors and cancellation, and on normal, failure, and cancellation alike reclaim with close then wait_closed limited by timeout. On a TimeoutError or ConnectionError while waiting for the close, call writer.transport.abort and do not hide the existing error. Do not close the DB connection. For an invalid timeout, the close wait uses 1 second.
The hook right before the ACK reproduces the failure window in which the save is finished but the publisher does not know.
Kill the receiver and pick up from the same DB
Implement async run_client(host,port,path,timeout,before_ack=None). Validate timeout first, then open in the order open_store(path), then asyncio.open_connection(host,port,limit=64) limited by timeout. Write the DB's last as b"RESUME last\n", and after a drain limited by timeout, pass it to consume_session and return its return value. Close the DB connection on every path, and for a failure before passing to consume_session, reclaim the socket yourself. Propagate errors and cancellation. The checker creates a temporary DB and a loopback server, terminates the first process between commit 1 and ACK 1, and then reconnects with a new process. It must not apply the resent 1 twice.
If you always send last as -1 or wipe the store every time, it shows up in the new process's actual request.