Packages
oban
2.16.1
2.23.0
2.22.1
2.22.0
2.21.1
2.21.0
2.20.3
2.20.2
2.20.1
2.20.0
2.19.4
2.19.3
2.19.2
2.19.1
2.19.0
2.18.3
2.18.2
2.18.1
2.18.0
2.17.12
2.17.11
2.17.10
2.17.9
2.17.8
2.17.7
2.17.6
2.17.5
2.17.4
2.17.3
2.17.2
2.17.1
2.17.0
2.16.3
2.16.2
2.16.1
2.16.0
2.15.4
2.15.3
2.15.2
2.15.1
2.15.0
2.14.2
2.14.1
2.14.0
2.13.6
2.13.5
2.13.4
2.13.3
2.13.2
2.13.1
2.13.0
2.12.1
2.12.0
2.11.3
2.11.2
2.11.1
2.11.0
2.10.1
2.10.0
retired
2.9.2
2.9.1
2.9.0
2.8.0
2.7.2
2.7.1
2.7.0
2.6.1
2.6.0
2.5.0
2.4.3
2.4.2
2.4.1
2.4.0
2.3.4
2.3.3
2.3.2
2.3.1
2.3.0
2.2.0
2.1.0
2.0.0
2.0.0-rc.3
2.0.0-rc.2
2.0.0-rc.1
2.0.0-rc.0
1.2.0
1.1.0
1.0.0
1.0.0-rc.2
1.0.0-rc.1
0.12.1
0.12.0
0.11.1
0.11.0
0.10.1
0.10.0
0.9.0
0.8.1
0.8.0
0.7.1
0.7.0
0.6.0
0.5.0
0.4.0
0.3.0
0.2.0
0.1.0
Robust job processing, backed by modern PostgreSQL, SQLite3, and MySQL.
Current section
Files
Jump to
Current section
Files
lib/oban/queue/drainer.ex
defmodule Oban.Queue.Drainer do
@moduledoc false
import Ecto.Query, only: [where: 3]
alias Oban.{Config, Job, Repo}
alias Oban.Queue.Executor
@infinite 100_000_000
def drain(%Config{} = conf, [_ | _] = opts) do
args =
opts
|> Map.new()
|> Map.put_new(:with_limit, @infinite)
|> Map.put_new(:with_recursion, false)
|> Map.put_new(:with_safety, true)
|> Map.put_new(:with_scheduled, false)
|> Map.update!(:queue, &to_string/1)
drain(conf, %{cancelled: 0, discard: 0, failure: 0, snoozed: 0, success: 0}, args)
end
defp stage_scheduled(conf, queue, with_scheduled) do
query =
Job
|> where([j], j.state in ["scheduled", "retryable"])
|> where([j], j.queue == ^queue)
query =
case with_scheduled do
true -> query
%DateTime{} -> where(query, [j], j.scheduled_at <= ^with_scheduled)
end
Repo.update_all(conf, query, set: [state: "available"])
end
defp drain(conf, old_acc, %{queue: queue} = args) do
if args.with_scheduled, do: stage_scheduled(conf, queue, args.with_scheduled)
new_acc =
conf
|> fetch_available(args)
|> Enum.reduce(old_acc, fn job, acc ->
result =
conf
|> Executor.new(job, safe: args.with_safety)
|> Executor.call()
|> case do
%{state: :exhausted} -> :discard
%{state: state} -> state
end
Map.update(acc, result, 1, &(&1 + 1))
end)
if args.with_recursion and old_acc != new_acc do
drain(conf, new_acc, args)
else
new_acc
end
end
defp fetch_available(conf, %{queue: queue, with_limit: limit}) do
{:ok, meta} = conf.engine.init(conf, queue: queue, limit: limit)
{:ok, {_meta, jobs}} = conf.engine.fetch_jobs(conf, meta, %{})
jobs
end
end