defmodule Phoenix.SocketClient.Socket do use GenServer alias Phoenix.SocketClient.ChannelManager alias Phoenix.SocketClient.Message alias Phoenix.SocketClient.Telemetry import Phoenix.SocketClient, only: [get_state: 2, put_state: 3, get_process_pid: 2] def start_link(opts) do GenServer.start_link(__MODULE__, opts) end @impl true def init(%{sup_pid: sup_pid} = _opts) do Process.flag(:trap_exit, true) {:ok, %{ sup_pid: sup_pid, status: :disconnected, transport_ref: nil, transport_pid: nil }, {:continue, :post_start}} end @impl true def handle_continue(:post_start, %{sup_pid: sup_pid} = state) do Process.sleep(1_000) if get_state(sup_pid, :auto_connect) do Process.send(self(), :connect, [:noconnect]) end {:noreply, state} end @impl true def handle_info(:disconnect, state) do {:noreply, close(:normal, state)} end def handle_info(:connect, %{sup_pid: sup_pid} = state) do case get_state(sup_pid, :url) do nil -> Telemetry.debug(self(), "Connection: no URL provided, skipping connection") {:noreply, state} url -> transport = get_state(sup_pid, :transport) transport_opts = get_state(sup_pid, :transport_opts) |> Keyword.put(:sender, self()) Telemetry.socket_connecting(self(), url) case transport.open(url, transport_opts) do {:ok, transport_pid} -> Telemetry.debug(self(), "Connection: transport started", url, %{ transport_pid: transport_pid }) transport_ref = Process.monitor(transport_pid) state = %{ state | transport_pid: transport_pid, status: :connecting, transport_ref: transport_ref } update_socket_state_status(sup_pid, :connecting) {:noreply, state} {:error, reason} -> Telemetry.debug(self(), "Connection: transport failed to start", url, %{ reason: reason }) Telemetry.socket_connection_error(self(), url, reason) {:noreply, close(reason, state)} end end end @impl true def handle_info({:connected, transport_pid}, %{sup_pid: sup_pid} = state) do url = get_state(sup_pid, :url) Telemetry.debug(self(), "Connection: connected", url, %{transport_pid: transport_pid}) Telemetry.socket_connected(self(), url) state = %{state | status: :connected, transport_pid: transport_pid} update_socket_state_status(sup_pid, :connected) {:noreply, state} end @impl true def handle_info({:disconnected, reason, _transport_pid}, %{sup_pid: sup_pid} = state) do url = get_state(sup_pid, :url) Telemetry.debug(self(), "Connection: disconnected", url, %{reason: reason}) Telemetry.socket_disconnected(self(), url, reason) cm_pid = get_process_pid(sup_pid, :channel_manager) ChannelManager.terminate(cm_pid) {:noreply, close(reason, state)} end @impl true def handle_info({:receive, message}, %{sup_pid: sup_pid} = state) do url = get_state(sup_pid, :url) Telemetry.debug(self(), "Connection: received", url, %{message: message}) transport_receive(message, state) {:noreply, state} end @impl true def handle_info(:flush, %{sup_pid: sup_pid} = state) do url = get_state(sup_pid, :url) Telemetry.debug(self(), "Connection: flushing messages", url) to_send = state[:to_send_r] || [] Enum.each(to_send, &transport_send(&1, state)) {:noreply, %{state | to_send_r: []}} end def handle_info(%Message{} = message, %{sup_pid: sup_pid} = state) do url = get_state(sup_pid, :url) Telemetry.debug(self(), "Connection: received message struct", url, %{message: message}) # Handle message structs that might be sent directly channels_pid = get_process_pid(sup_pid, :channel_manager) children = Supervisor.which_children(channels_pid) case find_channel(children, message.topic) do nil -> :noop {_id, channel_pid, _type, _modules} -> send(channel_pid, message) end {:noreply, state} end @impl true def handle_info({:closed, reason, _transport_pid}, %{sup_pid: sup_pid} = state) do url = get_state(sup_pid, :url) Telemetry.debug(self(), "Connection: closed", url, %{reason: reason}) {:noreply, close(reason, state)} end @impl true def handle_info( {:DOWN, ref, :process, _pid, reason}, %{sup_pid: sup_pid, transport_ref: transport_ref, transport_pid: transport_pid} = state ) do if ref == transport_ref do url = get_state(sup_pid, :url) Telemetry.debug(self(), "Connection process: DOWN", url, %{ reason: reason, transport_pid: transport_pid, transport_ref: transport_ref }) {:noreply, close(:shutdown, state)} else {:noreply, state} end end @impl true def handle_info( {:EXIT, _from_pid, reason}, state ) do terminate(reason, state) {:noreply, close(reason, state)} end @impl true def terminate(_reason, _state) do :ok end @impl true def handle_call(:get_status, _from, state) do {:reply, state[:status] || :disconnected, state} end @impl true def handle_call(:get_state, _from, state) do {:reply, state, state} end @impl true def handle_call({:push_message, message}, _from, %{sup_pid: sup_pid} = state) do transport_pid = state[:transport_pid] if transport_pid do protocol_vsn = get_state(sup_pid, :vsn) serializer = Phoenix.SocketClient.Message.serializer(protocol_vsn) json_library = get_state(sup_pid, :json_library) || Jason send(transport_pid, {:send, Message.encode!(serializer, message, json_library)}) {:reply, message, state} else {:reply, {:error, :not_connected}, state} end end @impl true def handle_call(not_matched, from, state) do IO.inspect({:not_matched_handle_call, not_matched, from, state}) {:noreply, state} end defp transport_receive(message, %{sup_pid: sup_pid} = _state) do protocol_vsn = get_state(sup_pid, :vsn) serializer = Phoenix.SocketClient.Message.serializer(protocol_vsn) json_library = get_state(sup_pid, :json_library) || Jason decoded = Message.decode!(serializer, message, json_library) channels_pid = get_process_pid(sup_pid, :channel_manager) children = Supervisor.which_children(channels_pid) case find_channel(children, decoded.topic) do nil -> :noop {_id, channel_pid, _type, _modules} -> send(channel_pid, decoded) end end defp find_channel(children, topic) do Enum.find(children, fn {_id, pid, _type, _modules} -> :sys.get_state(pid).topic == topic end) end defp transport_send(message, %{sup_pid: sup_pid} = state) do transport_pid = state.transport_pid if transport_pid do protocol_vsn = get_state(sup_pid, :vsn) serializer = Phoenix.SocketClient.Message.serializer(protocol_vsn) json_library = get_state(sup_pid, :json_library) || Jason send(transport_pid, {:send, Message.encode!(serializer, message, json_library)}) end end defp close(reason, %{sup_pid: sup_pid} = state) do url = get_state(sup_pid, :url) Telemetry.debug(self(), "Connection: closing connection", url, %{reason: reason}) # Update socket state status update_socket_state_status(sup_pid, :disconnected) reconnect = get_state(sup_pid, :reconnect?) if reconnect do reconnect_interval = get_state(sup_pid, :reconnect_interval) Telemetry.debug(self(), "Connection: reconnecting", url, %{interval: reconnect_interval}) Telemetry.reconnecting(self(), url, 1) Process.send_after(self(), :connect, reconnect_interval) end %{state | status: :disconnected} end defp update_socket_state_status(sup_pid, status) do put_state(sup_pid, :status, status) end end