Hooked it all up to a websocket

This commit is contained in:
Kieran Eglin 2024-05-17 10:40:03 -07:00
parent 93de2514e2
commit bd28e55342
No known key found for this signature in database
GPG key ID: 193984967FCF432D
7 changed files with 160 additions and 41 deletions

View file

@ -24,7 +24,7 @@ defmodule Pinchflat.Application do
PinchflatWeb.Endpoint PinchflatWeb.Endpoint
] ]
:ok = Oban.Telemetry.attach_default_logger() attach_oban_telemetry()
Logger.add_handlers(:pinchflat) Logger.add_handlers(:pinchflat)
# See https://hexdocs.pm/elixir/Supervisor.html # See https://hexdocs.pm/elixir/Supervisor.html
@ -40,4 +40,11 @@ defmodule Pinchflat.Application do
PinchflatWeb.Endpoint.config_change(changed, removed) PinchflatWeb.Endpoint.config_change(changed, removed)
:ok :ok
end end
defp attach_oban_telemetry do
events = [[:oban, :job, :start], [:oban, :job, :stop], [:oban, :job, :exception]]
:ok = Oban.Telemetry.attach_default_logger()
:telemetry.attach_many("job-telemetry-broadcast", events, &PinchflatWeb.Telemetry.job_state_change_broadcast/4, [])
end
end end

View file

@ -20,16 +20,14 @@ defmodule PinchflatWeb.Pages.PageController do
end end
defp render_home_page(conn) do defp render_home_page(conn) do
# TODO: revert
conn conn
|> render(:home, |> render(:home,
media_profile_count: 1 || Repo.aggregate(MediaProfile, :count, :id), media_profile_count: Repo.aggregate(MediaProfile, :count, :id),
source_count: 1 || Repo.aggregate(Source, :count, :id), source_count: Repo.aggregate(Source, :count, :id),
media_item_count: media_item_count:
1 || MediaQuery.new()
MediaQuery.new() |> MediaQuery.with_media_downloaded_at()
|> MediaQuery.with_media_downloaded_at() |> Repo.aggregate(:count, :id)
|> Repo.aggregate(:count, :id)
) )
end end

View file

@ -6,7 +6,7 @@ defmodule Pinchflat.Pages.HistoryTableLive do
alias Pinchflat.Utils.NumberUtils alias Pinchflat.Utils.NumberUtils
alias PinchflatWeb.CustomComponents.TextComponents alias PinchflatWeb.CustomComponents.TextComponents
@limit 10 @limit 5
def render(%{records: []} = assigns) do def render(%{records: []} = assigns) do
~H""" ~H"""

View file

@ -26,16 +26,15 @@
</div> </div>
<div class="rounded-sm border shadow-default border-strokedark bg-boxdark mt-4 p-5"> <div class="rounded-sm border shadow-default border-strokedark bg-boxdark mt-4 p-5">
<span class="text-2xl font-medium mb-4">Active Tasks</span> <span class="text-2xl font-medium mb-4">Media History</span>
<section class="mt-6"> <section class="mt-6">
<%= live_render(@conn, Pinchflat.Pages.JobTableLive) %> <%= live_render(@conn, Pinchflat.Pages.HistoryTableLive) %>
</section> </section>
</div> </div>
<div class="rounded-sm border shadow-default border-strokedark bg-boxdark mt-4 p-5"> <div class="rounded-sm border shadow-default border-strokedark bg-boxdark mt-4 p-5">
<span class="text-2xl font-medium mb-4">Media History</span> <span class="text-2xl font-medium mb-4">Active Tasks</span>
<section class="mt-6"> <section class="mt-6 min-h-80">
<%!-- TODO: revert --%> <%= live_render(@conn, Pinchflat.Pages.JobTableLive) %>
<%!-- <%= live_render(@conn, Pinchflat.Pages.HistoryTableLive) %> --%>
</section> </section>
</div> </div>

View file

@ -9,51 +9,50 @@ defmodule Pinchflat.Pages.JobTableLive do
def render(%{tasks: []} = assigns) do def render(%{tasks: []} = assigns) do
~H""" ~H"""
<div class="mb-4 flex items-center"> <div class="mb-4 flex items-center">
<p>Nothing here!</p> <p>Nothing Here!</p>
</div> </div>
""" """
end end
def render(assigns) do def render(assigns) do
~H""" ~H"""
<div> <div class="max-w-full overflow-x-auto">
<div class="max-w-full overflow-x-auto"> <.table rows={@tasks} table_class="text-white">
<.table rows={@tasks} table_class="text-white"> <:col :let={task} label="Task">
<:col :let={task} label="Task"> <%= worker_to_task_name(task.job.worker) %>
<%= worker_to_task_name(task.job.worker) %> </:col>
</:col> <:col :let={task} label="Subject">
<:col :let={task} label="Subject"> <.subtle_link href={task_to_link(task)}>
<.subtle_link href={task_to_link(task)}> <%= StringUtils.truncate(task_to_record_name(task), 35) %>
<%= StringUtils.truncate(task_to_record_name(task), 35) %> </.subtle_link>
</.subtle_link> </:col>
</:col> <:col :let={task} label="Attempt No.">
<:col :let={task} label="Attempt No."> <%= task.job.attempt %>
<%= task.job.attempt %> </:col>
</:col> <:col :let={task} label="Started At">
<:col :let={task} label="Started At"> <%= format_datetime(task.job.attempted_at) %>
<%= format_datetime(task.job.attempted_at) %> </:col>
</:col> </.table>
</.table>
</div>
</div> </div>
""" """
end end
def mount(_params, _session, socket) do def mount(_params, _session, socket) do
PinchflatWeb.Endpoint.subscribe("tasks:job_table_live") PinchflatWeb.Endpoint.subscribe("job:state")
{:ok, assign(socket, tasks: get_tasks())} {:ok, assign(socket, tasks: get_tasks())}
end end
def handle_info(%{topic: "tasks:job_table_live", event: "reload"}, socket) do def handle_info(%{topic: "job:state", event: "change"}, socket) do
{:noreply, assign(socket, tasks: get_tasks())} {:noreply, assign(socket, tasks: get_tasks())}
end end
defp get_tasks do defp get_tasks do
TasksQuery.new() TasksQuery.new()
|> TasksQuery.join_job() |> TasksQuery.join_job()
|> where(^TasksQuery.in_state(["executing", "completed"])) |> where(^TasksQuery.in_state("executing"))
|> where(^TasksQuery.has_tag("show_in_dashboard")) |> where(^TasksQuery.has_tag("show_in_dashboard"))
|> order_by([t, j], desc: j.attempted_at)
|> Repo.all() |> Repo.all()
|> Repo.preload([:media_item, :source]) |> Repo.preload([:media_item, :source])
end end
@ -68,10 +67,10 @@ defmodule Pinchflat.Pages.JobTableLive do
end end
defp map_worker_to_task_name("FastIndexingWorker"), do: "Fast Indexing Source" defp map_worker_to_task_name("FastIndexingWorker"), do: "Fast Indexing Source"
defp map_worker_to_task_name("MediaDownloadWorker"), do: "Download Media" defp map_worker_to_task_name("MediaDownloadWorker"), do: "Downloading Media"
defp map_worker_to_task_name("MediaCollectionIndexingWorker"), do: "Indexing Source" defp map_worker_to_task_name("MediaCollectionIndexingWorker"), do: "Indexing Source"
defp map_worker_to_task_name("MediaQualityUpgradeWorker"), do: "Upgrade Media Quality" defp map_worker_to_task_name("MediaQualityUpgradeWorker"), do: "Upgrading Media Quality"
defp map_worker_to_task_name("SourceMetadataStorageWorker"), do: "Fetch Source Metadata" defp map_worker_to_task_name("SourceMetadataStorageWorker"), do: "Fetching Source Metadata"
defp map_worker_to_task_name(other), do: other <> " (Report to Devs)" defp map_worker_to_task_name(other), do: other <> " (Report to Devs)"
defp task_to_record_name(%Task{} = task) do defp task_to_record_name(%Task{} = task) do

View file

@ -19,6 +19,11 @@ defmodule PinchflatWeb.Telemetry do
Supervisor.init(children, strategy: :one_for_one) Supervisor.init(children, strategy: :one_for_one)
end end
@doc false
def job_state_change_broadcast(_event, _measure, _meta, _config) do
PinchflatWeb.Endpoint.broadcast("job:state", "change", nil)
end
def metrics do def metrics do
[ [
# Phoenix Metrics # Phoenix Metrics

View file

@ -0,0 +1,111 @@
defmodule PinchflatWeb.Pages.JobTableLiveTest do
use PinchflatWeb.ConnCase
import Ecto.Query, warn: false
import Phoenix.LiveViewTest
import Pinchflat.MediaFixtures
import Pinchflat.SourcesFixtures
alias Pinchflat.Utils.StringUtils
alias Pinchflat.Pages.JobTableLive
alias Pinchflat.Downloading.MediaDownloadWorker
alias Pinchflat.FastIndexing.FastIndexingWorker
describe "initial rendering" do
test "shows message when no records", %{conn: conn} do
{:ok, _view, html} = live_isolated(conn, JobTableLive, session: %{})
assert html =~ "Nothing Here!"
refute html =~ "Subject"
end
test "shows records when present", %{conn: conn} do
{_source, _media_item, _task, _job} = create_media_item_job()
{:ok, _view, html} = live_isolated(conn, JobTableLive, session: %{})
assert html =~ "Subject"
end
test "doesn't show records when not in executing state", %{conn: conn} do
{_source, _media_item, _task, _job} = create_media_item_job(:scheduled)
{_source, _media_item, _task, _job} = create_media_item_job(:completed)
{:ok, _view, html} = live_isolated(conn, JobTableLive, session: %{})
assert html =~ "Nothing Here!"
refute html =~ "Subject"
end
end
describe "job rendering" do
test "shows worker name", %{conn: conn} do
{_source, _media_item, _task, _job} = create_media_item_job()
{:ok, _view, html} = live_isolated(conn, JobTableLive, session: %{})
assert html =~ "Downloading Media"
end
test "shows the media item title", %{conn: conn} do
{_source, media_item, _task, _job} = create_media_item_job()
{:ok, _view, html} = live_isolated(conn, JobTableLive, session: %{})
assert html =~ StringUtils.truncate(media_item.title, 35)
end
test "shows a media item link", %{conn: conn} do
{_source, media_item, _task, _job} = create_media_item_job()
{:ok, _view, html} = live_isolated(conn, JobTableLive, session: %{})
assert html =~ ~p"/sources/#{media_item.source_id}/media/#{media_item}"
end
test "shows the source custom name", %{conn: conn} do
{source, _task, _job} = create_source_job()
{:ok, _view, html} = live_isolated(conn, JobTableLive, session: %{})
assert html =~ StringUtils.truncate(source.custom_name, 35)
end
test "shows a source link", %{conn: conn} do
{source, _task, _job} = create_source_job()
{:ok, _view, html} = live_isolated(conn, JobTableLive, session: %{})
assert html =~ ~p"/sources/#{source.id}"
end
test "listens for job:state change events", %{conn: conn} do
{_source, _media_item, _task, _job} = create_media_item_job()
{:ok, _view, _html} = live_isolated(conn, JobTableLive, session: %{})
PinchflatWeb.Endpoint.broadcast("job:state", "change", nil)
assert_receive %Phoenix.Socket.Broadcast{topic: "job:state", event: "change", payload: nil}
end
end
defp create_media_item_job(job_state \\ :executing) do
source = source_fixture()
media_item = media_item_fixture(source_id: source.id)
{:ok, task} = MediaDownloadWorker.kickoff_with_task(media_item)
Oban.Job
|> where([j], j.id == ^task.job_id)
|> Repo.update_all(set: [state: to_string(job_state)])
job = Repo.get!(Oban.Job, task.job_id)
{source, media_item, task, job}
end
defp create_source_job(job_state \\ :executing) do
source = source_fixture()
{:ok, task} = FastIndexingWorker.kickoff_with_task(source)
Oban.Job
|> where([j], j.id == ^task.job_id)
|> Repo.update_all(set: [state: to_string(job_state)])
job = Repo.get!(Oban.Job, task.job_id)
{source, task, job}
end
end