-module(gaffer_postgres). -moduledoc "Pure SQL query builder and serializer for Postgres drivers.". % API % Migrations -export([migrations/1]). -export([migrate_up/1]). -export([migrate_down/1]). -export([ensure_migrations_table/0]). -export([applied_versions/0]). -export([history/0]). % Queues -export([queue_insert/1]). -export([queue_exists/1]). -export([queue_list/0]). -export([queue_delete/1]). % Introspection -export([info/1]). % Jobs -export([job_write/1]). -export([job_get/1]). -export([job_list/1]). -export([job_delete/1]). -export([job_claim/2]). -export([job_prune/2]). -doc """ A parameterized SQL query. The `QueryString` will contain `$1` etc. which each will correspond to the nth value in `Values`. """. -type query() :: {QueryString :: iodata(), Values :: list()}. -doc "A list of queries to be run in one transaction.". -type queries() :: [query()]. -export_type([query/0]). -export_type([queries/0]). %--- API ----------------------------------------------------------------------- % Migrations -doc "Sorted list of all schema migrations as `{Version, Up, Down}` tuples.". -spec migrations(Opts :: map()) -> [{Version :: pos_integer(), Up :: queries(), Down :: queries()}]. migrations(#{}) -> [ {1, queries([ % Queues ~""" CREATE TABLE IF NOT EXISTS gaffer_queues ( name TEXT PRIMARY KEY ) """, % Jobs ~""" CREATE TABLE IF NOT EXISTS gaffer_jobs ( id UUID PRIMARY KEY, queue TEXT NOT NULL REFERENCES gaffer_queues(name), state TEXT NOT NULL CHECK (state IN ('available', 'executing', 'completed', 'cancelled', 'failed')), payload JSONB NOT NULL, attempt INTEGER NOT NULL, max_attempts INTEGER NOT NULL, priority INTEGER NOT NULL, timeout INTEGER, backoff JSONB, shutdown_timeout INTEGER, result JSONB, errors JSONB NOT NULL, scheduled_at TIMESTAMPTZ, created_at TIMESTAMPTZ NOT NULL, attempted_at TIMESTAMPTZ, completed_at TIMESTAMPTZ, cancelled_at TIMESTAMPTZ, failed_at TIMESTAMPTZ ) """, % Query indexes ~""" CREATE INDEX IF NOT EXISTS idx_gaffer_jobs_claimable ON gaffer_jobs (queue, priority DESC, created_at ASC) WHERE state = 'available' """, ~""" CREATE INDEX IF NOT EXISTS idx_gaffer_jobs_queue_state ON gaffer_jobs (queue, state) """ ]), queries([ ~"DROP TABLE IF EXISTS gaffer_jobs", ~"DROP TABLE IF EXISTS gaffer_queues" ])} ]. -doc "Queries to apply a migration and record an `up` event.". -spec migrate_up({pos_integer(), queries(), _}) -> queries(). migrate_up({Version, UpQueries, _DownQueries}) -> UpQueries ++ [ { ~"INSERT INTO gaffer_schema_migrations (version, direction) VALUES ($1, 'up')", [Version] } ]. -doc "Queries to roll back a migration and record a `down` event.". -spec migrate_down({pos_integer(), _, queries()}) -> queries(). migrate_down({Version, _UpQueries, DownQueries}) -> DownQueries ++ [ { ~"INSERT INTO gaffer_schema_migrations (version, direction) VALUES ($1, 'down')", [Version] } ]. -doc "Queries to create the migrations history table if it does not exist.". -spec ensure_migrations_table() -> queries(). ensure_migrations_table() -> queries([ ~""" CREATE TABLE IF NOT EXISTS gaffer_schema_migrations ( id BIGSERIAL PRIMARY KEY, version BIGINT NOT NULL, direction TEXT NOT NULL CHECK (direction IN ('up', 'down')), created_at TIMESTAMPTZ NOT NULL DEFAULT NOW() ) """ ]). -doc "Query that fetches the currently applied migration versions.". -spec applied_versions() -> queries(). applied_versions() -> [ { ~""" SELECT version FROM ( SELECT DISTINCT ON (version) version, direction FROM gaffer_schema_migrations ORDER BY version, id DESC ) t WHERE direction = 'up' ORDER BY version """, [] } ]. -doc "Query that fetches the full migrations history, oldest first.". -spec history() -> queries(). history() -> SQL = [ ~"SELECT version, direction, ", ts_column(~"created_at", ~"created_at"), ~" FROM gaffer_schema_migrations ORDER BY id" ], [{SQL, []}]. % Queues -doc "Query to register a queue name. No-op if already registered.". -spec queue_insert(gaffer:queue()) -> queries(). queue_insert(Name) -> [ { ~"INSERT INTO gaffer_queues (name) VALUES ($1) ON CONFLICT (name) DO NOTHING", [atom_to_binary(Name)] } ]. -doc "Query to check whether a queue name exists.". -spec queue_exists(gaffer:queue()) -> queries(). queue_exists(Name) -> [{~"SELECT 1 FROM gaffer_queues WHERE name = $1", [atom_to_binary(Name)]}]. -doc "Query to list all registered queue names.". -spec queue_list() -> queries(). queue_list() -> [{~"SELECT name FROM gaffer_queues", []}]. -doc "Query to delete a queue by name.". -spec queue_delete(gaffer:queue()) -> queries(). queue_delete(Name) -> [{~"DELETE FROM gaffer_queues WHERE name = $1", [atom_to_binary(Name)]}]. % Introspection -doc "Query that aggregates job counts and timestamps per state.". -spec info(gaffer:queue()) -> queries(). info(Queue) -> TSCase = ts_case_for_state(), SQL = [ ~"SELECT state, COUNT(*) AS count, ", ts_column([~"MIN(", TSCase, ~")"], ~"oldest"), ~", ", ts_column([~"MAX(", TSCase, ~")"], ~"newest"), ~" FROM gaffer_jobs WHERE queue = $1 GROUP BY state" ], [{SQL, [atom_to_binary(Queue)]}]. % Jobs -doc "Returns upsert queries for a single encoded job.". -spec job_write(map()) -> queries(). job_write(Encoded) -> {Cols, Phs, Vals} = columns_and_values(Encoded), Mutable = [C || C <:- Cols, not lists:member(C, immutable_columns())], Sets = lists:join(~", ", [[C, ~" = EXCLUDED.", C] || C <:- Mutable]), SQL = [ ~"INSERT INTO gaffer_jobs (", lists:join(~", ", Cols), ~") VALUES (", lists:join(~", ", Phs), ~") ON CONFLICT (id) DO UPDATE SET ", Sets, ~" RETURNING ", job_columns() ], [{SQL, Vals}]. -doc "Query to fetch a job by ID.". -spec job_get(term()) -> queries(). job_get(ID) -> SQL = [ ~"SELECT ", job_columns(), ~" FROM gaffer_jobs WHERE id = $1" ], [{SQL, [ID]}]. -doc "Query to list jobs matching the given filters.". -spec job_list(map()) -> queries(). job_list(#{queue := Queue} = Opts) -> {SQL, Params} = case Opts of #{state := State} -> { [ ~"SELECT ", job_columns(), ~" FROM gaffer_jobs WHERE queue = $1 AND state = $2" ], [Queue, State] }; _ -> { [ ~"SELECT ", job_columns(), ~" FROM gaffer_jobs WHERE queue = $1" ], [Queue] } end, [{SQL, Params}]. -doc "Query to delete a job by ID.". -spec job_delete(term()) -> queries(). job_delete(ID) -> [{~"DELETE FROM gaffer_jobs WHERE id = $1", [ID]}]. job_columns() -> job_columns(~""). job_columns(Prefix) -> lists:join(~", ", [ [Prefix, ~"id"], [Prefix, ~"queue"], [Prefix, ~"state"], [Prefix, ~"payload"], [Prefix, ~"attempt"], [Prefix, ~"max_attempts"], [Prefix, ~"priority"], [Prefix, ~"timeout"], [Prefix, ~"backoff"], [Prefix, ~"shutdown_timeout"], [Prefix, ~"result"], [Prefix, ~"errors"] | [ ts_column([Prefix, C], C) || C <:- ts_column_names() ] ]). ts_column(Expr, Alias) -> [ ~"EXTRACT(EPOCH FROM date_trunc('second', ", Expr, ~"))::bigint * 1000000 + MOD(EXTRACT(MICROSECONDS FROM ", Expr, ~")::bigint, 1000000) AS ", Alias ]. ts_column_names() -> [ ~"scheduled_at", ~"created_at", ~"attempted_at", ~"completed_at", ~"cancelled_at", ~"failed_at" ]. % Job lifecycle -doc "Query to atomically claim available jobs for execution.". -spec job_claim(map(), map()) -> queries(). job_claim( #{queue := Queue, limit := Limit, global_max_workers := GlobalMax}, #{state := State, attempted_at := AttemptedAt} ) -> Now = AttemptedAt, {LimitClause, LimitParams} = job_claim_effective_limit(Limit, GlobalMax), SQL = [ ~""" WITH candidates AS ( SELECT id FROM gaffer_jobs WHERE queue = $1 AND state = 'available' AND (scheduled_at IS NULL OR scheduled_at <= to_timestamp($2::bigint / 1000000.0)) ORDER BY priority DESC, created_at ASC """, LimitClause, ~""" FOR UPDATE SKIP LOCKED ) UPDATE gaffer_jobs j SET state = $3, attempted_at = to_timestamp($4::bigint / 1000000.0) FROM candidates c WHERE j.id = c.id RETURNING """, ~" ", job_columns(~"j.") ], [{SQL, [Queue, Now, State, Now | LimitParams]}]. -doc "Query to delete jobs older than per-state cutoffs for a queue.". -spec job_prune(gaffer:queue(), map()) -> queries(). job_prune(Queue, Opts) -> {Clauses, Params, N} = maps:fold(fun prune_clause/3, {[], [], 1}, Opts), QueueParam = [~"$", integer_to_binary(N)], SQL = [ ~"DELETE FROM gaffer_jobs WHERE queue = ", QueueParam, ~" AND (", lists:join(~" OR ", lists:reverse(Clauses)), ~")", ~" RETURNING id" ], [{SQL, lists:reverse(Params, [atom_to_binary(Queue)])}]. prune_clause(State, all, {Cs, Ps, N}) -> {[[~"(state = '", atom_to_binary(State), ~"')"] | Cs], Ps, N}; prune_clause(State, Cutoff, {Cs, Ps, N}) -> Col = state_timestamp_column(State), StateName = atom_to_binary(State), C = [~"(state = '", StateName, ~"' AND ", Col, older_than(N), ~")"], {[C | Cs], [Cutoff | Ps], N + 1}. older_than(N) -> [~" < to_timestamp($", integer_to_binary(N), ~"::bigint / 1000000.0)"]. %--- Internal ------------------------------------------------------------------ ts_case_for_state() -> ~""" CASE state WHEN 'available' THEN created_at WHEN 'executing' THEN attempted_at WHEN 'completed' THEN completed_at WHEN 'cancelled' THEN cancelled_at WHEN 'failed' THEN failed_at END """. state_timestamp_column(available) -> ~"created_at"; state_timestamp_column(executing) -> ~"attempted_at"; state_timestamp_column(completed) -> ~"completed_at"; state_timestamp_column(cancelled) -> ~"cancelled_at"; state_timestamp_column(failed) -> ~"failed_at". immutable_columns() -> [~"id", ~"queue", ~"created_at"]. job_claim_effective_limit(infinity, infinity) -> {~"", []}; job_claim_effective_limit(Limit, infinity) -> {~"\n LIMIT $5::bigint", [Limit]}; job_claim_effective_limit(Limit, GlobalMax) -> { ~""" LIMIT GREATEST(0, ( SELECT LEAST($5::bigint, $6::bigint - count(*)) FROM gaffer_jobs WHERE queue = $1 AND state = 'executing' )) """, [Limit, GlobalMax] }. columns_and_values(Map) -> Pairs = maps:to_list(Map), Cols = [atom_to_binary(K) || {K, _} <:- Pairs], Vals = [V || {_, V} <:- Pairs], Phs = [ [~"$", integer_to_binary(I)] || I <:- lists:seq(1, length(Cols)) ], {Cols, Phs, Vals}. queries(SQLs) -> [{SQL, []} || SQL <:- SQLs].