ACKの前に止まった自販機を復旧する
目標
永続カーソル・アトミックな処理・再生範囲を実装し、実際のTCP再接続で重複処理を防ぎます。
なぜ重要なのか
接続の復旧と業務の復旧は別物です。保存と確認の間に落ちる受信側を直接検査して、2つの境界を学びます。Pythonの関数・例外・async/await、SQLの基礎、前のリアルタイム購読者のコースを前提知識として推奨します。インストールやインターネットは必要ありません。75分のラボなので、基本の60分セッションでは+時間で延長してください。セッションが終了するとすべてのファイルが消えるので、必要なコードは別に保管してください。
ステップ
- 途切れた最後の行をコマンドとして受け付けない: parse_line(line)を実装してください。bytesだけを受け付け、64バイト以下のASCII「EVENT 連番 変化量 LF」の1行を、(seq, delta)のタプルで返します。seqは0–2147483647、deltaは-1000–1000です。数字の0以外の先頭の0、+符号、-0、CRLF、改行の欠落、複数行、非ASCIIはValueErrorです。
- 再起動しても覚えているストアを開く: open_store(path)はsqlite3の接続を返します。checkpoint(id INTEGER PRIMARY KEY CHECK(id=1), last INTEGER NOT NULL, total INTEGER NOT NULL)とledger(seq INTEGER PRIMARY KEY, delta INTEGER NOT NULL)の2つのテーブルを、存在しないときにだけ作成し、checkpointの初期行(1,-1,0)を、存在しないときにだけ挿入します。state(con)は(last,total)のタプルです。明示的なSQLトランザクションを使えるように、isolation_level=None、busy timeout=1秒で接続してください。初期化に失敗したときは、接続を閉じてエラーを伝えます。
- 合計とカーソルを一緒に保存するか、一緒に取り消す: Exceptionのサブクラスとして、GapとConflictを宣言し、apply_event(con,seq,delta,fault=None)を実装してください。boolを除く、ステップ1の範囲のintだけを許可し、不正ならValueErrorです。seq=last+1なら、ledgerへの挿入→totalの更新→faultがあればfault()→lastの更新を、1つのトランザクションでコミットして、Trueを返します。seq<=lastなら、ledgerの同じdeltaならFalse、内容が違えばConflictです。それより大きな連番はGapです。失敗したときは、すべての変更をロールバックして元のエラーを伝え、トランザクションを残さないでください。
- 切り詰められた記録を、復旧の成功として隠さない: Exceptionのサブクラスとして、ResyncRequiredとreplay(events,last,limit)を実装してください。eventsは1–128個の(seq,delta)のタプルからなるlistで、すべての連番が連続・昇順であり、値はステップ1の範囲のintです。lastは-1–2147483647のint、limitは1–16のintです。bool・型・範囲のエラー、空のログ、last>最後の連番はValueErrorです。last<最初の連番-1ならResyncRequired、それ以外はseq>lastの項目を最大limit個、新しいlistとして返します。入力全体を先に検査し、変更はしません。
- コミットしたあとにだけ、業務のACKを作る: consume_line(con,line,fault=None)は、parse_lineとapply_eventを呼び出し、成功後にb"ACK 連番\n"を返します。faultをapply_eventに渡してください。内容が同じ重複にもACKを返しますが、パース・ギャップ・衝突・業務の失敗ではエラーを伝え、ACKは返しません。
- エラーやキャンセルでも接続を閉じる: async consume_session(reader,writer,con,timeout,before_ack=None)を実装してください。timeoutはboolを除く正の有限なint/floatで、そうでなければValueErrorです。readlineの待機を毎回timeoutで制限し、EOFならstateを返します。1行をconsume_lineでコミットしたあと、before_ackがあればbefore_ack(seq)を呼び出し、writer.write(ack)と、timeoutで制限したdrainを実行します。エラー・キャンセルは伝え、正常・失敗・キャンセルのいずれでも、close→timeoutで制限したwait_closedで回収します。終了待機のTimeoutError・ConnectionErrorにはwriter.transport.abortを呼び出し、元のエラーは隠しません。DB接続は閉じません。不正なtimeoutの終了待機には1秒を使います。
- 受信側を終了させて、同じDBから引き継ぐ: async run_client(host,port,path,timeout,before_ack=None)を実装してください。timeoutを先に検証し、open_store(path)、timeoutで制限したasyncio.open_connection(host,port,limit=64)の順に開きます。DBのlastをb"RESUME last\n"として書き込み、timeoutで制限したdrainのあと、consume_sessionに渡して、戻り値を返します。すべての経路でDB接続を閉じ、consume_sessionに渡す前に失敗したときは、直接ソケットを回収します。エラー・キャンセルは伝えます。採点ツールが一時DBとループバックサーバーを作り、最初のプロセスをコミット1/ACK 1の間に終了させたあと、新しいプロセスで再接続します。再送された1を重複して反映してはいけません。
参考
mkdir -p /root/resumeを実行したあと、/root/resume/client.pyに実装します。採点ツールは前のステップも実行し、一時DB・ループバックポート・受信プロセスを準備して回収します。学習者がサーバーやDBファイルを起動し続ける必要はありません。提出ファイルは変更しません。1ステップあたりの実行制限は15秒で、学習時間を制限するものではありません。最後の検査は実際のプロセスを終了させるので、後始末のfinallyが実行されるという前提を置いていません。本番用のTLS・認証・自動バックオフ・スナップショット復旧・分散exactly-onceを実装したラボではありません。
途切れた最後の行をコマンドとして受け付けない
parse_line(line)を実装してください。bytesだけを受け付け、64バイト以下のASCII「EVENT 連番 変化量 LF」の1行を、(seq, delta)のタプルで返します。seqは0–2147483647、deltaは-1000–1000です。数字の0以外の先頭の0、+符号、-0、CRLF、改行の欠落、複数行、非ASCIIはValueErrorです。
最後の改行をstripで取り除く前に、完全な1フレームかどうかを確認してください。
再起動しても覚えているストアを開く
open_store(path)はsqlite3の接続を返します。checkpoint(id INTEGER PRIMARY KEY CHECK(id=1), last INTEGER NOT NULL, total INTEGER NOT NULL)とledger(seq INTEGER PRIMARY KEY, delta INTEGER NOT NULL)の2つのテーブルを、存在しないときにだけ作成し、checkpointの初期行(1,-1,0)を、存在しないときにだけ挿入します。state(con)は(last,total)のタプルです。明示的なSQLトランザクションを使えるように、isolation_level=None、busy timeout=1秒で接続してください。初期化に失敗したときは、接続を閉じてエラーを伝えます。
すでに保存したカーソルを初期化すると、再接続ではなく最初からの重複処理になってしまいます。
合計とカーソルを一緒に保存するか、一緒に取り消す
Exceptionのサブクラスとして、GapとConflictを宣言し、apply_event(con,seq,delta,fault=None)を実装してください。boolを除く、ステップ1の範囲のintだけを許可し、不正ならValueErrorです。seq=last+1なら、ledgerへの挿入→totalの更新→faultがあればfault()→lastの更新を、1つのトランザクションでコミットして、Trueを返します。seq<=lastなら、ledgerの同じdeltaならFalse、内容が違えばConflictです。それより大きな連番はGapです。失敗したときは、すべての変更をロールバックして元のエラーを伝え、トランザクションを残さないでください。
カーソルを先にコミットしても、合計を先に別にコミットしても、障害のときに両者がずれます。
切り詰められた記録を、復旧の成功として隠さない
Exceptionのサブクラスとして、ResyncRequiredとreplay(events,last,limit)を実装してください。eventsは1–128個の(seq,delta)のタプルからなるlistで、すべての連番が連続・昇順であり、値はステップ1の範囲のintです。lastは-1–2147483647のint、limitは1–16のintです。bool・型・範囲のエラー、空のログ、last>最後の連番はValueErrorです。last<最初の連番-1ならResyncRequired、それ以外はseq>lastの項目を最大limit個、新しいlistとして返します。入力全体を先に検査し、変更はしません。
保管開始10とカーソル9はつながりますが、カーソル8には9が欠けています。
コミットしたあとにだけ、業務のACKを作る
consume_line(con,line,fault=None)は、parse_lineとapply_eventを呼び出し、成功後にb"ACK 連番\n"を返します。faultをapply_eventに渡してください。内容が同じ重複にもACKを返しますが、パース・ギャップ・衝突・業務の失敗ではエラーを伝え、ACKは返しません。
ACKの直後に、別のDB接続からも処理結果が見える必要があります。
エラーやキャンセルでも接続を閉じる
async consume_session(reader,writer,con,timeout,before_ack=None)を実装してください。timeoutはboolを除く正の有限なint/floatで、そうでなければValueErrorです。readlineの待機を毎回timeoutで制限し、EOFならstateを返します。1行をconsume_lineでコミットしたあと、before_ackがあればbefore_ack(seq)を呼び出し、writer.write(ack)と、timeoutで制限したdrainを実行します。エラー・キャンセルは伝え、正常・失敗・キャンセルのいずれでも、close→timeoutで制限したwait_closedで回収します。終了待機のTimeoutError・ConnectionErrorにはwriter.transport.abortを呼び出し、元のエラーは隠しません。DB接続は閉じません。不正なtimeoutの終了待機には1秒を使います。
ACK直前のフックは、保存は終わったもののパブリッシャーが知らない障害区間を再現します。
受信側を終了させて、同じDBから引き継ぐ
async run_client(host,port,path,timeout,before_ack=None)を実装してください。timeoutを先に検証し、open_store(path)、timeoutで制限したasyncio.open_connection(host,port,limit=64)の順に開きます。DBのlastをb"RESUME last\n"として書き込み、timeoutで制限したdrainのあと、consume_sessionに渡して、戻り値を返します。すべての経路でDB接続を閉じ、consume_sessionに渡す前に失敗したときは、直接ソケットを回収します。エラー・キャンセルは伝えます。採点ツールが一時DBとループバックサーバーを作り、最初のプロセスをコミット1/ACK 1の間に終了させたあと、新しいプロセスで再接続します。再送された1を重複して反映してはいけません。
lastを常に-1で送ったり、ストアを毎回消したりすると、新しいプロセスの実際のリクエストで明らかになります。