ローディングスピナーでは何も分かりません。AIタスクが数分間に及ぶ場合や、3回目のリトライのためにキューに戻る場合、ユーザーには現在の状態を見せる必要があります。Server-Sent Events(SSE)を使えば、WebSocketsのようなハンドシェイクのオーバーヘッドや、ロングポーリングのような複雑な制御を行うことなく、その可視性を実現できます。サーバーは単一のHTTPレスポンスを開いたままにし、状況の変化に応じてプレーンテキストのアップデートをプッシュします。クライアントはそれらが届くたびに読み取ります。
接続が切断されたとき、最初からやり直したいとは誰も思わないはずです。適切に構築されたSSEストリームは、中断した場所を記憶しています。Node.js 20と標準ライブラリだけで、これを実装できます。外部パッケージは必要ありません。
ワイヤーフォーマットの仕組み
SSEメッセージは単純なテキストです。サーバーは3つの要素を書き込みます。オプションのイベント名、必須の data フィールド、そしてセーブポイントとなる id フィールドです。各レコードは2つの改行文字(境界を示す空行)で終わります。
正常なストリームのワイヤー上での見え方は、以下のようになります。
id: 14
event: status
data: {"phase":"testing","progress":43}
id: 15
event: status
data: {"phase":"retrying","attempt":2}
ブラウザの EventSource クライアントは、これらの行を自動的に読み取ります。各ブロックに対してイベントを発生させ、最新の id を内部に保存します。TCP接続が不安定になった場合、クライアントは待機して再接続を行い、保存していた識別子を Last-Event-ID ヘッダーとしてサーバーに送り返します。このヘッダーこそが、このパターンが機能する最大の理由です。これがないと、永続的なカーソルを持つことができません。
Node.jsでのサーバー実装
Nodeの組み込み http モジュールで、これを直接扱うことができます。リクエストが来た際、クライアントがこれがページではなくストリームであることを認識できるように、正しいヘッダーを設定します。
Content-Type: text/event-stream
Cache-Control: no-cache
Connection: keep-alive
バッファリングを排除してください。プロキシやフレームワークは時としてレスポンスをバッチ処理(まとめて送信)することがあり、それがリアルタイム性を損なうため、チャンクごとに flush を行います。
まずIDを送り、次にイベントタイプ、ペイロードデータ、そして最後に終了を示す空行を送ります。順序が重要なのは、クライアントがキャプチャできるように、空行の前にIDが届かなければならないという点だけです。ネイティブの response.write() を使用する場合、出力は文字通り以下のようになります。
response.write(`id: ${cursor}\n`);
response.write(`event: ${eventName}\n`);
response.write(`data: ${JSON.stringify(payload)}\n\n`);
末尾の \n\n は単なる飾りではありません。SSEパーサーはこれをレコードの終端として扱います。これを見落とすと、クライアントはさらなるデータを待って停止してしまいます。
カーソルがすべてを決める
新しいHTTP接続が、必ずしも新しい状態を保証するわけではありません。クライアントが再接続したとき、Last-Event-ID ヘッダーは、そのクライアントが最後に受信したメッセージを教えてくれます。あなたの仕事は、最初からではなく、その次のメッセージから再開することです。
これは、サーバー側でイベントの順序付けられたログまたはジャーナルを維持することを意味します。デモであればメモリ上の配列で十分です。本番環境では、サーバーの再起動によって履歴が消え、すべてのクライアントがゼロから開始することを避けるため、データベースログ、Redisストリーム、あるいはライトアヘッドジャーナル(WAL)などの永続的な仕組みが必要です。
イベントは、単調増加する整数またはULIDでインデックスしてください。再接続が来たら、id > lastEventId となるイベントをクエリし、それらを順番に再生します。バックログされたメッセージが数百ある場合は、小さな人工的な遅延を入れたりバッチ化したりしてもよいですが、クライアントが時系列に沿って状態を再構築できるように、古いものから順に送信してください。
重複が発生することを想定する
ネットワークは信頼できるものではありません。サーバーがイベントを送信した後、TCPの確認応答(ACK)を失い、タイムアウト後に同じイベントを再度送信する可能性があります。最初から「最低1回の配信(at-least-once delivery)」を前提に設計してください。
クライアント側での重複排除は簡単です。イベントIDをキーとした Map を保持します。新しいイベントが届いたら、そのマップを確認します。もしIDが存在していれば、その重複分は静かに破棄します。サーバーが決定論的なIDを割り当てるため、重複が発生しても害はありません。マップを無限に肥大化させる必要はありません。イベントが安全に処理されたことを確認したら、古いIDを削除してください。ブラウザクライアントの場合、数百件程度のスライディングウィンドウがあれば通常は十分です。
カーソルが期限切れになったとき
最終的に、クライアントは何時間後、あるいは数日後に再接続することがあります。もし履歴バッファが直近の1,000イベント分しかなく、クライアントが2,000イベント分遅れている場合、欠落した部分を再生することは不可能です。
部分的な履歴をストリームしないでください。それではクライアントの状態が不整合になってしまいます。代わりに、カーソルが期限切れであることを検知し、次のイベントとして「フルスナップショット」を送信してください。スナップショットには、クライアントを現在の状態に固定するための新しいカーソルを含める必要があります。そこからは、通常通りライブな差分(delta)の送信を再開します。クライアント側のコードが、差分を追記するのではなくローカルモデルをリセットすべきタイミングを判断できるよう、プロトコル内でこの境界を明確に定義しておいてください。
ストリームを保護する
公開されたSSEエンドポイントは、攻撃の標的になりやすいものです。誰でも接続を維持し続けることができ、リクエストの再送によってストレージへの読み取り負荷が増幅される可能性があります。
エンドポイントを適切な認可で保護してください。ブラウザの EventSource はカスタムヘッダーをサポートしていないため、トークンをクエリ文字列で渡すか、厳格な SameSite ポリシーを設定したクッキーを使用してください。ストリームリソースを割り当てる前に、トークンを検証してください。
履歴の制限とユーザーごとのクォータを設定してください。タスクごとの保存イベント数に上限を設け、クライアントごとの同時接続数にも上限を設けます。切断やリプレイをログに記録することで、カーソルエンドポイントに過剰なリクエストを送り続ける不正なクライアントを特定できるようにします。
このパターンは応用が効く
このアプローチは HTTP 内だけに留まるものではありません。WebSockets、メッセージキュー、あるいはエージェント間インターフェースに移行する場合でも、同じルールが適用されます。トランスポート層は変わり(バイナリフレームやトピックのサブスクリプションを使用するかもしれません)、根本的な問題は同一のままです。カーソル、永続的なログ、少なくとも1回の配信(at-least-once)セマンティクス、クライアント側での重複排除、そしてカーソルが古くなった場合のフルスナップショットへのフォールバックが必要です。状態の収束(state convergence)を一度解決してしまえば、コアロジックを再設計することなく、TCP、WebSocket、あるいは RabbitMQ のようなブローカーを介して配信できるようになります。
シンプルさを保つ
Server-Sent Events が機能するのは、通常の HTTP 上で動作するからです。プロキシも理解できますし、ロードバランサーによるヘルスチェックも可能です。デバッグも curl を使うのと同じくらい簡単です。しかし、エッジケースを無視してしまうと、そのシンプルさは失われます。カーソルを構築し、リプレイを想定し、クライアント側で重複を排除し、履歴がなくなったらスナップショットを作成してください。そうすることで、不安定な Wi-Fi やサーバーの再起動、あるいは夜間のブラウザのスリープ状態であっても、長時間実行される AI タスクがその進捗を正確に報告できるようになります。
出典: Build a Reconnecting SSE Task Stream with Node.js
ディスカッションに参加する: GyaanSetu AI Community
