Packages
oban
2.18.2
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/postgres/v01.ex
defmodule Oban.Migrations.Postgres.V01 do
@moduledoc false
use Ecto.Migration
def up(%{create_schema: create?, prefix: prefix} = opts) do
%{escaped_prefix: escaped, quoted_prefix: quoted} = opts
if create?, do: execute("CREATE SCHEMA IF NOT EXISTS #{quoted}")
execute """
DO $$
BEGIN
IF NOT EXISTS (SELECT 1 FROM pg_type
WHERE typname = 'oban_job_state'
AND typnamespace = '#{escaped}'::regnamespace::oid) THEN
CREATE TYPE #{quoted}.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, :"#{quoted}.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: fragment("timezone('UTC', now())")
add :scheduled_at, :utc_datetime_usec,
null: false,
default: fragment("timezone('UTC', 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)
execute """
CREATE OR REPLACE FUNCTION #{quoted}.oban_jobs_notify() RETURNS trigger AS $$
DECLARE
channel text;
notice json;
BEGIN
IF (TG_OP = 'INSERT') THEN
channel = '#{escaped}.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 = '#{escaped}.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 #{quoted}.oban_jobs"
execute """
CREATE TRIGGER oban_notify
AFTER INSERT OR UPDATE OF state ON #{quoted}.oban_jobs
FOR EACH ROW EXECUTE PROCEDURE #{quoted}.oban_jobs_notify();
"""
end
def down(%{prefix: prefix, quoted_prefix: quoted}) do
execute "DROP TRIGGER IF EXISTS oban_notify ON #{quoted}.oban_jobs"
execute "DROP FUNCTION IF EXISTS #{quoted}.oban_jobs_notify()"
drop_if_exists table(:oban_jobs, prefix: prefix)
execute "DROP TYPE IF EXISTS #{quoted}.oban_job_state"
end
end