接続の復旧と業務処理の復旧は違う
一言でいうと
切れた接続はもう一度開けますが、すでに行った処理をやり直すかどうかを判断するには、業務の結果と処理カーソルを同じトランザクションで残す必要があります。
なぜ必要なのか
宇宙ステーションのお菓子の自販機がイベントを受け取ります。EVENT 0 7は在庫を7個増やすという意味です。受信側が在庫を保存した瞬間に電源が切れ、ACKを送れませんでした。パブリッシャーには「処理したが返事が消えた場合」と「まったく処理できなかった場合」が同じに見えます。TCPの再接続に成功しただけでは、この2つの状況を区別できません。不確かなイベントは再送しつつ、受信側がすでに処理したイベントを見分けられる設計が必要です。
前のコースのCursorは、実行中のオブジェクトの記憶です。プロセスを新しく起動すると、その記憶も消えます。今回はSQLiteファイルに連番と合計を残し、実際の受信プロセスをACKの直前に終了させます。2つ目のプロセスが同じファイルを開いて引き継げるかを確かめます。決済サーバーを作るラボではなく、1つのストリームの順序付き在庫変化量で復旧の境界を浮かび上がらせる実験です。
どう動くのか
まずバイト列を契約にする
練習用のプロトコルは、ASCIIのEVENT、空白、連番、空白、変化量、LF1個です。例えば、EVENT 0 7とそのあとの改行が1つのイベントです。連番は0から始まる連続した整数、変化量は-1000から1000までです。先頭の0、+符号、-0、CRLF、改行のない最後の断片は許可しません。厳格な文法は、「数字に変換できたからよい」と「約束したメッセージである」を区別させるために、著者が決めたものです。ほかのプロトコルがCRLFを許可するなら、その契約に従う必要があります。
このメッセージはWebSocketやSSEではありません。2つの標準は、それぞれ別のフレーム・再接続のルールを持っています。ここではTCPのバイトストリームの上に小さな行プロトコルを載せて、処理の意味に集中します。1行の上限は64バイトで、open_connectionのStreamReader limitも64に設定します。不正なメッセージを読んだあとで、そのまま次の行から続けると、実際に抜けたコマンドを隠してしまうおそれがあるため、接続を終了します。HTTP認証、TLS、ユーザー別の権限は実装していないので、インターネットにそのまま公開してはいけません。
業務上の効果とカーソルは1か所に
checkpointテーブルのただ1つのid=1の行は、lastとtotalを保存します。まだ処理したものがなければ、last=-1、total=0です。ledgerはseqを主キーにして、処理したdeltaを保管します。このラボではledgerを削除しません。処理履歴の長期保管・圧縮ポリシーは、別に設計する必要があります。
apply_eventはBEGIN IMMEDIATEで書き込みトランザクションを開き、現在のlastを読みます。seqがlast+1なら、ledgerに記録してtotalを変更し、lastを進めたあとでCOMMITします。途中でエラーが起きたらROLLBACKして、元のエラーを呼び出し元に返します。ラボのfaultフックは、totalを変更したあと、lastを変更する前に実行されます。採点ツールはこのとき、一部の変更を同じ接続で観察してエラーを注入し、そのあとで、別の接続では何の変更も残っていないことを確認します。別のプロセスをこのフックですぐに終了させて、finallyが実行されない場合も検査します。単に「例外を捕捉した」という出力を、原子性の証拠にすることはしません。
同じseqとdeltaが再び来たら、結果を加算せずにFalseを返します。同じseqでdeltaが違う場合はConflictです。イベントIDを再利用しながら内容を変えるパブリッシャーを、重複として黙って受け入れないためです。seqがlast+1より大きければGapです。4まで処理したのに6を受け取ってlast=6へ飛ばすと、あとから届いた5をすでに処理したものと誤解してしまいます。いったん止めて、抜けた区間を復旧する必要があります。
ここではsqlite3.connectのisolation_level=Noneで自動BEGINをオフにして、SQLでBEGIN・COMMIT・ROLLBACKを明示します。Python 3.12のautocommit=Falseの方式と混ぜて、BEGINをネストさせないでください。with conは接続そのものを閉じてくれる構文ではないので、所有者がfinallyでcloseします。open_storeは既存の行を初期値で上書きしません。再起動のたびにlast=-1へ初期化すると、ファイルを使う意味がなくなります。
ACKはコミットのあとに、再開はカーソルの次から
consume_lineは、パースとapply_eventが成功したあとで、ACKの連番とLFを返します。内容が同じ重複にもACKを返します。結果を再び加算せずに、パブリッシャーがリトライを終えられるようにするためです。まだコミットしていないのにACKを先に送って落ちると、パブリッシャーがイベントを消してしまうおそれがあります。逆に、コミットのあとでACKが消える場合は重複が起こりえますが、処理履歴でふるい落とせます。
run_clientはストアを開き、TCPで接続して、RESUME lastとLFを送ります。再接続は、新しいrun_clientの呼び出しで明示的に行います。無限の自動リトライループを隠して入れることはしません。最終検査では、最初の受信側が0と1を保存したあと、ACK 1の前にos._exitで終了します。2つ目のプロセスはRESUME 1を送り、サーバーがわざと再送した1と、新しいイベント2を受け取ります。サーバーで観察したリクエスト・ACKと、独立した接続で読んだledgerと合計を、あわせて突き合わせます。実行中のオブジェクトをすり替えた、疑似的な再起動ではありません。
保管されたログの範囲外では復旧方法が変わる
replayは、パブリッシャー側の有限なログから、カーソルより大きい項目を最大limit個返します。現在の保管開始が10なら、last=9は10から読めますが、last=8は必要な9がないのでResyncRequiredです。最新の項目だけを渡して成功とすると、静かな損失になります。実際の製品は、一貫したスナップショットとそのスナップショットのカーソルから再同期するか、明示的な復旧失敗を提供する必要があります。このラボは、スナップショットの作成・インストール自体は実装しません。
空のログは、保管情報がなければ「最初から空だった」と「すべて切り捨てられた」を区別できないため、このAPIではValueErrorです。入力は1–128個の連続した項目、バッチは1–16個に制限します。実際のログサービスには、保管の下限・上限のような別のメタデータが必要です。Redis XREADも指定したIDより後の項目を読みますが、このラボの整数の連番と例外を、Redisの実際の動作だと解釈してはいけません。
現場での姿
通知・共同編集の画面・AIストリーミングは、接続が長く維持されても、デプロイやネットワークの変更で切れます。再接続と再生、処理カーソル、保管ポリシーは、あわせて設計する必要があります。接続数が増えると、リトライに指数バックオフ・ジッター・最大回数、アカウント別のリソース制限も必要になります。今回は2回の接続と小さなログだけを検証するので、インターネットの遅延や大規模なスループットの証拠としては使いません。
同じSQLiteトランザクションに入った合計とカーソルの原子性が、外部のメールや決済APIまで束ねてくれるわけではありません。DBコミットとHTTPリクエストの間にも、新しい障害区間ができます。そのようなシステムでは、受信側の冪等キー、outboxのような設計と再調整の手順を、別に学ぶ必要があります。「ACKがあるから、分散システム全体でちょうど1回」という説明は誤りです。
ラボのファイルはプロセスの再起動中は残りますが、LabHubのセッションが終了すると消えます。今回のプロセス障害の検査は、ホストの電源喪失、ディスク障害、バックアップからの復旧を試したものではありません。SQLiteの耐久性も、ファイルシステムと同期動作についての前提を持っています。観察した障害の種類と、保証していない種類を区別して報告してください。
次のラボですること
フレームのパース → ストア → アトミックな処理 → 保管範囲の判断 → コミット後のACK → 接続の回収 → 実際のプロセスの再接続の順に、client.pyを完成させます。timeoutは各ネットワークawaitの待機上限であり、セッション全体やDB作業を含めた総期限ではありません。SQLiteの呼び出しは同期式で、この小さな実験は受信側が1つだけです。本番のイベントループで長いDB作業をそのまま実行する設計は、避ける必要があります。
公式ドキュメントでさらに読む
- Pythonのsqlite3トランザクション制御: 接続モードと明示的なコミットを区別します。
- SQLiteのトランザクション: BEGIN IMMEDIATEと同時書き込みの制約を確認します。
- SQLiteのアトミックコミット: 原子性の実装と、ストレージ装置に対する前提を読みます。
- Python asyncio Streams: 読み取りの制限、drain、closeとwait_closedの責務を確認します。
- Redis XREAD: 最後に読んだIDとそのあとの再生という、実際のAPIを比較します。