defmodule Observer.Web.Network.Page do
@moduledoc """
This is the live component responsible for the Network pillar: the busiest inet ports on a
selected node, ranked by received/sent bytes per refresh interval (deltas between samples,
cumulative on the first tick), plus the NIF-based `socket` module sockets that port listings
miss. A drill-down panel shows the selected port's full inet options.
Refresh ticks carry a generation counter: changing any control bumps the generation and starts
a new timer chain, so a tick from a cancelled chain that was already in flight is ignored
instead of spawning a second chain.
"""
@behaviour Observer.Web.Page
use Observer.Web, :live_component
alias Observer.Web.Components.Attention
alias Observer.Web.Components.Core
alias Observer.Web.Page
alias ObserverWeb.Network
@impl Phoenix.LiveComponent
def render(assigns) do
~H"""
<:inner_form>
<.form
for={@form}
id="network-update-form"
class="flex flex-col md:flex-row md:items-end shrink-0 ml-2 mr-2 py-2 text-xs text-center text-zinc-800 dark:text-white whitespace-nowrap gap-x-5 gap-y-1"
phx-change="form-update"
>
<:inner_button>
REFRESH
Could not sample {@form.params["service"]}: {inspect(@sample_error)}
Collecting network endpoints...
Inet Ports
No inet ports on this service.
<:col :let={{row, _index}} label="PORT">{inspect(row.port)}
<:col :let={{row, _index}} label="DRIVER">{row.name}
<:col :let={{row, _index}} label="OWNER">{row.owner_label}
<:col :let={{row, _index}} label="LOCAL">{row.local}
<:col :let={{row, _index}} label="REMOTE">{row.remote}
<:col :let={{row, _index}} label="RECV (Δ)">{format_bytes(row.recv_diff)}
<:col :let={{row, _index}} label="SENT (Δ)">{format_bytes(row.send_diff)}
<:col :let={{row, _index}} label="QUEUE">{row.queue_size}
<:col :let={{_row, index}} label="">
DETAILS
Sockets (NIF)
No `socket` module sockets on this service.
<:col :let={socket} label="ID">{socket.id_str}
<:col :let={socket} label="KIND">{inspect(socket.kind)}
<:col :let={socket} label="DOMAIN">{inspect(socket.domain)}
<:col :let={socket} label="TYPE">{inspect(socket.type)}
<:col :let={socket} label="PROTOCOL">{inspect(socket.protocol)}
<:col :let={socket} label="READ">{format_bytes(socket.read_bytes)}
<:col :let={socket} label="WRITE">{format_bytes(socket.write_bytes)}
Port {inspect(@details.port)} ({@details.local} → {@details.remote})
CLOSE
"""
end
@impl Page
def handle_mount(socket) when is_connected?(socket) do
:net_kernel.monitor_nodes(true)
send(self(), {:network_tick, 1})
socket
|> assign_defaults()
|> assign(:tick_gen, 1)
end
def handle_mount(socket) do
socket
|> assign_defaults()
|> assign(:tick_gen, 0)
end
defp assign_defaults(socket) do
socket
|> assign(:services, services())
|> assign(:rows, nil)
|> assign(:sockets, nil)
|> assign(:previous_counters, %{})
|> assign(:details, nil)
|> assign(:sample_error, nil)
|> assign(:form, to_form(default_form_options()))
end
@impl Page
def handle_params(params, _uri, socket) do
{:noreply, apply_action(socket, socket.assigns.live_action, params)}
end
defp apply_action(socket, :index, _params) do
assign(socket, :page_title, "Network")
end
@impl Page
def handle_parent_event("form-update", params, socket) do
service_changed? = params["service"] != socket.assigns.form.params["service"]
socket =
if service_changed? do
socket
|> assign(:previous_counters, %{})
|> assign(:details, nil)
|> assign(:rows, nil)
|> assign(:sockets, nil)
else
socket
end
{:noreply,
socket
|> assign(:form, to_form(params))
|> restart_tick_chain()}
end
def handle_parent_event("network-refresh", _params, socket) do
{:noreply, restart_tick_chain(socket)}
end
def handle_parent_event("network-select-row", %{"index" => index}, socket) do
case Enum.at(socket.assigns.rows || [], positive_int(index, 0)) do
nil -> {:noreply, socket}
row -> {:noreply, assign(socket, :details, row)}
end
end
def handle_parent_event("network-details-close", _params, socket) do
{:noreply, assign(socket, :details, nil)}
end
defp restart_tick_chain(socket) do
tick_gen = socket.assigns.tick_gen + 1
send(self(), {:network_tick, tick_gen})
assign(socket, :tick_gen, tick_gen)
end
@impl Page
def handle_info({:network_tick, gen}, %{assigns: %{tick_gen: gen}} = socket) do
socket =
case Network.sample(selected_service(socket)) do
{:ok, sample} -> assign_sample(socket, sample)
{:error, reason} -> assign(socket, :sample_error, reason)
end
refresh_seconds = positive_int(socket.assigns.form.params["refresh_seconds"], 0)
if refresh_seconds > 0 do
Process.send_after(self(), {:network_tick, gen}, refresh_seconds * 1_000)
end
{:noreply, socket}
end
def handle_info({:network_tick, _stale_gen}, socket) do
{:noreply, socket}
end
def handle_info({:nodeup, _node}, socket) do
{:noreply, assign(socket, :services, services())}
end
def handle_info({:nodedown, node}, %{assigns: %{form: form}} = socket) do
socket = assign(socket, :services, services())
if to_string(node) == form.params["service"] do
params = %{form.params | "service" => to_string(Node.self())}
{:noreply,
socket
|> assign(:form, to_form(params))
|> assign(:previous_counters, %{})
|> assign(:details, nil)
|> restart_tick_chain()}
else
{:noreply, socket}
end
end
defp assign_sample(socket, %{ports: ports, sockets: sockets}) do
%{form: form, previous_counters: previous_counters, details: details} = socket.assigns
sort_by = sort_by_atom(form.params["sort_by"])
rows = Network.rank_ports(ports, sort_by, 250, previous_counters)
details =
if details do
Enum.find(rows, &(&1.port == details.port))
end
socket
|> assign(:rows, rows)
|> assign(:sockets, Enum.sort_by(sockets, &(&1.read_bytes + &1.write_bytes), :desc))
|> assign(:previous_counters, Network.counters_by_port(ports))
|> assign(:details, details)
|> assign(:sample_error, nil)
end
defp details_rows(details) do
[
{"owner", details.owner_label},
{"driver", details.name},
{"recv (total)", format_bytes(details.recv_oct)},
{"sent (total)", format_bytes(details.send_oct)},
{"queue size", to_string(details.queue_size)},
{"memory", format_bytes(details.memory)},
{"statistics", inspect(Keyword.get(details.inet, :statistics, []), limit: 30)},
{"options", inspect(Keyword.get(details.inet, :options, []), limit: 50)}
]
end
defp services do
Enum.map([Node.self() | Node.list()], &{&1, to_string(&1)})
end
# The form's service value only becomes a node if it matches a currently known node - a
# crafted payload (or a node that just left the cluster) safely falls back to the local one.
defp selected_service(socket) do
service = socket.assigns.form.params["service"]
[Node.self() | Node.list()]
|> Enum.find(Node.self(), &(to_string(&1) == service))
end
defp sort_by_atom("recv"), do: :recv
defp sort_by_atom("send"), do: :send
defp sort_by_atom(_total), do: :total
defp positive_int(value, default) do
case Integer.parse(value || "") do
{int, ""} when int >= 0 -> int
_invalid -> default
end
end
defp default_form_options do
%{
"service" => to_string(Node.self()),
"sort_by" => "total",
"refresh_seconds" => "5"
}
end
defp attention_msg do
assigns = %{}
~H"""
The busiest inet ports on the selected service, ranked by bytes received/sent per refresh
interval, with the local and remote endpoints and the owning process - plus the NIF-based
socket module sockets that port listings miss. Click DETAILS for a port's full statistics
and options.
"""
end
defp format_bytes(bytes) when bytes >= 1_073_741_824,
do: "#{Float.round(bytes / 1_073_741_824, 1)} GB"
defp format_bytes(bytes) when bytes >= 1_048_576, do: "#{Float.round(bytes / 1_048_576, 1)} MB"
defp format_bytes(bytes) when bytes >= 1_024, do: "#{Float.round(bytes / 1_024, 1)} KB"
defp format_bytes(bytes), do: "#{bytes} B"
end