Server-Sent Events :id=server-sent-events
SSEミドルウェアは、Server-Sent Eventsプロトコルを使用して、サーバーからHTTPクライアントへイベントをストリーミングします。
利用できる仕組みは2つあります。HTTPハンドラからの直接ストリーミングと、sse_relayミドルウェアを介したプロセスベースのリレーです。
分類:部分的な統合レシピを含むプロトコルリファレンス。 リレーブロックは、HTTPサーバー、ルーター、プロセスホスト、対象プロセス、セキュリティコンテキストがすでに存在することを前提としています。アプリケーションのコールバックとクライアントの動作は、これらのスニペットの範囲外です。
直接ストリーミング
HTTPハンドラからSSEイベントを直接送信するには、res:write_event()を使用します。最初の呼び出しでレスポンスが自動的にSSEモードへ切り替わり、適切なヘッダーが設定されます。
local http = require("http")
local function handler()
local res, res_err = http.response()
if res_err then return nil, res_err end
local err = res:write_event({name = "status", data = {state = "started"}})
if err then return nil, err end
err = res:write_event({name = "progress", data = {percent = 50}})
if err then return nil, err end
err = res:write_event({name = "status", data = {state = "complete"}})
if err then return nil, err end
return true
end
各イベントにはnameフィールドとdataフィールドが必要です。dataの値は自動的にJSONエンコードされます。
SSEリレー
SSEリレーミドルウェアは、プロセスを基盤とする長時間のSSEストリームを作成します。WebSocketリレーと同じリレーパターンに従います。
仕組み
- HTTPハンドラが、JSON形式のリレー設定を
X-SSE-Relayヘッダーに設定します - ミドルウェアがレスポンスをインターセプトし、SSEセッションを作成します
- セッションが固有のPIDを持つプロセスとして登録されます
- セッションPIDへ送信されたメッセージがSSEイベントとしてクライアントに転送されます
プロセスのセマンティクス
SSEストリームは、固有のPIDを持つ完全なプロセスです。プロセスシステムと次のように統合されます:
- アドレス指定可能 — 任意のプロセスがストリームPIDへメッセージを送信できます
- 監視可能 — プロセスがSSEストリームの終了イベントを監視できます
- リンク可能 — SSEストリームを他のプロセスとリンクできます
- EXITイベント — ストリームが閉じると、監視プロセスは終了通知を受信します
-- Send event to SSE client from any process
local _, send_err = process.send(stream_pid, "sse.message", {event = "update", value = 42})
if send_err then return nil, send_err end
-- Monitor an SSE stream
local _, monitor_err = process.monitor(stream_pid)
if monitor_err then return nil, monitor_err end
設定
ルーターのマッチ後ミドルウェアとして追加します:
- name: sse_router
kind: http.router
meta:
server: gateway
prefix: /sse
post_middleware:
- sse_relay
post_options:
sserelay.allowed.origins: "https://app.example.com"
| オプション | 説明 |
|---|---|
sserelay.allowed.origins |
許可するオリジン(カンマ区切り、ワイルドカード対応) |
ハンドラのセットアップ
HTTPハンドラはプロセスを生成し、リレーを設定します:
local http = require("http")
local json = require("json")
local function handler()
local req, req_err = http.request()
if req_err then return nil, req_err end
local res, res_err = http.response()
if res_err then return nil, res_err end
local user_id, query_err = req:query("user_id")
if query_err then return nil, query_err end
-- Spawn handler process
local pid, spawn_err = process.spawn("app.sse:handler", "app:processes")
if spawn_err then return nil, spawn_err end
-- Configure relay
local relay_config, encode_err = json.encode({
target_pid = tostring(pid),
message_topic = "sse.message",
heartbeat_interval = "30s",
metadata = {
user_id = user_id
}
})
if encode_err then
local _, terminate_err = process.terminate(pid)
return nil, terminate_err or encode_err
end
local header_err = res:set_header("X-SSE-Relay", relay_config)
if header_err then
local _, terminate_err = process.terminate(pid)
return nil, terminate_err or header_err
end
end
リレー設定のフィールド
| フィールド | 型 | デフォルト | 説明 |
|---|---|---|---|
target_pid |
string | — | メッセージを受信するプロセスPID(デタッチモードでは省略) |
message_topic |
string | sse.message |
転送するイベントのトピックフィルター |
heartbeat_interval |
duration | 30s |
ハートビート間隔(30s、1mなど) |
idle_timeout |
duration | — | 非アクティブ状態が続いた場合にストリームを閉じるまでの時間 |
hard_timeout |
duration | — | 絶対経過時間に基づいてストリームを閉じるまでの時間 |
metadata |
object | — | join/leave/heartbeatメッセージに添付されるデータ |
管理モードとデタッチモード
管理モード
target_pidを設定すると、リレーは管理モードで動作します:
- 対象プロセスを監視します
- 接続時に
sse.join、切断時にsse.leaveを送信します - 対象が終了するとストリームを自動的に閉じます
デタッチモード
target_pidを省略すると、リレーはデタッチモードで開始します:
readyイベントを、stream_pidとmessage_topicを含めてクライアントへ送信します- 初期状態ではプロセスを監視しません
- プロセスは後から
sse.controlメッセージを送信してアタッチできます
jsonをインポートし、レスポンスオブジェクトをresとして取得済みのハンドラ内では、次のようにデタッチモードを設定し、両方の操作を確認します:
-- Detached setup: no target_pid
local relay_config, encode_err = json.encode({
heartbeat_interval = "30s"
})
if encode_err then return nil, encode_err end
local header_err = res:set_header("X-SSE-Relay", relay_config)
if header_err then return nil, header_err end
クライアントはreadyイベントを受信します:
{"stream_pid": "{n1@app:gateway|0x0002a}", "message_topic": "sse.message"}
メッセージトピック
リレーは、ストリームと対象プロセス間の通信に次のトピックを使用します:
| トピック | 方向 | タイミング | ペイロード |
|---|---|---|---|
sse.join |
stream → target | クライアント接続時 | client_pid、metadata |
sse.message |
target → stream | デフォルトのイベントトピック | SSE イベントとして転送 |
sse.heartbeat |
stream → target | 周期的(デフォルトは30秒ごと) | client_pid、uptime、message_count、metadata |
sse.leave |
stream → target | クライアント切断時 | client_pid、metadata |
sse.control |
any → stream | 制御コマンド | リレー設定フィールド |
sse.close |
any → stream | 強制クローズ | 任意の理由文字列 |
対象プロセスでの受信
local function handler()
local inbox = process.inbox()
while true do
local msg, ok = inbox:receive()
if not ok then break end
local topic = msg:topic()
local data, payload_err = msg:payload():data()
if payload_err then return nil, payload_err end
if topic == "sse.join" then
local client_pid = data.client_pid
elseif topic == "sse.heartbeat" then
-- Periodic health check
elseif topic == "sse.leave" then
-- Release application state associated with data.client_pid.
end
end
end
イベントの送信
ストリームPIDへメッセージを送信することで、クライアントへイベントを送ります:
-- Send on the default message topic
local _, send_err = process.send(stream_pid, "sse.message", {
event = "update",
value = 42
})
if send_err then return nil, send_err end
-- Force close the stream
local _, close_err = process.send(stream_pid, "sse.close", "session expired")
if close_err then return nil, close_err end
設定したmessage_topicで送信されたイベントは、SSEイベントとしてクライアントに転送されます。トピック名がSSEイベント名になります。
接続の移譲
対象プロセス、トピックフィルター、タイムアウトを動的に変更するには、制御メッセージを送信します:
local _, transfer_err = process.send(stream_pid, "sse.control", {
target_pid = tostring(new_pid),
message_topic = "custom.topic",
idle_timeout = "5m"
})
if transfer_err then return nil, transfer_err end
対象が変わると、リレーはまず新しい対象の監視を開始してsse.joinを送信し、その後で以前の対象の監視を停止してsse.leaveを送信します。再アタッチせずにデタッチするには、target_pidを空文字列に設定します。
関連項目
- ミドルウェア — ミドルウェア設定
- WebSocketリレー — WebSocketでの同等機能
- プロセス — プロセスメッセージング