確認・期限・後始末の責任を決める
一言でいうと
送信バッファが空になったという事実、業務が終わったという確認、接続を回収したという事実は、それぞれ別の証拠です。
なぜ必要なのか
サーバーがwriteを終えた後に成功ログを残したのに、ユーザーの画面は変わりません。バイトがオペレーティングシステムに預けられたことと、相手のプログラムが業務を終えたことは別だからです。逆に、相手は処理したのに、ACKが来る前に接続が切れたこともあります。リアルタイムのプログラムは、この不確実性を隠さず、メッセージの契約・期限・復旧位置で表現する必要があります。
今回の契約は、とても小さなものです。サーバーはASCIIでEVENT 7の後に改行を送り、クライアントは該当の業務を終えたと仮定して、ACK 7の後に改行を返します。連番は0から2147483647までの整数です。認証された業務プロトコルや永続メッセージブローカーではなく、確認と寿命を学べる教育用のプロトコルです。
どう動くのか
StreamWriter.writeは、データを書き込みバッファに入れます。await drainは、送信のフロー制御が許可するまで待ちます。drainの完了は、相手のアプリケーションの完了ではないので、続けてreader.readlineで、同じ連番のACKを待ちます。EOF、別の連番、誤った形式は、ConnectionErrorとして扱います。送信関数は、複数のイベントが使う永続的な接続を、勝手に閉じません。
時間予算は、イベント1つに1回与えます。全体で1秒なのにdrainで0.8秒を使ったなら、ACKには約0.2秒しか残りません。各awaitに新しく1秒を与えると、最大の待ち時間が作業の数に応じて増えます。開始時にmonotonicの時計で締切時刻を決め、待つ前に毎回、残りの時間を計算してください。最後の結果が戻った後も、期限を超えていないか確認します。壁時計の時刻の補正は、経過時間の判定に使いません。
キューのワーカーは、getで項目を1つ引き受けた後、send_oneを実行します。成功でも例外でも、引き受けた項目の後始末は、finallyでtask_doneとして終えます。ACKの前にtask_doneを呼ぶと、発行側のjoinが、実際の確認より先に解けます。逆に、例外のときにtask_doneを忘れると、作業は死んだのに、joinだけが待ち続けます。後始末と業務の成功は別の記録であるという点を、ここでもう一度適用します。
ワーカー関数が接続の所有権を引き受けたなら、空のキューでキャンセルされても、ACK待ちの間に失敗しても、writer.closeを呼ぶ必要があります。wait_closedにも有限な期限を置き、終了が終わらなければ、transport.abortで強制的に回収します。キューにまだ残っている項目の破棄・リトライは、購読者の管理者の責任です。ワーカー1つが自分の処理中の項目を後始末することと、キュー全体を空にすることは、混ぜません。
再接続には、最後に確実に処理した連番が必要です。Cursor(last=7)に8が来れば連続した進行で、7が再び来れば重複です。10が先に来れば、8と9が抜けているので、Gapを出し、lastを動かしません。ログの再生や信頼できるスナップショットで、欠落を復旧してから進める必要があります。ただ大きな連番を採用すると、欠落が永久に隠れます。
ここで作るCursorは、メモリの中で連続性を判断する小さな部品です。acceptがTrueなら、lastをすぐに進めるので、失敗しうる実際の決済関数を呼ぶ前に、そのまま使ってはいけません。業務の適用とチェックポイントをアトミックに保存する機能は、このコースにありません。プロセスが再起動するとメモリも消えます。重複判断の試験が通ったからといって、ちょうど1回の処理が保証されるわけではありません。
現場での姿
運用では、購読者の接続解除、パブリッシャーのキャンセル、プロセスの終了が、それぞれ別の時点で重なります。正常なACKの試験だけが成功すると、このような終了経路のリークを見逃します。繰り返し接続を切っても、ソケットとタスクが残らないかを確認する必要があります。リスナーを先に閉じて新しい接続を防ぎ、既存のタスクと接続を後始末した後で、サーバーの終了を待つ順序も重要です。Python 3.12のServer.wait_closedは、アクティブな接続の終了を待ちます。
最後の実際のTCPの試験は、2つの購読者の部分障害と回収を観測します。前のステップは、偽の時計で全体の期限を決定的に検査し、ACKが永遠に来ない場合も、実際の非同期の待機で検査します。2つの方式は、代替の関係ではありません。偽の時計は境界の計算を、実際の接続は配線と寿命を確認します。
次のラボですること
期限の試験では、1秒を実際に何回も待つ代わりに、注入した時計を動かして、境界を検査できます。しかし、偽のreaderがすぐ返すなら、無限の待機を断ち切るコードがなくても、その例だけは通ることがあります。そのため、永遠に返らないreaderも別に使います。失敗の試験が長くかかるからといって取り除くと、まさにその欠陥が、運用で無限の待機になります。計算の境界と実際のキャンセルの両方を残す理由です。
ステップ8で、キューのポリシー・カーソル・ACKの全体期限・ワーカー・broadcastを完成させます。すべてのコードは/root/realtime/delivery.pyに置き、サーバーは、検査器が一時ポートで準備します。正解を見る前に、失敗メッセージが、どの層の証拠を求めているかを先に読んでみてください。