Merge branch 'worker-messages' into 'develop'
[akkoma] / lib / pleroma / job_queue_monitor.ex
1 # Pleroma: A lightweight social networking server
2 # Copyright © 2017-2020 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, :failure], &handle_event/4, nil)
19 :telemetry.attach("oban-monitor-success", [:oban, :success], &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, status], %{duration: duration}, meta, _) do
29 GenServer.cast(__MODULE__, {:process_event, status, duration, meta})
30 end
31
32 @impl true
33 def handle_call(:stats, _from, state) do
34 {:reply, state, state}
35 end
36
37 @impl true
38 def handle_cast({:process_event, status, duration, meta}, state) do
39 state =
40 state
41 |> Map.update!(:workers, fn workers ->
42 workers
43 |> Map.put_new(meta.worker, %{})
44 |> Map.update!(meta.worker, &update_worker(&1, status, meta, duration))
45 end)
46 |> Map.update!(:queues, fn workers ->
47 workers
48 |> Map.put_new(meta.queue, @queue)
49 |> Map.update!(meta.queue, &update_queue(&1, status, meta, duration))
50 end)
51 |> Map.update!(:processed_jobs, &(&1 + 1))
52
53 {:noreply, state}
54 end
55
56 defp update_worker(worker, status, meta, duration) do
57 worker
58 |> Map.put_new(meta.args["op"], @operation)
59 |> Map.update!(meta.args["op"], &update_op(&1, status, meta, duration))
60 end
61
62 defp update_op(op, :enqueue, _meta, _duration) do
63 op
64 |> Map.update!(:enqueued, &(&1 + 1))
65 end
66
67 defp update_op(op, status, _meta, _duration) do
68 op
69 |> Map.update!(:processed_jobs, &(&1 + 1))
70 |> Map.update!(status, &(&1 + 1))
71 end
72
73 defp update_queue(queue, status, _meta, _duration) do
74 queue
75 |> Map.update!(:processed_jobs, &(&1 + 1))
76 |> Map.update!(status, &(&1 + 1))
77 end
78 end