cover/Elixir.WalletSettlement.JobQueue.html

1 defmodule WalletSettlement.JobQueue do
2 @moduledoc """
3 Lightweight in-process job queue for settlement and reconciliation workers.
4
5 Simulates Oban queue semantics (enqueue, process, retry/backoff, dead-letter)
6 using ETS + GenServer for CI operation where no database is available.
7
8 Queue semantics:
9 - Jobs have a queue name, attempt count, max attempts, and status.
10 - Statuses: :available -> :executing -> :completed | :retryable | :discarded
11 - Workers implement `perform/1` callback that returns `:ok` or `{:error, reason}`.
12 - On failure before max_attempts: job is re-queued as :retryable.
13 - On max_attempts exceeded: job is moved to :discarded (dead-letter).
14
15 For production Oban integration, replace this module with Oban configuration
16 as documented in `WalletSettlement.QueueConfig`.
17 """
18 use GenServer
19
20 alias WalletSettlement.QueueConfig
21
22 @table :wallet_settlement_jobs
23
24
:-(
defstruct [:job_id, :queue, :worker, :args, :status, :attempts, :max_attempts,
25 :error, :enqueued_at, :scheduled_at, :completed_at]
26
27
:-(
def start_link(opts), do: GenServer.start_link(__MODULE__, opts, name: __MODULE__)
28
29 @doc """
30 Enqueues a job for asynchronous execution.
31
32 `worker` must be a module implementing `perform/1`.
33 Returns `{:ok, job_id}`.
34 """
35 @spec enqueue(queue :: atom(), worker :: module(), args :: map()) :: {:ok, String.t()}
36
:-(
def enqueue(queue, worker, args \\ %{}) do
37 13 GenServer.call(__MODULE__, {:enqueue, queue, worker, args})
38 end
39
40 @doc "Returns all jobs for a given queue and status."
41 @spec list_jobs(queue :: atom(), status :: atom()) :: [map()]
42
:-(
def list_jobs(queue, status \\ nil) do
43 2 GenServer.call(__MODULE__, {:list_jobs, queue, status})
44 end
45
46 @doc "Returns a job by ID."
47 @spec get_job(String.t()) :: {:ok, map()} | {:error, :not_found}
48 1 def get_job(job_id), do: GenServer.call(__MODULE__, {:get_job, job_id})
49
50 @doc "Runs all available jobs in the given queue synchronously. For test use."
51 @spec drain_queue(queue :: atom()) :: %{completed: integer(), failed: integer(), discarded: integer()}
52 def drain_queue(queue) do
53 6 GenServer.call(__MODULE__, {:drain_queue, queue}, 30_000)
54 end
55
56 22 def reset, do: GenServer.call(__MODULE__, :reset)
57
58 @impl true
59 def init(_opts) do
60
:-(
:ets.new(@table, [:set, :protected, :named_table])
61 {:ok, %{}}
62 end
63
64 @impl true
65 def handle_call({:enqueue, queue, worker, args}, _from, state) do
66 13 {:ok, cfg} = QueueConfig.queue(queue)
67 13 job_id = WalletSharedKernel.TypedId.generate("job")
68 13 job = %__MODULE__{
69 job_id: job_id,
70 queue: queue,
71 worker: worker,
72 args: args,
73 status: :available,
74 attempts: 0,
75 13 max_attempts: cfg.max_attempts,
76 enqueued_at: DateTime.utc_now(),
77 scheduled_at: DateTime.utc_now()
78 }
79 13 :ets.insert(@table, {job_id, job})
80 13 {:reply, {:ok, job_id}, state}
81 end
82
83 @impl true
84 def handle_call({:list_jobs, queue, nil}, _from, state) do
85
:-(
jobs = :ets.tab2list(@table)
86
:-(
|> Enum.map(fn {_, j} -> j end)
87
:-(
|> Enum.filter(&(&1.queue == queue))
88
:-(
{:reply, jobs, state}
89 end
90
91 @impl true
92 def handle_call({:list_jobs, queue, status}, _from, state) do
93 2 jobs = :ets.tab2list(@table)
94 2 |> Enum.map(fn {_, j} -> j end)
95 2 |> Enum.filter(&(&1.queue == queue and &1.status == status))
96 2 {:reply, jobs, state}
97 end
98
99 @impl true
100 def handle_call({:get_job, job_id}, _from, state) do
101 1 result = case :ets.lookup(@table, job_id) do
102 1 [{_, j}] -> {:ok, j}
103
:-(
[] -> {:error, :not_found}
104 end
105 1 {:reply, result, state}
106 end
107
108 @impl true
109 def handle_call({:drain_queue, queue}, _from, state) do
110 6 available = :ets.tab2list(@table)
111 8 |> Enum.map(fn {_, j} -> j end)
112 8 |> Enum.filter(&(&1.queue == queue and &1.status == :available))
113
114 6 results = Enum.reduce(available, %{completed: 0, failed: 0, discarded: 0}, fn job, acc ->
115 7 job = %{job | status: :executing, attempts: job.attempts + 1}
116 7 :ets.insert(@table, {job.job_id, job})
117
118 7 case apply(job.worker, :perform, [job.args]) do
119 :ok ->
120 4 done = %{job | status: :completed, completed_at: DateTime.utc_now()}
121 4 :ets.insert(@table, {done.job_id, done})
122 4 %{acc | completed: acc.completed + 1}
123
124 {:error, reason} ->
125 3 if job.attempts >= job.max_attempts do
126 1 discarded = %{job | status: :discarded, error: inspect(reason)}
127 1 :ets.insert(@table, {discarded.job_id, discarded})
128 1 %{acc | discarded: acc.discarded + 1}
129 else
130 2 retry = %{job | status: :available, error: inspect(reason)}
131 2 :ets.insert(@table, {retry.job_id, retry})
132 2 %{acc | failed: acc.failed + 1}
133 end
134 end
135 end)
136
137 6 {:reply, results, state}
138 end
139
140 @impl true
141 def handle_call(:reset, _from, state) do
142 22 :ets.delete_all_objects(@table)
143 22 {:reply, :ok, state}
144 end
145 end
Line Hits Source