Packages
oban
2.16.3
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/engines/inline.ex
defmodule Oban.Engines.Inline do
@moduledoc """
A testing-specific engine that's used when Oban's started with `testing: :inline`.
## Usage
This is meant for testing and shouldn't be configured directly:
Oban.start_link(repo: MyApp.Repo, testing: :inline)
"""
@behaviour Oban.Engine
import DateTime, only: [utc_now: 0]
alias Ecto.Changeset
alias Oban.{Config, Engine, Job}
alias Oban.Queue.Executor
@impl Engine
def init(_conf, opts), do: {:ok, Map.new(opts)}
@impl Engine
def put_meta(_conf, meta, key, value), do: Map.put(meta, key, value)
@impl Engine
def check_meta(_conf, meta, _running), do: meta
@impl Engine
def refresh(_conf, meta), do: meta
@impl Engine
def shutdown(_conf, meta), do: meta
@impl Engine
def insert_job(%Config{} = conf, %Changeset{} = changeset, _opts) do
{:ok, execute_job(conf, changeset)}
end
@impl Engine
def insert_all_jobs(%Config{} = conf, changesets, _opts) do
changesets
|> expand()
|> Enum.map(&execute_job(conf, &1))
end
@impl Engine
def fetch_jobs(_conf, meta, _running), do: {:ok, {meta, []}}
@impl Engine
def complete_job(_conf, _job), do: :ok
@impl Engine
def discard_job(_conf, _job), do: :ok
@impl Engine
def error_job(_conf, _job, seconds) when is_integer(seconds), do: :ok
@impl Engine
def snooze_job(_conf, _job, seconds) when is_integer(seconds), do: :ok
@impl Engine
def cancel_job(_conf, _job), do: :ok
@impl Engine
def cancel_all_jobs(_conf, _queryable), do: {:ok, []}
@impl Engine
def retry_job(_conf, _job), do: :ok
@impl Engine
def retry_all_jobs(_conf, _queryable), do: {:ok, []}
# Changeset Helpers
defp expand(value), do: expand(value, %{})
defp expand(fun, changes) when is_function(fun, 1), do: expand(fun.(changes), changes)
defp expand(%{changesets: changesets}, _), do: expand(changesets, %{})
defp expand(changesets, _) when is_list(changesets), do: changesets
# Execution Helpers
defp execute_job(conf, changeset) do
changeset =
changeset
|> Changeset.put_change(:attempt, 1)
|> Changeset.put_change(:attempted_by, [conf.node])
|> Changeset.put_change(:attempted_at, utc_now())
|> Changeset.put_change(:scheduled_at, utc_now())
|> Changeset.update_change(:args, &json_encode_decode/1)
|> Changeset.update_change(:meta, &json_encode_decode/1)
case Changeset.apply_action(changeset, :insert) do
{:ok, job} ->
conf
|> Executor.new(job, safe: false)
|> Executor.call()
|> complete_job()
{:error, changeset} ->
raise Ecto.InvalidChangesetError, action: :insert, changeset: changeset
end
end
defp json_encode_decode(map) do
map
|> Jason.encode!()
|> Jason.decode!()
end
defp complete_job(%{job: job, state: :failure}) do
%Job{job | errors: [Job.format_attempt(job)], state: "retryable", scheduled_at: utc_now()}
end
defp complete_job(%{job: job, state: :cancelled}) do
%Job{job | errors: [Job.format_attempt(job)], state: "cancelled", cancelled_at: utc_now()}
end
defp complete_job(%{job: job, state: state}) when state in [:discard, :exhausted] do
%Job{job | errors: [Job.format_attempt(job)], state: "discarded", discarded_at: utc_now()}
end
defp complete_job(%{job: job, state: :success}) do
%Job{job | state: "completed", completed_at: utc_now()}
end
defp complete_job(%{job: job, result: {:snooze, snooze}, state: :snoozed}) do
%Job{job | state: "scheduled", scheduled_at: DateTime.add(utc_now(), snooze, :second)}
end
end