Packages
oban
1.2.0
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/migrations.ex
defmodule Oban.Migrations do
@moduledoc false
use Ecto.Migration
@initial_version 1
@current_version 8
@default_prefix "public"
def up(opts \\ []) when is_list(opts) do
prefix = Keyword.get(opts, :prefix, @default_prefix)
version = Keyword.get(opts, :version, @current_version)
initial = min(migrated_version(repo(), prefix) + 1, @current_version)
if initial <= version, do: change(prefix, initial..version, :up)
end
def down(opts \\ []) when is_list(opts) do
prefix = Keyword.get(opts, :prefix, @default_prefix)
version = Keyword.get(opts, :version, @initial_version)
initial = max(migrated_version(repo(), prefix), @initial_version)
if initial >= version, do: change(prefix, initial..version, :down)
end
def initial_version, do: @initial_version
def current_version, do: @current_version
def migrated_version(repo, prefix) do
query = """
SELECT description
FROM pg_class
LEFT JOIN pg_description ON pg_description.objoid = pg_class.oid
LEFT JOIN pg_namespace ON pg_namespace.oid = pg_class.relnamespace
WHERE pg_class.relname = 'oban_jobs'
AND pg_namespace.nspname = '#{prefix}'
"""
case repo.query(query) do
{:ok, %{rows: [[version]]}} when is_binary(version) -> String.to_integer(version)
_ -> 0
end
end
defp change(prefix, range, direction) do
for index <- range do
[__MODULE__, "V#{index}"]
|> Module.concat()
|> apply(direction, [prefix])
end
end
defmodule Helper do
@moduledoc false
defmacro now do
quote do
fragment("timezone('UTC', now())")
end
end
def record_version(prefix, version) do
execute "COMMENT ON TABLE #{prefix}.oban_jobs IS '#{version}'"
end
# Extracted into a helper function to facilitate sharing.
def v1_oban_notify(prefix) do
execute """
CREATE OR REPLACE FUNCTION #{prefix}.oban_jobs_notify() RETURNS trigger AS $$
DECLARE
channel text;
notice json;
BEGIN
IF (TG_OP = 'INSERT') THEN
channel = '#{prefix}.oban_insert';
notice = json_build_object('queue', NEW.queue, 'state', NEW.state);
-- No point triggering for a job that isn't scheduled to run now
IF NEW.scheduled_at IS NOT NULL AND NEW.scheduled_at > now() AT TIME ZONE 'utc' THEN
RETURN null;
END IF;
ELSE
channel = '#{prefix}.oban_update';
notice = json_build_object('queue', NEW.queue, 'new_state', NEW.state, 'old_state', OLD.state);
END IF;
PERFORM pg_notify(channel, notice::text);
RETURN NULL;
END;
$$ LANGUAGE plpgsql;
"""
execute "DROP TRIGGER IF EXISTS oban_notify ON #{prefix}.oban_jobs"
execute """
CREATE TRIGGER oban_notify
AFTER INSERT OR UPDATE OF state ON #{prefix}.oban_jobs
FOR EACH ROW EXECUTE PROCEDURE #{prefix}.oban_jobs_notify();
"""
end
def v2_oban_notify(prefix) do
execute """
CREATE OR REPLACE FUNCTION #{prefix}.oban_jobs_notify() RETURNS trigger AS $$
DECLARE
channel text;
notice json;
BEGIN
IF NEW.state = 'available' THEN
channel = '#{prefix}.oban_insert';
notice = json_build_object('queue', NEW.queue, 'state', NEW.state);
PERFORM pg_notify(channel, notice::text);
END IF;
RETURN NULL;
END;
$$ LANGUAGE plpgsql;
"""
execute "DROP TRIGGER oban_notify ON #{prefix}.oban_jobs"
execute """
CREATE TRIGGER oban_notify
AFTER INSERT ON #{prefix}.oban_jobs
FOR EACH ROW EXECUTE PROCEDURE #{prefix}.oban_jobs_notify();
"""
end
end
defmodule V1 do
@moduledoc false
use Ecto.Migration
import Oban.Migrations.Helper
def up(prefix) do
execute "CREATE SCHEMA IF NOT EXISTS #{prefix}"
execute """
DO $$
BEGIN
IF NOT EXISTS (SELECT 1 FROM pg_type
WHERE typname = 'oban_job_state'
AND typnamespace = '#{prefix}'::regnamespace::oid) THEN
CREATE TYPE #{prefix}.oban_job_state AS ENUM (
'available',
'scheduled',
'executing',
'retryable',
'completed',
'discarded'
);
END IF;
END$$;
"""
create_if_not_exists table(:oban_jobs, primary_key: false, prefix: prefix) do
add :id, :bigserial, primary_key: true
add :state, :"#{prefix}.oban_job_state", null: false, default: "available"
add :queue, :text, null: false, default: "default"
add :worker, :text, null: false
add :args, :map, null: false
add :errors, {:array, :map}, null: false, default: []
add :attempt, :integer, null: false, default: 0
add :max_attempts, :integer, null: false, default: 20
add :inserted_at, :utc_datetime_usec, null: false, default: now()
add :scheduled_at, :utc_datetime_usec, null: false, default: now()
add :attempted_at, :utc_datetime_usec
add :completed_at, :utc_datetime_usec
end
create_if_not_exists index(:oban_jobs, [:queue], prefix: prefix)
create_if_not_exists index(:oban_jobs, [:state], prefix: prefix)
create_if_not_exists index(:oban_jobs, [:scheduled_at], prefix: prefix)
v1_oban_notify(prefix)
record_version(prefix, 1)
end
def down(prefix) do
execute "DROP TRIGGER IF EXISTS oban_notify ON #{prefix}.oban_jobs"
execute "DROP FUNCTION IF EXISTS #{prefix}.oban_jobs_notify()"
drop_if_exists table(:oban_jobs, prefix: prefix)
execute "DROP TYPE IF EXISTS #{prefix}.oban_job_state"
end
end
defmodule V2 do
@moduledoc false
use Ecto.Migration
import Oban.Migrations.Helper
def up(prefix) do
# We only need the scheduled_at index for scheduled and available jobs
drop_if_exists index(:oban_jobs, [:scheduled_at], prefix: prefix)
state = "#{prefix}.oban_job_state"
create index(:oban_jobs, [:scheduled_at],
where: "state in ('available'::#{state}, 'scheduled'::#{state})",
prefix: prefix
)
create constraint(:oban_jobs, :worker_length,
check: "char_length(worker) > 0 AND char_length(worker) < 128",
prefix: prefix
)
create constraint(:oban_jobs, :queue_length,
check: "char_length(queue) > 0 AND char_length(queue) < 128",
prefix: prefix
)
execute """
CREATE OR REPLACE FUNCTION #{prefix}.oban_wrap_id(value bigint) RETURNS int AS $$
BEGIN
RETURN (CASE WHEN value > 2147483647 THEN mod(value, 2147483647) ELSE value END)::int;
END;
$$ LANGUAGE plpgsql IMMUTABLE;
"""
record_version(prefix, 2)
end
def down(prefix) do
drop_if_exists constraint(:oban_jobs, :queue_length, prefix: prefix)
drop_if_exists constraint(:oban_jobs, :worker_length, prefix: prefix)
drop_if_exists index(:oban_jobs, [:scheduled_at], prefix: prefix)
create index(:oban_jobs, [:scheduled_at], prefix: prefix)
execute("DROP FUNCTION IF EXISTS #{prefix}.oban_wrap_id(value bigint)")
record_version(prefix, 1)
end
end
defmodule V3 do
@moduledoc false
use Ecto.Migration
import Oban.Migrations.Helper
def up(prefix) do
alter table(:oban_jobs, prefix: prefix) do
add_if_not_exists(:attempted_by, {:array, :text})
end
create_if_not_exists table(:oban_beats, primary_key: false, prefix: prefix) do
add :node, :text, null: false
add :queue, :text, null: false
add :nonce, :text, null: false
add :limit, :integer, null: false
add :paused, :boolean, null: false, default: false
add :running, {:array, :integer}, null: false, default: []
add :inserted_at, :utc_datetime_usec, null: false, default: now()
add :started_at, :utc_datetime_usec, null: false
end
create_if_not_exists index(:oban_beats, [:inserted_at], prefix: prefix)
record_version(prefix, 3)
end
def down(prefix) do
alter table(:oban_jobs, prefix: prefix) do
remove_if_exists(:attempted_by, {:array, :text})
end
drop_if_exists table(:oban_beats, prefix: prefix)
record_version(prefix, 2)
end
end
defmodule V4 do
@moduledoc false
# Dropping the `oban_wrap_id` function is isolated to allow progressive rollout.
use Ecto.Migration
import Oban.Migrations.Helper
def up(prefix) do
execute("DROP FUNCTION IF EXISTS #{prefix}.oban_wrap_id(value bigint)")
record_version(prefix, 4)
end
def down(prefix) do
execute """
CREATE OR REPLACE FUNCTION #{prefix}.oban_wrap_id(value bigint) RETURNS int AS $$
BEGIN
RETURN (CASE WHEN value > 2147483647 THEN mod(value, 2147483647) ELSE value END)::int;
END;
$$ LANGUAGE plpgsql IMMUTABLE;
"""
record_version(prefix, 3)
end
end
defmodule V5 do
@moduledoc false
use Ecto.Migration
import Oban.Migrations.Helper
def up(prefix) do
drop_if_exists index(:oban_jobs, [:scheduled_at], prefix: prefix)
drop_if_exists index(:oban_jobs, [:queue], prefix: prefix)
drop_if_exists index(:oban_jobs, [:state], prefix: prefix)
create_if_not_exists index(:oban_jobs, [:queue, :state, :scheduled_at, :id], prefix: prefix)
record_version(prefix, 5)
end
def down(prefix) do
drop_if_exists index(:oban_jobs, [:queue, :state, :scheduled_at, :id], prefix: prefix)
state = "#{prefix}.oban_job_state"
create_if_not_exists index(:oban_jobs, [:queue], prefix: prefix)
create_if_not_exists index(:oban_jobs, [:state], prefix: prefix)
create index(:oban_jobs, [:scheduled_at],
where: "state in ('available'::#{state}, 'scheduled'::#{state})",
prefix: prefix
)
record_version(prefix, 4)
end
end
defmodule V6 do
@moduledoc false
use Ecto.Migration
import Oban.Migrations.Helper
def up(prefix) do
execute "ALTER TABLE #{prefix}.oban_beats ALTER COLUMN running TYPE bigint[]"
record_version(prefix, 6)
end
def down(prefix) do
execute "ALTER TABLE #{prefix}.oban_beats ALTER COLUMN running TYPE integer[]"
record_version(prefix, 5)
end
end
defmodule V7 do
@moduledoc false
use Ecto.Migration
import Oban.Migrations.Helper
def up(prefix) do
create_if_not_exists index(
:oban_jobs,
["attempted_at desc", :id],
where: "state in ('completed', 'discarded')",
prefix: prefix,
name: :oban_jobs_attempted_at_id_index
)
record_version(prefix, 7)
end
def down(prefix) do
drop_if_exists index(:oban_jobs, [:attempted_at, :id], prefix: prefix)
record_version(prefix, 6)
end
end
defmodule V8 do
@moduledoc false
use Ecto.Migration
import Oban.Migrations.Helper
def up(prefix) do
alter table(:oban_jobs, prefix: prefix) do
add_if_not_exists(:discarded_at, :utc_datetime_usec)
add_if_not_exists(:priority, :integer)
add_if_not_exists(:tags, {:array, :string})
end
alter table(:oban_jobs, prefix: prefix) do
modify :priority, :integer, default: 0
modify :tags, {:array, :string}, default: []
end
drop_if_exists index(:oban_jobs, [:queue, :state, :scheduled_at, :id], prefix: prefix)
create_if_not_exists index(:oban_jobs, [:queue, :state, :priority, :scheduled_at, :id],
prefix: prefix
)
v2_oban_notify(prefix)
record_version(prefix, 8)
end
def down(prefix) do
drop_if_exists index(:oban_jobs, [:queue, :state, :priority, :scheduled_at, :id],
prefix: prefix
)
create_if_not_exists index(:oban_jobs, [:queue, :state, :scheduled_at, :id], prefix: prefix)
alter table(:oban_jobs, prefix: prefix) do
remove_if_exists(:discarded_at, :utc_datetime_usec)
remove_if_exists(:priority, :integer)
remove_if_exists(:tags, {:array, :string})
end
v1_oban_notify(prefix)
record_version(prefix, 7)
end
end
end