Packages

phoenix_kit

1.7.49
1.7.207 1.7.206 1.7.205 1.7.204 1.7.203 1.7.202 1.7.201 1.7.200 1.7.199 1.7.198 1.7.197 1.7.196 1.7.194 1.7.193 1.7.192 1.7.191 1.7.190 1.7.189 1.7.187 1.7.186 1.7.185 1.7.184 1.7.183 1.7.182 1.7.181 1.7.180 1.7.179 1.7.178 1.7.177 1.7.176 1.7.175 1.7.174 1.7.173 1.7.172 1.7.171 1.7.170 1.7.169 1.7.168 1.7.167 1.7.166 1.7.165 1.7.164 1.7.162 1.7.161 1.7.160 1.7.159 1.7.157 1.7.156 1.7.155 1.7.154 1.7.153 1.7.152 1.7.151 1.7.150 1.7.149 1.7.146 1.7.145 1.7.144 1.7.143 1.7.138 1.7.133 1.7.132 1.7.131 1.7.130 1.7.128 1.7.126 1.7.125 1.7.121 1.7.120 1.7.119 1.7.118 1.7.117 1.7.116 1.7.115 1.7.114 1.7.113 1.7.112 1.7.111 1.7.110 1.7.109 1.7.108 1.7.107 1.7.106 1.7.105 1.7.104 1.7.103 1.7.102 1.7.101 1.7.100 1.7.99 1.7.98 1.7.97 1.7.96 1.7.95 1.7.94 1.7.93 1.7.92 1.7.91 1.7.90 1.7.89 1.7.88 1.7.87 1.7.86 1.7.85 1.7.84 1.7.83 1.7.82 1.7.81 1.7.80 1.7.79 1.7.78 1.7.77 1.7.76 1.7.75 1.7.74 1.7.71 1.7.70 1.7.69 1.7.66 1.7.65 1.7.64 1.7.63 1.7.62 1.7.61 1.7.59 1.7.58 1.7.57 1.7.56 1.7.55 1.7.54 1.7.53 1.7.52 1.7.51 1.7.49 1.7.44 1.7.43 1.7.42 1.7.41 1.7.39 1.7.38 1.7.37 1.7.36 1.7.34 1.7.33 1.7.31 1.7.30 1.7.29 1.7.28 1.7.27 1.7.26 1.7.25 1.7.24 1.7.23 1.7.22 1.7.21 1.7.20 1.7.19 1.7.18 1.7.17 1.7.16 1.7.15 1.7.14 1.7.13 1.7.12 1.7.11 1.7.10 1.7.9 1.7.8 1.7.7 1.7.6 1.7.5 1.7.4 1.7.3 1.7.2 1.7.1 1.7.0 1.6.20 1.6.19 1.6.18 1.6.17 1.6.16 1.6.15 1.6.14 1.6.13 1.6.12 1.6.11 1.6.10 1.6.9 1.6.8 1.6.7 1.6.6 1.6.5 1.6.4 1.6.3 1.5.2 1.5.1 1.5.0 1.4.9 1.4.8 1.4.7 1.4.6 1.4.5 1.4.4 1.4.3 1.4.2 1.4.1 1.4.0 1.3.2 1.3.1 1.3.0 1.2.10 1.2.9 1.2.8 1.2.7 1.2.5 1.2.4 1.2.2 1.2.1 1.2.0 1.1.0 1.0.0

A foundation for building Elixir Phoenix apps — SaaS, social networks, ERP systems, marketplaces, and more

Current section

Files

Jump to
phoenix_kit lib modules sync websocket_client.ex
Raw

lib/modules/sync/websocket_client.ex

defmodule PhoenixKit.Modules.Sync.WebSocketClient do
@moduledoc """
WebSocket client for Sync receiver connections.
Uses WebSockex to connect from the receiver site to the sender's
WebSocket endpoint. Sends requests and receives responses.
## Architecture
The receiver uses this client to:
1. Connect to the sender's hosted channel
2. Send requests for data (tables, schema, records)
3. Receive responses and notify the caller (LiveView)
## Usage
{:ok, pid} = WebSocketClient.start_link(
url: "https://sender-site.com",
code: "ABC12345",
caller: self()
)
# Request available tables
WebSocketClient.request_tables(pid)
# Receive: {:sync_client, {:tables, tables}}
# Request records
WebSocketClient.request_records(pid, "users", offset: 0, limit: 100)
# Receive: {:sync_client, {:records, "users", result}}
"""
use WebSockex
require Logger
@heartbeat_interval 30_000
@join_timeout 10_000
defstruct [
:url,
:code,
:caller,
:receiver_info,
:ref_counter,
:pending_refs,
:joined,
:heartbeat_ref
]
# ===========================================
# PUBLIC API
# ===========================================
@doc """
Starts the WebSocket client and connects to the sender.
"""
@spec start_link(keyword()) :: {:ok, pid()} | {:error, any()}
def start_link(opts) do
url = Keyword.fetch!(opts, :url)
code = Keyword.fetch!(opts, :code)
caller = Keyword.fetch!(opts, :caller)
receiver_info = Keyword.get(opts, :receiver_info, %{})
ws_url = build_websocket_url(url, code)
state = %__MODULE__{
url: url,
code: code,
caller: caller,
receiver_info: receiver_info,
ref_counter: 0,
pending_refs: %{},
joined: false,
heartbeat_ref: nil
}
WebSockex.start_link(ws_url, __MODULE__, state, [])
end
@doc """
Disconnects the WebSocket client.
"""
@spec disconnect(pid()) :: :ok
def disconnect(pid) do
WebSockex.cast(pid, :disconnect)
end
@doc """
Request server capabilities.
"""
@spec request_capabilities(pid()) :: :ok
def request_capabilities(pid) do
WebSockex.cast(pid, :request_capabilities)
end
@doc """
Request list of available tables from sender.
"""
@spec request_tables(pid()) :: :ok
def request_tables(pid) do
WebSockex.cast(pid, :request_tables)
end
@doc """
Request schema for a specific table.
"""
@spec request_schema(pid(), String.t()) :: :ok
def request_schema(pid, table) do
WebSockex.cast(pid, {:request_schema, table})
end
@doc """
Request record count for a table.
"""
@spec request_count(pid(), String.t()) :: :ok
def request_count(pid, table) do
WebSockex.cast(pid, {:request_count, table})
end
@doc """
Request records from a table with pagination.
"""
@spec request_records(pid(), String.t(), keyword()) :: :ok
def request_records(pid, table, opts \\ []) do
WebSockex.cast(pid, {:request_records, table, opts})
end
# ===========================================
# WEBSOCKEX CALLBACKS
# ===========================================
@impl true
def handle_connect(_conn, state) do
Logger.info("Sync.WebSocketClient: Connected to #{state.url}")
# Join the transfer channel with receiver info
join_payload = %{receiver_info: state.receiver_info}
join_msg = encode_message("transfer:#{state.code}", "phx_join", join_payload, make_ref(state))
WebSockex.cast(self(), {:send_raw, join_msg})
# Start join timeout
Process.send_after(self(), :join_timeout, @join_timeout)
{:ok, state}
end
@impl true
def handle_frame({:text, msg}, state) do
case Jason.decode(msg) do
{:ok, [_join_ref, ref, topic, event, payload]} ->
handle_phoenix_message(topic, event, payload, ref, state)
{:error, reason} ->
Logger.warning("Sync.WebSocketClient: Failed to decode message: #{inspect(reason)}")
{:ok, state}
end
end
def handle_frame(_frame, state) do
{:ok, state}
end
@impl true
def handle_cast({:send_raw, msg}, state) do
{:reply, {:text, msg}, state}
end
def handle_cast(:disconnect, state) do
notify_caller(state, :disconnected)
{:close, state}
end
def handle_cast(:request_capabilities, state) do
if state.joined do
{ref, state} = next_ref(state)
msg = encode_message("transfer:#{state.code}", "request:capabilities", %{ref: ref}, ref)
state = track_request(state, ref, :capabilities)
{:reply, {:text, msg}, state}
else
Logger.warning("Sync.WebSocketClient: Cannot send request - not joined")
{:ok, state}
end
end
def handle_cast(:request_tables, state) do
if state.joined do
{ref, state} = next_ref(state)
msg = encode_message("transfer:#{state.code}", "request:tables", %{ref: ref}, ref)
state = track_request(state, ref, :tables)
{:reply, {:text, msg}, state}
else
Logger.warning("Sync.WebSocketClient: Cannot send request - not joined")
{:ok, state}
end
end
def handle_cast({:request_schema, table}, state) do
if state.joined do
{ref, state} = next_ref(state)
msg =
encode_message("transfer:#{state.code}", "request:schema", %{table: table, ref: ref}, ref)
state = track_request(state, ref, {:schema, table})
{:reply, {:text, msg}, state}
else
Logger.warning("Sync.WebSocketClient: Cannot send request - not joined")
{:ok, state}
end
end
def handle_cast({:request_count, table}, state) do
if state.joined do
{ref, state} = next_ref(state)
msg =
encode_message("transfer:#{state.code}", "request:count", %{table: table, ref: ref}, ref)
state = track_request(state, ref, {:count, table})
{:reply, {:text, msg}, state}
else
Logger.warning("Sync.WebSocketClient: Cannot send request - not joined")
{:ok, state}
end
end
def handle_cast({:request_records, table, opts}, state) do
if state.joined do
{ref, state} = next_ref(state)
offset = Keyword.get(opts, :offset, 0)
limit = Keyword.get(opts, :limit, 100)
msg =
encode_message(
"transfer:#{state.code}",
"request:records",
%{
table: table,
offset: offset,
limit: limit,
ref: ref
},
ref
)
state = track_request(state, ref, {:records, table})
{:reply, {:text, msg}, state}
else
Logger.warning("Sync.WebSocketClient: Cannot send request - not joined")
{:ok, state}
end
end
@impl true
def handle_info(:heartbeat, state) do
if state.joined do
{ref, state} = next_ref(state)
msg = encode_message("phoenix", "heartbeat", %{}, ref)
state = %{state | heartbeat_ref: ref}
schedule_heartbeat()
{:reply, {:text, msg}, state}
else
{:ok, state}
end
end
def handle_info(:join_timeout, state) do
if state.joined do
{:ok, state}
else
Logger.warning("Sync.WebSocketClient: Join timeout")
notify_caller(state, {:error, :join_timeout})
{:close, state}
end
end
def handle_info(_msg, socket) do
{:ok, socket}
end
@impl true
def handle_disconnect(%{reason: reason}, state) do
Logger.info("Sync.WebSocketClient: Disconnected - #{inspect(reason)}")
notify_caller(state, {:disconnected, reason})
{:ok, state}
end
@impl true
def terminate(reason, state) do
Logger.info("Sync.WebSocketClient: Terminating - #{inspect(reason)}")
notify_caller(state, {:terminated, reason})
:ok
end
# ===========================================
# PHOENIX MESSAGE HANDLERS
# ===========================================
defp handle_phoenix_message(_topic, "phx_reply", %{"status" => "ok"} = payload, ref, state) do
cond do
# Join reply
not state.joined and Map.get(payload, "response") == %{} ->
Logger.info("Sync.WebSocketClient: Joined channel")
state = %{state | joined: true}
notify_caller(state, :connected)
schedule_heartbeat()
{:ok, state}
# Heartbeat reply
state.heartbeat_ref == ref ->
{:ok, %{state | heartbeat_ref: nil}}
true ->
{:ok, state}
end
end
defp handle_phoenix_message(_topic, "phx_reply", %{"status" => "error"} = payload, _ref, state) do
Logger.warning("Sync.WebSocketClient: Error response - #{inspect(payload)}")
notify_caller(state, {:error, payload})
{:ok, state}
end
defp handle_phoenix_message(_topic, "phx_error", payload, _ref, state) do
Logger.error("Sync.WebSocketClient: Channel error - #{inspect(payload)}")
notify_caller(state, {:error, payload})
{:ok, state}
end
defp handle_phoenix_message(_topic, "phx_close", _payload, _ref, state) do
Logger.info("Sync.WebSocketClient: Channel closed by server")
notify_caller(state, :channel_closed)
{:close, state}
end
# Handle response messages from sender
defp handle_phoenix_message(
_topic,
"response:capabilities",
%{"capabilities" => caps, "ref" => ref},
_msg_ref,
state
) do
Logger.debug("Sync.WebSocketClient: Received capabilities")
{request_type, state} = pop_request(state, ref)
if request_type == :capabilities do
notify_caller(state, {:capabilities, caps})
end
{:ok, state}
end
defp handle_phoenix_message(
_topic,
"response:tables",
%{"tables" => tables, "ref" => ref},
_msg_ref,
state
) do
Logger.debug("Sync.WebSocketClient: Received #{length(tables)} tables")
{request_type, state} = pop_request(state, ref)
if request_type == :tables do
notify_caller(state, {:tables, tables})
end
{:ok, state}
end
defp handle_phoenix_message(
_topic,
"response:schema",
%{"schema" => schema, "ref" => ref},
_msg_ref,
state
) do
Logger.info("Sync.WebSocketClient: Received schema response, ref: #{ref}")
{request_type, state} = pop_request(state, ref)
Logger.info("Sync.WebSocketClient: Request type for ref #{ref}: #{inspect(request_type)}")
case request_type do
{:schema, table} ->
Logger.info("Sync.WebSocketClient: Notifying caller with schema for #{table}")
notify_caller(state, {:schema, table, schema})
other ->
Logger.warning("Sync.WebSocketClient: Unexpected request type: #{inspect(other)}")
end
{:ok, state}
end
defp handle_phoenix_message(
_topic,
"response:count",
%{"count" => count, "ref" => ref},
_msg_ref,
state
) do
Logger.debug("Sync.WebSocketClient: Received count: #{count}")
{request_type, state} = pop_request(state, ref)
case request_type do
{:count, table} -> notify_caller(state, {:count, table, count})
_ -> :ok
end
{:ok, state}
end
defp handle_phoenix_message(_topic, "response:records", payload, _msg_ref, state) do
records = Map.get(payload, "records", [])
ref = Map.get(payload, "ref")
offset = Map.get(payload, "offset", 0)
has_more = Map.get(payload, "has_more", false)
Logger.debug("Sync.WebSocketClient: Received #{length(records)} records")
{request_type, state} = pop_request(state, ref)
case request_type do
{:records, table} ->
result = %{records: records, offset: offset, has_more: has_more}
notify_caller(state, {:records, table, result})
_ ->
:ok
end
{:ok, state}
end
defp handle_phoenix_message(
_topic,
"response:error",
%{"error" => error, "ref" => ref},
_msg_ref,
state
) do
Logger.warning("Sync.WebSocketClient: Error response - #{error}")
{request_type, state} = pop_request(state, ref)
notify_caller(state, {:request_error, request_type, error})
{:ok, state}
end
defp handle_phoenix_message(topic, event, payload, _ref, state) do
Logger.debug("Sync.WebSocketClient: Received #{event} on #{topic}: #{inspect(payload)}")
notify_caller(state, {:message, event, payload})
{:ok, state}
end
# ===========================================
# PRIVATE FUNCTIONS
# ===========================================
defp build_websocket_url(base_url, code) do
uri = URI.parse(base_url)
scheme =
case uri.scheme do
"https" -> "wss"
"http" -> "ws"
"wss" -> "wss"
"ws" -> "ws"
_ -> "wss"
end
path =
case uri.path do
nil ->
"/sync/websocket"
"" ->
"/sync/websocket"
# If path already ends with /sync/websocket, use as-is
# Otherwise append /sync/websocket to the prefix path
p ->
if String.ends_with?(p, "/sync/websocket") do
p
else
"#{String.trim_trailing(p, "/")}/sync/websocket"
end
end
query =
case uri.query do
nil -> "code=#{code}&vsn=2.0.0"
q -> "#{q}&code=#{code}&vsn=2.0.0"
end
URI.to_string(%{uri | scheme: scheme, path: path, query: query})
end
defp encode_message(topic, event, payload, ref) do
Jason.encode!([nil, ref, topic, event, payload])
end
defp next_ref(state) do
ref = to_string(state.ref_counter + 1)
{ref, %{state | ref_counter: state.ref_counter + 1}}
end
defp make_ref(%{ref_counter: counter}) do
to_string(counter + 1)
end
defp track_request(state, ref, request_type) do
pending = Map.put(state.pending_refs, ref, request_type)
%{state | pending_refs: pending}
end
defp pop_request(state, ref) do
{request_type, pending} = Map.pop(state.pending_refs, ref)
{request_type, %{state | pending_refs: pending}}
end
defp schedule_heartbeat do
Process.send_after(self(), :heartbeat, @heartbeat_interval)
end
defp notify_caller(%{caller: caller}, message) when is_pid(caller) do
if Process.alive?(caller) do
send(caller, {:sync_client, message})
end
end
defp notify_caller(_, _), do: :ok
end