685ba2eada6454b01b7e032c52c0f211cc9744ba
[akkoma] / lib / pleroma / job_queue_monitor.ex
1 # Pleroma: A lightweight social networking server
2 # Copyright © 2017-2019 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, enqueued: 0}
9 @queue %{processed_jobs: 0, success: 0, failure: 0, enqueued: 0}
10 @operation %{processed_jobs: 0, success: 0, failure: 0, enqueued: 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 enqueue({:ok, job}) do
29 meta = Map.take(job, [:args, :queue, :worker])
30 GenServer.cast(__MODULE__, {:process_enqueue, meta})
31
32 {:ok, job}
33 end
34
35 def enqueue(result), do: result
36
37 def handle_event([:oban, status], %{duration: duration}, meta, _) do
38 GenServer.cast(__MODULE__, {:process_event, status, duration, meta})
39 end
40
41 @impl true
42 def handle_call(:stats, _from, state) do
43 {:reply, state, state}
44 end
45
46 def handle_cast({:process_enqueue, meta}, state) do
47 state =
48 state
49 |> Map.update!(:workers, fn workers ->
50 workers
51 |> Map.put_new(meta.worker, %{})
52 |> Map.update!(meta.worker, &update_worker(&1, :enqueue, meta))
53 end)
54 |> Map.update!(:queues, fn workers ->
55 workers
56 |> Map.put_new(meta.queue, @queue)
57 |> Map.update!(meta.queue, fn queue -> Map.update!(queue, :enqueued, &(&1 + 1)) end)
58 end)
59 |> Map.update!(:enqueued, &(&1 + 1))
60
61 {:noreply, state}
62 end
63
64 @impl true
65 def handle_cast({:process_event, status, duration, meta}, state) do
66 state =
67 state
68 |> Map.update!(:workers, fn workers ->
69 workers
70 |> Map.put_new(meta.worker, %{})
71 |> Map.update!(meta.worker, &update_worker(&1, status, meta, duration))
72 end)
73 |> Map.update!(:queues, fn workers ->
74 workers
75 |> Map.put_new(meta.queue, @queue)
76 |> Map.update!(meta.queue, &update_queue(&1, status, meta, duration))
77 end)
78 |> Map.update!(:processed_jobs, &(&1 + 1))
79 |> decr_enqueued()
80
81 {:noreply, state}
82 end
83
84 defp update_worker(worker, status, meta, duration \\ 0) do
85 worker
86 |> Map.put_new(meta.args["op"], @operation)
87 |> Map.update!(meta.args["op"], &update_op(&1, status, meta, duration))
88 end
89
90 defp update_op(op, :enqueue, _meta, _duration) do
91 op
92 |> Map.update!(:enqueued, &(&1 + 1))
93 end
94
95 defp update_op(op, status, _meta, _duration) do
96 op
97 |> Map.update!(:processed_jobs, &(&1 + 1))
98 |> Map.update!(status, &(&1 + 1))
99 |> decr_enqueued()
100 end
101
102 defp update_queue(queue, status, _meta, _duration) do
103 queue
104 |> Map.update!(:processed_jobs, &(&1 + 1))
105 |> Map.update!(status, &(&1 + 1))
106 |> decr_enqueued()
107 end
108
109 defp decr_enqueued(map) do
110 Map.update!(map, :enqueued, fn
111 0 -> 0
112 enqueued -> enqueued - 1
113 end)
114 end
115 end