total rebase
[anni] / lib / pleroma / job_queue_monitor.ex
1 # Pleroma: A lightweight social networking server
2 # Copyright © 2017-2022 Pleroma Authors <https://pleroma.social/>
3 # SPDX-License-Identifier: AGPL-3.0-only
4
5 defmodule Pleroma.JobQueueMonitor do
6   use GenServer
7
8   @initial_state %{workers: %{}, queues: %{}, processed_jobs: 0}
9   @queue %{processed_jobs: 0, success: 0, failure: 0}
10   @operation %{processed_jobs: 0, success: 0, failure: 0}
11
12   def start_link(_) do
13     GenServer.start_link(__MODULE__, @initial_state, name: __MODULE__)
14   end
15
16   @impl true
17   def init(state) do
18     :telemetry.attach("oban-monitor-failure", [:oban, :job, :exception], &handle_event/4, nil)
19     :telemetry.attach("oban-monitor-success", [:oban, :job, :stop], &handle_event/4, nil)
20
21     {:ok, state}
22   end
23
24   def stats do
25     GenServer.call(__MODULE__, :stats)
26   end
27
28   def handle_event([:oban, :job, event], %{duration: duration}, meta, _) do
29     GenServer.cast(
30       __MODULE__,
31       {:process_event, mapping_status(event), duration, meta}
32     )
33   end
34
35   @impl true
36   def handle_call(:stats, _from, state) do
37     {:reply, state, state}
38   end
39
40   @impl true
41   def handle_cast({:process_event, status, duration, meta}, state) do
42     state =
43       state
44       |> Map.update!(:workers, fn workers ->
45         workers
46         |> Map.put_new(meta.worker, %{})
47         |> Map.update!(meta.worker, &update_worker(&1, status, meta, duration))
48       end)
49       |> Map.update!(:queues, fn workers ->
50         workers
51         |> Map.put_new(meta.queue, @queue)
52         |> Map.update!(meta.queue, &update_queue(&1, status, meta, duration))
53       end)
54       |> Map.update!(:processed_jobs, &(&1 + 1))
55
56     {:noreply, state}
57   end
58
59   defp update_worker(worker, status, meta, duration) do
60     worker
61     |> Map.put_new(meta.args["op"], @operation)
62     |> Map.update!(meta.args["op"], &update_op(&1, status, meta, duration))
63   end
64
65   defp update_op(op, :enqueue, _meta, _duration) do
66     op
67     |> Map.update!(:enqueued, &(&1 + 1))
68   end
69
70   defp update_op(op, status, _meta, _duration) do
71     op
72     |> Map.update!(:processed_jobs, &(&1 + 1))
73     |> Map.update!(status, &(&1 + 1))
74   end
75
76   defp update_queue(queue, status, _meta, _duration) do
77     queue
78     |> Map.update!(:processed_jobs, &(&1 + 1))
79     |> Map.update!(status, &(&1 + 1))
80   end
81
82   defp mapping_status(:stop), do: :success
83   defp mapping_status(:exception), do: :failure
84 end