defmodule Observer.Web.Processes.Page do @moduledoc """ This is the live component responsible for the Processes pillar: an etop-style, auto-refreshing table of the busiest processes on a selected node, ranked by reductions (delta per refresh interval), memory or message queue length, with a per-process drill-down panel. 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.Processes @impl Phoenix.LiveComponent def render(assigns) do ~H"""
<:inner_form> <.form for={@form} id="processes-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 processes...
Processes: {@summary.process_count} Run Queue: {@summary.run_queue} {key}: {format_bytes(Map.get(@summary.memory, key, 0))}
<:col :let={{row, _index}} label="NAME">{row.name} <:col :let={{row, _index}} label="PID">{inspect(row.pid)} <:col :let={{row, _index}} label="MEMORY">{format_bytes(row.memory)} <:col :let={{row, _index}} label={reductions_label(@form.params["sort_by"])}> {row.reductions_diff} <:col :let={{row, _index}} label="MSG QUEUE">{row.message_queue_len} <:col :let={{row, _index}} label="CURRENT FUNCTION">{row.current_function} <:col :let={{_row, index}} label="">

Process {inspect(@details.pid)}

The process has exited.
{key} {value}
""" end @impl Page def handle_mount(socket) when is_connected?(socket) do :net_kernel.monitor_nodes(true) send(self(), {:processes_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, []) |> assign(:summary, nil) |> assign(:previous_reductions, %{}) |> 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, "Processes") 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_reductions, %{}) |> assign(:details, nil) |> assign(:summary, nil) |> assign(:rows, []) else socket end {:noreply, socket |> assign(:form, to_form(params)) |> restart_tick_chain()} end def handle_parent_event("processes-refresh", _params, socket) do {:noreply, restart_tick_chain(socket)} end def handle_parent_event("processes-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, fetch_details(row.pid))} end end def handle_parent_event("processes-details-close", _params, socket) do {:noreply, assign(socket, :details, nil)} end # Bumps the generation and fires an immediate tick - any tick still in flight from the # previous chain is dropped by the generation guard in handle_info/2. defp restart_tick_chain(socket) do tick_gen = socket.assigns.tick_gen + 1 send(self(), {:processes_tick, tick_gen}) assign(socket, :tick_gen, tick_gen) end @impl Page def handle_info({:processes_tick, gen}, %{assigns: %{tick_gen: gen}} = socket) do socket = case Processes.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(), {:processes_tick, gen}, refresh_seconds * 1_000) end {:noreply, socket} end def handle_info({:processes_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_reductions, %{}) |> assign(:details, nil) |> restart_tick_chain()} else {:noreply, socket} end end defp assign_sample(socket, sample) do %{form: form, previous_reductions: previous_reductions, details: details} = socket.assigns sort_by = sort_by_atom(form.params["sort_by"]) limit = positive_int(form.params["limit"], 50) rows = Processes.rank(sample, sort_by, limit, previous_reductions) details = if details do fetch_details(details.pid) end socket |> assign(:rows, rows) |> assign( :summary, Map.take(sample, [:process_count, :run_queue, :memory]) ) |> assign(:previous_reductions, Processes.reductions_by_pid(sample)) |> assign(:details, details) |> assign(:sample_error, nil) end defp fetch_details(pid) do case Processes.details(pid) do {:ok, info} -> %{pid: pid, info: info} {:error, :not_found} -> %{pid: pid, info: :not_found} end 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("memory"), do: :memory defp sort_by_atom("message_queue_len"), do: :message_queue_len defp sort_by_atom(_reductions), do: :reductions 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" => "reductions", "limit" => "50", "refresh_seconds" => "5" } end defp reductions_label("reductions"), do: "REDUCTIONS (Δ)" defp reductions_label(_other), do: "REDUCTIONS" defp attention_msg do assigns = %{} ~H""" Lists the busiest processes on the selected service, ranked by reductions (delta per refresh interval), memory or message queue length - the same bounded collector etop uses. While this page is open the :scheduler_wall_time flag is enabled on the sampled node (etop behaves the same way); it is switched back off when the page is closed. """ 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