WebSocketリレー
websocket_relayミドルウェアはHTTP接続をアップグレードし、WebSocketメッセージを対象プロセスへ中継します。
分類:部分的な統合レシピを含むプロトコルリファレンス。 ブロックは、HTTPサーバー、ルーター、プロセスホスト、対象プロセス、セキュリティコンテキストが存在することを前提としています。アプリケーションのメッセージハンドラとクライアント状態のクリーンアップは、アプリケーション側が担当します。
仕組み
- HTTPハンドラが対象プロセスのPIDを
X-WS-Relayヘッダーに設定します - ミドルウェアが接続をWebSocketへアップグレードします
- リレーが対象プロセスへアタッチし、監視します
- クライアントとプロセスの間でメッセージが双方向に流れます
プロセスのセマンティクス
WebSocket接続は、固有のPIDを持つ完全なプロセスです。プロセスシステムと次のように統合されます:
- アドレス指定可能 → 任意のプロセスがWebSocket PIDへメッセージを送信できます
- 監視可能 → プロセスがWebSocket接続の終了イベントを監視できます
- リンク可能 → WebSocket接続を他のプロセスとリンクできます
- EXITイベント → 接続が閉じると、監視プロセスは終了通知を受信します
-- Monitor a WebSocket connection from another process
local _, monitor_err = process.monitor(websocket_pid)
if monitor_err then return nil, monitor_err end
-- Send a message to the WebSocket client from any process.
-- The relay wraps it as {topic, data} JSON; the topic name is arbitrary.
local _, send_err = process.send(websocket_pid, "update", "hello")
if send_err then return nil, send_err end
接続の移譲
制御メッセージを送信することで、接続を別のプロセスへ移譲できます:
local _, transfer_err = process.send(websocket_pid, "ws.control", {
target_pid = new_process_pid,
message_topic = "ws.message"
})
if transfer_err then return nil, transfer_err end
設定
ルーターのマッチ後ミドルウェアとして追加します:
- name: ws_router
kind: http.router
meta:
server: gateway
prefix: /ws
post_middleware:
- websocket_relay
post_options:
wsrelay.allowed.origins: "https://app.example.com"
| オプション | 説明 |
|---|---|
wsrelay.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
-- リレーを設定
res:set_header("X-WS-Relay", json.encode({
target_pid = tostring(pid),
message_topic = "ws.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-WS-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 | ws.message |
クライアントメッセージ用トピック |
heartbeat_interval |
duration | - | ハートビート頻度(例:30s) |
metadata |
object | - | join/leave/heartbeatメッセージに付与 |
メッセージトピック
リレーは対象プロセスへ次のメッセージを送信します:
| トピック | タイミング | ペイロード |
|---|---|---|
ws.join |
クライアント接続時 | JSON {client_pid, metadata} |
ws.message(または指定したmessage_topic) |
クライアントがメッセージを送信したとき | クライアントの生ペイロード(テキストフレーム → String形式、バイナリフレーム → Bytes形式)。どちらの形式でもpayload:data()はLua文字列を返し、送信元PIDはクライアントPIDです |
ws.heartbeat |
定期的(デフォルトは30秒ごと。heartbeat_intervalで変更可能) |
JSON {client_pid, uptime, message_count, metadata} |
ws.leave |
クライアント切断時 | JSON {client_pid, metadata} |
メッセージの受信
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 from = msg:from() -- client connection PID
if topic == "ws.join" then
-- Client connected — payload is {client_pid, metadata}
local data, payload_err = msg:payload():data()
if payload_err then return nil, payload_err end
local client_pid = data.client_pid
elseif topic == "ws.message" then
-- Raw client message; from() is the client PID
local incoming = msg:payload()
local frame_format = incoming:get_format() -- "text/plain" or "application/octet-stream"
local body, payload_err = incoming:data() -- Lua string in either case
if payload_err then return nil, payload_err end
-- Decode or dispatch `body` according to `frame_format` and the
-- application's protocol.
elseif topic == "ws.leave" then
-- Client disconnected — payload is {client_pid, metadata}
-- Release application state associated with `from`.
end
end
end
クライアントへの送信
クライアントPIDを使用してメッセージを送り返します。任意のトピックは {topic, data} JSONとしてラップされ、WebSocketに転送されます。サーバーからクライアントへのメッセージはすべて、{topic, data} JSONラッパーを含む単一のWebSocket TEXTフレームとして送信されます。バイナリペイロードは data フィールドにbase64エンコードされます。個別のバイナリフレームとしては送信されません。
-- Send a structured message (any topic name)
local _, send_err = process.send(client_pid, "update", {event = "update", value = 42})
if send_err then return nil, send_err end
-- 接続を閉じる (ペイロードはクローズ理由文字列)
process.send(client_pid, "ws.close", "Session ended")
サーバー → クライアント方向で予約されているトピックは、ws.control(リレーの再設定)とws.close(接続の終了)です。
ブロードキャスト
複数のクライアントへブロードキャストするには、クライアントPIDを追跡します:
local clients = {}
-- On join
clients[client_pid] = true
-- On leave
clients[client_pid] = nil
-- Broadcast
local function broadcast(message)
for pid, _ in pairs(clients) do
local _, send_err = process.send(pid, "broadcast", message)
if send_err then return nil, send_err end
end
return true
end
関連項目
- ミドルウェア - ミドルウェア設定
- プロセス - プロセスメッセージング
- WebSocketクライアント - 外向きWebSocket接続