Packages

phoenix_kit

1.7.66
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 web channel.ex
Raw

lib/modules/sync/web/channel.ex

defmodule PhoenixKitWeb.SyncChannel do
@moduledoc """
Channel for DB Sync protocol messages.
Handles communication between sender and receiver sites during
a data sync session.
## Architecture
The SENDER site hosts this channel (has data to share).
The RECEIVER connects via WebSocket to pull data.
**Data Flow:**
1. Receiver's WebSocketClient connects to this channel
2. Receiver sends requests (e.g., "request:tables")
3. Channel handles request by querying local database
4. Channel sends response back to Receiver
## Protocol Messages
### From Receiver (requests)
- `request:capabilities` - Get server capabilities/version
- `request:tables` - Request list of available tables
- `request:schema` - Request table schema
- `request:count` - Request record count for table
- `request:records` - Request paginated records
### To Receiver (responses)
- `response:capabilities` - Server capabilities
- `response:tables` - List of available tables
- `response:schema` - Table schema details
- `response:count` - Record count
- `response:records` - Paginated records
- `response:error` - Error response
"""
use Phoenix.Channel
require Logger
alias PhoenixKit.Modules.Sync
alias PhoenixKit.Modules.Sync.DataExporter
alias PhoenixKit.Modules.Sync.SchemaInspector
@impl true
def join("transfer:" <> code, _params, socket) do
# Verify the code matches the socket's session
if socket.assigns.session_code == code do
# Update session with channel PID so sender's LiveView can track
Sync.update_session(code, %{channel_pid: self()})
# Notify the sender's LiveView that a receiver has joined
send_to_sender(socket.assigns.session, {:receiver_joined, self()})
Logger.info("Sync: Receiver joined channel for code #{code}")
{:ok, socket}
else
Logger.warning(
"Sync: Channel join mismatch - expected #{socket.assigns.session_code}, got #{code}"
)
{:error, %{reason: "code_mismatch"}}
end
end
# ===========================================
# INCOMING REQUESTS FROM RECEIVER
# ===========================================
@impl true
def handle_in("request:capabilities", %{"ref" => ref}, socket) do
Logger.debug("Sync.Channel: Capabilities requested")
capabilities = %{
version: "1.0.0",
phoenix_kit_version: Application.spec(:phoenix_kit, :vsn) |> to_string(),
features: ["list_tables", "get_schema", "fetch_records"]
}
push(socket, "response:capabilities", %{capabilities: capabilities, ref: ref})
{:noreply, socket}
end
def handle_in("request:tables", %{"ref" => ref}, socket) do
Logger.debug("Sync.Channel: Tables requested")
case SchemaInspector.list_tables() do
{:ok, tables} ->
push(socket, "response:tables", %{tables: tables, ref: ref})
{:error, reason} ->
push(socket, "response:error", %{
error: "Failed to list tables: #{inspect(reason)}",
ref: ref
})
end
{:noreply, socket}
end
def handle_in("request:schema", %{"table" => table, "ref" => ref}, socket) do
Logger.info("Sync.Channel: Schema requested for #{table}")
case SchemaInspector.get_schema(table) do
{:ok, schema} ->
Logger.info("Sync.Channel: Schema found for #{table}, columns: #{length(schema.columns)}")
push(socket, "response:schema", %{schema: schema, ref: ref})
{:error, :not_found} ->
Logger.warning("Sync.Channel: Table not found: #{table}")
push(socket, "response:error", %{error: "Table not found: #{table}", ref: ref})
{:error, reason} ->
Logger.error("Sync.Channel: Failed to get schema for #{table}: #{inspect(reason)}")
push(socket, "response:error", %{
error: "Failed to get schema: #{inspect(reason)}",
ref: ref
})
end
{:noreply, socket}
end
def handle_in("request:count", %{"table" => table, "ref" => ref}, socket) do
Logger.debug("Sync.Channel: Count requested for #{table}")
case DataExporter.get_count(table) do
{:ok, count} ->
push(socket, "response:count", %{count: count, ref: ref})
{:error, reason} ->
push(socket, "response:error", %{
error: "Failed to get count: #{inspect(reason)}",
ref: ref
})
end
{:noreply, socket}
end
def handle_in("request:records", payload, socket) do
table = Map.fetch!(payload, "table")
ref = Map.fetch!(payload, "ref")
offset = Map.get(payload, "offset", 0)
limit = Map.get(payload, "limit", 100)
Logger.debug(
"Sync.Channel: Records requested for #{table} (offset: #{offset}, limit: #{limit})"
)
case DataExporter.fetch_records(table, offset: offset, limit: limit) do
{:ok, records} ->
push(socket, "response:records", %{
records: records,
offset: offset,
has_more: length(records) == limit,
ref: ref
})
{:error, reason} ->
push(socket, "response:error", %{
error: "Failed to fetch records: #{inspect(reason)}",
ref: ref
})
end
{:noreply, socket}
end
def handle_in(event, payload, socket) do
Logger.warning("Sync: Unknown event #{event} with payload #{inspect(payload)}")
{:reply, {:error, %{message: "Unknown event: #{event}"}}, socket}
end
# ===========================================
# HANDLE INFO (from Sender's LiveView)
# ===========================================
@impl true
def handle_info(_msg, socket) do
{:noreply, socket}
end
# ===========================================
# TERMINATE
# ===========================================
@impl true
def terminate(reason, socket) do
Logger.info(
"Sync: Channel terminated for code #{socket.assigns.session_code}, reason: #{inspect(reason)}"
)
# Notify the sender's LiveView that receiver has disconnected (with PID for multi-receiver support)
send_to_sender(socket.assigns.session, {:receiver_disconnected, self()})
:ok
end
# ===========================================
# PRIVATE FUNCTIONS
# ===========================================
defp send_to_sender(%{owner_pid: pid}, message) when is_pid(pid) do
if Process.alive?(pid) do
send(pid, {:sync, message})
end
end
defp send_to_sender(_session, _message), do: :ok
end