From ef5202a0cfacde0fa76401beb5faa92258c56ca3 Mon Sep 17 00:00:00 2001 From: nicocsh <119978472+nicocsh@users.noreply.github.com> Date: Wed, 22 Jul 2026 19:58:20 +0200 Subject: [PATCH 1/3] Add WebRTC Star and Mesh --- .../lib/game_server_web/application.ex | 12 + .../channels/signaling_broker.ex | 380 ++++++++++++++++++ .../channels/signaling_channel.ex | 308 ++++++++++++++ .../game_server_web/channels/user_socket.ex | 3 +- .../lib/game_server_web/signaling.ex | 15 + apps/game_server_web/mix.exs | 1 + modules/plugins/.gitignore | 2 + modules/plugins/handle_webrtc/.formatter.exs | 3 + modules/plugins/handle_webrtc/.gitignore | 3 + modules/plugins/handle_webrtc/README.md | 35 ++ .../plugins/handle_webrtc/config/config.exs | 3 + .../example_hook/v1/example_hook.pb.ex.old | 35 ++ .../game_server/modules/handle_webrtc_hook.ex | 15 + modules/plugins/handle_webrtc/mix.exs | 33 ++ modules/plugins/handle_webrtc/mix.lock | 13 + .../handle_webrtc/proto/example_hook.proto | 26 ++ sdk/lib/game_server/signaling.ex | 8 + 17 files changed, 894 insertions(+), 1 deletion(-) create mode 100644 apps/game_server_web/lib/game_server_web/application.ex create mode 100644 apps/game_server_web/lib/game_server_web/channels/signaling_broker.ex create mode 100644 apps/game_server_web/lib/game_server_web/channels/signaling_channel.ex create mode 100644 apps/game_server_web/lib/game_server_web/signaling.ex create mode 100644 modules/plugins/.gitignore create mode 100644 modules/plugins/handle_webrtc/.formatter.exs create mode 100644 modules/plugins/handle_webrtc/.gitignore create mode 100644 modules/plugins/handle_webrtc/README.md create mode 100644 modules/plugins/handle_webrtc/config/config.exs create mode 100644 modules/plugins/handle_webrtc/lib/example_hook/v1/example_hook.pb.ex.old create mode 100644 modules/plugins/handle_webrtc/lib/game_server/modules/handle_webrtc_hook.ex create mode 100644 modules/plugins/handle_webrtc/mix.exs create mode 100644 modules/plugins/handle_webrtc/mix.lock create mode 100644 modules/plugins/handle_webrtc/proto/example_hook.proto create mode 100644 sdk/lib/game_server/signaling.ex diff --git a/apps/game_server_web/lib/game_server_web/application.ex b/apps/game_server_web/lib/game_server_web/application.ex new file mode 100644 index 000000000..56abc3288 --- /dev/null +++ b/apps/game_server_web/lib/game_server_web/application.ex @@ -0,0 +1,12 @@ +defmodule GameServerWeb.Application do + use Application + + @impl true + def start(_type, _args) do + children = [ + GameServerWeb.SignalingBroker + ] + + Supervisor.start_link(children, strategy: :one_for_one, name: GameServerWeb.Supervisor) + end +end diff --git a/apps/game_server_web/lib/game_server_web/channels/signaling_broker.ex b/apps/game_server_web/lib/game_server_web/channels/signaling_broker.ex new file mode 100644 index 000000000..0d3efc8ae --- /dev/null +++ b/apps/game_server_web/lib/game_server_web/channels/signaling_broker.ex @@ -0,0 +1,380 @@ +defmodule GameServerWeb.SignalingBroker do + @moduledoc """ + Signaling relay for WebRTC peer-to-peer and client-server topologies. + + Rooms are created explicitly by a worker process (e.g. a lobby worker) and + are keyed by the lobby id. Each room stores its topology and, for :star, + the designated host user id. The broker validates membership and topology + rules on every relay. + + Does not create PeerConnections or handle media; only routes SDP offers, + answers, and ICE candidates between registered peers in a room. + + ## Topologies + + * `:mesh` — any member may send an offer/answer/ICE to any other member. + * `:star` — one host peer (the Godot headless server) and client peers. + Clients may only signal to the host; the host may signal to any client. + Non-host peers cannot exchange messages directly. + + Each peer is monitored via `Process.monitor/1`. When a peer crashes or + disconnects it is automatically removed and the remaining peers are notified. + """ + + use GenServer + require Logger + + defstruct rooms: %{}, refs: %{} + + # ── Public API ─────────────────────────────────────────────────────────── + + def start_link(opts) do + name = Keyword.get(opts, :name, __MODULE__) + GenServer.start_link(__MODULE__, opts, name: name) + end + + @doc """ + Creates a signaling room. `room_id` is typically the lobby id. + + For `:star` topology `host_user_id` is required and designates the user + that will act as the authoritative server peer. + """ + def create_room(room_id, topology, opts \\ []) when topology in [:mesh, :star] do + host_user_id = if topology == :star, do: Keyword.fetch!(opts, :host_user_id), else: nil + GenServer.call(__MODULE__, {:create_room, room_id, topology, host_user_id}) + end + + @doc """ + Closes a signaling room. Existing peers are notified with a `room_closed` + event so their channels can stop gracefully. + """ + def close_room(room_id) do + GenServer.call(__MODULE__, {:close_room, room_id}) + end + + def room_exists?(room_id) do + GenServer.call(__MODULE__, {:room_exists, room_id}) + end + + @doc """ + Registers a peer in a room. + + Returns `{:ok, role}` where `role` is derived from the room topology and + the provided `user_id`. Returns `{:error, :room_not_found}` if the room + does not exist, and `{:error, :duplicate_peer}` if `peer_id` is already + present. + """ + def join(room_id, peer_id, pid, user_id, metadata \\ %{}) + when is_binary(room_id) and is_binary(peer_id) and is_pid(pid) and is_binary(user_id) do + GenServer.call(__MODULE__, {:join, room_id, peer_id, pid, user_id, metadata}) + end + + def leave(room_id, peer_id) when is_binary(room_id) and is_binary(peer_id) do + GenServer.call(__MODULE__, {:leave, room_id, peer_id}) + end + + @doc """ + Routes a signaling message from one peer to a specific target. + + Enforces topology rules: in `:star` mode a non-host peer may only relay + to the host. + """ + def relay(room_id, from_peer_id, to_peer_id, type, payload) do + GenServer.call(__MODULE__, {:relay, room_id, from_peer_id, to_peer_id, type, payload}) + end + + @doc """ + Broadcasts a signaling message to every other peer in the room. + + In `:star` mode only the host may broadcast. + """ + def broadcast(room_id, from_peer_id, type, payload) do + GenServer.call(__MODULE__, {:broadcast, room_id, from_peer_id, type, payload}) + end + + def list_peers(room_id) do + GenServer.call(__MODULE__, {:list_peers, room_id}) + end + + def is_host?(room_id, user_id) do + GenServer.call(__MODULE__, {:is_host, room_id, user_id}) + end + + # ── GenServer callbacks ────────────────────────────────────────────────── + + @impl true + def init(_opts) do + {:ok, %__MODULE__{}} + end + + @impl true + def handle_call({:create_room, room_id, topology, host_user_id}, _from, state) do + if Map.has_key?(state.rooms, room_id) do + Logger.warning("SignalingBroker: room already exists room=#{room_id}") + {:reply, {:error, :already_exists}, state} + else + room = %{ + topology: topology, + host_user_id: host_user_id, + peers: %{} + } + + Logger.info("SignalingBroker: room created room=#{room_id} topology=#{topology} host_user_id=#{host_user_id || "none"}") + {:reply, :ok, %{state | rooms: Map.put(state.rooms, room_id, room)}} + end + end + + @impl true + def handle_call({:close_room, room_id}, _from, state) do + case Map.pop(state.rooms, room_id) do + {nil, _} -> + Logger.warning("SignalingBroker: close_room for non-existent room=#{room_id}") + {:reply, {:error, :room_not_found}, state} + + {room, rooms} -> + peer_count = map_size(room.peers) + Logger.info("SignalingBroker: closing room=#{room_id} topology=#{room.topology} evicting=#{peer_count}") + + for {peer_id, %{pid: pid}} <- room.peers do + Logger.debug("SignalingBroker: sending room_closed to peer=#{peer_id} pid=#{inspect(pid)}") + send(pid, {:signaling_relay, :room_closed, nil, %{}}) + end + + refs = Enum.reject(state.refs, fn {_ref, {r, _p}} -> r == room_id end) |> Map.new() + + {:reply, :ok, %{state | rooms: rooms, refs: refs}} + end + end + + @impl true + def handle_call({:room_exists, room_id}, _from, state) do + {:reply, Map.has_key?(state.rooms, room_id), state} + end + + @impl true + def handle_call({:join, room_id, peer_id, pid, user_id, metadata}, _from, state) do + case Map.get(state.rooms, room_id) do + nil -> + Logger.warning("SignalingBroker: join failed room_not_found room=#{room_id} user=#{user_id}") + {:reply, {:error, :room_not_found}, state} + + room -> + if Map.has_key?(room.peers, peer_id) do + Logger.warning("SignalingBroker: join failed duplicate_peer room=#{room_id} peer=#{peer_id}") + {:reply, {:error, :duplicate_peer}, state} + else + ref = Process.monitor(pid) + + role = + case room.topology do + :mesh -> :peer + :star -> if user_id == room.host_user_id, do: :host, else: :client + end + + peer = %{ + pid: pid, + user_id: user_id, + role: role, + metadata: metadata, + joined_at: System.monotonic_time(:second) + } + + peers = Map.put(room.peers, peer_id, peer) + room = %{room | peers: peers} + rooms = Map.put(state.rooms, room_id, room) + refs = Map.put(state.refs, ref, {room_id, peer_id}) + + peer_count = map_size(peers) + Logger.info("SignalingBroker: peer joined room=#{room_id} peer=#{peer_id} user=#{user_id} role=#{role} total_peers=#{peer_count}") + + for {other_id, %{pid: other_pid}} <- room.peers, other_id != peer_id do + Logger.debug("SignalingBroker: notifying peer=#{other_id} of peer_joined peer=#{peer_id}") + send(other_pid, {:signaling_relay, :peer_joined, peer_id, %{ + peer_id: peer_id, + role: role, + user_id: user_id + }}) + end + + {:reply, {:ok, role}, %{state | rooms: rooms, refs: refs}} + end + end + end + + @impl true + def handle_call({:leave, room_id, peer_id}, _from, state) do + case Map.get(state.rooms, room_id) do + nil -> + Logger.warning("SignalingBroker: leave failed room_not_found room=#{room_id} peer=#{peer_id}") + {:reply, {:error, :room_not_found}, state} + + room -> + case Map.pop(room.peers, peer_id) do + {nil, _} -> + Logger.warning("SignalingBroker: leave failed peer_not_found room=#{room_id} peer=#{peer_id}") + {:reply, {:error, :peer_not_found}, state} + + {_peer, peers} -> + peer_count = map_size(peers) + Logger.info("SignalingBroker: peer leaving room=#{room_id} peer=#{peer_id} user=#{_peer.user_id} role=#{_peer.role} remaining_peers=#{peer_count}") + + for {other_id, %{pid: other_pid}} <- peers do + send(other_pid, {:signaling_relay, :peer_left, peer_id, %{peer_id: peer_id}}) + end + + room = %{room | peers: peers} + + rooms = + if map_size(peers) == 0 do + Logger.info("SignalingBroker: room empty, removing room=#{room_id}") + Map.delete(state.rooms, room_id) + else + Map.put(state.rooms, room_id, room) + end + + ref_entry = + Enum.find(state.refs, fn {_ref, {r, p}} -> r == room_id and p == peer_id end) + + refs = + if ref_entry do + Map.delete(state.refs, elem(ref_entry, 0)) + else + state.refs + end + + {:reply, :ok, %{state | rooms: rooms, refs: refs}} + end + end + end + + @impl true + def handle_call({:relay, room_id, from, to, type, payload}, _from, state) do + case Map.get(state.rooms, room_id) do + nil -> + Logger.warning("SignalingBroker: relay failed room_not_found room=#{room_id} from=#{from} to=#{to} type=#{type}") + {:reply, {:error, :room_not_found}, state} + + room -> + from_peer = Map.get(room.peers, from) + to_peer = Map.get(room.peers, to) + + cond do + is_nil(from_peer) -> + Logger.warning("SignalingBroker: relay failed peer_not_found room=#{room_id} from=#{from} to=#{to} type=#{type}") + {:reply, {:error, :peer_not_found}, state} + + is_nil(to_peer) -> + Logger.warning("SignalingBroker: relay failed peer_not_found room=#{room_id} from=#{from} to=#{to} type=#{type}") + {:reply, {:error, :peer_not_found}, state} + + room.topology == :star and from_peer.role != :host and to_peer.role != :host -> + Logger.warning("SignalingBroker: relay failed not_allowed room=#{room_id} from=#{from} role=#{from_peer.role} to=#{to} role=#{to_peer.role}") + {:reply, {:error, :not_allowed}, state} + + true -> + Logger.debug("SignalingBroker: relaying room=#{room_id} type=#{type} from=#{from} to=#{to}") + send(to_peer.pid, {:signaling_relay, type, from, payload}) + {:reply, :ok, state} + end + end + end + + @impl true + def handle_call({:broadcast, room_id, from, type, payload}, _from, state) do + case Map.get(state.rooms, room_id) do + nil -> + Logger.warning("SignalingBroker: broadcast failed room_not_found room=#{room_id} from=#{from} type=#{type}") + {:reply, {:error, :room_not_found}, state} + + room -> + from_peer = Map.get(room.peers, from) + + if is_nil(from_peer) do + Logger.warning("SignalingBroker: broadcast failed peer_not_found room=#{room_id} from=#{from}") + {:reply, {:error, :peer_not_found}, state} + else + if room.topology == :star and from_peer.role != :host do + Logger.warning("SignalingBroker: broadcast failed not_allowed room=#{room_id} from=#{from} role=#{from_peer.role}") + {:reply, {:error, :not_allowed}, state} + else + targets = Enum.filter(room.peers, fn {peer_id, _} -> peer_id != from end) |> Enum.map(fn {id, _} -> id end) + Logger.debug("SignalingBroker: broadcasting room=#{room_id} type=#{type} from=#{from} targets=#{length(targets)}") + for {peer_id, %{pid: pid}} <- room.peers, peer_id != from do + send(pid, {:signaling_relay, type, from, payload}) + end + + {:reply, :ok, state} + end + end + end + end + + @impl true + def handle_call({:list_peers, room_id}, _from, state) do + case Map.get(state.rooms, room_id) do + nil -> + Logger.warning("SignalingBroker: list_peers failed room_not_found room=#{room_id}") + {:reply, {:error, :room_not_found}, state} + + room -> + peers = + Map.new(room.peers, fn {id, peer} -> + {id, %{user_id: peer.user_id, role: peer.role, metadata: peer.metadata}} + end) + + {:reply, peers, state} + end + end + + def handle_call({:is_host, room_id, user_id}, _from, state) do + case Map.get(state.rooms, room_id) do + nil -> + {:reply, false, state} + + room -> + {:reply, room.host_user_id == user_id, state} + end + end + + @impl true + def handle_info({:DOWN, ref, :process, pid, reason}, state) do + case Map.pop(state.refs, ref) do + {nil, _} -> + Logger.debug("SignalingBroker: DOWN from unknown pid=#{inspect(pid)} reason=#{inspect(reason)}") + {:noreply, state} + + {{room_id, peer_id}, refs} -> + case Map.get(state.rooms, room_id) do + nil -> + Logger.warning("SignalingBroker: DOWN for removed room room=#{room_id} peer=#{peer_id} pid=#{inspect(pid)} reason=#{inspect(reason)}") + {:noreply, %{state | refs: refs}} + + room -> + {_peer, peers} = Map.pop(room.peers, peer_id) + remaining = map_size(peers) + Logger.info("SignalingBroker: peer DOWN room=#{room_id} peer=#{peer_id} pid=#{inspect(pid)} reason=#{inspect(reason)} remaining_peers=#{remaining}") + + for {other_id, %{pid: pid}} <- peers, other_id != peer_id do + send(pid, {:signaling_relay, :peer_left, peer_id, %{peer_id: peer_id}}) + end + + room = %{room | peers: peers} + + rooms = + if map_size(peers) == 0 do + Logger.info("SignalingBroker: room empty after DOWN, removing room=#{room_id}") + Map.delete(state.rooms, room_id) + else + Map.put(state.rooms, room_id, room) + end + + {:noreply, %{state | rooms: rooms, refs: refs}} + end + end + end + + @impl true + def handle_info(_msg, state) do + {:noreply, state} + end +end diff --git a/apps/game_server_web/lib/game_server_web/channels/signaling_channel.ex b/apps/game_server_web/lib/game_server_web/channels/signaling_channel.ex new file mode 100644 index 000000000..05a6a6407 --- /dev/null +++ b/apps/game_server_web/lib/game_server_web/channels/signaling_channel.ex @@ -0,0 +1,308 @@ +defmodule GameServerWeb.SignalingChannel do + @moduledoc """ + Channel for WebRTC signaling relay. + + Topic: `signaling:` + + Rooms are created by a worker process (e.g. a lobby worker) through + `SignalingBroker.create_room/3`. Clients may only join if the room exists + and they are a member of the corresponding lobby. The topology and host + are fixed at room creation; clients cannot choose their role. + + ## Lifecycle + + On join a unique `peer_id` is generated. The broker assigns the role + (`:host` or `:client` for `:star`, `:peer` for `:mesh`) based on the + room's configuration and the authenticated user id. + + ## Messages + + Inbound events (from client): + + push("offer", %{target: "peer-uuid", sdp: "..."}) + push("answer", %{target: "peer-uuid", sdp: "..."}) + push("ice", %{target: "peer-uuid", candidate: "..."}) + push("broadcast_offer", %{sdp: "..."}) + + Outbound events (to client): + + "offer" — %{sdp: "...", from_peer_id: "..."} + "answer" — %{sdp: "...", from_peer_id: "..."} + "ice" — %{candidate: "...", from_peer_id: "..."} + "peer_joined" — %{peer_id: "...", role: :host | :client | :peer, user_id: "..."} + "peer_left" — %{peer_id: "..."} + "room_closed" — %{} + """ + + use Phoenix.Channel + + import GameServerWeb.ChannelPush + require Logger + + alias GameServerWeb.SignalingBroker + + # WebSocket message rate limits (per user) — defaults, overridden by config + @default_ws_rate_limit 300 + @default_ws_rate_window :timer.seconds(10) + + # Separate ICE candidate budget — prevents ICE flooding from starving + # other channel events. A typical WebRTC session sends 5–30 candidates. + @default_ice_rate_limit 150 + @default_ice_rate_window :timer.seconds(30) + + @impl true + def join("signaling:" <> room_id, _payload, socket) do + user_id = socket.assigns.current_scope.user.id + + if is_nil(user_id) do + Logger.warning("SignalingChannel: unauthorized join attempt room=#{room_id} missing user_id") + {:error, %{reason: "unauthorized"}} + else + user = GameServer.Accounts.get_user(user_id) + is_host = SignalingBroker.is_host?(room_id, user_id) + Logger.warning(is_host) + + if not is_host and (is_nil(user.lobby_id) or user.lobby_id != room_id) do + Logger.warning("SignalingChannel: join rejected not_lobby_member room=#{room_id} user=#{user_id} user_lobby=#{user.lobby_id || "nil"}") + {:error, %{reason: "not_lobby_member"}} + else + peer_id = Ecto.UUID.generate() + + case SignalingBroker.join(room_id, peer_id, self(), user_id, %{}) do + {:ok, role} -> + Logger.info("SignalingChannel: join ok room=#{room_id} peer=#{peer_id} user=#{user_id} role=#{role}") + {:ok, %{peer_id: peer_id, role: role}, + assign(socket, + signaling_room: room_id, + signaling_peer_id: peer_id, + signaling_role: role + )} + + {:error, :room_not_found} -> + Logger.warning("SignalingChannel: join failed room_not_found room=#{room_id} user=#{user_id}") + {:error, %{reason: "room_not_found"}} + + {:error, :duplicate_peer} -> + Logger.warning("SignalingChannel: join failed duplicate_peer room=#{room_id} user=#{user_id}") + {:error, %{reason: "duplicate_peer"}} + end + end + end + end + + # ── Signaling relay ────────────────────────────────────────────────────── + + @impl true + def handle_in("offer", %{"target" => target, "sdp" => sdp}, socket) do + with :ok <- check_ws_rate_limit(socket) do + room = socket.assigns.signaling_room + from = socket.assigns.signaling_peer_id + + case SignalingBroker.relay(room, from, target, :offer, %{sdp: sdp}) do + :ok -> + {:reply, {:ok, %{}}, socket} + + {:error, :peer_not_found} -> + Logger.warning("SignalingChannel: offer failed peer_not_found room=#{room} from=#{from} target=#{target}") + {:reply, {:error, %{error: "peer_not_found"}}, socket} + + {:error, :not_allowed} -> + Logger.warning("SignalingChannel: offer failed not_allowed room=#{room} from=#{from} target=#{target}") + {:reply, {:error, %{error: "not_allowed"}}, socket} + + {:error, :room_not_found} -> + Logger.warning("SignalingChannel: offer failed room_not_found room=#{room} from=#{from}") + {:stop, :normal, {:error, %{error: "room_not_found"}}, socket} + end + end + end + + @impl true + def handle_in("answer", %{"target" => target, "sdp" => sdp}, socket) do + with :ok <- check_ws_rate_limit(socket) do + room = socket.assigns.signaling_room + from = socket.assigns.signaling_peer_id + + case SignalingBroker.relay(room, from, target, :answer, %{sdp: sdp}) do + :ok -> + {:reply, {:ok, %{}}, socket} + + {:error, :peer_not_found} -> + Logger.warning("SignalingChannel: answer failed peer_not_found room=#{room} from=#{from} target=#{target}") + {:reply, {:error, %{error: "peer_not_found"}}, socket} + + {:error, :not_allowed} -> + Logger.warning("SignalingChannel: answer failed not_allowed room=#{room} from=#{from} target=#{target}") + {:reply, {:error, %{error: "not_allowed"}}, socket} + + {:error, :room_not_found} -> + Logger.warning("SignalingChannel: answer failed room_not_found room=#{room} from=#{from}") + {:stop, :normal, {:error, %{error: "room_not_found"}}, socket} + end + end + end + + @impl true + def handle_in("ice", %{"target" => target, "candidate" => candidate}, socket) do + with :ok <- check_ice_rate_limit(socket) do + room = socket.assigns.signaling_room + from = socket.assigns.signaling_peer_id + + case SignalingBroker.relay(room, from, target, :ice, %{candidate: candidate}) do + :ok -> + {:reply, {:ok, %{}}, socket} + + {:error, :peer_not_found} -> + Logger.warning("SignalingChannel: ice failed peer_not_found room=#{room} from=#{from} target=#{target}") + {:reply, {:error, %{error: "peer_not_found"}}, socket} + + {:error, :not_allowed} -> + Logger.warning("SignalingChannel: ice failed not_allowed room=#{room} from=#{from} target=#{target}") + {:reply, {:error, %{error: "not_allowed"}}, socket} + + {:error, :room_not_found} -> + Logger.warning("SignalingChannel: ice failed room_not_found room=#{room} from=#{from}") + {:stop, :normal, {:error, %{error: "room_not_found"}}, socket} + end + end + end + + @impl true + def handle_in("broadcast_offer", %{"sdp" => sdp}, socket) do + with :ok <- check_ws_rate_limit(socket) do + room = socket.assigns.signaling_room + from = socket.assigns.signaling_peer_id + + case SignalingBroker.broadcast(room, from, :offer, %{sdp: sdp, from_peer_id: from}) do + :ok -> + {:reply, {:ok, %{}}, socket} + + {:error, :not_allowed} -> + Logger.warning("SignalingChannel: broadcast_offer failed not_allowed room=#{room} from=#{from}") + {:reply, {:error, %{error: "not_allowed"}}, socket} + + {:error, :room_not_found} -> + Logger.warning("SignalingChannel: broadcast_offer failed room_not_found room=#{room} from=#{from}") + {:stop, :normal, {:error, %{error: "room_not_found"}}, socket} + end + end + end + + @impl true + def handle_in("list_peers", _payload, socket) do + with :ok <- check_ws_rate_limit(socket) do + room = socket.assigns.signaling_room + + case SignalingBroker.list_peers(room) do + peers when is_map(peers) -> + {:reply, {:ok, %{peers: peers}}, socket} + + {:error, :room_not_found} -> + Logger.warning("SignalingChannel: list_peers failed room_not_found room=#{room}") + {:stop, :normal, {:error, %{error: "room_not_found"}}, socket} + end + end + end + + @impl true + def handle_in(event, payload, socket) do + Logger.warning("SignalingChannel: unknown event=#{event} room=#{socket.assigns[:signaling_room] || "nil"} peer=#{socket.assigns[:signaling_peer_id] || "nil"}") + {:reply, {:error, %{error: "unknown_event"}}, socket} + end + + # ── Broker relay messages ──────────────────────────────────────────────── + + @impl true + def handle_info({:signaling_relay, :room_closed, nil, payload}, socket) do + Logger.info("SignalingChannel: room_closed received, stopping room=#{socket.assigns.signaling_room} peer=#{socket.assigns.signaling_peer_id}") + push_event(socket, "room_closed", payload) + {:stop, :normal, socket} + end + + @impl true + def handle_info({:signaling_relay, type, from_peer_id, payload}, socket) do + event_name = relay_event_name(type) + payload = if is_nil(from_peer_id), do: payload, else: Map.put(payload, :from_peer_id, from_peer_id) + push_event(socket, event_name, payload) + {:noreply, socket} + end + + @impl true + def handle_info({:channel_updates_flush, _}, socket) do + {:noreply, socket} + end + + @impl true + def handle_info(msg, socket) do + Logger.debug("SignalingChannel: unexpected msg=#{inspect(msg)} room=#{socket.assigns[:signaling_room] || "nil"} peer=#{socket.assigns[:signaling_peer_id] || "nil"}") + {:noreply, socket} + end + + @impl true + def terminate(reason, socket) do + room_id = socket.assigns[:signaling_room] + peer_id = socket.assigns[:signaling_peer_id] + + if room_id && peer_id do + Logger.info("SignalingChannel: terminating reason=#{inspect(reason)} room=#{room_id} peer=#{peer_id}") + SignalingBroker.leave(room_id, peer_id) + else + Logger.debug("SignalingChannel: terminating without room/peer reason=#{inspect(reason)}") + end + + :ok + end + + # ── Private helpers ─────────────────────────────────────────────────────── + + defp relay_event_name(:offer), do: "offer" + defp relay_event_name(:answer), do: "answer" + defp relay_event_name(:ice), do: "ice" + defp relay_event_name(:peer_joined), do: "peer_joined" + defp relay_event_name(:peer_left), do: "peer_left" + defp relay_event_name(:room_closed), do: "room_closed" + + # ── WebSocket rate limiting ───────────────────────────────────────────── + + defp check_ws_rate_limit(socket) do + config = Application.get_env(:game_server_web, GameServerWeb.Plugs.RateLimiter, []) + + if Keyword.get(config, :enabled, true) do + user_id = socket.assigns.current_scope.user.id + limit = Keyword.get(config, :signaling_ws_limit, @default_ws_rate_limit) + window = Keyword.get(config, :signaling_ws_window, @default_ws_rate_window) + + case GameServerWeb.RateLimit.hit("signaling_ws:#{user_id}", window, limit) do + {:allow, _count} -> + :ok + + {:deny, _retry_after} -> + Logger.warning("SignalingChannel: rate limit exceeded user=#{user_id} room=#{socket.assigns[:signaling_room] || "nil"}") + {:stop, :normal, {:error, %{error: "rate_limited"}}, socket} + end + else + :ok + end + end + + defp check_ice_rate_limit(socket) do + config = Application.get_env(:game_server_web, GameServerWeb.Plugs.RateLimiter, []) + + if Keyword.get(config, :enabled, true) do + user_id = socket.assigns.current_scope.user.id + limit = Keyword.get(config, :signaling_ice_limit, @default_ice_rate_limit) + window = Keyword.get(config, :signaling_ice_window, @default_ice_rate_window) + + case GameServerWeb.RateLimit.hit("signaling_ice:#{user_id}", window, limit) do + {:allow, _count} -> + :ok + + {:deny, _retry_after} -> + Logger.warning("SignalingChannel: ICE rate limit exceeded user=#{user_id} room=#{socket.assigns[:signaling_room] || "nil"}") + {:reply, {:error, %{error: "ice_rate_limited"}}, socket} + end + else + :ok + end + end +end diff --git a/apps/game_server_web/lib/game_server_web/channels/user_socket.ex b/apps/game_server_web/lib/game_server_web/channels/user_socket.ex index bd0c2a5a3..134d5243b 100644 --- a/apps/game_server_web/lib/game_server_web/channels/user_socket.ex +++ b/apps/game_server_web/lib/game_server_web/channels/user_socket.ex @@ -23,7 +23,8 @@ defmodule GameServerWeb.UserSocket do "Global lobby list: created/updated/deleted, membership counts"}, {"group:*", GameServerWeb.GroupChannel, "One group: state, members, join requests, chat"}, {"groups", GameServerWeb.GroupsChannel, "Global group list: created/updated/deleted"}, - {"party:*", GameServerWeb.PartyChannel, "One party: state, members, chat, disband"} + {"party:*", GameServerWeb.PartyChannel, "One party: state, members, chat, disband"}, + {"signaling:*", GameServerWeb.SignalingChannel, "Signaling channel for WebRTC"} ] for {pattern, module, _description} <- @channels do diff --git a/apps/game_server_web/lib/game_server_web/signaling.ex b/apps/game_server_web/lib/game_server_web/signaling.ex new file mode 100644 index 000000000..1030c6abb --- /dev/null +++ b/apps/game_server_web/lib/game_server_web/signaling.ex @@ -0,0 +1,15 @@ +defmodule GameServer.Signaling do + @moduledoc """ + Public API for the WebRTC signaling broker. + + This module is exposed to plugins through the SDK. All operations are + delegated to the internal `GameServerWeb.SignalingBroker` process. + """ + + alias GameServerWeb.SignalingBroker + + defdelegate create_room(room_id, topology, opts \\ []), to: SignalingBroker + defdelegate close_room(room_id), to: SignalingBroker + defdelegate room_exists?(room_id), to: SignalingBroker + defdelegate list_peers(room_id), to: SignalingBroker +end diff --git a/apps/game_server_web/mix.exs b/apps/game_server_web/mix.exs index ead49cced..e2fddcf42 100644 --- a/apps/game_server_web/mix.exs +++ b/apps/game_server_web/mix.exs @@ -24,6 +24,7 @@ defmodule GameServerWeb.MixProject do def application do [ + mod: {GameServerWeb.Application, []}, extra_applications: [:logger] ] end diff --git a/modules/plugins/.gitignore b/modules/plugins/.gitignore new file mode 100644 index 000000000..e9d9cf718 --- /dev/null +++ b/modules/plugins/.gitignore @@ -0,0 +1,2 @@ +ebin/ +_build/ diff --git a/modules/plugins/handle_webrtc/.formatter.exs b/modules/plugins/handle_webrtc/.formatter.exs new file mode 100644 index 000000000..d304ff320 --- /dev/null +++ b/modules/plugins/handle_webrtc/.formatter.exs @@ -0,0 +1,3 @@ +[ + inputs: ["{mix,.formatter}.exs", "{config,lib,test}/**/*.{ex,exs}"] +] diff --git a/modules/plugins/handle_webrtc/.gitignore b/modules/plugins/handle_webrtc/.gitignore new file mode 100644 index 000000000..914532603 --- /dev/null +++ b/modules/plugins/handle_webrtc/.gitignore @@ -0,0 +1,3 @@ +/_build/ +/deps/ +/ebin/ diff --git a/modules/plugins/handle_webrtc/README.md b/modules/plugins/handle_webrtc/README.md new file mode 100644 index 000000000..0a19e8d57 --- /dev/null +++ b/modules/plugins/handle_webrtc/README.md @@ -0,0 +1,35 @@ +# Example hook plugin + +This is a minimal OTP hook plugin example. + +## Build + +From the repo root: + +- `cd modules/plugins_examples/example_hook` +- `mix deps.get` +- `mix compile` + +This compiles the plugin into `_build/dev/lib/example_hook/ebin/` (and generates the `.app` file with `hooks_module`). + +## Package (plugin bundle) + +The server loads plugins from a directory containing an `ebin/` folder. +To create a bundle directory you can drop into `modules/plugins/`: + +- `mix plugin.bundle` + +This also copies transitive runtime dependencies, including: + +- compiled dependency BEAMs into `deps//ebin` +- plugin/runtime `priv/` directories into `priv` and `deps//priv` + +The `priv/` copy is important for dependencies that ship NIFs or other runtime assets. + +## Run locally + +- Set `GAME_SERVER_PLUGINS_DIR=modules/plugins_examples` (or copy the built bundle into `modules/plugins/example_hook`) +- Open the Admin Config page and click **Reload plugins** +- Call a function via the “Hooks - Test RPC” form using: + - `plugin`: `example_hook` + - `fn`: `hello` (or `set_current_user_meta`) diff --git a/modules/plugins/handle_webrtc/config/config.exs b/modules/plugins/handle_webrtc/config/config.exs new file mode 100644 index 000000000..5b66c9106 --- /dev/null +++ b/modules/plugins/handle_webrtc/config/config.exs @@ -0,0 +1,3 @@ +import Config + +# No runtime configuration required for this example plugin. diff --git a/modules/plugins/handle_webrtc/lib/example_hook/v1/example_hook.pb.ex.old b/modules/plugins/handle_webrtc/lib/example_hook/v1/example_hook.pb.ex.old new file mode 100644 index 000000000..74d5868f7 --- /dev/null +++ b/modules/plugins/handle_webrtc/lib/example_hook/v1/example_hook.pb.ex.old @@ -0,0 +1,35 @@ +defmodule ExampleHook.V1.HelloProtoRequest do + @moduledoc false + + use Protobuf, + full_name: "example_hook.v1.HelloProtoRequest", + protoc_gen_elixir_version: "0.17.0", + syntax: :proto3 + + field :name, 1, type: :string + field :repeat, 2, type: :uint32 +end + +defmodule ExampleHook.V1.HelloProtoReply do + @moduledoc false + + use Protobuf, + full_name: "example_hook.v1.HelloProtoReply", + protoc_gen_elixir_version: "0.17.0", + syntax: :proto3 + + field :greeting, 1, type: :string + field :name_length, 2, type: :uint32, json_name: "nameLength" +end + +defmodule ExampleHook.V1.ExampleLoadout do + @moduledoc false + + use Protobuf, + full_name: "example_hook.v1.ExampleLoadout", + protoc_gen_elixir_version: "0.17.0", + syntax: :proto3 + + field :weapon_id, 1, type: :uint32, json_name: "weaponId" + field :perk_ids, 2, repeated: true, type: :uint32, json_name: "perkIds" +end diff --git a/modules/plugins/handle_webrtc/lib/game_server/modules/handle_webrtc_hook.ex b/modules/plugins/handle_webrtc/lib/game_server/modules/handle_webrtc_hook.ex new file mode 100644 index 000000000..af8161240 --- /dev/null +++ b/modules/plugins/handle_webrtc/lib/game_server/modules/handle_webrtc_hook.ex @@ -0,0 +1,15 @@ +defmodule GameServer.Modules.HandleWebRTCHook do + use GameServer.Hooks + + require Logger + + alias GameServer.Signaling + + @impl true + def after_matchmaking_matched(tickets, lobby_id) do + server_user_id = "019f8a7d-485e-7000-bd09-6743f91a74e3" + Logger.warning("Creating WebRTC room") + :ok = Signaling.create_room(lobby_id, :star, host_user_id: server_user_id) + + end +end diff --git a/modules/plugins/handle_webrtc/mix.exs b/modules/plugins/handle_webrtc/mix.exs new file mode 100644 index 000000000..00ece4748 --- /dev/null +++ b/modules/plugins/handle_webrtc/mix.exs @@ -0,0 +1,33 @@ +defmodule HandleWebRTC.MixProject do + use Mix.Project + + def project do + [ + app: :handle_webrtc, + version: "0.1.1", + elixir: "~> 1.20", + start_permanent: Mix.env() == :prod, + deps: deps() + ] + end + + def application do + [ + extra_applications: [:logger], + env: [hooks_module: GameServer.Modules.HandleWebRTCHook] + ] + end + + # NOTE: This example lives inside the main server repo, so we depend on the + # in-repo SDK via a path dependency. + defp deps do + [ + {:game_server_sdk, path: "../../../sdk", runtime: false, optional: true}, + {:game_server_plugin_tools, path: "../../../sdk_tools", runtime: false}, + {:bunt, "~> 1.0"}, + {:phoenix, "~> 1.8.3"}, + # Typed hook payloads (see proto/example_hook.proto). + #{:protobuf, "~> 0.17"} + ] + end +end diff --git a/modules/plugins/handle_webrtc/mix.lock b/modules/plugins/handle_webrtc/mix.lock new file mode 100644 index 000000000..20a9d4a30 --- /dev/null +++ b/modules/plugins/handle_webrtc/mix.lock @@ -0,0 +1,13 @@ +%{ + "bunt": {:hex, :bunt, "1.0.0", "081c2c665f086849e6d57900292b3a161727ab40431219529f13c4ddcf3e7a44", [:mix], [], "hexpm", "dc5f86aa08a5f6fa6b8096f0735c4e76d54ae5c9fa2c143e5a1fc7c1cd9bb6b5"}, + "mime": {:hex, :mime, "2.0.7", "b8d739037be7cd402aee1ba0306edfdef982687ee7e9859bee6198c1e7e2f128", [:mix], [], "hexpm", "6171188e399ee16023ffc5b76ce445eb6d9672e2e241d2df6050f3c771e80ccd"}, + "phoenix": {:hex, :phoenix, "1.8.9", "a63ed0962ed5b903b146dab0ae8eb8387fe478f8171a5e26d56a165f35996fe1", [:mix], [{:bandit, "~> 1.0", [hex: :bandit, repo: "hexpm", optional: true]}, {:jason, "~> 1.0", [hex: :jason, repo: "hexpm", optional: true]}, {:phoenix_pubsub, "~> 2.1", [hex: :phoenix_pubsub, repo: "hexpm", optional: false]}, {:phoenix_template, "~> 1.0", [hex: :phoenix_template, repo: "hexpm", optional: false]}, {:phoenix_view, "~> 2.0", [hex: :phoenix_view, repo: "hexpm", optional: true]}, {:plug, "~> 1.14", [hex: :plug, repo: "hexpm", optional: false]}, {:plug_cowboy, "~> 2.7", [hex: :plug_cowboy, repo: "hexpm", optional: true]}, {:plug_crypto, "~> 1.2 or ~> 2.0", [hex: :plug_crypto, repo: "hexpm", optional: false]}, {:telemetry, "~> 0.4 or ~> 1.0", [hex: :telemetry, repo: "hexpm", optional: false]}, {:websock_adapter, "~> 0.5", [hex: :websock_adapter, repo: "hexpm", optional: false]}], "hexpm", "3477e2dd5a4f61820341169031bdfe21275f659923bea9c5c0ea2aa1c3fcc046"}, + "phoenix_pubsub": {:hex, :phoenix_pubsub, "2.2.0", "ff3a5616e1bed6804de7773b92cbccfc0b0f473faf1f63d7daf1206c7aeaaa6f", [:mix], [], "hexpm", "adc313a5bf7136039f63cfd9668fde73bba0765e0614cba80c06ac9460ff3e96"}, + "phoenix_template": {:hex, :phoenix_template, "1.0.4", "e2092c132f3b5e5b2d49c96695342eb36d0ed514c5b252a77048d5969330d639", [:mix], [{:phoenix_html, "~> 2.14.2 or ~> 3.0 or ~> 4.0", [hex: :phoenix_html, repo: "hexpm", optional: true]}], "hexpm", "2c0c81f0e5c6753faf5cca2f229c9709919aba34fab866d3bc05060c9c444206"}, + "plug": {:hex, :plug, "1.20.3", "56c480c633ec2ce10140e236e15233bf576e1d323887d7c96711bd02ab5160db", [:mix], [{:mime, "~> 1.0 or ~> 2.0", [hex: :mime, repo: "hexpm", optional: false]}, {:plug_crypto, "~> 1.1.1 or ~> 1.2 or ~> 2.0", [hex: :plug_crypto, repo: "hexpm", optional: false]}, {:telemetry, "~> 0.4.3 or ~> 1.0", [hex: :telemetry, repo: "hexpm", optional: false]}], "hexpm", "be266aee1b8536ef6409d58cf39a3121319f0ec47cfa1b24024485aa0e76ad76"}, + "plug_crypto": {:hex, :plug_crypto, "2.1.1", "19bda8184399cb24afa10be734f84a16ea0a2bc65054e23a62bb10f06bc89491", [:mix], [], "hexpm", "6470bce6ffe41c8bd497612ffde1a7e4af67f36a15eea5f921af71cf3e11247c"}, + "protobuf": {:hex, :protobuf, "0.17.0", "39e24e43c9648e148feba16ed51100b5b2028ea900b55460377b0476f6e10613", [:mix], [{:jason, "~> 1.2", [hex: :jason, repo: "hexpm", optional: true]}], "hexpm", "ca6c91f6f63e2c147b47f03eefd10b80538aa6fc55ff4b12b795efb786b0152f"}, + "telemetry": {:hex, :telemetry, "1.4.2", "a0cb522801dffb1c49fe6e30561badffc7b6d0e180db1300df759faa22062855", [:rebar3], [], "hexpm", "928f6495066506077862c0d1646609eed891a4326bee3126ba54b60af61febb1"}, + "websock": {:hex, :websock, "0.5.3", "2f69a6ebe810328555b6fe5c831a851f485e303a7c8ce6c5f675abeb20ebdadc", [:mix], [], "hexpm", "6105453d7fac22c712ad66fab1d45abdf049868f253cf719b625151460b8b453"}, + "websock_adapter": {:hex, :websock_adapter, "0.6.0", "73db5ab8aaefd1a876a97ce3e6afc96562625de69ef17a4e04426e034849d0b8", [:mix], [{:bandit, ">= 0.6.0", [hex: :bandit, repo: "hexpm", optional: true]}, {:plug, "~> 1.14", [hex: :plug, repo: "hexpm", optional: false]}, {:plug_cowboy, "~> 2.6", [hex: :plug_cowboy, repo: "hexpm", optional: true]}, {:websock, "~> 0.5", [hex: :websock, repo: "hexpm", optional: false]}], "hexpm", "50021a85bce8f203b086705d9e0c5415e2c7eb05d319111b0428fe71f9934617"}, +} diff --git a/modules/plugins/handle_webrtc/proto/example_hook.proto b/modules/plugins/handle_webrtc/proto/example_hook.proto new file mode 100644 index 000000000..9c178d67f --- /dev/null +++ b/modules/plugins/handle_webrtc/proto/example_hook.proto @@ -0,0 +1,26 @@ +// Typed hook payloads for the example plugin. +// +// Convention: a typed hook `foo_bar` names its messages FooBarRequest and +// FooBarReply. The gamend server relays these bytes untouched (RtcEnvelope +// args_raw / data_raw) — only this plugin and the game client know the +// schema, and both sides use their own generated bindings. +syntax = "proto3"; + +package example_hook.v1; + +message HelloProtoRequest { + string name = 1; + uint32 repeat = 2; +} + +message HelloProtoReply { + string greeting = 1; + uint32 name_length = 2; +} + +// KV data schema example, registered for the "pb_loadout" key via +// kv_schemas/0 in the hooks module. +message ExampleLoadout { + uint32 weapon_id = 1; + repeated uint32 perk_ids = 2; +} diff --git a/sdk/lib/game_server/signaling.ex b/sdk/lib/game_server/signaling.ex new file mode 100644 index 000000000..e219c2579 --- /dev/null +++ b/sdk/lib/game_server/signaling.ex @@ -0,0 +1,8 @@ +defmodule GameServer.Signaling do + @moduledoc "SDK stub for GameServer.Signaling." + + def create_room(_room_id, _topology, _opts \\ []), do: :ok + def close_room(_room_id), do: :ok + def room_exists?(_room_id), do: false + def list_peers(_room_id), do: %{} +end From c794c35a7b9196bd641799c328cdd5bc148fed4b Mon Sep 17 00:00:00 2001 From: nicocsh <119978472+nicocsh@users.noreply.github.com> Date: Thu, 23 Jul 2026 20:02:38 +0200 Subject: [PATCH 2/3] New Socket User ID fix and fixed some warnings --- .../lib/game_server_web/channels/signaling_broker.ex | 4 ++-- .../lib/game_server_web/channels/signaling_channel.ex | 10 +++++----- 2 files changed, 7 insertions(+), 7 deletions(-) diff --git a/apps/game_server_web/lib/game_server_web/channels/signaling_broker.ex b/apps/game_server_web/lib/game_server_web/channels/signaling_broker.ex index 0d3efc8ae..2eeee00ec 100644 --- a/apps/game_server_web/lib/game_server_web/channels/signaling_broker.ex +++ b/apps/game_server_web/lib/game_server_web/channels/signaling_broker.ex @@ -214,9 +214,9 @@ defmodule GameServerWeb.SignalingBroker do Logger.warning("SignalingBroker: leave failed peer_not_found room=#{room_id} peer=#{peer_id}") {:reply, {:error, :peer_not_found}, state} - {_peer, peers} -> + {peer, peers} -> peer_count = map_size(peers) - Logger.info("SignalingBroker: peer leaving room=#{room_id} peer=#{peer_id} user=#{_peer.user_id} role=#{_peer.role} remaining_peers=#{peer_count}") + Logger.info("SignalingBroker: peer leaving room=#{room_id} peer=#{peer_id} user=#{peer.user_id} role=#{peer.role} remaining_peers=#{peer_count}") for {other_id, %{pid: other_pid}} <- peers do send(other_pid, {:signaling_relay, :peer_left, peer_id, %{peer_id: peer_id}}) diff --git a/apps/game_server_web/lib/game_server_web/channels/signaling_channel.ex b/apps/game_server_web/lib/game_server_web/channels/signaling_channel.ex index 05a6a6407..df39c031c 100644 --- a/apps/game_server_web/lib/game_server_web/channels/signaling_channel.ex +++ b/apps/game_server_web/lib/game_server_web/channels/signaling_channel.ex @@ -52,7 +52,7 @@ defmodule GameServerWeb.SignalingChannel do @impl true def join("signaling:" <> room_id, _payload, socket) do - user_id = socket.assigns.current_scope.user.id + user_id = socket.assigns.current_scope.user_id if is_nil(user_id) do Logger.warning("SignalingChannel: unauthorized join attempt room=#{room_id} missing user_id") @@ -60,7 +60,7 @@ defmodule GameServerWeb.SignalingChannel do else user = GameServer.Accounts.get_user(user_id) is_host = SignalingBroker.is_host?(room_id, user_id) - Logger.warning(is_host) + #Logger.warning(is_host) if not is_host and (is_nil(user.lobby_id) or user.lobby_id != room_id) do Logger.warning("SignalingChannel: join rejected not_lobby_member room=#{room_id} user=#{user_id} user_lobby=#{user.lobby_id || "nil"}") @@ -205,7 +205,7 @@ defmodule GameServerWeb.SignalingChannel do end @impl true - def handle_in(event, payload, socket) do + def handle_in(event, _payload, socket) do Logger.warning("SignalingChannel: unknown event=#{event} room=#{socket.assigns[:signaling_room] || "nil"} peer=#{socket.assigns[:signaling_peer_id] || "nil"}") {:reply, {:error, %{error: "unknown_event"}}, socket} end @@ -268,7 +268,7 @@ defmodule GameServerWeb.SignalingChannel do config = Application.get_env(:game_server_web, GameServerWeb.Plugs.RateLimiter, []) if Keyword.get(config, :enabled, true) do - user_id = socket.assigns.current_scope.user.id + user_id = socket.assigns.current_scope.user_id limit = Keyword.get(config, :signaling_ws_limit, @default_ws_rate_limit) window = Keyword.get(config, :signaling_ws_window, @default_ws_rate_window) @@ -289,7 +289,7 @@ defmodule GameServerWeb.SignalingChannel do config = Application.get_env(:game_server_web, GameServerWeb.Plugs.RateLimiter, []) if Keyword.get(config, :enabled, true) do - user_id = socket.assigns.current_scope.user.id + user_id = socket.assigns.current_scope.user_id limit = Keyword.get(config, :signaling_ice_limit, @default_ice_rate_limit) window = Keyword.get(config, :signaling_ice_window, @default_ice_rate_window) From 88db5f1a97729115b33eab3ef0add522e4325f4d Mon Sep 17 00:00:00 2001 From: nicocsh <119978472+nicocsh@users.noreply.github.com> Date: Sun, 26 Jul 2026 20:41:39 +0200 Subject: [PATCH 3/3] New module to make WebRTC Signaling 1:1 to Lobbies --- .../channels/signaling_broker.ex | 371 ++++++++++++++---- .../channels/signaling_channel.ex | 119 +++--- .../lib/game_server_web/signaling.ex | 6 +- modules/plugins/handle_webrtc/README.md | 35 -- .../example_hook/v1/example_hook.pb.ex.old | 35 -- .../game_server/modules/handle_webrtc_hook.ex | 15 - .../handle_webrtc/proto/example_hook.proto | 26 -- .../.formatter.exs | 0 .../.gitignore | 0 modules/plugins/webrtc_lobby_hook/README.md | 85 ++++ .../config/config.exs | 0 .../game_server/modules/webrtc_lobby_hook.ex | 239 +++++++++++ .../mix.exs | 6 +- .../mix.lock | 0 sdk/lib/game_server/signaling.ex | 4 + 15 files changed, 679 insertions(+), 262 deletions(-) delete mode 100644 modules/plugins/handle_webrtc/README.md delete mode 100644 modules/plugins/handle_webrtc/lib/example_hook/v1/example_hook.pb.ex.old delete mode 100644 modules/plugins/handle_webrtc/lib/game_server/modules/handle_webrtc_hook.ex delete mode 100644 modules/plugins/handle_webrtc/proto/example_hook.proto rename modules/plugins/{handle_webrtc => webrtc_lobby_hook}/.formatter.exs (100%) rename modules/plugins/{handle_webrtc => webrtc_lobby_hook}/.gitignore (100%) create mode 100644 modules/plugins/webrtc_lobby_hook/README.md rename modules/plugins/{handle_webrtc => webrtc_lobby_hook}/config/config.exs (100%) create mode 100644 modules/plugins/webrtc_lobby_hook/lib/game_server/modules/webrtc_lobby_hook.ex rename modules/plugins/{handle_webrtc => webrtc_lobby_hook}/mix.exs (83%) rename modules/plugins/{handle_webrtc => webrtc_lobby_hook}/mix.lock (100%) diff --git a/apps/game_server_web/lib/game_server_web/channels/signaling_broker.ex b/apps/game_server_web/lib/game_server_web/channels/signaling_broker.ex index 2eeee00ec..4816c5416 100644 --- a/apps/game_server_web/lib/game_server_web/channels/signaling_broker.ex +++ b/apps/game_server_web/lib/game_server_web/channels/signaling_broker.ex @@ -3,8 +3,8 @@ defmodule GameServerWeb.SignalingBroker do Signaling relay for WebRTC peer-to-peer and client-server topologies. Rooms are created explicitly by a worker process (e.g. a lobby worker) and - are keyed by the lobby id. Each room stores its topology and, for :star, - the designated host user id. The broker validates membership and topology + are keyed by the lobby id. Each room stores its topology and, for :star, + the designated host user id. The broker validates membership and topology rules on every relay. Does not create PeerConnections or handle media; only routes SDP offers, @@ -17,8 +17,9 @@ defmodule GameServerWeb.SignalingBroker do Clients may only signal to the host; the host may signal to any client. Non-host peers cannot exchange messages directly. - Each peer is monitored via `Process.monitor/1`. When a peer crashes or - disconnects it is automatically removed and the remaining peers are notified. + Each peer is monitored via `Process.monitor/1`. When a peer crashes or + disconnects it enters a grace period so that reconnections keep the same + user_id. If the grace period expires, the remaining peers are notified. """ use GenServer @@ -34,18 +35,23 @@ defmodule GameServerWeb.SignalingBroker do end @doc """ - Creates a signaling room. `room_id` is typically the lobby id. + Creates a signaling room. `room_id` is typically the lobby id. For `:star` topology `host_user_id` is required and designates the user that will act as the authoritative server peer. + + `allowed_users` is a map of `user_id => role` populated by the lobby hook. + `late_join` controls whether users not in the initial list may join later. + `reconnect_timeout` is the grace period in milliseconds before a + disconnected user is removed. """ def create_room(room_id, topology, opts \\ []) when topology in [:mesh, :star] do host_user_id = if topology == :star, do: Keyword.fetch!(opts, :host_user_id), else: nil - GenServer.call(__MODULE__, {:create_room, room_id, topology, host_user_id}) + GenServer.call(__MODULE__, {:create_room, room_id, topology, host_user_id, opts}) end @doc """ - Closes a signaling room. Existing peers are notified with a `room_closed` + Closes a signaling room. Existing peers are notified with a `room_closed` event so their channels can stop gracefully. """ def close_room(room_id) do @@ -57,20 +63,37 @@ defmodule GameServerWeb.SignalingBroker do end @doc """ - Registers a peer in a room. + Registers a peer in a room using the authenticated `user_id`. Returns `{:ok, role}` where `role` is derived from the room topology and - the provided `user_id`. Returns `{:error, :room_not_found}` if the room - does not exist, and `{:error, :duplicate_peer}` if `peer_id` is already - present. + the provided `user_id`. Returns `{:error, :room_not_found}` if the room + does not exist, `{:error, :not_allowed}` if the user is not in the + allowed list and late join is disabled, and `{:ok, role}` on + reconnection. + """ + def join(room_id, user_id, pid, metadata \\ %{}) + when is_binary(room_id) and is_binary(user_id) and is_pid(pid) do + GenServer.call(__MODULE__, {:join, room_id, user_id, pid, metadata}) + end + + def leave(room_id, user_id) when is_binary(room_id) and is_binary(user_id) do + GenServer.call(__MODULE__, {:leave, room_id, user_id}) + end + + @doc """ + Allows a user to join a room after it has been created (late join). + Called by the lobby hook when a new user joins the lobby. """ - def join(room_id, peer_id, pid, user_id, metadata \\ %{}) - when is_binary(room_id) and is_binary(peer_id) and is_pid(pid) and is_binary(user_id) do - GenServer.call(__MODULE__, {:join, room_id, peer_id, pid, user_id, metadata}) + def allow_user(room_id, user_id, role \\ :peer) do + GenServer.call(__MODULE__, {:allow_user, room_id, user_id, role}) end - def leave(room_id, peer_id) when is_binary(room_id) and is_binary(peer_id) do - GenServer.call(__MODULE__, {:leave, room_id, peer_id}) + @doc """ + Removes a user from the allowed list and kicks them if connected. + Called by the lobby hook when a user leaves the lobby. + """ + def disallow_user(room_id, user_id) do + GenServer.call(__MODULE__, {:disallow_user, room_id, user_id}) end @doc """ @@ -79,8 +102,8 @@ defmodule GameServerWeb.SignalingBroker do Enforces topology rules: in `:star` mode a non-host peer may only relay to the host. """ - def relay(room_id, from_peer_id, to_peer_id, type, payload) do - GenServer.call(__MODULE__, {:relay, room_id, from_peer_id, to_peer_id, type, payload}) + def relay(room_id, from_user_id, to_user_id, type, payload) do + GenServer.call(__MODULE__, {:relay, room_id, from_user_id, to_user_id, type, payload}) end @doc """ @@ -88,8 +111,8 @@ defmodule GameServerWeb.SignalingBroker do In `:star` mode only the host may broadcast. """ - def broadcast(room_id, from_peer_id, type, payload) do - GenServer.call(__MODULE__, {:broadcast, room_id, from_peer_id, type, payload}) + def broadcast(room_id, from_user_id, type, payload) do + GenServer.call(__MODULE__, {:broadcast, room_id, from_user_id, type, payload}) end def list_peers(room_id) do @@ -100,6 +123,14 @@ defmodule GameServerWeb.SignalingBroker do GenServer.call(__MODULE__, {:is_host, room_id, user_id}) end + def get_room(room_id) do + GenServer.call(__MODULE__, {:get_room, room_id}) + end + + def update_room_host(room_id, new_host_user_id) do + GenServer.call(__MODULE__, {:update_room_host, room_id, new_host_user_id}) + end + # ── GenServer callbacks ────────────────────────────────────────────────── @impl true @@ -108,18 +139,25 @@ defmodule GameServerWeb.SignalingBroker do end @impl true - def handle_call({:create_room, room_id, topology, host_user_id}, _from, state) do + def handle_call({:create_room, room_id, topology, host_user_id, opts}, _from, state) do if Map.has_key?(state.rooms, room_id) do Logger.warning("SignalingBroker: room already exists room=#{room_id}") {:reply, {:error, :already_exists}, state} else + allowed_users = Keyword.get(opts, :allowed_users, %{}) + late_join = Keyword.get(opts, :late_join, true) + reconnect_timeout = Keyword.get(opts, :reconnect_timeout, 30_000) + room = %{ topology: topology, host_user_id: host_user_id, - peers: %{} + allowed_users: allowed_users, + peers: %{}, + late_join: late_join, + reconnect_timeout: reconnect_timeout } - Logger.info("SignalingBroker: room created room=#{room_id} topology=#{topology} host_user_id=#{host_user_id || "none"}") + Logger.info("SignalingBroker: room created room=#{room_id} topology=#{topology} host_user_id=#{host_user_id || "none"} allowed_users=#{map_size(allowed_users)} late_join=#{late_join}") {:reply, :ok, %{state | rooms: Map.put(state.rooms, room_id, room)}} end end @@ -135,12 +173,12 @@ defmodule GameServerWeb.SignalingBroker do peer_count = map_size(room.peers) Logger.info("SignalingBroker: closing room=#{room_id} topology=#{room.topology} evicting=#{peer_count}") - for {peer_id, %{pid: pid}} <- room.peers do - Logger.debug("SignalingBroker: sending room_closed to peer=#{peer_id} pid=#{inspect(pid)}") + for {user_id, %{pid: pid}} <- room.peers do + Logger.debug("SignalingBroker: sending room_closed to user=#{user_id} pid=#{inspect(pid)}") send(pid, {:signaling_relay, :room_closed, nil, %{}}) end - refs = Enum.reject(state.refs, fn {_ref, {r, _p}} -> r == room_id end) |> Map.new() + refs = Enum.reject(state.refs, fn {_ref, {r, _u}} -> r == room_id end) |> Map.new() {:reply, :ok, %{state | rooms: rooms, refs: refs}} end @@ -152,74 +190,167 @@ defmodule GameServerWeb.SignalingBroker do end @impl true - def handle_call({:join, room_id, peer_id, pid, user_id, metadata}, _from, state) do + def handle_call({:allow_user, room_id, user_id, role}, _from, state) do case Map.get(state.rooms, room_id) do nil -> - Logger.warning("SignalingBroker: join failed room_not_found room=#{room_id} user=#{user_id}") + Logger.warning("SignalingBroker: allow_user failed room_not_found room=#{room_id} user=#{user_id}") {:reply, {:error, :room_not_found}, state} room -> - if Map.has_key?(room.peers, peer_id) do - Logger.warning("SignalingBroker: join failed duplicate_peer room=#{room_id} peer=#{peer_id}") - {:reply, {:error, :duplicate_peer}, state} - else - ref = Process.monitor(pid) + allowed_users = Map.put(room.allowed_users, user_id, role) + room = %{room | allowed_users: allowed_users} + rooms = Map.put(state.rooms, room_id, room) + + Logger.info("SignalingBroker: allowed user room=#{room_id} user=#{user_id} role=#{role}") + {:reply, :ok, %{state | rooms: rooms}} + end + end + + @impl true + def handle_call({:disallow_user, room_id, user_id}, _from, state) do + case Map.get(state.rooms, room_id) do + nil -> + {:reply, {:error, :room_not_found}, state} + + room -> + allowed_users = Map.delete(room.allowed_users, user_id) + room = %{room | allowed_users: allowed_users} + rooms = Map.put(state.rooms, room_id, room) + + {room, rooms, refs} = + if peer = Map.get(room.peers, user_id) do + if peer.disconnect_timer, do: Process.cancel_timer(peer.disconnect_timer) + + {peer, peers} = Map.pop(room.peers, user_id) + send(peer.pid, {:signaling_relay, :room_closed, nil, %{reason: "removed_from_lobby"}}) - role = - case room.topology do - :mesh -> :peer - :star -> if user_id == room.host_user_id, do: :host, else: :client + for {other_id, %{pid: other_pid}} <- peers, other_id != user_id do + send(other_pid, {:signaling_relay, :peer_left, user_id, %{user_id: user_id}}) end - peer = %{ - pid: pid, - user_id: user_id, - role: role, - metadata: metadata, - joined_at: System.monotonic_time(:second) - } + room = %{room | peers: peers} + rooms = Map.put(rooms, room_id, room) - peers = Map.put(room.peers, peer_id, peer) - room = %{room | peers: peers} - rooms = Map.put(state.rooms, room_id, room) - refs = Map.put(state.refs, ref, {room_id, peer_id}) + ref_entry = Enum.find(state.refs, fn {_ref, {r, u}} -> r == room_id and u == user_id end) + refs = if ref_entry, do: Map.delete(state.refs, elem(ref_entry, 0)), else: state.refs + + {room, rooms, refs} + else + {room, rooms, state.refs} + end + + Logger.info("SignalingBroker: disallowed user room=#{room_id} user=#{user_id}") + {:reply, :ok, %{state | rooms: rooms, refs: refs}} + end + end + + @impl true + def handle_call({:join, room_id, user_id, pid, metadata}, _from, state) do + case Map.get(state.rooms, room_id) do + nil -> + Logger.warning("SignalingBroker: join failed room_not_found room=#{room_id} user=#{user_id}") + {:reply, {:error, :room_not_found}, state} + + room -> + allowed_role = Map.get(room.allowed_users, user_id) + + cond do + is_nil(allowed_role) and not room.late_join -> + Logger.warning("SignalingBroker: join failed not_allowed room=#{room_id} user=#{user_id}") + {:reply, {:error, :not_allowed}, state} + + peer = Map.get(room.peers, user_id) -> + # Reconnection: same user_id reconnecting before the grace period expires. + if peer.disconnect_timer, do: Process.cancel_timer(peer.disconnect_timer) - peer_count = map_size(peers) - Logger.info("SignalingBroker: peer joined room=#{room_id} peer=#{peer_id} user=#{user_id} role=#{role} total_peers=#{peer_count}") + ref = Process.monitor(pid) + peer = %{peer | pid: pid, ref: ref, status: :connected, disconnect_timer: nil} + peers = Map.put(room.peers, user_id, peer) + room = %{room | peers: peers} + rooms = Map.put(state.rooms, room_id, room) + refs = Map.put(state.refs, ref, {room_id, user_id}) + + peer_count = map_size(peers) + Logger.info("SignalingBroker: user reconnected room=#{room_id} user=#{user_id} role=#{peer.role} total_peers=#{peer_count}") + + for {other_id, %{pid: other_pid}} <- peers, other_id != user_id do + Logger.debug("SignalingBroker: notifying user=#{other_id} of peer_rejoined user=#{user_id}") + send(other_pid, {:signaling_relay, :peer_rejoined, user_id, %{ + user_id: user_id, + role: peer.role + }}) + end + + {:reply, {:ok, peer.role}, %{state | rooms: rooms, refs: refs}} + + true -> + role = allowed_role || default_role(room, user_id) + ref = Process.monitor(pid) - for {other_id, %{pid: other_pid}} <- room.peers, other_id != peer_id do - Logger.debug("SignalingBroker: notifying peer=#{other_id} of peer_joined peer=#{peer_id}") - send(other_pid, {:signaling_relay, :peer_joined, peer_id, %{ - peer_id: peer_id, + peer = %{ + pid: pid, + ref: ref, + user_id: user_id, role: role, - user_id: user_id - }}) - end + metadata: metadata, + status: :connected, + joined_at: System.monotonic_time(:second), + disconnect_timer: nil + } + + peers = Map.put(room.peers, user_id, peer) + room = %{room | peers: peers} + rooms = Map.put(state.rooms, room_id, room) + refs = Map.put(state.refs, ref, {room_id, user_id}) + + peer_count = map_size(peers) + Logger.info("SignalingBroker: user joined room=#{room_id} user=#{user_id} role=#{role} total_peers=#{peer_count}") + + # Notify existing peers about the newcomer. + for {other_id, %{pid: other_pid}} <- peers, other_id != user_id do + Logger.debug("SignalingBroker: notifying user=#{other_id} of peer_joined user=#{user_id}") + send(other_pid, {:signaling_relay, :peer_joined, user_id, %{ + user_id: user_id, + role: role + }}) + end + + # Notify the newly joined peer about existing peers so it can initiate + # connections (e.g. star clients connecting to the host). + for {other_id, %{role: other_role}} <- peers, other_id != user_id do + Logger.debug("SignalingBroker: seeding existing peer to new user=#{user_id} other=#{other_id} role=#{other_role}") + send(pid, {:signaling_relay, :peer_joined, other_id, %{ + user_id: other_id, + role: other_role + }}) + end - {:reply, {:ok, role}, %{state | rooms: rooms, refs: refs}} + {:reply, {:ok, role}, %{state | rooms: rooms, refs: refs}} end end end @impl true - def handle_call({:leave, room_id, peer_id}, _from, state) do + def handle_call({:leave, room_id, user_id}, _from, state) do case Map.get(state.rooms, room_id) do nil -> - Logger.warning("SignalingBroker: leave failed room_not_found room=#{room_id} peer=#{peer_id}") + Logger.warning("SignalingBroker: leave failed room_not_found room=#{room_id} user=#{user_id}") {:reply, {:error, :room_not_found}, state} room -> - case Map.pop(room.peers, peer_id) do + case Map.pop(room.peers, user_id) do {nil, _} -> - Logger.warning("SignalingBroker: leave failed peer_not_found room=#{room_id} peer=#{peer_id}") + Logger.warning("SignalingBroker: leave failed peer_not_found room=#{room_id} user=#{user_id}") {:reply, {:error, :peer_not_found}, state} {peer, peers} -> + if peer.disconnect_timer, do: Process.cancel_timer(peer.disconnect_timer) + peer_count = map_size(peers) - Logger.info("SignalingBroker: peer leaving room=#{room_id} peer=#{peer_id} user=#{peer.user_id} role=#{peer.role} remaining_peers=#{peer_count}") + Logger.info("SignalingBroker: user leaving room=#{room_id} user=#{user_id} role=#{peer.role} remaining_peers=#{peer_count}") for {other_id, %{pid: other_pid}} <- peers do - send(other_pid, {:signaling_relay, :peer_left, peer_id, %{peer_id: peer_id}}) + send(other_pid, {:signaling_relay, :peer_left, user_id, %{user_id: user_id}}) end room = %{room | peers: peers} @@ -232,8 +363,7 @@ defmodule GameServerWeb.SignalingBroker do Map.put(state.rooms, room_id, room) end - ref_entry = - Enum.find(state.refs, fn {_ref, {r, p}} -> r == room_id and p == peer_id end) + ref_entry = Enum.find(state.refs, fn {_ref, {r, u}} -> r == room_id and u == user_id end) refs = if ref_entry do @@ -297,9 +427,10 @@ defmodule GameServerWeb.SignalingBroker do Logger.warning("SignalingBroker: broadcast failed not_allowed room=#{room_id} from=#{from} role=#{from_peer.role}") {:reply, {:error, :not_allowed}, state} else - targets = Enum.filter(room.peers, fn {peer_id, _} -> peer_id != from end) |> Enum.map(fn {id, _} -> id end) + targets = Enum.filter(room.peers, fn {user_id, _} -> user_id != from end) |> Enum.map(fn {id, _} -> id end) Logger.debug("SignalingBroker: broadcasting room=#{room_id} type=#{type} from=#{from} targets=#{length(targets)}") - for {peer_id, %{pid: pid}} <- room.peers, peer_id != from do + + for {user_id, %{pid: pid}} <- room.peers, user_id != from do send(pid, {:signaling_relay, type, from, payload}) end @@ -318,8 +449,8 @@ defmodule GameServerWeb.SignalingBroker do room -> peers = - Map.new(room.peers, fn {id, peer} -> - {id, %{user_id: peer.user_id, role: peer.role, metadata: peer.metadata}} + Map.new(room.peers, fn {user_id, peer} -> + {user_id, %{user_id: peer.user_id, role: peer.role, metadata: peer.metadata}} end) {:reply, peers, state} @@ -336,6 +467,33 @@ defmodule GameServerWeb.SignalingBroker do end end + def handle_call({:get_room, room_id}, _from, state) do + case Map.get(state.rooms, room_id) do + nil -> + {:reply, {:error, :room_not_found}, state} + + room -> + {:reply, {:ok, %{topology: room.topology, host_user_id: room.host_user_id}}, state} + end + end + + def handle_call({:update_room_host, room_id, new_host_user_id}, _from, state) do + case Map.get(state.rooms, room_id) do + nil -> + {:reply, {:error, :room_not_found}, state} + + room -> + if room.topology != :star do + {:reply, {:error, :not_star}, state} + else + room = %{room | host_user_id: new_host_user_id} + rooms = Map.put(state.rooms, room_id, room) + Logger.info("SignalingBroker: updated host room=#{room_id} host=#{new_host_user_id}") + {:reply, :ok, %{state | rooms: rooms}} + end + end + end + @impl true def handle_info({:DOWN, ref, :process, pid, reason}, state) do case Map.pop(state.refs, ref) do @@ -343,32 +501,62 @@ defmodule GameServerWeb.SignalingBroker do Logger.debug("SignalingBroker: DOWN from unknown pid=#{inspect(pid)} reason=#{inspect(reason)}") {:noreply, state} - {{room_id, peer_id}, refs} -> + {{room_id, user_id}, refs} -> case Map.get(state.rooms, room_id) do nil -> - Logger.warning("SignalingBroker: DOWN for removed room room=#{room_id} peer=#{peer_id} pid=#{inspect(pid)} reason=#{inspect(reason)}") + Logger.warning("SignalingBroker: DOWN for removed room room=#{room_id} user=#{user_id} pid=#{inspect(pid)} reason=#{inspect(reason)}") {:noreply, %{state | refs: refs}} room -> - {_peer, peers} = Map.pop(room.peers, peer_id) - remaining = map_size(peers) - Logger.info("SignalingBroker: peer DOWN room=#{room_id} peer=#{peer_id} pid=#{inspect(pid)} reason=#{inspect(reason)} remaining_peers=#{remaining}") - - for {other_id, %{pid: pid}} <- peers, other_id != peer_id do - send(pid, {:signaling_relay, :peer_left, peer_id, %{peer_id: peer_id}}) + peer = Map.get(room.peers, user_id) + + if peer do + # Start grace period instead of removing immediately so the same + # user_id can reconnect and keep its role. + timer = Process.send_after(self(), {:reconnect_timeout, room_id, user_id}, room.reconnect_timeout) + peer = %{peer | status: :disconnected, disconnect_timer: timer} + peers = Map.put(room.peers, user_id, peer) + room = %{room | peers: peers} + rooms = Map.put(state.rooms, room_id, room) + + Logger.info("SignalingBroker: user disconnected room=#{room_id} user=#{user_id} grace=#{room.reconnect_timeout}ms pid=#{inspect(pid)} reason=#{inspect(reason)}") + {:noreply, %{state | rooms: rooms, refs: refs}} + else + {:noreply, %{state | refs: refs}} end + end + end + end - room = %{room | peers: peers} + @impl true + def handle_info({:reconnect_timeout, room_id, user_id}, state) do + case Map.get(state.rooms, room_id) do + nil -> + {:noreply, state} - rooms = - if map_size(peers) == 0 do - Logger.info("SignalingBroker: room empty after DOWN, removing room=#{room_id}") - Map.delete(state.rooms, room_id) - else - Map.put(state.rooms, room_id, room) - end + room -> + peer = Map.get(room.peers, user_id) + + if peer && peer.status == :disconnected do + {_peer, peers} = Map.pop(room.peers, user_id) + remaining = map_size(peers) + Logger.info("SignalingBroker: reconnect timeout expired room=#{room_id} user=#{user_id} remaining_peers=#{remaining}") - {:noreply, %{state | rooms: rooms, refs: refs}} + for {other_id, %{pid: other_pid}} <- peers, other_id != user_id do + send(other_pid, {:signaling_relay, :peer_left, user_id, %{user_id: user_id}}) + end + + rooms = + if map_size(peers) == 0 do + Logger.info("SignalingBroker: room empty after timeout, removing room=#{room_id}") + Map.delete(state.rooms, room_id) + else + Map.put(state.rooms, room_id, %{room | peers: peers}) + end + + {:noreply, %{state | rooms: rooms}} + else + {:noreply, state} end end end @@ -377,4 +565,13 @@ defmodule GameServerWeb.SignalingBroker do def handle_info(_msg, state) do {:noreply, state} end + + # ── Private helpers ───────────────────────────────────────────────────── + + defp default_role(room, user_id) do + case room.topology do + :mesh -> :peer + :star -> if user_id == room.host_user_id, do: :host, else: :client + end + end end diff --git a/apps/game_server_web/lib/game_server_web/channels/signaling_channel.ex b/apps/game_server_web/lib/game_server_web/channels/signaling_channel.ex index df39c031c..82d014a2b 100644 --- a/apps/game_server_web/lib/game_server_web/channels/signaling_channel.ex +++ b/apps/game_server_web/lib/game_server_web/channels/signaling_channel.ex @@ -4,34 +4,37 @@ defmodule GameServerWeb.SignalingChannel do Topic: `signaling:` - Rooms are created by a worker process (e.g. a lobby worker) through - `SignalingBroker.create_room/3`. Clients may only join if the room exists - and they are a member of the corresponding lobby. The topology and host - are fixed at room creation; clients cannot choose their role. + Rooms are created by the `WebRTCLobbyHook` through `SignalingBroker.create_room/3`. + The allowed-user list is populated by the hook, so this channel does not + need to query the lobby system. The topology and host are fixed at room + creation; clients cannot choose their role. ## Lifecycle - On join a unique `peer_id` is generated. The broker assigns the role - (`:host` or `:client` for `:star`, `:peer` for `:mesh`) based on the - room's configuration and the authenticated user id. + On join the authenticated `user_id` is used directly as the peer identity. + The broker assigns the role (`:host` or `:client` for `:star`, `:peer` for + `:mesh`) based on the room's configuration. If the same user_id reconnects + within the configured grace period, the existing peer is preserved and a + `peer_rejoined` event is broadcast. ## Messages Inbound events (from client): - push("offer", %{target: "peer-uuid", sdp: "..."}) - push("answer", %{target: "peer-uuid", sdp: "..."}) - push("ice", %{target: "peer-uuid", candidate: "..."}) + push("offer", %{target: "user-uuid", sdp: "..."}) + push("answer", %{target: "user-uuid", sdp: "..."}) + push("ice", %{target: "user-uuid", candidate: "..."}) push("broadcast_offer", %{sdp: "..."}) Outbound events (to client): - "offer" — %{sdp: "...", from_peer_id: "..."} - "answer" — %{sdp: "...", from_peer_id: "..."} - "ice" — %{candidate: "...", from_peer_id: "..."} - "peer_joined" — %{peer_id: "...", role: :host | :client | :peer, user_id: "..."} - "peer_left" — %{peer_id: "..."} - "room_closed" — %{} + "offer" — %{sdp: "...", from_user_id: "..."} + "answer" — %{sdp: "...", from_user_id: "..."} + "ice" — %{candidate: "...", from_user_id: "..."} + "peer_joined" — %{user_id: "...", role: :host | :client | :peer} + "peer_rejoined" — %{user_id: "...", role: :host | :client | :peer} + "peer_left" — %{user_id: "..."} + "room_closed" — %{} """ use Phoenix.Channel @@ -46,7 +49,7 @@ defmodule GameServerWeb.SignalingChannel do @default_ws_rate_window :timer.seconds(10) # Separate ICE candidate budget — prevents ICE flooding from starving - # other channel events. A typical WebRTC session sends 5–30 candidates. + # other channel events. A typical WebRTC session sends 5–30 candidates. @default_ice_rate_limit 150 @default_ice_rate_window :timer.seconds(30) @@ -58,34 +61,27 @@ defmodule GameServerWeb.SignalingChannel do Logger.warning("SignalingChannel: unauthorized join attempt room=#{room_id} missing user_id") {:error, %{reason: "unauthorized"}} else - user = GameServer.Accounts.get_user(user_id) - is_host = SignalingBroker.is_host?(room_id, user_id) - #Logger.warning(is_host) - - if not is_host and (is_nil(user.lobby_id) or user.lobby_id != room_id) do - Logger.warning("SignalingChannel: join rejected not_lobby_member room=#{room_id} user=#{user_id} user_lobby=#{user.lobby_id || "nil"}") - {:error, %{reason: "not_lobby_member"}} - else - peer_id = Ecto.UUID.generate() - - case SignalingBroker.join(room_id, peer_id, self(), user_id, %{}) do - {:ok, role} -> - Logger.info("SignalingChannel: join ok room=#{room_id} peer=#{peer_id} user=#{user_id} role=#{role}") - {:ok, %{peer_id: peer_id, role: role}, - assign(socket, - signaling_room: room_id, - signaling_peer_id: peer_id, - signaling_role: role - )} - - {:error, :room_not_found} -> - Logger.warning("SignalingChannel: join failed room_not_found room=#{room_id} user=#{user_id}") - {:error, %{reason: "room_not_found"}} - - {:error, :duplicate_peer} -> - Logger.warning("SignalingChannel: join failed duplicate_peer room=#{room_id} user=#{user_id}") - {:error, %{reason: "duplicate_peer"}} - end + case SignalingBroker.join(room_id, user_id, self(), %{}) do + {:ok, role} -> + Logger.info("SignalingChannel: join ok room=#{room_id} user=#{user_id} role=#{role}") + {:ok, %{user_id: user_id, role: role}, + assign(socket, + signaling_room: room_id, + signaling_user_id: user_id, + signaling_role: role + )} + + {:error, :room_not_found} -> + Logger.warning("SignalingChannel: join failed room_not_found room=#{room_id} user=#{user_id}") + {:error, %{reason: "room_not_found"}} + + {:error, :not_allowed} -> + Logger.warning("SignalingChannel: join failed not_allowed room=#{room_id} user=#{user_id}") + {:error, %{reason: "not_allowed"}} + + {:error, reason} -> + Logger.warning("SignalingChannel: join failed reason=#{reason} room=#{room_id} user=#{user_id}") + {:error, %{reason: to_string(reason)}} end end end @@ -96,7 +92,7 @@ defmodule GameServerWeb.SignalingChannel do def handle_in("offer", %{"target" => target, "sdp" => sdp}, socket) do with :ok <- check_ws_rate_limit(socket) do room = socket.assigns.signaling_room - from = socket.assigns.signaling_peer_id + from = socket.assigns.signaling_user_id case SignalingBroker.relay(room, from, target, :offer, %{sdp: sdp}) do :ok -> @@ -121,7 +117,7 @@ defmodule GameServerWeb.SignalingChannel do def handle_in("answer", %{"target" => target, "sdp" => sdp}, socket) do with :ok <- check_ws_rate_limit(socket) do room = socket.assigns.signaling_room - from = socket.assigns.signaling_peer_id + from = socket.assigns.signaling_user_id case SignalingBroker.relay(room, from, target, :answer, %{sdp: sdp}) do :ok -> @@ -146,7 +142,7 @@ defmodule GameServerWeb.SignalingChannel do def handle_in("ice", %{"target" => target, "candidate" => candidate}, socket) do with :ok <- check_ice_rate_limit(socket) do room = socket.assigns.signaling_room - from = socket.assigns.signaling_peer_id + from = socket.assigns.signaling_user_id case SignalingBroker.relay(room, from, target, :ice, %{candidate: candidate}) do :ok -> @@ -171,9 +167,9 @@ defmodule GameServerWeb.SignalingChannel do def handle_in("broadcast_offer", %{"sdp" => sdp}, socket) do with :ok <- check_ws_rate_limit(socket) do room = socket.assigns.signaling_room - from = socket.assigns.signaling_peer_id + from = socket.assigns.signaling_user_id - case SignalingBroker.broadcast(room, from, :offer, %{sdp: sdp, from_peer_id: from}) do + case SignalingBroker.broadcast(room, from, :offer, %{sdp: sdp, from_user_id: from}) do :ok -> {:reply, {:ok, %{}}, socket} @@ -206,7 +202,7 @@ defmodule GameServerWeb.SignalingChannel do @impl true def handle_in(event, _payload, socket) do - Logger.warning("SignalingChannel: unknown event=#{event} room=#{socket.assigns[:signaling_room] || "nil"} peer=#{socket.assigns[:signaling_peer_id] || "nil"}") + Logger.warning("SignalingChannel: unknown event=#{event} room=#{socket.assigns[:signaling_room] || "nil"} user=#{socket.assigns[:signaling_user_id] || "nil"}") {:reply, {:error, %{error: "unknown_event"}}, socket} end @@ -214,15 +210,15 @@ defmodule GameServerWeb.SignalingChannel do @impl true def handle_info({:signaling_relay, :room_closed, nil, payload}, socket) do - Logger.info("SignalingChannel: room_closed received, stopping room=#{socket.assigns.signaling_room} peer=#{socket.assigns.signaling_peer_id}") + Logger.info("SignalingChannel: room_closed received, stopping room=#{socket.assigns.signaling_room} user=#{socket.assigns.signaling_user_id}") push_event(socket, "room_closed", payload) {:stop, :normal, socket} end @impl true - def handle_info({:signaling_relay, type, from_peer_id, payload}, socket) do + def handle_info({:signaling_relay, type, from_user_id, payload}, socket) do event_name = relay_event_name(type) - payload = if is_nil(from_peer_id), do: payload, else: Map.put(payload, :from_peer_id, from_peer_id) + payload = if is_nil(from_user_id), do: payload, else: Map.put(payload, :from_user_id, from_user_id) push_event(socket, event_name, payload) {:noreply, socket} end @@ -234,20 +230,22 @@ defmodule GameServerWeb.SignalingChannel do @impl true def handle_info(msg, socket) do - Logger.debug("SignalingChannel: unexpected msg=#{inspect(msg)} room=#{socket.assigns[:signaling_room] || "nil"} peer=#{socket.assigns[:signaling_peer_id] || "nil"}") + Logger.debug("SignalingChannel: unexpected msg=#{inspect(msg)} room=#{socket.assigns[:signaling_room] || "nil"} user=#{socket.assigns[:signaling_user_id] || "nil"}") {:noreply, socket} end @impl true def terminate(reason, socket) do room_id = socket.assigns[:signaling_room] - peer_id = socket.assigns[:signaling_peer_id] + user_id = socket.assigns[:signaling_user_id] - if room_id && peer_id do - Logger.info("SignalingChannel: terminating reason=#{inspect(reason)} room=#{room_id} peer=#{peer_id}") - SignalingBroker.leave(room_id, peer_id) + if room_id && user_id do + Logger.info("SignalingChannel: terminating reason=#{inspect(reason)} room=#{room_id} user=#{user_id}") + # Do NOT call SignalingBroker.leave here. The broker's DOWN handler + # starts a grace period so the same user_id can reconnect and keep + # its role. Explicit leave is only used for intentional removal. else - Logger.debug("SignalingChannel: terminating without room/peer reason=#{inspect(reason)}") + Logger.debug("SignalingChannel: terminating without room/user reason=#{inspect(reason)}") end :ok @@ -259,6 +257,7 @@ defmodule GameServerWeb.SignalingChannel do defp relay_event_name(:answer), do: "answer" defp relay_event_name(:ice), do: "ice" defp relay_event_name(:peer_joined), do: "peer_joined" + defp relay_event_name(:peer_rejoined), do: "peer_rejoined" defp relay_event_name(:peer_left), do: "peer_left" defp relay_event_name(:room_closed), do: "room_closed" diff --git a/apps/game_server_web/lib/game_server_web/signaling.ex b/apps/game_server_web/lib/game_server_web/signaling.ex index 1030c6abb..5b053a6d6 100644 --- a/apps/game_server_web/lib/game_server_web/signaling.ex +++ b/apps/game_server_web/lib/game_server_web/signaling.ex @@ -2,7 +2,7 @@ defmodule GameServer.Signaling do @moduledoc """ Public API for the WebRTC signaling broker. - This module is exposed to plugins through the SDK. All operations are + This module is exposed to plugins through the SDK. All operations are delegated to the internal `GameServerWeb.SignalingBroker` process. """ @@ -11,5 +11,9 @@ defmodule GameServer.Signaling do defdelegate create_room(room_id, topology, opts \\ []), to: SignalingBroker defdelegate close_room(room_id), to: SignalingBroker defdelegate room_exists?(room_id), to: SignalingBroker + defdelegate allow_user(room_id, user_id, role \\ :peer), to: SignalingBroker + defdelegate disallow_user(room_id, user_id), to: SignalingBroker defdelegate list_peers(room_id), to: SignalingBroker + defdelegate get_room(room_id), to: SignalingBroker + defdelegate update_room_host(room_id, new_host_user_id), to: SignalingBroker end diff --git a/modules/plugins/handle_webrtc/README.md b/modules/plugins/handle_webrtc/README.md deleted file mode 100644 index 0a19e8d57..000000000 --- a/modules/plugins/handle_webrtc/README.md +++ /dev/null @@ -1,35 +0,0 @@ -# Example hook plugin - -This is a minimal OTP hook plugin example. - -## Build - -From the repo root: - -- `cd modules/plugins_examples/example_hook` -- `mix deps.get` -- `mix compile` - -This compiles the plugin into `_build/dev/lib/example_hook/ebin/` (and generates the `.app` file with `hooks_module`). - -## Package (plugin bundle) - -The server loads plugins from a directory containing an `ebin/` folder. -To create a bundle directory you can drop into `modules/plugins/`: - -- `mix plugin.bundle` - -This also copies transitive runtime dependencies, including: - -- compiled dependency BEAMs into `deps//ebin` -- plugin/runtime `priv/` directories into `priv` and `deps//priv` - -The `priv/` copy is important for dependencies that ship NIFs or other runtime assets. - -## Run locally - -- Set `GAME_SERVER_PLUGINS_DIR=modules/plugins_examples` (or copy the built bundle into `modules/plugins/example_hook`) -- Open the Admin Config page and click **Reload plugins** -- Call a function via the “Hooks - Test RPC” form using: - - `plugin`: `example_hook` - - `fn`: `hello` (or `set_current_user_meta`) diff --git a/modules/plugins/handle_webrtc/lib/example_hook/v1/example_hook.pb.ex.old b/modules/plugins/handle_webrtc/lib/example_hook/v1/example_hook.pb.ex.old deleted file mode 100644 index 74d5868f7..000000000 --- a/modules/plugins/handle_webrtc/lib/example_hook/v1/example_hook.pb.ex.old +++ /dev/null @@ -1,35 +0,0 @@ -defmodule ExampleHook.V1.HelloProtoRequest do - @moduledoc false - - use Protobuf, - full_name: "example_hook.v1.HelloProtoRequest", - protoc_gen_elixir_version: "0.17.0", - syntax: :proto3 - - field :name, 1, type: :string - field :repeat, 2, type: :uint32 -end - -defmodule ExampleHook.V1.HelloProtoReply do - @moduledoc false - - use Protobuf, - full_name: "example_hook.v1.HelloProtoReply", - protoc_gen_elixir_version: "0.17.0", - syntax: :proto3 - - field :greeting, 1, type: :string - field :name_length, 2, type: :uint32, json_name: "nameLength" -end - -defmodule ExampleHook.V1.ExampleLoadout do - @moduledoc false - - use Protobuf, - full_name: "example_hook.v1.ExampleLoadout", - protoc_gen_elixir_version: "0.17.0", - syntax: :proto3 - - field :weapon_id, 1, type: :uint32, json_name: "weaponId" - field :perk_ids, 2, repeated: true, type: :uint32, json_name: "perkIds" -end diff --git a/modules/plugins/handle_webrtc/lib/game_server/modules/handle_webrtc_hook.ex b/modules/plugins/handle_webrtc/lib/game_server/modules/handle_webrtc_hook.ex deleted file mode 100644 index af8161240..000000000 --- a/modules/plugins/handle_webrtc/lib/game_server/modules/handle_webrtc_hook.ex +++ /dev/null @@ -1,15 +0,0 @@ -defmodule GameServer.Modules.HandleWebRTCHook do - use GameServer.Hooks - - require Logger - - alias GameServer.Signaling - - @impl true - def after_matchmaking_matched(tickets, lobby_id) do - server_user_id = "019f8a7d-485e-7000-bd09-6743f91a74e3" - Logger.warning("Creating WebRTC room") - :ok = Signaling.create_room(lobby_id, :star, host_user_id: server_user_id) - - end -end diff --git a/modules/plugins/handle_webrtc/proto/example_hook.proto b/modules/plugins/handle_webrtc/proto/example_hook.proto deleted file mode 100644 index 9c178d67f..000000000 --- a/modules/plugins/handle_webrtc/proto/example_hook.proto +++ /dev/null @@ -1,26 +0,0 @@ -// Typed hook payloads for the example plugin. -// -// Convention: a typed hook `foo_bar` names its messages FooBarRequest and -// FooBarReply. The gamend server relays these bytes untouched (RtcEnvelope -// args_raw / data_raw) — only this plugin and the game client know the -// schema, and both sides use their own generated bindings. -syntax = "proto3"; - -package example_hook.v1; - -message HelloProtoRequest { - string name = 1; - uint32 repeat = 2; -} - -message HelloProtoReply { - string greeting = 1; - uint32 name_length = 2; -} - -// KV data schema example, registered for the "pb_loadout" key via -// kv_schemas/0 in the hooks module. -message ExampleLoadout { - uint32 weapon_id = 1; - repeated uint32 perk_ids = 2; -} diff --git a/modules/plugins/handle_webrtc/.formatter.exs b/modules/plugins/webrtc_lobby_hook/.formatter.exs similarity index 100% rename from modules/plugins/handle_webrtc/.formatter.exs rename to modules/plugins/webrtc_lobby_hook/.formatter.exs diff --git a/modules/plugins/handle_webrtc/.gitignore b/modules/plugins/webrtc_lobby_hook/.gitignore similarity index 100% rename from modules/plugins/handle_webrtc/.gitignore rename to modules/plugins/webrtc_lobby_hook/.gitignore diff --git a/modules/plugins/webrtc_lobby_hook/README.md b/modules/plugins/webrtc_lobby_hook/README.md new file mode 100644 index 000000000..bfffb2dc2 --- /dev/null +++ b/modules/plugins/webrtc_lobby_hook/README.md @@ -0,0 +1,85 @@ +# WebRTC Lobby Hook + +Keeps a Phoenix WebRTC signaling room in sync with a game lobby. When a lobby has WebRTC enabled in its metadata, this hook automatically creates, updates and tears down a matching signaling room managed by the SignalingBroker. + +## What it does + +- Creates a signaling room automatically when a WebRTC-enabled lobby is created. +- Closes the signaling room when the lobby is deleted. +- Syncs the allowed-user list with lobby joins and leaves. +- Notifies the designated star-topology host so a headless server can join the signaling channel automatically. +- Supports late joins and role assignment (host, client, peer). + +## Supported topologies + +| Topology | Behaviour | +|---|---| +| star | One host and many clients. The host is notified via user: with webrtc:room_ready. Clients send offers to the host. | +| mesh | All participants are peers. No automatic host notification. | + +## Lobby metadata configuration + +Set metadata.webrtc on the lobby: + +```Elixir +%{ + "webrtc" => %{ + "enabled" => true, + "topology" => "star", + "host_user_id" => "star-topology-server-user-id", + "late_join" => true, + "reconnect_timeout" => 30000 + } +} +``` + +### Options + +- enabled — Set to true to create the signaling room. +- topology — star or mesh. +- host_user_id — Optional. In star mode this user is assigned the :host role. Falls back to lobby.host_id. +- late_join — Allow users who join the lobby later to enter the signaling room. Defaults to true. +- reconnect_timeout — Grace period for disconnected peers in milliseconds. Defaults to 30000. + +## Hook callbacks + +| Callback | Purpose | +|---|---| +| before_lobby_create/1 | Injects default WebRTC metadata into lobby attributes if missing. | +| after_lobby_create/1 | Creates the signaling room and seeds the allowed-user list from lobby members. | +| after_lobby_updated/1 | Creates or closes the room if WebRTC is toggled. | +| after_lobby_deleted/1 | Closes the signaling room. | +| after_lobby_join/2 | Allows the joining user into the signaling room with the correct role. | +| after_lobby_leave/2 | Removes the user from the signaling room. | +| after_lobby_host_change/2 | Updates the star host and notifies the new host. | + +## Star topology specifics + +1. The host is resolved in this order: metadata.webrtc.host_user_id -> lobby.host_id. +2. The host is always added to the allowed-user list even if it is not a lobby member (useful for headless servers). +3. When the room is created, the host receives a broadcast on user:: + +```Elixir +%{ + "lobby_id" => lobby_id, + "topology" => "star", + "host_user_id" => host_user_id, + "signaling_topic" => "signaling:#{lobby_id}" +} +``` + +## Role assignment + +Roles are assigned automatically: + +- star + matching host ID -> :host +- star + anyone else -> :client +- mesh -> :peer + +The SignalingBroker uses these roles to enforce who can send offers and answers. + +## Notes + +- The hook normalizes all metadata keys to strings so that Phoenix changesets and Ecto interop cleanly. +- The signaling room ID is the same as the lobby ID, so they are always 1-to-1. +- If the signaling room already exists when a lobby is updated, the hook leaves it alone. diff --git a/modules/plugins/handle_webrtc/config/config.exs b/modules/plugins/webrtc_lobby_hook/config/config.exs similarity index 100% rename from modules/plugins/handle_webrtc/config/config.exs rename to modules/plugins/webrtc_lobby_hook/config/config.exs diff --git a/modules/plugins/webrtc_lobby_hook/lib/game_server/modules/webrtc_lobby_hook.ex b/modules/plugins/webrtc_lobby_hook/lib/game_server/modules/webrtc_lobby_hook.ex new file mode 100644 index 000000000..c2b779636 --- /dev/null +++ b/modules/plugins/webrtc_lobby_hook/lib/game_server/modules/webrtc_lobby_hook.ex @@ -0,0 +1,239 @@ +defmodule GameServer.Modules.WebRTCLobbyHook do + @moduledoc """ + Keeps a WebRTC signaling room in sync with a lobby. + + This is the only module that connects the lobby system to the WebRTC + signaling layer. When a lobby has `metadata.webrtc.enabled = true`, a + signaling room with the same id as the lobby is created automatically. + The room is closed when the lobby is deleted, and the allowed-user list + is kept in sync with lobby joins and leaves. + + When a star-topology room is created, the designated host is notified on + its user channel (`user:`) with a `webrtc:room_ready` event + so a headless server can connect automatically. + + Configuration is read from `lobby.metadata.webrtc`: + + %{ + "enabled" => true, + "topology" => "star" | "mesh", + "late_join" => true, + "reconnect_timeout" => 30000, + "host_user_id" => "optional-server-user-id" + } + + In `:star` mode the host is resolved in this order: + 1. `metadata.webrtc.host_user_id` + 2. `lobby.host_id` + + The allowed-user list is seeded from the lobby members at creation time. + Late joiners are added via `after_lobby_join/2`. + """ + + use GameServer.Hooks + + require Logger + + alias GameServer.Signaling + alias GameServer.Lobbies + alias GameServerWeb.Endpoint + + # Force WebRTC Star for all lobbies. + @impl true + def before_lobby_create(attrs) do + metadata = Map.get(attrs, :metadata) || Map.get(attrs, "metadata") || %{} + + metadata = + Map.new(metadata, fn {k, v} -> + {to_string(k), v} + end) + + webrtc_meta = %{ + "webrtc" => %{ + "enabled" => true, + "topology" => "star", + "late_join" => true, + "reconnect_timeout" => 30000, + "host_user_id" => "example_host_id" + } + } + + new_metadata = Map.merge(metadata, webrtc_meta) + + metadata_key = + cond do + Map.has_key?(attrs, "metadata") -> "metadata" + Map.has_key?(attrs, :metadata) -> :metadata + true -> "metadata" + end + + new_attrs = Map.put(attrs, metadata_key, new_metadata) + + {:ok, new_attrs} + end + + @impl true + def after_lobby_create(lobby) do + ensure_room(lobby) + end + + @impl true + def after_lobby_updated(lobby) do + # If WebRTC is enabled later, create the room. If disabled, close it. + ensure_room(lobby) + end + + @impl true + def after_lobby_deleted(lobby) do + Logger.info("WebRTC: closing signaling room for deleted lobby=#{lobby.id}") + Signaling.close_room(lobby.id) + end + + @impl true + def after_lobby_join(user, lobby) do + # Late join: allow the user into the signaling room. + role = role_for(user.id, lobby) + Logger.info("WebRTC: late join allowed lobby=#{lobby.id} user=#{user.id} role=#{role}") + Signaling.allow_user(lobby.id, user.id, role) + end + + @impl true + def after_lobby_leave(user, lobby) do + Logger.info("WebRTC: user left lobby, removing from signaling room lobby=#{lobby.id} user=#{user.id}") + Signaling.disallow_user(lobby.id, user.id) + end + + @impl true + def after_lobby_host_change(lobby, new_host_id) do + # In star topology, update the host user id when the lobby host changes. + with {:ok, %{topology: :star}} <- Signaling.get_room(lobby.id) do + Logger.info("WebRTC: updating star host lobby=#{lobby.id} host=#{new_host_id}") + Signaling.allow_user(lobby.id, new_host_id, :host) + Signaling.update_room_host(lobby.id, new_host_id) + notify_host_ready(lobby.id, :star, new_host_id) + else + _ -> :ok + end + end + + # ── Private helpers ───────────────────────────────────────────────────── + + defp ensure_room(lobby) do + with %{"webrtc" => %{"enabled" => true, "topology" => topology}} <- lobby.metadata, + topology_atom <- parse_topology(topology), + {:ok, host_user_id} <- resolve_host(lobby, topology_atom) do + if Signaling.room_exists?(lobby.id) do + :ok + else + allowed_users = build_allowed_users(lobby, host_user_id, topology_atom) + late_join = get_in(lobby.metadata, ["webrtc", "late_join"]) || true + reconnect_timeout = get_in(lobby.metadata, ["webrtc", "reconnect_timeout"]) || 30_000 + + Logger.info("WebRTC: creating signaling room lobby=#{lobby.id} topology=#{topology} host=#{host_user_id} allowed_users=#{map_size(allowed_users)}") + + :ok = Signaling.create_room(lobby.id, topology_atom, + host_user_id: host_user_id, + allowed_users: allowed_users, + late_join: late_join, + reconnect_timeout: reconnect_timeout + ) + + # Notify the host so a headless server can join automatically. + if topology_atom == :star do + notify_host_ready(lobby.id, topology_atom, host_user_id) + end + + :ok + end + else + _ -> + # WebRTC not enabled or invalid config; close room if it exists. + if Signaling.room_exists?(lobby.id) do + Signaling.close_room(lobby.id) + end + + :ok + end + end + + # Broadcasts a notification to the host's user channel so the headless + # server can connect to the signaling room automatically. + defp notify_host_ready(lobby_id, topology, host_user_id) do + Logger.info("WebRTC: notifying host user=#{host_user_id} of ready room=#{lobby_id}") + + Endpoint.broadcast("user:#{host_user_id}", "webrtc:room_ready", %{ + "lobby_id" => lobby_id, + "topology" => to_string(topology), + "host_user_id" => host_user_id, + "signaling_topic" => "signaling:#{lobby_id}" + }) + + Logger.info("WebRTC: notifying host=#{host_user_id} about signaling room lobby=#{lobby_id}") + end + + defp build_allowed_users(lobby, host_user_id, topology) do + # Returns lobby members to build the list of allowed users. + members = Lobbies.get_lobby_members(lobby) + + # Convert everything to strings so Ecto UUIDs, binaries and atoms all + # interop cleanly. The host_user_id is always forced to :host. + host_id = to_string(host_user_id) + + users = + Map.new(members, fn member -> + id = to_string(member.id) + + role = + cond do + topology == :star and id == host_id -> :host + topology == :star -> :client + true -> :peer + end + + {id, role} + end) + + # Ensure the dedicated host is allowed even if it is not a lobby member. + # This is critical for headless servers that never join the lobby itself. + users = + if topology == :star and is_binary(host_user_id) and host_user_id != "" do + Map.put(users, host_id, :host) + else + users + end + + Logger.info("WebRTC: build_allowed_users host=#{host_id} users=#{inspect(users)}") + users + end + + defp role_for(user_id, lobby) do + webrtc = lobby.metadata["webrtc"] || %{} + topology = parse_topology(webrtc["topology"] || "mesh") + host_user_id = webrtc["host_user_id"] || lobby.host_id + + cond do + topology == :star and to_string(user_id) == to_string(host_user_id) -> :host + topology == :star -> :client + true -> :peer + end + end + + defp resolve_host(%{metadata: %{"webrtc" => %{"host_user_id" => host_id}}}, :star) + when is_binary(host_id) and host_id != "", + do: {:ok, host_id} + + defp resolve_host(%{host_id: host_id}, :star) + when is_binary(host_id) and host_id != "", + do: {:ok, host_id} + + defp resolve_host(_, :star), do: {:error, :no_host_for_star} + defp resolve_host(_, :mesh), do: {:ok, nil} + + defp parse_topology("star"), do: :star + defp parse_topology("mesh"), do: :mesh + + defp parse_topology(other) do + Logger.warning("WebRTC: unknown topology=#{other}, defaulting to mesh") + :mesh + end +end diff --git a/modules/plugins/handle_webrtc/mix.exs b/modules/plugins/webrtc_lobby_hook/mix.exs similarity index 83% rename from modules/plugins/handle_webrtc/mix.exs rename to modules/plugins/webrtc_lobby_hook/mix.exs index 00ece4748..d149b6937 100644 --- a/modules/plugins/handle_webrtc/mix.exs +++ b/modules/plugins/webrtc_lobby_hook/mix.exs @@ -1,9 +1,9 @@ -defmodule HandleWebRTC.MixProject do +defmodule HandleWebRTC.WebRTCLobbyHook do use Mix.Project def project do [ - app: :handle_webrtc, + app: :webrtc_lobby_hook, version: "0.1.1", elixir: "~> 1.20", start_permanent: Mix.env() == :prod, @@ -14,7 +14,7 @@ defmodule HandleWebRTC.MixProject do def application do [ extra_applications: [:logger], - env: [hooks_module: GameServer.Modules.HandleWebRTCHook] + env: [hooks_module: GameServer.Modules.WebRTCLobbyHook] ] end diff --git a/modules/plugins/handle_webrtc/mix.lock b/modules/plugins/webrtc_lobby_hook/mix.lock similarity index 100% rename from modules/plugins/handle_webrtc/mix.lock rename to modules/plugins/webrtc_lobby_hook/mix.lock diff --git a/sdk/lib/game_server/signaling.ex b/sdk/lib/game_server/signaling.ex index e219c2579..a39858b16 100644 --- a/sdk/lib/game_server/signaling.ex +++ b/sdk/lib/game_server/signaling.ex @@ -4,5 +4,9 @@ defmodule GameServer.Signaling do def create_room(_room_id, _topology, _opts \\ []), do: :ok def close_room(_room_id), do: :ok def room_exists?(_room_id), do: false + def allow_user(_room_id, _user_id, _role \\ :peer), do: :ok + def disallow_user(_room_id, _user_id), do: :ok def list_peers(_room_id), do: %{} + def get_room(_room_id), do: {:error, :room_not_found} + def update_room_host(_room_id, _new_host_user_id), do: :ok end