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>
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"> <: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="">

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})

{key} {value}
""" 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