diff --git a/apps/game_server_core/lib/game_server/signaling.ex b/apps/game_server_core/lib/game_server/signaling.ex new file mode 100644 index 000000000..e05402116 --- /dev/null +++ b/apps/game_server_core/lib/game_server/signaling.ex @@ -0,0 +1,63 @@ +defmodule GameServer.Signaling do + @moduledoc """ + Public API for the WebRTC signaling server. + + This module delegates to `GameServer.Signaling.Server`, which is the + GenServer that actually manages rooms and relays messages. + """ + + alias GameServer.Signaling.Server + + def create_room(room_id, topology, opts \\ []) when topology in [:mesh, :star] do + Server.create_room(room_id, topology, opts) + end + + def close_room(room_id) do + Server.close_room(room_id) + end + + def exists_room?(room_id) do + Server.exists_room?(room_id) + end + + def join_room(room_id, user_id, pid, metadata \\ %{}) + when is_binary(room_id) and is_binary(user_id) and is_pid(pid) do + Server.join_room(room_id, user_id, pid, metadata) + end + + def leave_room(room_id, user_id) when is_binary(room_id) and is_binary(user_id) do + Server.leave_room(room_id, user_id) + end + + def allow_user(room_id, user_id, role \\ :user) do + Server.allow_user(room_id, user_id, role) + end + + def disallow_user(room_id, user_id) do + Server.disallow_user(room_id, user_id) + end + + def relay_message(room_id, from_user_id, to_user_id, type, payload) do + Server.relay_message(room_id, from_user_id, to_user_id, type, payload) + end + + def broadcast_message(room_id, from_user_id, type, payload) do + Server.broadcast_message(room_id, from_user_id, type, payload) + end + + def list_users(room_id) do + Server.list_users(room_id) + end + + def room_host?(room_id, user_id) do + Server.room_host?(room_id, user_id) + end + + def get_room(room_id) do + Server.get_room(room_id) + end + + def update_room_host(room_id, new_host_user_id) do + Server.update_room_host(room_id, new_host_user_id) + end +end diff --git a/apps/game_server_core/lib/game_server/signaling/server.ex b/apps/game_server_core/lib/game_server/signaling/server.ex new file mode 100644 index 000000000..eb04d044d --- /dev/null +++ b/apps/game_server_core/lib/game_server/signaling/server.ex @@ -0,0 +1,692 @@ +defmodule GameServer.Signaling.Server do + @moduledoc """ + Signaling relay for WebRTC user-to-user 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 server 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 users in a room. + + ## Topologies + + * `:mesh` — any member may send an offer/answer/ICE to any other member. + * `:star` — one host user (the Godot headless server) and client users. + Clients may only signal to the host; the host may signal to any client. + Non-host users cannot exchange messages directly. + + Each user is monitored via `Process.monitor/1`. When a user crashes or + disconnects it enters a grace period so that reconnections keep the same + user_id. If the grace period expires, the remaining users 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 + + @impl true + def init(_opts) do + if enabled?() do + Logger.info("SignalingServer: started") + else + Logger.info("SignalingServer: disabled, idle") + end + + {:ok, %__MODULE__{}} + end + + defp enabled? do + Application.get_env(:game_server_core, __MODULE__, [])[:enabled] != false + 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 user. + + `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, opts}) + end + + @doc """ + Closes a signaling room. Existing users 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 exists_room?(room_id) do + GenServer.call(__MODULE__, {:room_exists, room_id}) + end + + @doc """ + Registers a user 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, `{: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(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, room_id, user_id, pid, metadata}) + end + + def leave_room(room_id, user_id) when is_binary(room_id) and is_binary(user_id) do + GenServer.call(__MODULE__, {:leave_room, 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 allow_user(room_id, user_id, role \\ :user) do + GenServer.call(__MODULE__, {:allow_user, room_id, user_id, role}) + end + + @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 """ + Routes a signaling message from one user to a specific target. + + Enforces topology rules: in `:star` mode a non-host user may only relay + to the host. + """ + def relay_message(room_id, from_user_id, to_user_id, type, payload) do + GenServer.call(__MODULE__, {:relay_message, room_id, from_user_id, to_user_id, type, payload}) + end + + @doc """ + Broadcasts a signaling message to every other user in the room. + + In `:star` mode only the host may broadcast. + """ + def broadcast_message(room_id, from_user_id, type, payload) do + GenServer.call(__MODULE__, {:broadcast_message, room_id, from_user_id, type, payload}) + end + + def list_users(room_id) do + GenServer.call(__MODULE__, {:list_users, room_id}) + end + + def room_host?(room_id, user_id) do + GenServer.call(__MODULE__, {:room_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 + 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("SignalingServer: 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, + allowed_users: allowed_users, + users: %{}, + late_join: late_join, + reconnect_timeout: reconnect_timeout + } + + Logger.info( + "SignalingServer: 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 + + @impl true + def handle_call({:close_room, room_id}, _from, state) do + case Map.pop(state.rooms, room_id) do + {nil, _} -> + Logger.warning("SignalingServer: close_room for non-existent room=#{room_id}") + {:reply, {:error, :room_not_found}, state} + + {room, rooms} -> + user_count = map_size(room.users) + + Logger.info( + "SignalingServer: closing room=#{room_id} topology=#{room.topology} evicting=#{user_count}" + ) + + for {user_id, %{pid: pid}} <- room.users do + Logger.debug( + "SignalingServer: 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, _u}} -> 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({:allow_user, room_id, user_id, role}, _from, state) do + case Map.get(state.rooms, room_id) do + nil -> + Logger.warning( + "SignalingServer: allow_user failed room_not_found room=#{room_id} user=#{user_id}" + ) + + {:reply, {:error, :room_not_found}, state} + + room -> + 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("SignalingServer: 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 user = Map.get(room.users, user_id) do + if user.disconnect_timer, do: Process.cancel_timer(user.disconnect_timer) + + {user, users} = Map.pop(room.users, user_id) + send(user.pid, {:signaling_relay, :room_closed, nil, %{reason: "removed_from_lobby"}}) + + for {other_id, %{pid: other_pid}} <- users, other_id != user_id do + send(other_pid, {:signaling_relay, :user_left, user_id, %{user_id: user_id}}) + end + + room = %{room | users: users} + rooms = Map.put(rooms, room_id, room) + + 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("SignalingServer: disallowed user room=#{room_id} user=#{user_id}") + {:reply, :ok, %{state | rooms: rooms, refs: refs}} + end + end + + @impl true + def handle_call({:join_room, room_id, user_id, pid, metadata}, _from, state) do + case Map.get(state.rooms, room_id) do + nil -> + Logger.warning( + "SignalingServer: join_room 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( + "SignalingServer: join_room failed not_allowed room=#{room_id} user=#{user_id}" + ) + + {:reply, {:error, :not_allowed}, state} + + user = Map.get(room.users, user_id) -> + # Reconnection: same user_id reconnecting before the grace period expires. + if user.disconnect_timer, do: Process.cancel_timer(user.disconnect_timer) + + ref = Process.monitor(pid) + user = %{user | pid: pid, ref: ref, status: :connected, disconnect_timer: nil} + users = Map.put(room.users, user_id, user) + room = %{room | users: users} + rooms = Map.put(state.rooms, room_id, room) + refs = Map.put(state.refs, ref, {room_id, user_id}) + + user_count = map_size(users) + + Logger.info( + "SignalingServer: user reconnected room=#{room_id} user=#{user_id} role=#{user.role} total_users=#{user_count}" + ) + + for {other_id, %{pid: other_pid}} <- users, other_id != user_id do + Logger.debug( + "SignalingServer: notifying user=#{other_id} of user_rejoined user=#{user_id}" + ) + + send( + other_pid, + {:signaling_relay, :user_rejoined, user_id, + %{ + user_id: user_id, + role: user.role + }} + ) + end + + {:reply, {:ok, user.role}, %{state | rooms: rooms, refs: refs}} + + true -> + role = allowed_role || default_role(room, user_id) + ref = Process.monitor(pid) + + user = %{ + pid: pid, + ref: ref, + user_id: user_id, + role: role, + metadata: metadata, + status: :connected, + joined_at: System.monotonic_time(:second), + disconnect_timer: nil + } + + users = Map.put(room.users, user_id, user) + room = %{room | users: users} + rooms = Map.put(state.rooms, room_id, room) + refs = Map.put(state.refs, ref, {room_id, user_id}) + + user_count = map_size(users) + + Logger.info( + "SignalingServer: user joined room=#{room_id} user=#{user_id} role=#{role} total_users=#{user_count}" + ) + + # Notify existing peers about the newcomer. + for {other_id, %{pid: other_pid}} <- users, other_id != user_id do + Logger.debug( + "SignalingServer: notifying user=#{other_id} of user_joined user=#{user_id}" + ) + + send( + other_pid, + {:signaling_relay, :user_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}} <- users, other_id != user_id do + Logger.debug( + "SignalingServer: seeding existing user to new user=#{user_id} other=#{other_id} role=#{other_role}" + ) + + send( + pid, + {:signaling_relay, :user_joined, other_id, %{user_id: other_id, role: other_role}} + ) + end + + {:reply, {:ok, role}, %{state | rooms: rooms, refs: refs}} + end + end + end + + @impl true + def handle_call({:leave_room, room_id, user_id}, _from, state) do + case Map.get(state.rooms, room_id) do + nil -> + Logger.warning( + "SignalingServer: leave_room failed room_not_found room=#{room_id} user=#{user_id}" + ) + + {:reply, {:error, :room_not_found}, state} + + room -> + case Map.pop(room.users, user_id) do + {nil, _} -> + Logger.warning( + "SignalingServer: leave_room failed user_not_found room=#{room_id} user=#{user_id}" + ) + + {:reply, {:error, :user_not_found}, state} + + {user, users} -> + if user.disconnect_timer, do: Process.cancel_timer(user.disconnect_timer) + + user_count = map_size(users) + + Logger.info( + "SignalingServer: user leaving room=#{room_id} user=#{user_id} role=#{user.role} remaining_users=#{user_count}" + ) + + for {_other_id, %{pid: other_pid}} <- users do + send(other_pid, {:signaling_relay, :user_left, user_id, %{user_id: user_id}}) + end + + room = %{room | users: users} + + rooms = + if map_size(users) == 0 do + Logger.info("SignalingServer: 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, 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 + end + + {:reply, :ok, %{state | rooms: rooms, refs: refs}} + end + end + end + + @impl true + def handle_call({:relay_message, room_id, from, to, type, payload}, _from, state) do + case Map.get(state.rooms, room_id) do + nil -> + Logger.warning( + "SignalingServer: relay_message failed room_not_found room=#{room_id} from=#{from} to=#{to} type=#{type}" + ) + + {:reply, {:error, :room_not_found}, state} + + room -> + from_user = Map.get(room.users, from) + to_user = Map.get(room.users, to) + + cond do + is_nil(from_user) -> + Logger.warning( + "SignalingServer: relay_message failed user_not_found room=#{room_id} from=#{from} to=#{to} type=#{type}" + ) + + {:reply, {:error, :user_not_found}, state} + + is_nil(to_user) -> + Logger.warning( + "SignalingServer: relay_message failed user_not_found room=#{room_id} from=#{from} to=#{to} type=#{type}" + ) + + {:reply, {:error, :user_not_found}, state} + + room.topology == :star and from_user.role != :host and to_user.role != :host -> + Logger.warning( + "SignalingServer: relay_message failed not_allowed room=#{room_id} from=#{from} role=#{from_user.role} to=#{to} role=#{to_user.role}" + ) + + {:reply, {:error, :not_allowed}, state} + + true -> + Logger.debug( + "SignalingServer: relaying room=#{room_id} type=#{type} from=#{from} to=#{to}" + ) + + send(to_user.pid, {:signaling_relay, type, from, payload}) + {:reply, :ok, state} + end + end + end + + @impl true + def handle_call({:broadcast_message, room_id, from, type, payload}, _from, state) do + case Map.get(state.rooms, room_id) do + nil -> + Logger.warning( + "SignalingServer: broadcast_message failed room_not_found room=#{room_id} from=#{from} type=#{type}" + ) + + {:reply, {:error, :room_not_found}, state} + + room -> + case Map.get(room.users, from) do + nil -> + Logger.warning( + "SignalingServer: broadcast_message failed user_not_found room=#{room_id} from=#{from}" + ) + + {:reply, {:error, :user_not_found}, state} + + from_user -> + if broadcast_allowed?(room, from_user) do + broadcast_to_room(room_id, room, from, type, payload) + {:reply, :ok, state} + else + Logger.warning( + "SignalingServer: broadcast_message failed not_allowed room=#{room.topology} from=#{from} role=#{from_user.role}" + ) + + {:reply, {:error, :not_allowed}, state} + end + end + end + end + + @impl true + def handle_call({:list_users, room_id}, _from, state) do + case Map.get(state.rooms, room_id) do + nil -> + Logger.warning("SignalingServer: list_users failed room_not_found room=#{room_id}") + {:reply, {:error, :room_not_found}, state} + + room -> + users = + Map.new(room.users, fn {user_id, user} -> + {user_id, %{user_id: user.user_id, role: user.role, metadata: user.metadata}} + end) + + {:reply, users, state} + end + end + + def handle_call({:room_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 + + 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("SignalingServer: 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 + {nil, _} -> + Logger.debug( + "SignalingServer: DOWN from unknown pid=#{inspect(pid)} reason=#{inspect(reason)}" + ) + + {:noreply, state} + + {{room_id, user_id}, refs} -> + case Map.get(state.rooms, room_id) do + nil -> + Logger.warning( + "SignalingServer: DOWN for removed room room=#{room_id} user=#{user_id} pid=#{inspect(pid)} reason=#{inspect(reason)}" + ) + + {:noreply, %{state | refs: refs}} + + room -> + user = Map.get(room.users, user_id) + + if user 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 + ) + + user = %{user | status: :disconnected, disconnect_timer: timer} + users = Map.put(room.users, user_id, user) + room = %{room | users: users} + rooms = Map.put(state.rooms, room_id, room) + + Logger.info( + "SignalingServer: 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 + + @impl true + def handle_info({:reconnect_timeout, room_id, user_id}, state) do + case Map.get(state.rooms, room_id) do + nil -> + {:noreply, state} + + room -> + user = Map.get(room.users, user_id) + + if user && user.status == :disconnected do + {_user, users} = Map.pop(room.users, user_id) + remaining = map_size(users) + + Logger.info( + "SignalingServer: reconnect timeout expired room=#{room_id} user=#{user_id} remaining_users=#{remaining}" + ) + + for {other_id, %{pid: other_pid}} <- users, other_id != user_id do + send(other_pid, {:signaling_relay, :user_left, user_id, %{user_id: user_id}}) + end + + rooms = + if map_size(users) == 0 do + Logger.info("SignalingServer: room empty after timeout, removing room=#{room_id}") + Map.delete(state.rooms, room_id) + else + Map.put(state.rooms, room_id, %{room | users: users}) + end + + {:noreply, %{state | rooms: rooms}} + else + {:noreply, state} + end + end + end + + @impl true + def handle_info(_msg, state) do + {:noreply, state} + end + + # ── Private helpers ───────────────────────────────────────────────────── + + defp default_role(room, user_id) do + case room.topology do + :mesh -> :user + :star -> if user_id == room.host_user_id, do: :host, else: :client + end + end + + defp broadcast_allowed?(room, from_user) do + room.topology != :star or from_user.role == :host + end + + defp broadcast_to_room(room_id, room, from, type, payload) do + targets = + Enum.filter(room.users, fn {user_id, _} -> user_id != from end) + |> Enum.map(fn {id, _} -> id end) + + Logger.debug( + "SignalingServer: broadcasting room=#{room_id} type=#{type} from=#{from} targets=#{length(targets)}" + ) + + for {user_id, %{pid: pid}} <- room.users, user_id != from do + send(pid, {:signaling_relay, type, from, payload}) + end + + :ok + end +end diff --git a/apps/game_server_web/config/test.exs b/apps/game_server_web/config/test.exs index d281bd07d..9a131ad84 100644 --- a/apps/game_server_web/config/test.exs +++ b/apps/game_server_web/config/test.exs @@ -84,6 +84,7 @@ config :game_server_core, GameServer.Accounts.StalePresenceSweeper, enabled: fal # ("database is locked"). Tests drive tick/0 and sweep/0 directly. config :game_server_core, GameServer.Tournaments.Ticker, enabled: false config :game_server_core, GameServer.Matchmaking.Worker, enabled: false +config :game_server_core, GameServer.Signaling.Server, enabled: false # NOTE: deliberately NOT setting `async_inline: true` here, unlike the root # config/test.exs. Payments call GameServer.Async.run/1 from inside a 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..5040daa71 --- /dev/null +++ b/apps/game_server_web/lib/game_server_web/channels/signaling_channel.ex @@ -0,0 +1,371 @@ +defmodule GameServerWeb.SignalingChannel do + @moduledoc """ + Channel for WebRTC signaling relay. + + Topic: `signaling:` + + Rooms are created by the `WebRTCLobbyHook` through `Server.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 the authenticated `user_id` is used directly as the user identity. + The server assigns the role (`:host` or `:client` for `:star`, `:user` for + `:mesh`) based on the room's configuration. If the same user_id reconnects + within the configured grace period, the existing user is preserved and a + `user_rejoined` event is broadcast. + + ## Messages + + Inbound events (from client): + + 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_user_id: "..."} + "answer" — %{sdp: "...", from_user_id: "..."} + "ice" — %{candidate: "...", from_user_id: "..."} + "user_joined" — %{user_id: "...", role: :host | :client | :user} + "user_rejoined" — %{user_id: "...", role: :host | :client | :user} + "user_left" — %{user_id: "..."} + "room_closed" — %{} + """ + + use Phoenix.Channel + + import GameServerWeb.ChannelPush + require Logger + + alias GameServer.Signaling.Server + + # 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 + case Server.join_room(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 + + # ── 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_user_id + + case Server.relay_message(room, from, target, :offer, %{sdp: sdp}) do + :ok -> + {:reply, {:ok, %{}}, socket} + + {:error, :user_not_found} -> + Logger.warning( + "SignalingChannel: offer failed user_not_found room=#{room} from=#{from} target=#{target}" + ) + + {:reply, {:error, %{error: "user_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_user_id + + case Server.relay_message(room, from, target, :answer, %{sdp: sdp}) do + :ok -> + {:reply, {:ok, %{}}, socket} + + {:error, :user_not_found} -> + Logger.warning( + "SignalingChannel: answer failed user_not_found room=#{room} from=#{from} target=#{target}" + ) + + {:reply, {:error, %{error: "user_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_user_id + + case Server.relay_message(room, from, target, :ice, %{candidate: candidate}) do + :ok -> + {:reply, {:ok, %{}}, socket} + + {:error, :user_not_found} -> + Logger.warning( + "SignalingChannel: ice failed user_not_found room=#{room} from=#{from} target=#{target}" + ) + + {:reply, {:error, %{error: "user_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_user_id + + case Server.broadcast_message(room, from, :offer, %{sdp: sdp, from_user_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_users", _payload, socket) do + with :ok <- check_ws_rate_limit(socket) do + room = socket.assigns.signaling_room + + case Server.list_users(room) do + users when is_map(users) -> + {:reply, {:ok, %{users: users}}, socket} + + {:error, :room_not_found} -> + Logger.warning("SignalingChannel: list_users 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"} user=#{socket.assigns[:signaling_user_id] || "nil"}" + ) + + {:reply, {:error, %{error: "unknown_event"}}, socket} + end + + # ── Server 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} 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_user_id, payload}, socket) do + event_name = relay_event_name(type) + + 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 + + @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"} user=#{socket.assigns[:signaling_user_id] || "nil"}" + ) + + {:noreply, socket} + end + + @impl true + def terminate(reason, socket) do + room_id = socket.assigns[:signaling_room] + user_id = socket.assigns[:signaling_user_id] + + if room_id && user_id do + Logger.info( + "SignalingChannel: terminating reason=#{inspect(reason)} room=#{room_id} user=#{user_id}" + ) + + # Do NOT call Server.leave here. The server'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/user 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(:user_joined), do: "user_joined" + defp relay_event_name(:user_rejoined), do: "user_rejoined" + defp relay_event_name(:user_left), do: "user_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 506cb23b9..3b9617e62 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/host_supervision.ex b/apps/game_server_web/lib/game_server_web/host_supervision.ex index 88aca9895..d1c9c1d98 100644 --- a/apps/game_server_web/lib/game_server_web/host_supervision.ex +++ b/apps/game_server_web/lib/game_server_web/host_supervision.ex @@ -123,7 +123,9 @@ defmodule GameServerWeb.HostSupervision do GameServer.Matchmaking.Worker, # Buffers lobby snapshots/events and assigns seq. :global-registered, so # only one node runs it and start_link returns :ignore on the others. - GameServer.LobbySnapshots.Writer + GameServer.LobbySnapshots.Writer, + # Signaling relay for WebRTC user-to-user and client-server topologies + GameServer.Signaling.Server ] ++ extra end diff --git a/apps/game_server_web/lib/game_server_web/realtime_events.ex b/apps/game_server_web/lib/game_server_web/realtime_events.ex index b183f83b8..f5b7ea844 100644 --- a/apps/game_server_web/lib/game_server_web/realtime_events.ex +++ b/apps/game_server_web/lib/game_server_web/realtime_events.ex @@ -26,6 +26,7 @@ defmodule GameServerWeb.RealtimeEvents do @group "group:*" @groups "groups" @party "party:*" + @signaling "signaling:*" @events [ # ── user:* ────────────────────────────────────────────────────────── @@ -133,7 +134,18 @@ defmodule GameServerWeb.RealtimeEvents do {@party, "disbanded", true, "party ref", "The party was disbanded"}, {@party, "chat_message_created", true, "chat message", "Party chat message"}, {@party, "chat_message_updated", true, "chat message", "Party chat message edited"}, - {@party, "chat_message_deleted", true, "message id", "Party chat message deleted"} + {@party, "chat_message_deleted", true, "message id", "Party chat message deleted"}, + + # ── signaling:* ───────────────────────────────────────────────────── + {@signaling, "user_joined", false, "user id + role", "A user joined the signaling room"}, + {@signaling, "user_rejoined", false, "user id + role", + "A user rejoined the signaling room after a transient disconnect"}, + {@signaling, "user_left", false, "user id", "A user left the signaling room"}, + {@signaling, "offer", false, "sdp + from user id", "WebRTC offer relayed to this peer"}, + {@signaling, "answer", false, "sdp + from user id", "WebRTC answer relayed to this peer"}, + {@signaling, "ice", false, "candidate + from user id", + "WebRTC ICE candidate relayed to this peer"}, + {@signaling, "room_closed", false, "empty", "The signaling room was closed"} ] @doc "Every server→client event as a list of maps." diff --git a/config/test.exs b/config/test.exs index 2736cea37..3c338df88 100644 --- a/config/test.exs +++ b/config/test.exs @@ -83,6 +83,10 @@ config :game_server_core, GameServer.Tournaments.Ticker, enabled: false # Tests drive GameServer.Matchmaking.Worker.sweep/0 directly. config :game_server_core, GameServer.Matchmaking.Worker, enabled: false +# The signaling server is disabled in tests by default. If a test needs it, +# start it manually with start_supervised!(GameServer.Signaling.Server). +config :game_server_core, GameServer.Signaling.Server, enabled: false + # Disable app-level caching in tests to avoid stale reads across assertions. # Still provide the multilevel configuration so the cache can start. config :game_server_core, GameServer.Cache, diff --git a/mix.lock b/mix.lock index 03f584e9a..4fb9aa92b 100644 --- a/mix.lock +++ b/mix.lock @@ -1,5 +1,5 @@ %{ - "bandit": {:hex, :bandit, "1.12.3", "23f49ae03d86365b0caff3b3a1eba5d981066b1ecafe2b19c0bef5ad36c211ce", [:mix], [{:hpax, "~> 1.0", [hex: :hpax, repo: "hexpm", optional: false]}, {:plug, "~> 1.18", [hex: :plug, repo: "hexpm", optional: false]}, {:telemetry, "~> 0.4 or ~> 1.0", [hex: :telemetry, repo: "hexpm", optional: false]}, {:thousand_island, "~> 1.5", [hex: :thousand_island, repo: "hexpm", optional: false]}, {:websock, "~> 0.5", [hex: :websock, repo: "hexpm", optional: false]}], "hexpm", "a253ec03f391755b2126e4181ee2fee05c75b712b407aa399de1831b5088c58e"}, + "bandit": {:hex, :bandit, "1.12.4", "10bbab488edf8162318d736c19c5837077b8fca2bf5d95b07b33830387124f62", [:mix], [{:hpax, "~> 1.0", [hex: :hpax, repo: "hexpm", optional: false]}, {:plug, "~> 1.18", [hex: :plug, repo: "hexpm", optional: false]}, {:telemetry, "~> 0.4 or ~> 1.0", [hex: :telemetry, repo: "hexpm", optional: false]}, {:thousand_island, "~> 1.5", [hex: :thousand_island, repo: "hexpm", optional: false]}, {:websock, "~> 0.5", [hex: :websock, repo: "hexpm", optional: false]}], "hexpm", "84513318c5752a2a8017664450f889b47fae5d53d64698ddf1e4fb09a7449e8d"}, "bcrypt_elixir": {:hex, :bcrypt_elixir, "3.3.2", "d50091e3c9492d73e17fc1e1619a9b09d6a5ef99160eb4d736926fd475a16ca3", [:make, :mix], [{:comeonin, "~> 5.3", [hex: :comeonin, repo: "hexpm", optional: false]}, {:elixir_make, "~> 0.6", [hex: :elixir_make, repo: "hexpm", optional: false]}], "hexpm", "471be5151874ae7931911057d1467d908955f93554f7a6cd1b7d804cac8cef53"}, "bunch": {:hex, :bunch, "1.6.4", "b23ecefdcbc1d8b9723044055daa2ad8573a0ebe6ef35c0ac4f3bad243273ba3", [:mix], [], "hexpm", "cacbe7437dd5299c9e665011d32127e0ce92b26ab4b9f64ebd2389f57cc7c141"}, "bunch_native": {:hex, :bunch_native, "0.5.2", "706723acd1644fb6f8e7c8f3864172a7eee7664250e790696253825e3b65e18a", [:mix], [{:bundlex, "~> 1.0", [hex: :bundlex, repo: "hexpm", optional: false]}], "hexpm", "23f9e036c34510d8ace14fc47553fd7672e2330e769ab84fe4b54fac80095a4b"}, @@ -93,8 +93,8 @@ "phoenix_ecto": {:hex, :phoenix_ecto, "4.7.0", "75c4b9dfb3efdc42aec2bd5f8bccd978aca0651dbcbc7a3f362ea5d9d43153c6", [:mix], [{:ecto, "~> 3.5", [hex: :ecto, repo: "hexpm", optional: false]}, {:phoenix_html, "~> 2.14.2 or ~> 3.0 or ~> 4.1", [hex: :phoenix_html, repo: "hexpm", optional: true]}, {:plug, "~> 1.9", [hex: :plug, repo: "hexpm", optional: false]}, {:postgrex, "~> 0.16 or ~> 1.0", [hex: :postgrex, repo: "hexpm", optional: true]}], "hexpm", "1d75011e4254cb4ddf823e81823a9629559a1be93b4321a6a5f11a5306fbf4cc"}, "phoenix_html": {:hex, :phoenix_html, "4.3.0", "d3577a5df4b6954cd7890c84d955c470b5310bb49647f0a114a6eeecc850f7ad", [:mix], [], "hexpm", "3eaa290a78bab0f075f791a46a981bbe769d94bc776869f4f3063a14f30497ad"}, "phoenix_live_dashboard": {:hex, :phoenix_live_dashboard, "0.8.7", "405880012cb4b706f26dd1c6349125bfc903fb9e44d1ea668adaf4e04d4884b7", [:mix], [{:ecto, "~> 3.6.2 or ~> 3.7", [hex: :ecto, repo: "hexpm", optional: true]}, {:ecto_mysql_extras, "~> 0.5", [hex: :ecto_mysql_extras, repo: "hexpm", optional: true]}, {:ecto_psql_extras, "~> 0.7", [hex: :ecto_psql_extras, repo: "hexpm", optional: true]}, {:ecto_sqlite3_extras, "~> 1.1.7 or ~> 1.2.0", [hex: :ecto_sqlite3_extras, repo: "hexpm", optional: true]}, {:mime, "~> 1.6 or ~> 2.0", [hex: :mime, repo: "hexpm", optional: false]}, {:phoenix_live_view, "~> 0.19 or ~> 1.0", [hex: :phoenix_live_view, repo: "hexpm", optional: false]}, {:telemetry_metrics, "~> 0.6 or ~> 1.0", [hex: :telemetry_metrics, repo: "hexpm", optional: false]}], "hexpm", "3a8625cab39ec261d48a13b7468dc619c0ede099601b084e343968309bd4d7d7"}, - "phoenix_live_reload": {:hex, :phoenix_live_reload, "1.6.2", "b18b0773a1ba77f28c52decbb0f10fd1ac4d3ae5b8632399bbf6986e3b665f62", [:mix], [{:file_system, "~> 0.2.10 or ~> 1.0", [hex: :file_system, repo: "hexpm", optional: false]}, {:phoenix, "~> 1.4", [hex: :phoenix, repo: "hexpm", optional: false]}], "hexpm", "d1f89c18114c50d394721365ffb428cce24f1c13de0467ffa773e2ff4a30d5b9"}, - "phoenix_live_view": {:hex, :phoenix_live_view, "1.2.7", "d0f20871681216598e78baccd4f66d8686fdb8007821989bab508d293a3ed9ad", [:mix], [{:igniter, ">= 0.6.16 and < 1.0.0-0", [hex: :igniter, repo: "hexpm", optional: true]}, {:jason, "~> 1.0", [hex: :jason, repo: "hexpm", optional: true]}, {:lazy_html, "~> 0.1.0", [hex: :lazy_html, repo: "hexpm", optional: true]}, {:phoenix, "~> 1.6.15 or ~> 1.7.0 or ~> 1.8.0", [hex: :phoenix, repo: "hexpm", optional: false]}, {:phoenix_html, "~> 3.3 or ~> 4.0", [hex: :phoenix_html, 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.15", [hex: :plug, repo: "hexpm", optional: false]}, {:telemetry, "~> 0.4.2 or ~> 1.0", [hex: :telemetry, repo: "hexpm", optional: false]}], "hexpm", "61e97938a4fcca6d6f2c836925623abf2f52a572cc8c6085e4074f3f6337e0eb"}, + "phoenix_live_reload": {:hex, :phoenix_live_reload, "1.7.0", "fb1e429f6d8778ce3a6962debdc5e555428a05a6e7b058d6dbad13d281a2c31f", [:mix], [{:file_system, "~> 0.2.10 or ~> 1.0", [hex: :file_system, repo: "hexpm", optional: false]}, {:phoenix, "~> 1.4", [hex: :phoenix, repo: "hexpm", optional: false]}], "hexpm", "dc9f44271aa6fc4ab7797f2aa374ba096ef2c87520586280eb095626b7387a68"}, + "phoenix_live_view": {:hex, :phoenix_live_view, "1.2.8", "5006fd7b429c42489600fbc1600c750d0f0e5b5ea4965d9758e429b979392991", [:mix], [{:igniter, ">= 0.6.16 and < 1.0.0-0", [hex: :igniter, repo: "hexpm", optional: true]}, {:jason, "~> 1.0", [hex: :jason, repo: "hexpm", optional: true]}, {:lazy_html, "~> 0.1.0", [hex: :lazy_html, repo: "hexpm", optional: true]}, {:phoenix, "~> 1.6.15 or ~> 1.7.0 or ~> 1.8.0", [hex: :phoenix, repo: "hexpm", optional: false]}, {:phoenix_html, "~> 3.3 or ~> 4.0", [hex: :phoenix_html, 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.15", [hex: :plug, repo: "hexpm", optional: false]}, {:telemetry, "~> 0.4.2 or ~> 1.0", [hex: :telemetry, repo: "hexpm", optional: false]}], "hexpm", "b05ffe21f43c0ff219da62948b482c324aa5b8873e17b0c0cac58289a178af38"}, "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"}, "pigeon": {:git, "https://github.com/codedge-llc/pigeon.git", "712d5c2b20100d56bed08efed42e4eed924c422a", [ref: "712d5c2b20100d56bed08efed42e4eed924c422a"]}, @@ -106,16 +106,16 @@ "protobuf": {:hex, :protobuf, "0.17.0", "39e24e43c9648e148feba16ed51100b5b2028ea900b55460377b0476f6e10613", [:mix], [{:jason, "~> 1.2", [hex: :jason, repo: "hexpm", optional: true]}], "hexpm", "ca6c91f6f63e2c147b47f03eefd10b80538aa6fc55ff4b12b795efb786b0152f"}, "qex": {:hex, :qex, "0.5.2", "a0c861a2de2380314c23ef592349824ca9016c5845380667ff1d9a22a8796f9b", [:mix], [], "hexpm", "6fb81bf3ae354a9abb471b9561538ea3e8540125d803b00f45cbccff52f00496"}, "quic": {:hex, :quic, "1.7.1", "1ac5331739e17928a797d5ff103629a9faa91167481c326d7043b323e3bdc3b1", [:rebar3], [], "hexpm", "9b85784f1d78f5c3adfee5f52c3448cac3b03f5f0403d5dc782231c3d3ad709d"}, - "ranch": {:hex, :ranch, "2.2.0", "25528f82bc8d7c6152c57666ca99ec716510fe0925cb188172f41ce93117b1b0", [:make, :rebar3], [], "hexpm", "fa0b99a1780c80218a4197a59ea8d3bdae32fbff7e88527d7d8a4787eff4f8e7"}, + "ranch": {:hex, :ranch, "2.2.1", "fc2bb0e800efbeab8ff82eb325450b9c08653b14c3f69e8b09bbea304973105d", [:make, :rebar3], [], "hexpm", "55f05cce20ec2da1d90de5d5981afb93dbfc01325fc7e933189aa5c62f037c24"}, "redix": {:hex, :redix, "1.6.0", "694179c7a3c71bffac8848fcbfe16f6f6e2a1e020f00a2ad01bc6484815470d1", [:mix], [{:castore, "~> 0.1.0 or ~> 1.0", [hex: :castore, repo: "hexpm", optional: true]}, {:nimble_options, "~> 0.5.0 or ~> 1.0", [hex: :nimble_options, repo: "hexpm", optional: false]}, {:telemetry, "~> 0.4.0 or ~> 1.0", [hex: :telemetry, repo: "hexpm", optional: false]}], "hexpm", "b2eccb05e02f21c0c3ca57513e6bacb4dd48e6406dadbd7ff9fbe07bd6745999"}, - "req": {:hex, :req, "0.6.3", "7fe5e68792ff0546e45d5919104fa1764a13694cfe3e48c8a0f32ad051ae77e4", [:mix], [{:brotli, "~> 0.3.1", [hex: :brotli, repo: "hexpm", optional: true]}, {:ezstd, "~> 1.0", [hex: :ezstd, repo: "hexpm", optional: true]}, {:finch, "~> 0.21", [hex: :finch, repo: "hexpm", optional: false]}, {:jason, "~> 1.0", [hex: :jason, repo: "hexpm", optional: false]}, {:mime, "~> 2.0.6 or ~> 2.1", [hex: :mime, repo: "hexpm", optional: false]}, {:nimble_csv, "~> 1.0", [hex: :nimble_csv, repo: "hexpm", optional: true]}, {:plug, "~> 1.0", [hex: :plug, repo: "hexpm", optional: true]}], "hexpm", "e85b5c6c990e6c3f52bbba68e6f099118f2b8252825f96c7c3636b97a3de307d"}, + "req": {:hex, :req, "0.7.1", "86271f9e29dca83d382222428ced09ba03bad21e24888c50dd4527553d85f642", [:mix], [{:brotli, "~> 0.3.1", [hex: :brotli, repo: "hexpm", optional: true]}, {:finch, "~> 0.21", [hex: :finch, repo: "hexpm", optional: false]}, {:jason, "~> 1.0", [hex: :jason, repo: "hexpm", optional: false]}, {:mime, "~> 2.0.6 or ~> 2.1", [hex: :mime, repo: "hexpm", optional: false]}, {:nimble_csv, "~> 1.0", [hex: :nimble_csv, repo: "hexpm", optional: true]}, {:plug, "~> 1.0", [hex: :plug, repo: "hexpm", optional: true]}], "hexpm", "254638b15ceb9a2624d15aff13bf7903ea2e95cd4b9c1aa18da1fb06e1086b50"}, "rustler": {:hex, :rustler, "0.36.2", "6c2142f912166dfd364017ab2bf61242d4a5a3c88e7b872744642ae004b82501", [:mix], [{:jason, "~> 1.0", [hex: :jason, repo: "hexpm", optional: false]}, {:toml, "~> 0.7", [hex: :toml, repo: "hexpm", optional: false]}], "hexpm", "93832a6dbc1166739a19cd0c25e110e4cf891f16795deb9361dfcae95f6c88fe"}, "rustler_precompiled": {:hex, :rustler_precompiled, "0.9.0", "3a052eda09f3d2436364645cc1f13279cf95db310eb0c17b0d8f25484b233aa0", [:mix], [{:rustler, "~> 0.23", [hex: :rustler, repo: "hexpm", optional: true]}], "hexpm", "471d97315bd3bf7b64623418b3693eedd8e47de3d1cb79a0ac8f9da7d770d94c"}, "shmex": {:hex, :shmex, "0.5.2", "5e807c4ec0320fdd3b224d3f8052d501908217a066b67c0433f4cfa369a2b07c", [:mix], [{:bunch_native, "~> 0.5.0", [hex: :bunch_native, repo: "hexpm", optional: false]}, {:bundlex, "~> 1.0", [hex: :bundlex, repo: "hexpm", optional: false]}], "hexpm", "7325f40a7308fecaaab1b19790903034926e5fb202a48fcf51dfd2c4ba97f861"}, "ssl_verify_fun": {:hex, :ssl_verify_fun, "1.1.7", "354c321cf377240c7b8716899e182ce4890c5938111a1296add3ec74cf1715df", [:make, :mix, :rebar3], [], "hexpm", "fe4c190e8f37401d30167c8c405eda19469f34577987c76dde613e838bbc67f8"}, "stripity_stripe": {:hex, :stripity_stripe, "3.3.2", "795fb7d56baa2fb0d80a44d4fe23b4e947809e8d61e80d9f5b5eba64610f63e2", [:mix], [{:hackney, "~> 4.0", [hex: :hackney, repo: "hexpm", optional: false]}, {:jason, "~> 1.1", [hex: :jason, repo: "hexpm", optional: false]}, {:plug, "~> 1.14", [hex: :plug, repo: "hexpm", optional: true]}, {:telemetry, "~> 1.1", [hex: :telemetry, repo: "hexpm", optional: false]}, {:uri_query, "~> 0.2.0", [hex: :uri_query, repo: "hexpm", optional: false]}], "hexpm", "73db5ab782f1eec9dac0c9cc7b5391452c704960ddfff42acdba83dbe37b5011"}, "sweet_xml": {:hex, :sweet_xml, "0.7.5", "803a563113981aaac202a1dbd39771562d0ad31004ddbfc9b5090bdcd5605277", [:mix], [], "hexpm", "193b28a9b12891cae351d81a0cead165ffe67df1b73fe5866d10629f4faefb12"}, - "swoosh": {:hex, :swoosh, "1.26.3", "9d8b60077305ce259298d9a1102e5be67cd3c41d1ea930c29e9288af195ca017", [:mix], [{:bandit, ">= 1.0.0", [hex: :bandit, repo: "hexpm", optional: true]}, {:cowboy, "~> 1.1 or ~> 2.4", [hex: :cowboy, repo: "hexpm", optional: true]}, {:ex_aws, "~> 2.1", [hex: :ex_aws, repo: "hexpm", optional: true]}, {:finch, "~> 0.6", [hex: :finch, repo: "hexpm", optional: true]}, {:gen_smtp, "~> 0.13 or ~> 1.0", [hex: :gen_smtp, repo: "hexpm", optional: true]}, {:hackney, ">= 1.9.0 and < 5.0.0", [hex: :hackney, repo: "hexpm", optional: true]}, {:idna, ">= 6.0.0 and < 8.0.0", [hex: :idna, repo: "hexpm", optional: false]}, {:jason, "~> 1.0", [hex: :jason, repo: "hexpm", optional: false]}, {:mail, "~> 0.2", [hex: :mail, repo: "hexpm", optional: true]}, {:mime, "~> 1.1 or ~> 2.0", [hex: :mime, repo: "hexpm", optional: false]}, {:mua, "~> 0.2.3", [hex: :mua, repo: "hexpm", optional: true]}, {:multipart, "~> 0.4", [hex: :multipart, repo: "hexpm", optional: true]}, {:plug, "~> 1.9", [hex: :plug, repo: "hexpm", optional: true]}, {:plug_cowboy, ">= 1.0.0", [hex: :plug_cowboy, repo: "hexpm", optional: true]}, {:req, "~> 0.5.10 or ~> 0.6 or ~> 1.0", [hex: :req, repo: "hexpm", optional: true]}, {:telemetry, "~> 0.4.2 or ~> 1.0", [hex: :telemetry, repo: "hexpm", optional: false]}], "hexpm", "c7683d070fe8f8aa9d174e61b01f2d527be73cd8ac40037b7109184941eb569f"}, + "swoosh": {:hex, :swoosh, "1.27.0", "2df57c342f3854f2d429e6afebf19a03a179976a64006a93750ed81639453f75", [:mix], [{:bandit, ">= 1.0.0", [hex: :bandit, repo: "hexpm", optional: true]}, {:cowboy, "~> 1.1 or ~> 2.4", [hex: :cowboy, repo: "hexpm", optional: true]}, {:ex_aws, "~> 2.1", [hex: :ex_aws, repo: "hexpm", optional: true]}, {:finch, "~> 0.6", [hex: :finch, repo: "hexpm", optional: true]}, {:gen_smtp, "~> 0.13 or ~> 1.0", [hex: :gen_smtp, repo: "hexpm", optional: true]}, {:hackney, ">= 1.9.0 and < 5.0.0", [hex: :hackney, repo: "hexpm", optional: true]}, {:idna, ">= 6.0.0 and < 8.0.0", [hex: :idna, repo: "hexpm", optional: false]}, {:jason, "~> 1.0", [hex: :jason, repo: "hexpm", optional: false]}, {:mail, "~> 0.2", [hex: :mail, repo: "hexpm", optional: true]}, {:mime, "~> 1.1 or ~> 2.0", [hex: :mime, repo: "hexpm", optional: false]}, {:mua, "~> 0.2.3", [hex: :mua, repo: "hexpm", optional: true]}, {:multipart, "~> 0.4", [hex: :multipart, repo: "hexpm", optional: true]}, {:plug, "~> 1.9", [hex: :plug, repo: "hexpm", optional: true]}, {:plug_cowboy, ">= 1.0.0", [hex: :plug_cowboy, repo: "hexpm", optional: true]}, {:req, "~> 0.5.10 or ~> 0.6 or ~> 1.0", [hex: :req, repo: "hexpm", optional: true]}, {:telemetry, "~> 0.4.2 or ~> 1.0", [hex: :telemetry, repo: "hexpm", optional: false]}], "hexpm", "5da7d3b11de5d61327275ed599cb311942ea6e23cdbb411981f46ee150cddf76"}, "tailwind": {:hex, :tailwind, "0.5.1", "35435b13158c90d37da11e1cfc808755fca1d7b6c5ab87b1b19c5de87e2f0a10", [:mix], [], "hexpm", "c4e26302a59fec72abc5610ecb6ad2116d9aa31f31aab2d4b8eb6e95d25a689c"}, "telemetry": {:hex, :telemetry, "1.4.2", "a0cb522801dffb1c49fe6e30561badffc7b6d0e180db1300df759faa22062855", [:rebar3], [], "hexpm", "928f6495066506077862c0d1646609eed891a4326bee3126ba54b60af61febb1"}, "telemetry_metrics": {:hex, :telemetry_metrics, "1.1.0", "5bd5f3b5637e0abea0426b947e3ce5dd304f8b3bc6617039e2b5a008adc02f8f", [:mix], [{:telemetry, "~> 0.4 or ~> 1.0", [hex: :telemetry, repo: "hexpm", optional: false]}], "hexpm", "e7b79e8ddfde70adb6db8a6623d1778ec66401f366e9a8f5dd0955c56bc8ce67"}, 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/webrtc_lobby_hook/.formatter.exs b/modules/plugins/webrtc_lobby_hook/.formatter.exs new file mode 100644 index 000000000..d304ff320 --- /dev/null +++ b/modules/plugins/webrtc_lobby_hook/.formatter.exs @@ -0,0 +1,3 @@ +[ + inputs: ["{mix,.formatter}.exs", "{config,lib,test}/**/*.{ex,exs}"] +] diff --git a/modules/plugins/webrtc_lobby_hook/.gitignore b/modules/plugins/webrtc_lobby_hook/.gitignore new file mode 100644 index 000000000..914532603 --- /dev/null +++ b/modules/plugins/webrtc_lobby_hook/.gitignore @@ -0,0 +1,3 @@ +/_build/ +/deps/ +/ebin/ 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/webrtc_lobby_hook/config/config.exs b/modules/plugins/webrtc_lobby_hook/config/config.exs new file mode 100644 index 000000000..5b66c9106 --- /dev/null +++ b/modules/plugins/webrtc_lobby_hook/config/config.exs @@ -0,0 +1,3 @@ +import Config + +# No runtime configuration required for this example plugin. 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..57ec0cebe --- /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.exists_room?(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.exists_room?(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 -> :user + 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 -> :user + 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/webrtc_lobby_hook/mix.exs b/modules/plugins/webrtc_lobby_hook/mix.exs new file mode 100644 index 000000000..b74a6119a --- /dev/null +++ b/modules/plugins/webrtc_lobby_hook/mix.exs @@ -0,0 +1,39 @@ +defmodule HandleWebRTC.WebRTCLobbyHook do + use Mix.Project + + def project do + [ + app: :webrtc_lobby_hook, + 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.WebRTCLobbyHook] + ] + 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 + [ + shared_dep(:game_server_sdk, "../../../sdk"), + shared_dep(:game_server_plugin_tools, "../../../sdk_tools"), + {:bunt, "~> 1.0"}, + {:phoenix, "~> 1.8.3"}, + ] + end + + defp shared_dep(app, local_path) do + if File.dir?(local_path) do + {app, path: local_path, runtime: false} + else + {app, github: "appsinacup/game_server", override: true, runtime: false} + end + end +end diff --git a/modules/plugins/webrtc_lobby_hook/mix.lock b/modules/plugins/webrtc_lobby_hook/mix.lock new file mode 100644 index 000000000..20a9d4a30 --- /dev/null +++ b/modules/plugins/webrtc_lobby_hook/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/sdk/lib/game_server/signaling.ex b/sdk/lib/game_server/signaling.ex new file mode 100644 index 000000000..c9be19a23 --- /dev/null +++ b/sdk/lib/game_server/signaling.ex @@ -0,0 +1,12 @@ +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 exists_room?(_room_id), do: false + def allow_user(_room_id, _user_id, _role \\ :user), do: :ok + def disallow_user(_room_id, _user_id), do: :ok + def list_users(_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