defmodule MicrowavepropWeb.BackfillLive do @moduledoc false use MicrowavepropWeb, :live_view import Ecto.Query alias Microwaveprop.Radio.Contact alias Microwaveprop.Repo alias Microwaveprop.Workers.ContactWeatherEnqueueWorker @refresh_interval 2_000 @impl true def mount(_params, _session, socket) do if connected?(socket) do Process.send_after(self(), :refresh_stats, @refresh_interval) end stats = fetch_stats() unprocessed = count_unprocessed() {:ok, assign(socket, page_title: "Backfill", limit: 500, stats: stats, unprocessed: unprocessed, last_enqueued: nil, enqueuing: false )} end @impl true def handle_event("enqueue", %{"limit" => limit_str}, socket) do limit = String.to_integer(limit_str) # Run in a task to not block the LiveView pid = self() Task.start(fn -> count = enqueue_batch(limit) send(pid, {:enqueued, count}) end) {:noreply, assign(socket, enqueuing: true)} end @impl true def handle_info({:enqueued, count}, socket) do {:noreply, assign(socket, last_enqueued: count, enqueuing: false, stats: fetch_stats(), unprocessed: count_unprocessed() )} end def handle_info(:refresh_stats, socket) do Process.send_after(self(), :refresh_stats, @refresh_interval) {:noreply, assign(socket, stats: fetch_stats(), unprocessed: count_unprocessed() )} end defp enqueue_batch(limit) do contacts = from(c in Contact, where: not is_nil(c.pos1) and (c.hrrr_queued == false or c.weather_queued == false or c.terrain_queued == false or c.iemre_queued == false), order_by: [desc: c.qso_timestamp], limit: ^limit ) |> Repo.all() Enum.each(contacts, &ContactWeatherEnqueueWorker.enqueue_for_contact/1) length(contacts) end defp count_unprocessed do Repo.one( from(c in Contact, where: not is_nil(c.pos1) and (c.hrrr_queued == false or c.weather_queued == false or c.terrain_queued == false or c.iemre_queued == false), select: count(c.id) ) ) end defp fetch_stats do jobs = Repo.all( from(j in "oban_jobs", where: j.state in ["available", "executing", "scheduled", "retryable"], group_by: [j.worker, j.state], select: {j.worker, j.state, count(j.id)} ) ) by_worker = Enum.group_by(jobs, &elem(&1, 0), fn {_, state, count} -> {state, count} end) |> Enum.map(fn {worker, states} -> short_name = worker |> String.split(".") |> List.last() total = Enum.reduce(states, 0, fn {_, c}, acc -> acc + c end) executing = Enum.find_value(states, 0, fn {s, c} -> if s == "executing", do: c end) available = Enum.find_value(states, 0, fn {s, c} -> if s == "available", do: c end) retryable = Enum.find_value(states, 0, fn {s, c} -> if s == "retryable", do: c end) %{ worker: short_name, total: total, executing: executing, available: available, retryable: retryable } end) |> Enum.sort_by(& &1.total, :desc) completed_1h = Repo.one( from(j in "oban_jobs", where: j.state == "completed" and j.completed_at > ago(1, "hour"), select: count(j.id) ) ) %{by_worker: by_worker, completed_1h: completed_1h} end @impl true def render(assigns) do ~H""" <.header> Backfill Dashboard <:subtitle>Enqueue and monitor contact enrichment jobs
Contacts Needing Enrichment
{@unprocessed}
Active/Queued Jobs
{Enum.reduce(@stats.by_worker, 0, fn w, acc -> acc + w.total end)}
Completed (last hour)
{@stats.completed_1h}
<%= if @last_enqueued do %> Enqueued {@last_enqueued} contacts <% end %>

Job Queue Status

<%= if @stats.by_worker == [] do %>

No active jobs.

<% else %>
<%= for w <- @stats.by_worker do %> <% end %>
Worker Executing Available Retryable Total
{w.worker} <%= if w.executing > 0 do %> {w.executing} <% else %> 0 <% end %> {w.available} <%= if w.retryable > 0 do %> {w.retryable} <% else %> 0 <% end %> {w.total}
<% end %>

Auto-refreshes every 2 seconds.

""" end end