リアルタイム通信 — WebSocket・gRPC ストリーミング・WebRTC
長く繋がる WebSocket ハブを運用する
目標
websocketsで、発行・購読ハブと再接続クライアントを作り、遅い購読者・アイドルタイムアウト・黙って消えた相手・デプロイによる一斉切断に、順番に耐えられるようにします。
なぜ重要なのか
WebSocketサービスの障害は、ほとんどがプロトコルではなく、時間から来ます。1人が遅くなればハブのメモリがいっぱいになり、静かな接続はロードバランサーが先に切り、携帯電話がトンネルに入ると相手はcloseなしで消え、デプロイ1回で数万が同時に再接続します。このラボの5つのルール、つまり、購読者ごとの上限、1013での切断、アイドル上限より短いping、終了も一緒に待つこと、ジッター付きの再接続とseqでの続きの受け取りは、チャットでも相場でも、音声AIの途中結果でも、同じように必要です。
ステップ
- 再接続の間隔にジッターを混ぜる: /root/rt/wsops/client.pyにbackoff(attempt, base=0.2, cap=5.0, rng=random)を作ってください。attempt回目のリトライの前に待つ秒数を返します。値は、rng.uniform(0, min(cap, base * 2 ** attempt))を1回呼んで決めます(full jitter)。rngには、randomモジュールまたはrandom.Randomオブジェクトが来ます。
- クローズコードを見て再接続するか決める: /root/rt/wsops/client.pyにshould_reconnect(code)を追加してください。1001(去っていく)・1006(closeフレームなしで切断)・1011(サーバーの内部エラー)・1012(サービスの再起動)・1013(あとでもう一度)ならTrue、それ以外のコードはすべてFalseを返します。
- ハブが1人の発言を全員に運ぶ: /root/rt/wsops/hub.pyにrun(host, port, ring=1000, queue=100, ping_interval=1.0, ping_timeout=2.0)コルーチンを作ってください。websockets.asyncio.server.serveで、3つのパスを受け付けます。/pubに入ってきたテキストメッセージごとに、1から増えるseqを付けて、{"seq": n, "data": メッセージ}のJSON文字列にし、/subにつながったすべての購読者に送ります。/statsは、{"subscribers": 購読者数, "dropped": 切断した購読者数, "seq": 最後のseq}のJSONを1つ送って終わります。購読者ごとに、asyncio.Queueと、そのキューを空にしながら送るループを、別に持たせます。
- 遅い購読者1人だけを切る: 購読者キューのサイズをqueueで制限してください。キューがいっぱいで入れられなければ、その購読者をリストから外し、droppedを1上げて、クローズコード1013と理由"slow consumer"で閉じます。採点ツールは、読まない購読者1人とよく読む購読者1人をつないだまま、8000バイトのメッセージ8000個を発行し、よく読む側がすべて受け取るか、ハブのメモリがどれだけ増えたか、読まなかった側が1013を受け取るかを確認します。
- 静かな接続がロードバランサーに切られないようにする: serveに、ping_intervalとping_timeoutを、runの引数のまま渡してください。採点ツールは、2.5秒のあいだ静かな接続を切る中継器をハブの前に置き、pingを自分では送らない購読者で6秒待ってから、メッセージを1つ発行します。そのメッセージが購読者に届かなければなりません。
- 黙って消えた購読者を片づける: 購読者ループが、キューだけを待たずに、接続が閉じることも一緒に待つように直してください。ws.wait_closed()とq.get()をasyncio.wait(..., return_when=FIRST_COMPLETED)で一緒に待ち、接続が先に閉じたら、購読者リストから外してループを終えます。serveには、close_timeout=1.0も渡してください。採点ツールは、ハンドシェイクだけを行ってpingに答えない購読者をつないだあと、6秒以内に/statsのsubscribersが0に戻るかを確認します。
- 切れた位置から続きを受け取る: ハブに、最近のring個のメッセージを入れるリングバッファーを置き、/sub?last=Nでつながったら、seqがNより大きいメッセージを先に再送してから、リアルタイムで続けて送るようにしてください。N+1がバッファーからすでに押し出されていたら、{"reset": true, "seq": 最後のseq}を先に送ります。そして、/root/rt/wsops/client.pyにconsume(url, out_path, stop_after, base=0.1, cap=1.0)コルーチンを追加してください。urlにつないで受け取ったメッセージのseqを、1行に1つずつout_pathに追記し、接続が切れたら、should_reconnectとbackoffで待ってから、最後に書き込んだseqをlastとして渡して再接続します。seqがstop_afterに達したら終わります。採点ツールは、発行の途中で接続を2回切り、ファイルに1からstop_afterまでが、抜けなく1回ずつあるかを確認します。
参考
- 作業フォルダは/root/rt/wsopsです。mkdir -p /root/rt/wsopsで、先に作ってください。
- ハブを自分で起動してみるには、cd /root/rt/wsops && /opt/rt-lab/bin/python -c "import asyncio, hub; asyncio.run(hub.run('127.0.0.1', 9002))"を使ってください。採点ツールは、空いているポートを選んで、別に起動します。
- アイドルタイムアウトと接続の切断は、/opt/fixtures/rt/rtnet.pyのIdleProxyとRelay.cut()で作ります。どう切っているのか、読んでみてもかまいません。
- よくあるミスが2つあります。発行ループで購読者ごとにawait sendを呼んで、1人が全体を止めてしまうことと、受け取るとすぐlastを進めてしまい、書き込む前に切れたメッセージを失うことです。
- Pythonは、必ず/opt/rt-lab/bin/pythonで実行します。このラボのライブラリは、その仮想環境にだけ入っていて、普通のpython3で実行すると、ModuleNotFoundErrorが出ます。alias rpy=/opt/rt-lab/bin/pythonのように短くしておくと便利です。
- ラボのPodは、外へ出る接続が塞がれています。すべての通信は、同じPodの中の127.0.0.1で行われ、インストールやダウンロードは必要ありません。
- 採点ツールは、コードを別プロセスで読み込んで、実際に接続を張ってみます。例のファイルは関数の枠にすぎないので、そのままでは合格しません。前のステップで完成させた関数は、消さないでください。
- ラボのセッションが終わると、/rootのファイルは残りません。必要なコードは、終える前に別に保管してください。
再接続の間隔にジッターを混ぜる
/root/rt/wsops/client.pyにbackoff(attempt, base=0.2, cap=5.0, rng=random)を作ってください。attempt回目のリトライの前に待つ秒数を返します。値は、rng.uniform(0, min(cap, base * 2 ** attempt))を1回呼んで決めます(full jitter)。rngには、randomモジュールまたはrandom.Randomオブジェクトが来ます。
ジッターがないと、一斉に切れたクライアント数万が、まったく同じ時刻に再び押し寄せて、生き返ったばかりのサーバーをまた倒します。上限(cap)がないと、10回目のリトライは数分後になります。採点ツールは、シードを決めたrandom.Randomを渡して、同じ値を計算してみます。
クローズコードを見て再接続するか決める
/root/rt/wsops/client.pyにshould_reconnect(code)を追加してください。1001(去っていく)・1006(closeフレームなしで切断)・1011(サーバーの内部エラー)・1012(サービスの再起動)・1013(あとでもう一度)ならTrue、それ以外のコードはすべてFalseを返します。
1008(ポリシー違反)や1002(プロトコルエラー)で切れた接続に再接続すると、同じ理由でまた切れます。リトライに意味があるのは、相手側の事情が変わりうる場合だけです。1000は、誰かが意図的に閉じたという意味です。
ハブが1人の発言を全員に運ぶ
/root/rt/wsops/hub.pyにrun(host, port, ring=1000, queue=100, ping_interval=1.0, ping_timeout=2.0)コルーチンを作ってください。websockets.asyncio.server.serveで、3つのパスを受け付けます。/pubに入ってきたテキストメッセージごとに、1から増えるseqを付けて、{"seq": n, "data": メッセージ}のJSON文字列にし、/subにつながったすべての購読者に送ります。/statsは、{"subscribers": 購読者数, "dropped": 切断した購読者数, "seq": 最後のseq}のJSONを1つ送って終わります。購読者ごとに、asyncio.Queueと、そのキューを空にしながら送るループを、別に持たせます。
発行する側が購読者ごとにawait sendを順番に呼ぶと、1人が遅いとき、その後ろの全員が待ちます。購読者ごとのキューは、その待ちを人ごとに切り離す仕掛けです。発行は、キューに入れるだけで、待ちません。
遅い購読者1人だけを切る
購読者キューのサイズをqueueで制限してください。キューがいっぱいで入れられなければ、その購読者をリストから外し、droppedを1上げて、クローズコード1013と理由"slow consumer"で閉じます。採点ツールは、読まない購読者1人とよく読む購読者1人をつないだまま、8000バイトのメッセージ8000個を発行し、よく読む側がすべて受け取るか、ハブのメモリがどれだけ増えたか、読まなかった側が1013を受け取るかを確認します。
キューに上限がないと、遅い購読者1人の分のメッセージが、ハブのメモリに際限なくたまります。メッセージを捨てることと接続を切ることのどちらがよいかは、データの性質が決めます。順序と欠落が重要なストリームなら、穴が空いたまま送り続けるより、切って再接続させるほうが正直です。
静かな接続がロードバランサーに切られないようにする
serveに、ping_intervalとping_timeoutを、runの引数のまま渡してください。採点ツールは、2.5秒のあいだ静かな接続を切る中継器をハブの前に置き、pingを自分では送らない購読者で6秒待ってから、メッセージを1つ発行します。そのメッセージが購読者に届かなければなりません。
websocketsの既定のping間隔は20秒なので、アイドル上限がそれより短い機器の後ろでは、静かな接続が先に切れます。よくあるAWS ALBの既定のアイドル上限は60秒で、社内プロキシはもっと短い場合が多いです。間隔は、最も短いアイドル上限より短くする必要があります。
黙って消えた購読者を片づける
購読者ループが、キューだけを待たずに、接続が閉じることも一緒に待つように直してください。ws.wait_closed()とq.get()をasyncio.wait(..., return_when=FIRST_COMPLETED)で一緒に待ち、接続が先に閉じたら、購読者リストから外してループを終えます。serveには、close_timeout=1.0も渡してください。採点ツールは、ハンドシェイクだけを行ってpingに答えない購読者をつないだあと、6秒以内に/statsのsubscribersが0に戻るかを確認します。
pingのタイムアウトでライブラリが接続を閉じても、キューを待つコルーチンは、新しいメッセージが来るまで起きません。そのあいだ、購読者リストには死んだ接続が残って、数字を膨らませ、メッセージを受け取ってためます。そして、答えのない相手にcloseを送ったあと、答えを待つ時間の既定値は10秒です。
切れた位置から続きを受け取る
ハブに、最近のring個のメッセージを入れるリングバッファーを置き、/sub?last=Nでつながったら、seqがNより大きいメッセージを先に再送してから、リアルタイムで続けて送るようにしてください。N+1がバッファーからすでに押し出されていたら、{"reset": true, "seq": 最後のseq}を先に送ります。そして、/root/rt/wsops/client.pyにconsume(url, out_path, stop_after, base=0.1, cap=1.0)コルーチンを追加してください。urlにつないで受け取ったメッセージのseqを、1行に1つずつout_pathに追記し、接続が切れたら、should_reconnectとbackoffで待ってから、最後に書き込んだseqをlastとして渡して再接続します。seqがstop_afterに達したら終わります。採点ツールは、発行の途中で接続を2回切り、ファイルに1からstop_afterまでが、抜けなく1回ずつあるかを確認します。
「受け取ったもの」ではなく、「処理して書き残したもの」を基準に再接続する必要があります。受け取っただけで書く前に切れると、そのメッセージは永遠に消えます。同じseqが2回来たら、1回だけ書いてください。