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リレーミドルウェアは、プロセスを基盤とする長時間のSSEストリームを作成します。WebSocketリレーと同じリレーパターンに従います。

仕組み

  1. HTTPハンドラが、JSON形式のリレー設定をX-SSE-Relayヘッダーに設定します
  2. ミドルウェアがレスポンスをインターセプトし、SSEセッションを作成します
  3. セッションが固有のPIDを持つプロセスとして登録されます
  4. セッション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
リレーはターゲットプロセスをモニタします。ターゲットが終了すると、SSE ストリームは自動的に閉じられ、クライアントは `done` イベントを受け取ります。

設定

ルーターのマッチ後ミドルウェアとして追加します:

- 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を空文字列に設定します。

関連項目