Packages
oban
2.17.12
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/v08.ex
defmodule Oban.Migrations.Postgres.V08 do
@moduledoc false
use Ecto.Migration
def up(%{escaped_prefix: escaped, prefix: prefix, quoted_prefix: quoted}) 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, [:state, :queue, :priority, :scheduled_at, :id],
prefix: prefix
)
execute """
CREATE OR REPLACE FUNCTION #{quoted}.oban_jobs_notify() RETURNS trigger AS $$
DECLARE
channel text;
notice json;
BEGIN
IF NEW.state = 'available' THEN
channel = '#{escaped}.oban_insert';
notice = json_build_object('queue', NEW.queue);
PERFORM pg_notify(channel, notice::text);
END IF;
RETURN NULL;
END;
$$ LANGUAGE plpgsql;
"""
execute "DROP TRIGGER IF EXISTS oban_notify ON #{quoted}.oban_jobs"
execute """
CREATE TRIGGER oban_notify
AFTER INSERT ON #{quoted}.oban_jobs
FOR EACH ROW EXECUTE PROCEDURE #{quoted}.oban_jobs_notify();
"""
end
def down(%{escaped_prefix: escaped, prefix: prefix, quoted_prefix: quoted}) 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
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
end