packages feed

pgmq-migration (empty) → 0.1.0.0

raw patch · 33 files changed

+5538/−0 lines, 33 filesdep +basedep +bytestringdep +ephemeral-pg

Dependencies added: base, bytestring, ephemeral-pg, file-embed, hasql, hasql-migration, hasql-transaction, pgmq-migration, tasty, tasty-hunit, text, transformers

Files

+ CHANGELOG.md view
@@ -0,0 +1,8 @@+# Changelog for pgmq-migration++## 0.1.0.0 -- 2026-02-21++* Initial release+* Support for PGMQ v1.9.0 schema installation+* Support for PGMQ v1.10.0 schema installation+* Incremental migration support (v1.9.0 to v1.10.0 upgrade path)
+ LICENSE view
@@ -0,0 +1,20 @@+Copyright (c) 2025 Nadeem Bitar++Permission is hereby granted, free of charge, to any person obtaining+a copy of this software and associated documentation files (the+"Software"), to deal in the Software without restriction, including+without limitation the rights to use, copy, modify, merge, publish,+distribute, sublicense, and/or sell copies of the Software, and to+permit persons to whom the Software is furnished to do so, subject to+the following conditions:++The above copyright notice and this permission notice shall be included+in all copies or substantial portions of the Software.++THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND,+EXPRESS OR IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF+MERCHANTABILITY, FITNESS FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT.+IN NO EVENT SHALL THE AUTHORS OR COPYRIGHT HOLDERS BE LIABLE FOR ANY+CLAIM, DAMAGES OR OTHER LIABILITY, WHETHER IN AN ACTION OF CONTRACT,+TORT OR OTHERWISE, ARISING FROM, OUT OF OR IN CONNECTION WITH THE+SOFTWARE OR THE USE OR OTHER DEALINGS IN THE SOFTWARE.
+ database/migrations/v1.9.0_to_v1.10.0.sql view
@@ -0,0 +1,594 @@+-- Add last_read_at column to all existing queue tables (q_*) and archive tables (a_*)+DO $$+DECLARE+    queue_record RECORD;+    qtable TEXT;+    atable TEXT;+BEGIN+    FOR queue_record IN SELECT queue_name FROM pgmq.meta LOOP+        qtable := 'q_' || queue_record.queue_name;+        atable := 'a_' || queue_record.queue_name;++        -- Add last_read_at to queue table if it doesn't exist+        IF NOT EXISTS (+            SELECT 1 FROM information_schema.columns+            WHERE table_schema = 'pgmq'+            AND table_name = qtable+            AND column_name = 'last_read_at'+        ) THEN+            EXECUTE FORMAT('ALTER TABLE pgmq.%I ADD COLUMN last_read_at TIMESTAMP WITH TIME ZONE', qtable);+        END IF;++        -- Add last_read_at to archive table if it doesn't exist+        IF EXISTS (+            SELECT 1 FROM information_schema.tables+            WHERE table_schema = 'pgmq'+            AND table_name = atable+        ) AND NOT EXISTS (+            SELECT 1 FROM information_schema.columns+            WHERE table_schema = 'pgmq'+            AND table_name = atable+            AND column_name = 'last_read_at'+        ) THEN+            EXECUTE FORMAT('ALTER TABLE pgmq.%I ADD COLUMN last_read_at TIMESTAMP WITH TIME ZONE', atable);+        END IF;+    END LOOP;+END;+$$;++-- The functions that use this type are being dropped and recreated below+DROP TYPE IF EXISTS pgmq.message_record CASCADE;++CREATE TYPE pgmq.message_record AS (+    msg_id BIGINT,+    read_ct INTEGER,+    enqueued_at TIMESTAMP WITH TIME ZONE,+    last_read_at TIMESTAMP WITH TIME ZONE,+    vt TIMESTAMP WITH TIME ZONE,+    message JSONB,+    headers JSONB+);++------------------------------------------------------------+-- Migration: Updated functions+------------------------------------------------------------++DROP FUNCTION IF EXISTS pgmq.read_grouped_rr(+queue_name TEXT,+    vt INTEGER,+    qty INTEGER+);+CREATE FUNCTION pgmq.read_grouped_rr(+    queue_name TEXT,+    vt INTEGER,+    qty INTEGER+)+RETURNS SETOF pgmq.message_record AS $$+DECLARE+    sql TEXT;+    qtable TEXT := pgmq.format_table_name(queue_name, 'q');+BEGIN+    sql := FORMAT(+        $QUERY$+        WITH fifo_groups AS (+            -- Determine the absolute head (oldest) message id per FIFO group, regardless of visibility+            SELECT+                COALESCE(headers->>'x-pgmq-group', '_default_fifo_group') AS fifo_key,+                MIN(msg_id) AS head_msg_id+            FROM pgmq.%1$I+            GROUP BY COALESCE(headers->>'x-pgmq-group', '_default_fifo_group')+        ),+        eligible_groups AS (+            -- Only groups whose head message is currently visible+            -- Acquire a transaction-level advisory lock per group to prevent concurrent selection+            SELECT+                g.fifo_key,+                g.head_msg_id,+                ROW_NUMBER() OVER (ORDER BY g.head_msg_id) AS group_priority+            FROM fifo_groups g+            JOIN pgmq.%2$I h ON h.msg_id = g.head_msg_id+            WHERE h.vt <= clock_timestamp()+              AND pg_try_advisory_xact_lock(pg_catalog.hashtextextended(g.fifo_key, 0))+        ),+        available_messages AS (+            -- All currently visible messages starting at the head for each eligible group+            SELECT+                m.msg_id,+                eg.group_priority,+                ROW_NUMBER() OVER (+                    PARTITION BY eg.fifo_key+                    ORDER BY m.msg_id+                ) AS msg_rank_in_group+            FROM pgmq.%3$I m+            JOIN eligible_groups eg+              ON COALESCE(m.headers->>'x-pgmq-group', '_default_fifo_group') = eg.fifo_key+            WHERE m.vt <= clock_timestamp()+              AND m.msg_id >= eg.head_msg_id+        ),+        ordered_messages AS (+            -- Layered round-robin: take rank 1 of all groups by group_priority, then rank 2, etc.+            -- Assign selection order before locking+            SELECT msg_id, ROW_NUMBER() OVER (ORDER BY msg_rank_in_group, group_priority) as selection_order+            FROM available_messages+        ),+        selected_messages AS (+            -- Lock the messages in the correct order, preserving selection_order+            SELECT om.msg_id, om.selection_order+            FROM ordered_messages om+            JOIN pgmq.%4$I m ON m.msg_id = om.msg_id+            WHERE om.selection_order <= $1+            ORDER BY om.selection_order+            FOR UPDATE OF m SKIP LOCKED+        ),+        updated_messages AS (+            UPDATE pgmq.%5$I m+            SET+                vt = clock_timestamp() + %6$L,+                read_ct = read_ct + 1,+                last_read_at = clock_timestamp()+            FROM selected_messages sm+            WHERE m.msg_id = sm.msg_id+              AND m.vt <= clock_timestamp() -- final guard to avoid duplicate reads under races+            RETURNING m.msg_id, m.read_ct, m.enqueued_at, m.last_read_at, m.vt, m.message, m.headers, sm.selection_order+        )+        SELECT msg_id, read_ct, enqueued_at, last_read_at, vt, message, headers+        FROM updated_messages+        ORDER BY selection_order;+        $QUERY$,+        qtable, qtable, qtable, qtable, qtable, make_interval(secs => vt)+    );+    RETURN QUERY EXECUTE sql USING qty;+END;+$$ LANGUAGE plpgsql;+++DROP FUNCTION IF EXISTS pgmq.read_grouped_rr_with_poll(+    queue_name TEXT,+    vt INTEGER,+    qty INTEGER,+    max_poll_seconds INTEGER,+    poll_interval_ms INTEGER+);+CREATE FUNCTION pgmq.read_grouped_rr_with_poll(+    queue_name TEXT,+    vt INTEGER,+    qty INTEGER,+    max_poll_seconds INTEGER DEFAULT 5,+    poll_interval_ms INTEGER DEFAULT 100+)+RETURNS SETOF pgmq.message_record AS $$+DECLARE+    r pgmq.message_record;+    stop_at TIMESTAMP;+BEGIN+    stop_at := clock_timestamp() + make_interval(secs => max_poll_seconds);+    LOOP+      IF (SELECT clock_timestamp() >= stop_at) THEN+        RETURN;+      END IF;++      FOR r IN+        SELECT * FROM pgmq.read_grouped_rr(queue_name, vt, qty)+      LOOP+        RETURN NEXT r;+      END LOOP;+      IF FOUND THEN+        RETURN;+      ELSE+        PERFORM pg_sleep(poll_interval_ms::numeric / 1000);+      END IF;+    END LOOP;+END;+$$ LANGUAGE plpgsql;+++DROP FUNCTION IF EXISTS pgmq.read(+    queue_name TEXT,+    vt INTEGER,+    qty INTEGER,+    conditional JSONB+);+CREATE FUNCTION pgmq.read(+    queue_name TEXT,+    vt INTEGER,+    qty INTEGER,+    conditional JSONB DEFAULT '{}'+)+RETURNS SETOF pgmq.message_record AS $$+DECLARE+    sql TEXT;+    qtable TEXT := pgmq.format_table_name(queue_name, 'q');+BEGIN+    sql := FORMAT(+        $QUERY$+        WITH cte AS+        (+            SELECT msg_id+            FROM pgmq.%I+            WHERE vt <= clock_timestamp() AND CASE+                WHEN %L != '{}'::jsonb THEN (message @> %2$L)::integer+                ELSE 1+            END = 1+            ORDER BY msg_id ASC+            LIMIT $1+            FOR UPDATE SKIP LOCKED+        )+        UPDATE pgmq.%I m+        SET+            last_read_at = clock_timestamp(),+            vt = clock_timestamp() + %L,+            read_ct = read_ct + 1+        FROM cte+        WHERE m.msg_id = cte.msg_id+        RETURNING m.*;+        $QUERY$,+        qtable, conditional, qtable, make_interval(secs => vt)+    );+    RETURN QUERY EXECUTE sql USING qty;+END;+$$ LANGUAGE plpgsql;++DROP FUNCTION IF EXISTS pgmq.read_grouped(+    queue_name TEXT,+    vt INTEGER,+    qty INTEGER+);+CREATE FUNCTION pgmq.read_grouped(+    queue_name TEXT,+    vt INTEGER,+    qty INTEGER+)+RETURNS SETOF pgmq.message_record AS $$+DECLARE+    sql TEXT;+    qtable TEXT := pgmq.format_table_name(queue_name, 'q');+BEGIN+    sql := FORMAT(+        $QUERY$+        WITH fifo_groups AS (+            -- Find the minimum msg_id for each FIFO group that's ready to be processed+            SELECT+                COALESCE(headers->>'x-pgmq-group', '_default_fifo_group') as fifo_key,+                MIN(msg_id) as min_msg_id+            FROM pgmq.%I+            WHERE vt <= clock_timestamp()+            GROUP BY COALESCE(headers->>'x-pgmq-group', '_default_fifo_group')+        ),+        locked_groups AS (+            -- Lock the first available message in each FIFO group+            SELECT+                m.msg_id,+                fg.fifo_key+            FROM pgmq.%I m+            INNER JOIN fifo_groups fg ON+                COALESCE(m.headers->>'x-pgmq-group', '_default_fifo_group') = fg.fifo_key+                AND m.msg_id = fg.min_msg_id+            WHERE m.vt <= clock_timestamp()+            ORDER BY m.msg_id ASC+            FOR UPDATE SKIP LOCKED+        ),+        group_priorities AS (+            -- Assign priority to groups based on their oldest message+            SELECT+                fifo_key,+                msg_id as min_msg_id,+                ROW_NUMBER() OVER (ORDER BY msg_id) as group_priority+            FROM locked_groups+        ),+        available_messages AS (+            -- Get messages prioritizing filling batch from earliest group first+            SELECT+                m.msg_id,+                gp.group_priority,+                ROW_NUMBER() OVER (PARTITION BY gp.fifo_key ORDER BY m.msg_id) as msg_rank_in_group+            FROM pgmq.%I m+            INNER JOIN group_priorities gp ON+                COALESCE(m.headers->>'x-pgmq-group', '_default_fifo_group') = gp.fifo_key+            WHERE m.vt <= clock_timestamp()+            AND m.msg_id >= gp.min_msg_id  -- Only messages from min_msg_id onwards in each group+            AND NOT EXISTS (+                -- Ensure no earlier message in this group is currently being processed+                SELECT 1+                FROM pgmq.%I m2+                WHERE COALESCE(m2.headers->>'x-pgmq-group', '_default_fifo_group') =+                      COALESCE(m.headers->>'x-pgmq-group', '_default_fifo_group')+                AND m2.vt > clock_timestamp()+                AND m2.msg_id < m.msg_id+            )+        ),+        batch_selection AS (+            -- Select messages to fill batch, prioritizing earliest group+            SELECT+                msg_id,+                ROW_NUMBER() OVER (ORDER BY group_priority, msg_rank_in_group) as overall_rank+            FROM available_messages+        ),+        selected_messages AS (+            -- Limit to requested quantity+            SELECT msg_id+            FROM batch_selection+            WHERE overall_rank <= $1+            ORDER BY msg_id+            FOR UPDATE SKIP LOCKED+        )+        UPDATE pgmq.%I m+        SET+            vt = clock_timestamp() + %L,+            read_ct = read_ct + 1,+            last_read_at = clock_timestamp()+        FROM selected_messages sm+        WHERE m.msg_id = sm.msg_id+        RETURNING m.msg_id, m.read_ct, m.enqueued_at, m.last_read_at, m.vt, m.message, m.headers;+        $QUERY$,+        qtable, qtable, qtable, qtable, qtable, make_interval(secs => vt)+    );+    RETURN QUERY EXECUTE sql USING qty;+END;+$$ LANGUAGE plpgsql;+++DROP FUNCTION IF EXISTS pgmq.read_grouped_with_poll(+    queue_name TEXT,+    vt INTEGER,+    qty INTEGER,+    max_poll_seconds INTEGER,+    poll_interval_ms INTEGER+);+CREATE FUNCTION pgmq.read_grouped_with_poll(+    queue_name TEXT,+    vt INTEGER,+    qty INTEGER,+    max_poll_seconds INTEGER DEFAULT 5,+    poll_interval_ms INTEGER DEFAULT 100+)+RETURNS SETOF pgmq.message_record AS $$+DECLARE+    r pgmq.message_record;+    stop_at TIMESTAMP;+BEGIN+    stop_at := clock_timestamp() + make_interval(secs => max_poll_seconds);+    LOOP+      IF (SELECT clock_timestamp() >= stop_at) THEN+        RETURN;+      END IF;++      FOR r IN+        SELECT * FROM pgmq.read_grouped(queue_name, vt, qty)+      LOOP+        RETURN NEXT r;+      END LOOP;+      IF FOUND THEN+        RETURN;+      ELSE+        PERFORM pg_sleep(poll_interval_ms::numeric / 1000);+      END IF;+    END LOOP;+END;+$$ LANGUAGE plpgsql;++DROP FUNCTION IF EXISTS pgmq.read_with_poll(+    queue_name TEXT,+    vt INTEGER,+    qty INTEGER,+    max_poll_seconds INTEGER,+    poll_interval_ms INTEGER,+    conditional JSONB+);+CREATE FUNCTION pgmq.read_with_poll(+    queue_name TEXT,+    vt INTEGER,+    qty INTEGER,+    max_poll_seconds INTEGER DEFAULT 5,+    poll_interval_ms INTEGER DEFAULT 100,+    conditional JSONB DEFAULT '{}'+)+RETURNS SETOF pgmq.message_record AS $$+DECLARE+    r pgmq.message_record;+    stop_at TIMESTAMP;+    sql TEXT;+    qtable TEXT := pgmq.format_table_name(queue_name, 'q');+BEGIN+    stop_at := clock_timestamp() + make_interval(secs => max_poll_seconds);+    LOOP+      IF (SELECT clock_timestamp() >= stop_at) THEN+        RETURN;+      END IF;++      sql := FORMAT(+          $QUERY$+          WITH cte AS+          (+              SELECT msg_id+              FROM pgmq.%I+              WHERE vt <= clock_timestamp() AND CASE+                  WHEN %L != '{}'::jsonb THEN (message @> %2$L)::integer+                  ELSE 1+              END = 1+              ORDER BY msg_id ASC+              LIMIT $1+              FOR UPDATE SKIP LOCKED+          )+          UPDATE pgmq.%I m+          SET+              last_read_at = clock_timestamp(),+              vt = clock_timestamp() + %L,+              read_ct = read_ct + 1+          FROM cte+          WHERE m.msg_id = cte.msg_id+          RETURNING m.*;+          $QUERY$,+          qtable, conditional, qtable, make_interval(secs => vt)+      );++      FOR r IN+        EXECUTE sql USING qty+      LOOP+        RETURN NEXT r;+      END LOOP;+      IF FOUND THEN+        RETURN;+      ELSE+        PERFORM pg_sleep(poll_interval_ms::numeric / 1000);+      END IF;+    END LOOP;+END;+$$ LANGUAGE plpgsql;++-- Update archive functions to include last_read_at column+CREATE OR REPLACE FUNCTION pgmq.archive(+    queue_name TEXT,+    msg_id BIGINT+)+RETURNS BOOLEAN AS $$+DECLARE+    sql TEXT;+    result BIGINT;+    qtable TEXT := pgmq.format_table_name(queue_name, 'q');+    atable TEXT := pgmq.format_table_name(queue_name, 'a');+BEGIN+    sql := FORMAT(+        $QUERY$+        WITH archived AS (+            DELETE FROM pgmq.%I+            WHERE msg_id = $1+            RETURNING msg_id, vt, read_ct, enqueued_at, last_read_at, message, headers+        )+        INSERT INTO pgmq.%I (msg_id, vt, read_ct, enqueued_at, last_read_at, message, headers)+        SELECT msg_id, vt, read_ct, enqueued_at, last_read_at, message, headers+        FROM archived+        RETURNING msg_id;+        $QUERY$,+        qtable, atable+    );+    EXECUTE sql USING msg_id INTO result;+    RETURN NOT (result IS NULL);+END;+$$ LANGUAGE plpgsql;++CREATE OR REPLACE FUNCTION pgmq.archive(+    queue_name TEXT,+    msg_ids BIGINT[]+)+RETURNS SETOF BIGINT AS $$+DECLARE+    sql TEXT;+    qtable TEXT := pgmq.format_table_name(queue_name, 'q');+    atable TEXT := pgmq.format_table_name(queue_name, 'a');+BEGIN+    sql := FORMAT(+        $QUERY$+        WITH archived AS (+            DELETE FROM pgmq.%I+            WHERE msg_id = ANY($1)+            RETURNING msg_id, vt, read_ct, enqueued_at, last_read_at, message, headers+        )+        INSERT INTO pgmq.%I (msg_id, vt, read_ct, enqueued_at, last_read_at, message, headers)+        SELECT msg_id, vt, read_ct, enqueued_at, last_read_at, message, headers+        FROM archived+        RETURNING msg_id;+        $QUERY$,+        qtable, atable+    );+    RETURN QUERY EXECUTE sql USING msg_ids;+END;+$$ LANGUAGE plpgsql;++DROP FUNCTION IF EXISTS pgmq.pop(+    queue_name TEXT,+    qty INTEGER+);+CREATE FUNCTION pgmq.pop(queue_name TEXT, qty INTEGER DEFAULT 1)+RETURNS SETOF pgmq.message_record AS $$+DECLARE+    sql TEXT;+    result pgmq.message_record;+    qtable TEXT := pgmq.format_table_name(queue_name, 'q');+BEGIN+    sql := FORMAT(+        $QUERY$+        WITH cte AS+            (+                SELECT msg_id+                FROM pgmq.%I+                WHERE vt <= clock_timestamp()+                ORDER BY msg_id ASC+                LIMIT $1+                FOR UPDATE SKIP LOCKED+            )+        DELETE from pgmq.%I+        WHERE msg_id IN (select msg_id from cte)+        RETURNING *;+        $QUERY$,+        qtable, qtable+    );+    RETURN QUERY EXECUTE sql USING qty;+END;+$$ LANGUAGE plpgsql;++-- Drop old set_vt overloads (they'll be recreated below)+DROP FUNCTION IF EXISTS pgmq.set_vt(TEXT, BIGINT, INTEGER);+DROP FUNCTION IF EXISTS pgmq.set_vt(TEXT, BIGINT, TIMESTAMP WITH TIME ZONE);+DROP FUNCTION IF EXISTS pgmq.set_vt(TEXT, BIGINT[], INTEGER);+DROP FUNCTION IF EXISTS pgmq.set_vt(TEXT, BIGINT[], TIMESTAMP WITH TIME ZONE);++-- Sets timestamp vt of a message, returns it (base implementation)+CREATE FUNCTION pgmq.set_vt(queue_name TEXT, msg_id BIGINT, vt TIMESTAMP WITH TIME ZONE)+RETURNS SETOF pgmq.message_record AS $$+DECLARE+    sql TEXT;+    result pgmq.message_record;+    qtable TEXT := pgmq.format_table_name(queue_name, 'q');+BEGIN+    sql := FORMAT(+        $QUERY$+        UPDATE pgmq.%I+        SET vt = $1+        WHERE msg_id = $2+        RETURNING *;+        $QUERY$,+        qtable+    );+    RETURN QUERY EXECUTE sql USING vt, msg_id;+END;+$$ LANGUAGE plpgsql;++-- Sets integer vt of a message, returns it+CREATE FUNCTION pgmq.set_vt(queue_name TEXT, msg_id BIGINT, vt INTEGER)+RETURNS SETOF pgmq.message_record AS $$+    SELECT * FROM pgmq.set_vt(queue_name, msg_id, clock_timestamp() + make_interval(secs => vt));+$$ LANGUAGE sql;++-- Sets timestamp vt of multiple messages, returns them (base implementation)+CREATE FUNCTION pgmq.set_vt(+    queue_name TEXT,+    msg_ids BIGINT[],+    vt TIMESTAMP WITH TIME ZONE+)+RETURNS SETOF pgmq.message_record AS $$+DECLARE+    sql TEXT;+    qtable TEXT := pgmq.format_table_name(queue_name, 'q');+BEGIN+    sql := FORMAT(+        $QUERY$+        UPDATE pgmq.%I+        SET vt = $1+        WHERE msg_id = ANY($2)+        RETURNING *;+        $QUERY$,+        qtable+    );+    RETURN QUERY EXECUTE sql USING vt, msg_ids;+END;+$$ LANGUAGE plpgsql;++-- Sets integer vt of multiple messages, returns them+CREATE FUNCTION pgmq.set_vt(+    queue_name TEXT,+    msg_ids BIGINT[],+    vt INTEGER+)+RETURNS SETOF pgmq.message_record AS $$+    SELECT * FROM pgmq.set_vt(queue_name, msg_ids, clock_timestamp() + make_interval(secs => vt));+$$ LANGUAGE sql;
+ database/v1.10.0/01_schema.sql view
@@ -0,0 +1,11 @@+------------------------------------------------------------+-- Schema creation and grants (without extension)+------------------------------------------------------------+CREATE SCHEMA IF NOT EXISTS pgmq;++-- Grant permission to pg_monitor to all tables and sequences+GRANT USAGE ON SCHEMA pgmq TO pg_monitor;+GRANT SELECT ON ALL TABLES IN SCHEMA pgmq TO pg_monitor;+GRANT SELECT ON ALL SEQUENCES IN SCHEMA pgmq TO pg_monitor;+ALTER DEFAULT PRIVILEGES IN SCHEMA pgmq GRANT SELECT ON TABLES TO pg_monitor;+ALTER DEFAULT PRIVILEGES IN SCHEMA pgmq GRANT SELECT ON SEQUENCES TO pg_monitor;
+ database/v1.10.0/02_tables.sql view
@@ -0,0 +1,25 @@+------------------------------------------------------------+-- Tables: meta and notify_insert_throttle+------------------------------------------------------------++-- Table where queues and metadata about them is stored+CREATE TABLE IF NOT EXISTS pgmq.meta (+    queue_name VARCHAR UNIQUE NOT NULL,+    is_partitioned BOOLEAN NOT NULL,+    is_unlogged BOOLEAN NOT NULL,+    created_at TIMESTAMP WITH TIME ZONE DEFAULT now() NOT NULL+);++-- Table to track notification throttling for queues+CREATE UNLOGGED TABLE IF NOT EXISTS pgmq.notify_insert_throttle (+    queue_name           VARCHAR UNIQUE NOT NULL -- Queue name (without 'q_' prefix)+       CONSTRAINT notify_insert_throttle_meta_queue_name_fk+            REFERENCES pgmq.meta (queue_name)+            ON DELETE CASCADE,+    throttle_interval_ms INTEGER NOT NULL DEFAULT 0, -- Min milliseconds between notifications (0 = no throttling)+    last_notified_at     TIMESTAMP WITH TIME ZONE NOT NULL DEFAULT to_timestamp(0) -- Timestamp of last sent notification+);++CREATE INDEX IF NOT EXISTS idx_notify_throttle_active+    ON pgmq.notify_insert_throttle (queue_name, last_notified_at)+    WHERE throttle_interval_ms > 0;
+ database/v1.10.0/03_types.sql view
@@ -0,0 +1,34 @@+------------------------------------------------------------+-- Custom types: message_record, queue_record, metrics_result+------------------------------------------------------------++-- This type has the shape of a message in a queue, and is often returned by+-- pgmq functions that return messages+-- Note: last_read_at field added in pgmq 1.10.0+CREATE TYPE pgmq.message_record AS (+    msg_id BIGINT,+    read_ct INTEGER,+    enqueued_at TIMESTAMP WITH TIME ZONE,+    last_read_at TIMESTAMP WITH TIME ZONE,+    vt TIMESTAMP WITH TIME ZONE,+    message JSONB,+    headers JSONB+);++CREATE TYPE pgmq.queue_record AS (+    queue_name VARCHAR,+    is_partitioned BOOLEAN,+    is_unlogged BOOLEAN,+    created_at TIMESTAMP WITH TIME ZONE+);++-- returned by pgmq.metrics() and pgmq.metrics_all+CREATE TYPE pgmq.metrics_result AS (+    queue_name text,+    queue_length bigint,+    newest_msg_age_sec int,+    oldest_msg_age_sec int,+    total_messages bigint,+    scrape_time timestamp with time zone,+    queue_visible_length bigint+);
+ database/v1.10.0/04_core_functions.sql view
@@ -0,0 +1,69 @@+------------------------------------------------------------+-- Core utility functions+------------------------------------------------------------++-- prevents race conditions during queue creation by acquiring a transaction-level advisory lock+-- uses a transaction advisory lock maintain the lock until transaction commit+-- a race condition would still exist if lock was released before commit+CREATE FUNCTION pgmq.acquire_queue_lock(queue_name TEXT)+RETURNS void AS $$+BEGIN+  PERFORM pg_advisory_xact_lock(hashtext('pgmq.queue_' || queue_name));+END;+$$ LANGUAGE plpgsql;++-- a helper to format table names and check for invalid characters+CREATE FUNCTION pgmq.format_table_name(queue_name text, prefix text)+RETURNS TEXT AS $$+BEGIN+    IF queue_name ~ '\$|;|--|'''+    THEN+        RAISE EXCEPTION 'queue name contains invalid characters: $, ;, --, or \''';+    END IF;+    RETURN lower(prefix || '_' || queue_name);+END;+$$ LANGUAGE plpgsql;++CREATE FUNCTION pgmq.validate_queue_name(queue_name TEXT)+RETURNS void AS $$+BEGIN+  IF length(queue_name) > 47 THEN+    -- complete table identifier must be <= 63+    -- https://www.postgresql.org/docs/17/sql-syntax-lexical.html#SQL-SYNTAX-IDENTIFIERS+    -- e.g. template_pgmq_q_my_queue is an identifier for my_queue when partitioned+    -- template_pgmq_q_ (16) + <a max length queue name> (47) = 63+    RAISE EXCEPTION 'queue name is too long, maximum length is 47 characters';+  END IF;+END;+$$ LANGUAGE plpgsql;++CREATE FUNCTION pgmq._belongs_to_pgmq(table_name TEXT)+RETURNS BOOLEAN AS $$+DECLARE+    sql TEXT;+    result BOOLEAN;+BEGIN+  SELECT EXISTS (+    SELECT 1+    FROM pg_depend+    WHERE refobjid = (SELECT oid FROM pg_extension WHERE extname = 'pgmq')+    AND objid = (+        SELECT oid+        FROM pg_class+        WHERE relname = table_name+    )+  ) INTO result;+  RETURN result;+END;+$$ LANGUAGE plpgsql;++CREATE FUNCTION pgmq._extension_exists(extension_name TEXT)+    RETURNS BOOLEAN+    LANGUAGE SQL+AS $$+SELECT EXISTS (+    SELECT 1+    FROM pg_extension+    WHERE extname = extension_name+)+$$;
+ database/v1.10.0/05_queue_management.sql view
@@ -0,0 +1,265 @@+------------------------------------------------------------+-- Queue management functions+------------------------------------------------------------++CREATE FUNCTION pgmq.create_non_partitioned(queue_name TEXT)+RETURNS void AS $$+DECLARE+  qtable TEXT := pgmq.format_table_name(queue_name, 'q');+  qtable_seq TEXT := qtable || '_msg_id_seq';+  atable TEXT := pgmq.format_table_name(queue_name, 'a');+BEGIN+  PERFORM pgmq.validate_queue_name(queue_name);+  PERFORM pgmq.acquire_queue_lock(queue_name);++  EXECUTE FORMAT(+    $QUERY$+    CREATE TABLE IF NOT EXISTS pgmq.%I (+        msg_id BIGINT PRIMARY KEY GENERATED ALWAYS AS IDENTITY,+        read_ct INT DEFAULT 0 NOT NULL,+        enqueued_at TIMESTAMP WITH TIME ZONE DEFAULT now() NOT NULL,+        last_read_at TIMESTAMP WITH TIME ZONE,+        vt TIMESTAMP WITH TIME ZONE NOT NULL,+        message JSONB,+        headers JSONB+    )+    $QUERY$,+    qtable+  );++  EXECUTE FORMAT(+    $QUERY$+    CREATE TABLE IF NOT EXISTS pgmq.%I (+      msg_id BIGINT PRIMARY KEY,+      read_ct INT DEFAULT 0 NOT NULL,+      enqueued_at TIMESTAMP WITH TIME ZONE DEFAULT now() NOT NULL,+      last_read_at TIMESTAMP WITH TIME ZONE,+      archived_at TIMESTAMP WITH TIME ZONE DEFAULT now() NOT NULL,+      vt TIMESTAMP WITH TIME ZONE NOT NULL,+      message JSONB,+      headers JSONB+    );+    $QUERY$,+    atable+  );++  EXECUTE FORMAT(+    $QUERY$+    CREATE INDEX IF NOT EXISTS %I ON pgmq.%I (vt ASC);+    $QUERY$,+    qtable || '_vt_idx', qtable+  );++  EXECUTE FORMAT(+    $QUERY$+    CREATE INDEX IF NOT EXISTS %I ON pgmq.%I (archived_at);+    $QUERY$,+    'archived_at_idx_' || queue_name, atable+  );++  EXECUTE FORMAT(+    $QUERY$+    INSERT INTO pgmq.meta (queue_name, is_partitioned, is_unlogged)+    VALUES (%L, false, false)+    ON CONFLICT+    DO NOTHING;+    $QUERY$,+    queue_name+  );++END;+$$ LANGUAGE plpgsql;++CREATE FUNCTION pgmq.create_unlogged(queue_name TEXT)+RETURNS void AS $$+DECLARE+  qtable TEXT := pgmq.format_table_name(queue_name, 'q');+  qtable_seq TEXT := qtable || '_msg_id_seq';+  atable TEXT := pgmq.format_table_name(queue_name, 'a');+BEGIN+  PERFORM pgmq.validate_queue_name(queue_name);+  PERFORM pgmq.acquire_queue_lock(queue_name);++  EXECUTE FORMAT(+    $QUERY$+    CREATE UNLOGGED TABLE IF NOT EXISTS pgmq.%I (+        msg_id BIGINT PRIMARY KEY GENERATED ALWAYS AS IDENTITY,+        read_ct INT DEFAULT 0 NOT NULL,+        enqueued_at TIMESTAMP WITH TIME ZONE DEFAULT now() NOT NULL,+        last_read_at TIMESTAMP WITH TIME ZONE,+        vt TIMESTAMP WITH TIME ZONE NOT NULL,+        message JSONB,+        headers JSONB+    )+    $QUERY$,+    qtable+  );++  EXECUTE FORMAT(+    $QUERY$+    CREATE TABLE IF NOT EXISTS pgmq.%I (+      msg_id BIGINT PRIMARY KEY,+      read_ct INT DEFAULT 0 NOT NULL,+      enqueued_at TIMESTAMP WITH TIME ZONE DEFAULT now() NOT NULL,+      last_read_at TIMESTAMP WITH TIME ZONE,+      archived_at TIMESTAMP WITH TIME ZONE DEFAULT now() NOT NULL,+      vt TIMESTAMP WITH TIME ZONE NOT NULL,+      message JSONB,+      headers JSONB+    );+    $QUERY$,+    atable+  );++  EXECUTE FORMAT(+    $QUERY$+    CREATE INDEX IF NOT EXISTS %I ON pgmq.%I (vt ASC);+    $QUERY$,+    qtable || '_vt_idx', qtable+  );++  EXECUTE FORMAT(+    $QUERY$+    CREATE INDEX IF NOT EXISTS %I ON pgmq.%I (archived_at);+    $QUERY$,+    'archived_at_idx_' || queue_name, atable+  );++  EXECUTE FORMAT(+    $QUERY$+    INSERT INTO pgmq.meta (queue_name, is_partitioned, is_unlogged)+    VALUES (%L, false, true)+    ON CONFLICT+    DO NOTHING;+    $QUERY$,+    queue_name+  );+END;+$$ LANGUAGE plpgsql;++CREATE FUNCTION pgmq.create(queue_name TEXT)+RETURNS void AS $$+BEGIN+    PERFORM pgmq.create_non_partitioned(queue_name);+END;+$$ LANGUAGE plpgsql;++CREATE FUNCTION pgmq.drop_queue(queue_name TEXT, partitioned BOOLEAN)+RETURNS BOOLEAN AS $$+DECLARE+    qtable TEXT := pgmq.format_table_name(queue_name, 'q');+    fq_qtable TEXT := 'pgmq.' || qtable;+    atable TEXT := pgmq.format_table_name(queue_name, 'a');+    fq_atable TEXT := 'pgmq.' || atable;+BEGIN+    RAISE WARNING 'drop_queue(queue_name, partitioned) is deprecated and will be removed in PGMQ v2.0. Use drop_queue(queue_name) instead';++    PERFORM pgmq.drop_queue(queue_name);++    RETURN TRUE;+END;+$$ LANGUAGE plpgsql;++CREATE FUNCTION pgmq.drop_queue(queue_name TEXT)+RETURNS BOOLEAN AS $$+DECLARE+    qtable TEXT := pgmq.format_table_name(queue_name, 'q');+    qtable_seq TEXT := qtable || '_msg_id_seq';+    fq_qtable TEXT := 'pgmq.' || qtable;+    atable TEXT := pgmq.format_table_name(queue_name, 'a');+    fq_atable TEXT := 'pgmq.' || atable;+    partitioned BOOLEAN;+BEGIN+    PERFORM pgmq.acquire_queue_lock(queue_name);+    EXECUTE FORMAT(+        $QUERY$+        SELECT is_partitioned FROM pgmq.meta WHERE queue_name = %L+        $QUERY$,+        queue_name+    ) INTO partitioned;++    -- check if the queue exists+    IF NOT EXISTS (+        SELECT 1+        FROM information_schema.tables+        WHERE table_name = qtable and table_schema = 'pgmq'+    ) THEN+        RAISE NOTICE 'pgmq queue `%` does not exist', queue_name;+        RETURN FALSE;+    END IF;++    EXECUTE FORMAT(+        $QUERY$+        DROP TABLE IF EXISTS pgmq.%I+        $QUERY$,+        qtable+    );++    EXECUTE FORMAT(+        $QUERY$+        DROP TABLE IF EXISTS pgmq.%I+        $QUERY$,+        atable+    );++     IF EXISTS (+          SELECT 1+          FROM information_schema.tables+          WHERE table_name = 'meta' and table_schema = 'pgmq'+     ) THEN+        EXECUTE FORMAT(+            $QUERY$+            DELETE FROM pgmq.meta WHERE queue_name = %L+            $QUERY$,+            queue_name+        );+     END IF;++     IF partitioned THEN+        EXECUTE FORMAT(+          $QUERY$+          DELETE FROM %I.part_config where parent_table in (%L, %L)+          $QUERY$,+          pgmq._get_pg_partman_schema(), fq_qtable, fq_atable+        );+     END IF;++    RETURN TRUE;+END;+$$ LANGUAGE plpgsql;++-- list queues+CREATE FUNCTION pgmq."list_queues"()+RETURNS SETOF pgmq.queue_record AS $$+BEGIN+  RETURN QUERY SELECT * FROM pgmq.meta;+END+$$ LANGUAGE plpgsql;++-- purge queue, deleting all entries in it.+CREATE OR REPLACE FUNCTION pgmq."purge_queue"(queue_name TEXT)+RETURNS BIGINT AS $$+DECLARE+  deleted_count INTEGER;+  qtable TEXT := pgmq.format_table_name(queue_name, 'q');+BEGIN+  -- Get the row count before truncating+  EXECUTE format('SELECT count(*) FROM pgmq.%I', qtable) INTO deleted_count;++  -- Use TRUNCATE for better performance on large tables+  EXECUTE format('TRUNCATE TABLE pgmq.%I', qtable);++  -- Return the number of purged rows+  RETURN deleted_count;+END+$$ LANGUAGE plpgsql;++-- unassign archive, so it can be kept when a queue is deleted+CREATE FUNCTION pgmq."detach_archive"(queue_name TEXT)+RETURNS VOID AS $$+DECLARE+  atable TEXT := pgmq.format_table_name(queue_name, 'a');+BEGIN+  RAISE WARNING 'detach_archive(queue_name) is deprecated and is a no-op. It will be removed in PGMQ v2.0. Archive tables are no longer member objects.';+END+$$ LANGUAGE plpgsql;
+ database/v1.10.0/06_message_ops.sql view
@@ -0,0 +1,676 @@+------------------------------------------------------------+-- Message operations: send, read, delete, archive, pop, set_vt+------------------------------------------------------------++-- send: 2 args, no delay+CREATE FUNCTION pgmq.send(+    queue_name TEXT,+    msg JSONB+)+RETURNS SETOF BIGINT AS $$+BEGIN+    RETURN QUERY SELECT * FROM pgmq.send(queue_name, msg, NULL, clock_timestamp());+END;+$$ LANGUAGE plpgsql;++-- send: 3 args with headers+CREATE FUNCTION pgmq.send(+    queue_name TEXT,+    msg JSONB,+    headers JSONB+)+RETURNS SETOF BIGINT AS $$+BEGIN+    RETURN QUERY SELECT * FROM pgmq.send(queue_name, msg, headers, clock_timestamp());+END;+$$ LANGUAGE plpgsql;++-- send: 3 args with delay+CREATE FUNCTION pgmq.send(+    queue_name TEXT,+    msg JSONB,+    delay INTEGER+)+RETURNS SETOF BIGINT AS $$+BEGIN+    RETURN QUERY SELECT * FROM pgmq.send(queue_name, msg, NULL, clock_timestamp() + make_interval(secs => delay));+END;+$$ LANGUAGE plpgsql;++-- send: 4 args with headers and delay+CREATE FUNCTION pgmq.send(+    queue_name TEXT,+    msg JSONB,+    headers JSONB,+    delay INTEGER+)+RETURNS SETOF BIGINT AS $$+BEGIN+    RETURN QUERY SELECT * FROM pgmq.send(queue_name, msg, headers, clock_timestamp() + make_interval(secs => delay));+END;+$$ LANGUAGE plpgsql;++-- send: actual implementation+CREATE FUNCTION pgmq.send(+    queue_name TEXT,+    msg JSONB,+    headers JSONB,+    vt TIMESTAMP WITH TIME ZONE+)+RETURNS SETOF BIGINT AS $$+DECLARE+    sql TEXT;+    qtable TEXT := pgmq.format_table_name(queue_name, 'q');+BEGIN+    sql := FORMAT(+        $QUERY$+        INSERT INTO pgmq.%I (vt, message, headers)+        VALUES ($1, $2, $3)+        RETURNING msg_id;+        $QUERY$,+        qtable+    );+    RETURN QUERY EXECUTE sql USING vt, msg, headers;+END;+$$ LANGUAGE plpgsql;++-- send_batch: no delay+CREATE FUNCTION pgmq.send_batch(+    queue_name TEXT,+    msgs JSONB[]+)+RETURNS SETOF BIGINT AS $$+BEGIN+    RETURN QUERY SELECT * FROM pgmq.send_batch(queue_name, msgs, NULL, clock_timestamp());+END;+$$ LANGUAGE plpgsql;++-- send_batch: with delay+CREATE FUNCTION pgmq.send_batch(+    queue_name TEXT,+    msgs JSONB[],+    delay INTEGER+)+RETURNS SETOF BIGINT AS $$+BEGIN+    RETURN QUERY SELECT * FROM pgmq.send_batch(queue_name, msgs, NULL, clock_timestamp() + make_interval(secs => delay));+END;+$$ LANGUAGE plpgsql;++-- send_batch: with headers+CREATE FUNCTION pgmq.send_batch(+    queue_name TEXT,+    msgs JSONB[],+    headers JSONB[]+)+RETURNS SETOF BIGINT AS $$+BEGIN+    RETURN QUERY SELECT * FROM pgmq.send_batch(queue_name, msgs, headers, clock_timestamp());+END;+$$ LANGUAGE plpgsql;++-- send_batch: with headers and delay+CREATE FUNCTION pgmq.send_batch(+    queue_name TEXT,+    msgs JSONB[],+    headers JSONB[],+    delay INTEGER+)+RETURNS SETOF BIGINT AS $$+BEGIN+    RETURN QUERY SELECT * FROM pgmq.send_batch(queue_name, msgs, headers, clock_timestamp() + make_interval(secs => delay));+END;+$$ LANGUAGE plpgsql;++-- send_batch: actual implementation+CREATE FUNCTION pgmq.send_batch(+    queue_name TEXT,+    msgs JSONB[],+    headers JSONB[],+    vt TIMESTAMP WITH TIME ZONE+)+RETURNS SETOF BIGINT AS $$+DECLARE+    sql TEXT;+    qtable TEXT := pgmq.format_table_name(queue_name, 'q');+BEGIN+    sql := FORMAT(+        $QUERY$+        INSERT INTO pgmq.%I (vt, message, headers)+        SELECT $1, unnest($2), unnest(coalesce($3, array_fill(NULL::jsonb, array[cardinality($2)])))+        RETURNING msg_id;+        $QUERY$,+        qtable+    );+    RETURN QUERY EXECUTE sql USING vt, msgs, headers;+END;+$$ LANGUAGE plpgsql;++---- read+---- reads a number of messages from a queue, setting a visibility timeout on them+---- Note: last_read_at is set on read (pgmq 1.10.0+)+CREATE FUNCTION pgmq.read(+    queue_name TEXT,+    vt INTEGER,+    qty INTEGER+)+RETURNS SETOF pgmq.message_record AS $$+DECLARE+    sql TEXT;+    qtable TEXT := pgmq.format_table_name(queue_name, 'q');+BEGIN+    sql := FORMAT(+        $QUERY$+        WITH cte AS+        (+            SELECT msg_id+            FROM pgmq.%I+            WHERE vt <= clock_timestamp()+            ORDER BY msg_id ASC+            LIMIT $1+            FOR UPDATE SKIP LOCKED+        )+        UPDATE pgmq.%I m+        SET+            vt = clock_timestamp() + %L,+            read_ct = read_ct + 1,+            last_read_at = clock_timestamp()+        FROM cte+        WHERE m.msg_id = cte.msg_id+        RETURNING m.msg_id, m.read_ct, m.enqueued_at, m.last_read_at, m.vt, m.message, m.headers;+        $QUERY$,+        qtable, qtable, make_interval(secs => vt)+    );+    RETURN QUERY EXECUTE sql USING qty;+END;+$$ LANGUAGE plpgsql;++---- read+---- reads a number of messages from a queue, setting a visibility timeout on them+CREATE FUNCTION pgmq.read(+    queue_name TEXT,+    vt INTEGER,+    qty INTEGER,+    conditional JSONB+)+RETURNS SETOF pgmq.message_record AS $$+DECLARE+    sql TEXT;+    qtable TEXT := pgmq.format_table_name(queue_name, 'q');+BEGIN+    sql := FORMAT(+        $QUERY$+        WITH cte AS+        (+            SELECT msg_id+            FROM pgmq.%I+            WHERE vt <= clock_timestamp() AND message @> $2+            ORDER BY msg_id ASC+            LIMIT $1+            FOR UPDATE SKIP LOCKED+        )+        UPDATE pgmq.%I m+        SET+            vt = clock_timestamp() + %L,+            read_ct = read_ct + 1,+            last_read_at = clock_timestamp()+        FROM cte+        WHERE m.msg_id = cte.msg_id+        RETURNING m.msg_id, m.read_ct, m.enqueued_at, m.last_read_at, m.vt, m.message, m.headers;+        $QUERY$,+        qtable, qtable, make_interval(secs => vt)+    );+    RETURN QUERY EXECUTE sql USING qty, conditional;+END;+$$ LANGUAGE plpgsql;++---- read_with_poll+---- reads a number of messages from a queue, setting a visibility timeout on them+---- 5-arg version without conditional filter+CREATE FUNCTION pgmq.read_with_poll(+    queue_name TEXT,+    vt INTEGER,+    qty INTEGER,+    max_poll_seconds INTEGER DEFAULT 5,+    poll_interval_ms INTEGER DEFAULT 100+)+RETURNS SETOF pgmq.message_record AS $$+DECLARE+    r pgmq.message_record;+    stop_at TIMESTAMP;+    sql TEXT;+    qtable TEXT := pgmq.format_table_name(queue_name, 'q');+BEGIN+    stop_at := clock_timestamp() + make_interval(secs => max_poll_seconds);+    LOOP+      IF clock_timestamp() >= stop_at THEN+        RETURN;+      END IF;++      sql := FORMAT(+          $QUERY$+          WITH cte AS+          (+              SELECT msg_id+              FROM pgmq.%I+              WHERE vt <= clock_timestamp()+              ORDER BY msg_id ASC+              LIMIT $1+              FOR UPDATE SKIP LOCKED+          )+          UPDATE pgmq.%I m+          SET+              vt = clock_timestamp() + %L,+              read_ct = read_ct + 1,+              last_read_at = clock_timestamp()+          FROM cte+          WHERE m.msg_id = cte.msg_id+          RETURNING m.msg_id, m.read_ct, m.enqueued_at, m.last_read_at, m.vt, m.message, m.headers;+          $QUERY$,+          qtable, qtable, make_interval(secs => vt)+      );++      FOR r IN+        EXECUTE sql USING qty+      LOOP+        RETURN NEXT r;+      END LOOP;+      IF FOUND THEN+        RETURN;+      ELSE+        PERFORM pg_sleep(poll_interval_ms::numeric / 1000);+      END IF;+    END LOOP;+END;+$$ LANGUAGE plpgsql;++---- read_with_poll+---- 6-arg version with conditional as 6th parameter (matches pgmq-hasql expected order)+---- This is the overload used by pgmq-hasql: read_with_poll($1,$2,$3,$4,$5,$6)+CREATE FUNCTION pgmq.read_with_poll(+    queue_name TEXT,+    vt INTEGER,+    qty INTEGER,+    max_poll_seconds INTEGER,+    poll_interval_ms INTEGER,+    conditional JSONB+)+RETURNS SETOF pgmq.message_record AS $$+DECLARE+    r pgmq.message_record;+    stop_at TIMESTAMP;+    sql TEXT;+    qtable TEXT := pgmq.format_table_name(queue_name, 'q');+BEGIN+    stop_at := clock_timestamp() + make_interval(secs => max_poll_seconds);+    LOOP+      IF clock_timestamp() >= stop_at THEN+        RETURN;+      END IF;++      -- If conditional is NULL or empty, don't filter+      IF conditional IS NULL OR conditional = '{}'::jsonb THEN+          sql := FORMAT(+              $QUERY$+              WITH cte AS+              (+                  SELECT msg_id+                  FROM pgmq.%I+                  WHERE vt <= clock_timestamp()+                  ORDER BY msg_id ASC+                  LIMIT $1+                  FOR UPDATE SKIP LOCKED+              )+              UPDATE pgmq.%I m+              SET+                  vt = clock_timestamp() + %L,+                  read_ct = read_ct + 1,+                  last_read_at = clock_timestamp()+              FROM cte+              WHERE m.msg_id = cte.msg_id+              RETURNING m.msg_id, m.read_ct, m.enqueued_at, m.last_read_at, m.vt, m.message, m.headers;+              $QUERY$,+              qtable, qtable, make_interval(secs => vt)+          );++          FOR r IN+            EXECUTE sql USING qty+          LOOP+            RETURN NEXT r;+          END LOOP;+      ELSE+          sql := FORMAT(+              $QUERY$+              WITH cte AS+              (+                  SELECT msg_id+                  FROM pgmq.%I+                  WHERE vt <= clock_timestamp() AND message @> $2+                  ORDER BY msg_id ASC+                  LIMIT $1+                  FOR UPDATE SKIP LOCKED+              )+              UPDATE pgmq.%I m+              SET+                  vt = clock_timestamp() + %L,+                  read_ct = read_ct + 1,+                  last_read_at = clock_timestamp()+              FROM cte+              WHERE m.msg_id = cte.msg_id+              RETURNING m.msg_id, m.read_ct, m.enqueued_at, m.last_read_at, m.vt, m.message, m.headers;+              $QUERY$,+              qtable, qtable, make_interval(secs => vt)+          );++          FOR r IN+            EXECUTE sql USING qty, conditional+          LOOP+            RETURN NEXT r;+          END LOOP;+      END IF;++      IF FOUND THEN+        RETURN;+      ELSE+        PERFORM pg_sleep(poll_interval_ms::numeric / 1000);+      END IF;+    END LOOP;+END;+$$ LANGUAGE plpgsql;++---- read_with_poll+---- 6-arg version with conditional as 4th parameter (legacy order)+CREATE FUNCTION pgmq.read_with_poll(+    queue_name TEXT,+    vt INTEGER,+    qty INTEGER,+    conditional JSONB,+    max_poll_seconds INTEGER DEFAULT 5,+    poll_interval_ms INTEGER DEFAULT 100+)+RETURNS SETOF pgmq.message_record AS $$+DECLARE+    r pgmq.message_record;+    stop_at TIMESTAMP;+    sql TEXT;+    qtable TEXT := pgmq.format_table_name(queue_name, 'q');+BEGIN+    stop_at := clock_timestamp() + make_interval(secs => max_poll_seconds);+    LOOP+      IF clock_timestamp() >= stop_at THEN+        RETURN;+      END IF;++      sql := FORMAT(+          $QUERY$+          WITH cte AS+          (+              SELECT msg_id+              FROM pgmq.%I+              WHERE vt <= clock_timestamp() AND message @> $2+              ORDER BY msg_id ASC+              LIMIT $1+              FOR UPDATE SKIP LOCKED+          )+          UPDATE pgmq.%I m+          SET+              vt = clock_timestamp() + %L,+              read_ct = read_ct + 1,+              last_read_at = clock_timestamp()+          FROM cte+          WHERE m.msg_id = cte.msg_id+          RETURNING m.msg_id, m.read_ct, m.enqueued_at, m.last_read_at, m.vt, m.message, m.headers;+          $QUERY$,+          qtable, qtable, make_interval(secs => vt)+      );++      FOR r IN+        EXECUTE sql USING qty, conditional+      LOOP+        RETURN NEXT r;+      END LOOP;+      IF FOUND THEN+        RETURN;+      ELSE+        PERFORM pg_sleep(poll_interval_ms::numeric / 1000);+      END IF;+    END LOOP;+END;+$$ LANGUAGE plpgsql;++-- delete+-- deletes a message id from the queue permanently+CREATE FUNCTION pgmq.delete(+    queue_name TEXT,+    msg_id BIGINT+)+RETURNS BOOLEAN AS $$+DECLARE+    sql TEXT;+    result BIGINT;+    qtable TEXT := pgmq.format_table_name(queue_name, 'q');+BEGIN+    sql := FORMAT(+        $QUERY$+        DELETE FROM pgmq.%I+        WHERE msg_id = $1+        RETURNING msg_id+        $QUERY$,+        qtable+    );+    EXECUTE sql USING msg_id INTO result;+    RETURN NOT (result IS NULL);+END;+$$ LANGUAGE plpgsql;++-- delete+-- deletes an array of message ids from the queue permanently+CREATE FUNCTION pgmq.delete(+    queue_name TEXT,+    msg_ids BIGINT[]+)+RETURNS SETOF BIGINT AS $$+DECLARE+    sql TEXT;+    qtable TEXT := pgmq.format_table_name(queue_name, 'q');+BEGIN+    sql := FORMAT(+        $QUERY$+        DELETE FROM pgmq.%I+        WHERE msg_id = ANY($1)+        RETURNING msg_id+        $QUERY$,+        qtable+    );+    RETURN QUERY EXECUTE sql USING msg_ids;+END;+$$ LANGUAGE plpgsql;++-- archive+-- removes a message from the queue, and sends it to the archive+-- Note: last_read_at is preserved during archival (pgmq 1.10.0+)+CREATE FUNCTION pgmq.archive(+    queue_name TEXT,+    msg_id BIGINT+)+RETURNS BOOLEAN AS $$+DECLARE+    sql TEXT;+    result BIGINT;+    qtable TEXT := pgmq.format_table_name(queue_name, 'q');+    atable TEXT := pgmq.format_table_name(queue_name, 'a');+BEGIN+    sql := FORMAT(+        $QUERY$+        WITH archived AS (+            DELETE FROM pgmq.%I+            WHERE msg_id = $1+            RETURNING msg_id, vt, read_ct, enqueued_at, last_read_at, message, headers+        )+        INSERT INTO pgmq.%I (msg_id, vt, read_ct, enqueued_at, last_read_at, message, headers)+        SELECT msg_id, vt, read_ct, enqueued_at, last_read_at, message, headers+        FROM archived+        RETURNING msg_id;+        $QUERY$,+        qtable, atable+    );+    EXECUTE sql USING msg_id INTO result;+    RETURN NOT (result IS NULL);+END;+$$ LANGUAGE plpgsql;++-- archive+-- removes an array of message ids from the queue, and sends it to the archive+CREATE FUNCTION pgmq.archive(+    queue_name TEXT,+    msg_ids BIGINT[]+)+RETURNS SETOF BIGINT AS $$+DECLARE+    sql TEXT;+    qtable TEXT := pgmq.format_table_name(queue_name, 'q');+    atable TEXT := pgmq.format_table_name(queue_name, 'a');+BEGIN+    sql := FORMAT(+        $QUERY$+        WITH archived AS (+            DELETE FROM pgmq.%I+            WHERE msg_id = ANY($1)+            RETURNING msg_id, vt, read_ct, enqueued_at, last_read_at, message, headers+        )+        INSERT INTO pgmq.%I (msg_id, vt, read_ct, enqueued_at, last_read_at, message, headers)+        SELECT msg_id, vt, read_ct, enqueued_at, last_read_at, message, headers+        FROM archived+        RETURNING msg_id;+        $QUERY$,+        qtable, atable+    );+    RETURN QUERY EXECUTE sql USING msg_ids;+END;+$$ LANGUAGE plpgsql;++-- pop messages from queue (atomic read + delete)+-- qty parameter added in pgmq 1.7.0+-- Note: Returns message with last_read_at field (pgmq 1.10.0+)+CREATE FUNCTION pgmq.pop(+    queue_name TEXT,+    qty INTEGER DEFAULT 1+)+RETURNS SETOF pgmq.message_record AS $$+DECLARE+    sql TEXT;+    qtable TEXT := pgmq.format_table_name(queue_name, 'q');+BEGIN+    sql := FORMAT(+        $QUERY$+        WITH cte AS+            (+                SELECT msg_id+                FROM pgmq.%I+                WHERE vt <= clock_timestamp()+                ORDER BY msg_id ASC+                LIMIT $1+                FOR UPDATE SKIP LOCKED+            )+        DELETE from pgmq.%I+        WHERE msg_id IN (select msg_id from cte)+        RETURNING *;+        $QUERY$,+        qtable, qtable+    );+    RETURN QUERY EXECUTE sql USING qty;+END;+$$ LANGUAGE plpgsql;++-- Sets vt of a single message, returns it+-- Note: Returns message with last_read_at field (pgmq 1.10.0+)+CREATE FUNCTION pgmq.set_vt(+    queue_name TEXT,+    msg_id BIGINT,+    vt_offset INTEGER+)+RETURNS SETOF pgmq.message_record AS $$+DECLARE+    sql TEXT;+    qtable TEXT := pgmq.format_table_name(queue_name, 'q');+BEGIN+    sql := FORMAT(+        $QUERY$+        UPDATE pgmq.%I+        SET vt = (clock_timestamp() + %L)+        WHERE msg_id = %L+        RETURNING *;+        $QUERY$,+        qtable, make_interval(secs => vt_offset), msg_id+    );+    RETURN QUERY EXECUTE sql;+END;+$$ LANGUAGE plpgsql;++-- Batch set_vt for multiple messages (pgmq 1.8.0+)+CREATE FUNCTION pgmq.set_vt(+    queue_name TEXT,+    msg_ids BIGINT[],+    vt_offset INTEGER+)+RETURNS SETOF pgmq.message_record AS $$+DECLARE+    sql TEXT;+    qtable TEXT := pgmq.format_table_name(queue_name, 'q');+BEGIN+    sql := FORMAT(+        $QUERY$+        UPDATE pgmq.%I+        SET vt = (clock_timestamp() + %L)+        WHERE msg_id = ANY($1)+        RETURNING *;+        $QUERY$,+        qtable, make_interval(secs => vt_offset)+    );+    RETURN QUERY EXECUTE sql USING msg_ids;+END;+$$ LANGUAGE plpgsql;++-- Sets timestamp vt of a message, returns it (pgmq 1.10.0+)+CREATE FUNCTION pgmq.set_vt(queue_name TEXT, msg_id BIGINT, vt TIMESTAMP WITH TIME ZONE)+RETURNS SETOF pgmq.message_record AS $$+DECLARE+    sql TEXT;+    qtable TEXT := pgmq.format_table_name(queue_name, 'q');+BEGIN+    sql := FORMAT(+        $QUERY$+        UPDATE pgmq.%I+        SET vt = $1+        WHERE msg_id = $2+        RETURNING *;+        $QUERY$,+        qtable+    );+    RETURN QUERY EXECUTE sql USING vt, msg_id;+END;+$$ LANGUAGE plpgsql;++-- Sets timestamp vt of multiple messages, returns them (pgmq 1.10.0+)+CREATE FUNCTION pgmq.set_vt(+    queue_name TEXT,+    msg_ids BIGINT[],+    vt TIMESTAMP WITH TIME ZONE+)+RETURNS SETOF pgmq.message_record AS $$+DECLARE+    sql TEXT;+    qtable TEXT := pgmq.format_table_name(queue_name, 'q');+BEGIN+    sql := FORMAT(+        $QUERY$+        UPDATE pgmq.%I+        SET vt = $1+        WHERE msg_id = ANY($2)+        RETURNING *;+        $QUERY$,+        qtable+    );+    RETURN QUERY EXECUTE sql USING vt, msg_ids;+END;+$$ LANGUAGE plpgsql;
+ database/v1.10.0/07_metrics.sql view
@@ -0,0 +1,83 @@+------------------------------------------------------------+-- Metrics functions+------------------------------------------------------------++CREATE FUNCTION pgmq.metrics(queue_name TEXT)+RETURNS pgmq.metrics_result AS $$+DECLARE+    result_row pgmq.metrics_result;+    query TEXT;+    qtable TEXT := pgmq.format_table_name(queue_name, 'q');+BEGIN+    query := FORMAT(+        $QUERY$+        WITH q_summary AS (+            SELECT+                count(*) as queue_length,+                EXTRACT(epoch FROM (NOW() - max(enqueued_at)))::int as newest_msg_age_sec,+                EXTRACT(epoch FROM (NOW() - min(enqueued_at)))::int as oldest_msg_age_sec,+                NOW() as scrape_time,+                count(*) FILTER (WHERE vt <= NOW()) AS queue_visible_length+            FROM pgmq.%I+        ),+        all_metrics AS (+            SELECT CASE+                WHEN is_partitioned THEN+                    %L || '_part_' || pgmq._get_partition_col(%L)+                ELSE+                    %L+            END as queue_name+            FROM pgmq.meta+            WHERE meta.queue_name = %L+        )+        SELECT+            m.queue_name,+            q_summary.queue_length,+            q_summary.newest_msg_age_sec,+            q_summary.oldest_msg_age_sec,+            pgmq._get_pg_stat_get_xact_tuples_inserted(%L)::bigint,+            q_summary.scrape_time,+            q_summary.queue_visible_length+        FROM q_summary+        CROSS JOIN all_metrics m+        $QUERY$,+        qtable, queue_name, queue_name, queue_name, queue_name, 'pgmq.' || qtable+    );+    EXECUTE query INTO result_row;+    RETURN result_row;+END;+$$ LANGUAGE plpgsql;++CREATE FUNCTION pgmq.metrics_all()+RETURNS SETOF pgmq.metrics_result AS $$+DECLARE+    row_name RECORD;+    result_row pgmq.metrics_result;+BEGIN+    FOR row_name IN SELECT queue_name FROM pgmq.meta LOOP+        result_row := pgmq.metrics(row_name.queue_name);+        RETURN NEXT result_row;+    END LOOP;+END;+$$ LANGUAGE plpgsql;++CREATE FUNCTION pgmq._get_pg_stat_get_xact_tuples_inserted(table_name TEXT)+RETURNS bigint AS $$+DECLARE+    result bigint;+BEGIN+    SELECT pg_stat_get_xact_tuples_inserted(oid)+    INTO result+    FROM pg_class+    WHERE relname = table_name;++    IF result IS NULL THEN+        SELECT pg_stat_get_xact_tuples_inserted(oid)+        INTO result+        FROM pg_class+        WHERE relname = SUBSTRING(table_name FROM 6);  -- Remove 'pgmq.' prefix (5 chars + 1 for the dot)+    END IF;++    RETURN COALESCE(result, 0);+END;+$$ LANGUAGE plpgsql;
+ database/v1.10.0/08_partitioning.sql view
@@ -0,0 +1,268 @@+------------------------------------------------------------+-- Partitioning support functions+------------------------------------------------------------++CREATE FUNCTION pgmq._get_pg_partman_schema()+RETURNS TEXT AS $$+    SELECT extnamespace::regnamespace::text+    FROM pg_extension+    WHERE extname = 'pg_partman';+$$ LANGUAGE sql STABLE;++CREATE FUNCTION pgmq._get_partition_col(queue_name TEXT)+RETURNS TEXT AS $$+DECLARE+    _partition_col TEXT;+    _table_name TEXT := pgmq.format_table_name(queue_name, 'q');+BEGIN+    SELECT part_config.partition_interval INTO _partition_col+    FROM pg_class+    JOIN pg_namespace ON pg_class.relnamespace = pg_namespace.oid+    JOIN (SELECT parent_table, partition_interval FROM pgmq._get_partman_config()) AS part_config+      ON part_config.parent_table = pg_namespace.nspname || '.' || pg_class.relname+    WHERE pg_class.relname = _table_name+      AND pg_namespace.nspname = 'pgmq';++    RETURN _partition_col;+END;+$$ LANGUAGE plpgsql;++CREATE FUNCTION pgmq._get_partman_config()+RETURNS TABLE (parent_table text, partition_interval text) AS $$+DECLARE+    schema_name TEXT;+BEGIN+    SELECT pgmq._get_pg_partman_schema() INTO schema_name;+    IF schema_name IS NULL THEN+        RETURN QUERY SELECT NULL::text, NULL::text WHERE FALSE;+        RETURN;+    END IF;+    RETURN QUERY EXECUTE FORMAT(+        $QUERY$+        SELECT parent_table, partition_interval+        FROM %I.part_config+        $QUERY$,+        schema_name+    );+END;+$$ LANGUAGE plpgsql;++-- partitioned, with partition interval+-- Note: last_read_at column added (pgmq 1.10.0+)+CREATE FUNCTION pgmq.create_partitioned(+  queue_name TEXT,+  partition_interval TEXT DEFAULT '10000',+  retention_interval TEXT DEFAULT '100000'+)+RETURNS void AS $$+DECLARE+  qtable TEXT := pgmq.format_table_name(queue_name, 'q');+  qtable_seq TEXT := qtable || '_msg_id_seq';+  atable TEXT := pgmq.format_table_name(queue_name, 'a');+  atable_seq TEXT := atable || '_msg_id_seq';+BEGIN+  PERFORM pgmq.validate_queue_name(queue_name);+  PERFORM pgmq.acquire_queue_lock(queue_name);++  EXECUTE FORMAT(+    $QUERY$+    CREATE TABLE IF NOT EXISTS pgmq.%I (+        msg_id BIGINT GENERATED ALWAYS AS IDENTITY,+        read_ct INT DEFAULT 0 NOT NULL,+        enqueued_at TIMESTAMP WITH TIME ZONE DEFAULT now() NOT NULL,+        last_read_at TIMESTAMP WITH TIME ZONE,+        vt TIMESTAMP WITH TIME ZONE NOT NULL,+        message JSONB,+        headers JSONB+    ) PARTITION BY RANGE (msg_id)+    $QUERY$,+    qtable+  );++  EXECUTE FORMAT(+    $QUERY$+    CREATE TABLE IF NOT EXISTS pgmq.%I (+      msg_id BIGINT NOT NULL,+      read_ct INT DEFAULT 0 NOT NULL,+      enqueued_at TIMESTAMP WITH TIME ZONE DEFAULT now() NOT NULL,+      last_read_at TIMESTAMP WITH TIME ZONE,+      archived_at TIMESTAMP WITH TIME ZONE DEFAULT now() NOT NULL,+      vt TIMESTAMP WITH TIME ZONE NOT NULL,+      message JSONB,+      headers JSONB+    ) PARTITION BY RANGE (msg_id);+    $QUERY$,+    atable+  );++  IF NOT pgmq._extension_exists('pg_partman') THEN+    RAISE EXCEPTION 'pg_partman is required for partitioned queues';+  END IF;++  EXECUTE FORMAT(+    $QUERY$+    SELECT %I.create_parent(+        p_parent_table := 'pgmq.%I',+        p_control := 'msg_id',+        p_type := 'range',+        p_interval := %L,+        p_premake := 1+    )+    $QUERY$,+    pgmq._get_pg_partman_schema(), qtable, partition_interval+  );++  EXECUTE FORMAT(+    $QUERY$+    SELECT %I.create_parent(+        p_parent_table := 'pgmq.%I',+        p_control := 'msg_id',+        p_type := 'range',+        p_interval := %L,+        p_premake := 1+    )+    $QUERY$,+    pgmq._get_pg_partman_schema(), atable, partition_interval+  );++  EXECUTE FORMAT(+    $QUERY$+    UPDATE %I.part_config+    SET+        retention = %L,+        retention_keep_table = false,+        retention_keep_index = true,+        automatic_maintenance = 'on'+    WHERE parent_table = 'pgmq.%I';+    $QUERY$,+    pgmq._get_pg_partman_schema(), retention_interval, qtable+  );++  EXECUTE FORMAT(+    $QUERY$+    UPDATE %I.part_config+    SET+        retention = %L,+        retention_keep_table = false,+        retention_keep_index = true,+        automatic_maintenance = 'on'+    WHERE parent_table = 'pgmq.%I';+    $QUERY$,+    pgmq._get_pg_partman_schema(), retention_interval, atable+  );++  EXECUTE FORMAT(+    $QUERY$+    CREATE INDEX IF NOT EXISTS %I ON pgmq.%I (vt ASC);+    $QUERY$,+    qtable || '_vt_idx', qtable+  );++  EXECUTE FORMAT(+    $QUERY$+    CREATE INDEX IF NOT EXISTS %I ON pgmq.%I (archived_at);+    $QUERY$,+    'archived_at_idx_' || queue_name, atable+  );++  EXECUTE FORMAT(+    $QUERY$+    INSERT INTO pgmq.meta (queue_name, is_partitioned, is_unlogged)+    VALUES (%L, true, false)+    ON CONFLICT+    DO NOTHING;+    $QUERY$,+    queue_name+  );++END;+$$ LANGUAGE plpgsql;++-- convert_archive_partitioned+-- Note: last_read_at column added (pgmq 1.10.0+)+CREATE FUNCTION pgmq.convert_archive_partitioned(+  queue_name TEXT,+  partition_interval TEXT DEFAULT '10000',+  retention_interval TEXT DEFAULT '100000',+  leading_partition INT DEFAULT 10+)+RETURNS void AS $$+DECLARE+  atable TEXT := pgmq.format_table_name(queue_name, 'a');+BEGIN+  EXECUTE FORMAT(+    $QUERY$+    ALTER TABLE IF EXISTS pgmq.%I RENAME TO %I+    $QUERY$,+    atable, atable || '_old'+  );++  EXECUTE FORMAT(+    $QUERY$+    CREATE TABLE IF NOT EXISTS pgmq.%I (+      msg_id BIGINT NOT NULL,+      read_ct INT DEFAULT 0 NOT NULL,+      enqueued_at TIMESTAMP WITH TIME ZONE DEFAULT now() NOT NULL,+      last_read_at TIMESTAMP WITH TIME ZONE,+      archived_at TIMESTAMP WITH TIME ZONE DEFAULT now() NOT NULL,+      vt TIMESTAMP WITH TIME ZONE NOT NULL,+      message JSONB,+      headers JSONB+    ) PARTITION BY RANGE (msg_id);+    $QUERY$,+    atable+  );++  IF NOT pgmq._extension_exists('pg_partman') THEN+    RAISE EXCEPTION 'pg_partman is required for partitioned queues';+  END IF;++  EXECUTE FORMAT(+    $QUERY$+    SELECT %I.create_parent(+        p_parent_table := 'pgmq.%I',+        p_control := 'msg_id',+        p_type := 'range',+        p_interval := %L,+        p_premake := %L+    )+    $QUERY$,+    pgmq._get_pg_partman_schema(), atable, partition_interval, leading_partition+  );++  EXECUTE FORMAT(+    $QUERY$+    UPDATE %I.part_config+    SET+        retention = %L,+        retention_keep_table = false,+        retention_keep_index = true,+        automatic_maintenance = 'on'+    WHERE parent_table = 'pgmq.%I';+    $QUERY$,+    pgmq._get_pg_partman_schema(), retention_interval, atable+  );++  EXECUTE FORMAT(+    $QUERY$+    INSERT INTO pgmq.%I (msg_id, read_ct, enqueued_at, last_read_at, archived_at, vt, message, headers)+    SELECT msg_id, read_ct, enqueued_at, last_read_at, archived_at, vt, message, headers FROM pgmq.%I+    $QUERY$,+    atable, atable || '_old'+  );++  EXECUTE FORMAT(+    $QUERY$+    CREATE INDEX IF NOT EXISTS %I ON pgmq.%I (archived_at);+    $QUERY$,+    'archived_at_idx_' || queue_name, atable+  );++  EXECUTE FORMAT(+    $QUERY$+    DROP TABLE IF EXISTS pgmq.%I+    $QUERY$,+    atable || '_old'+  );+END;+$$ LANGUAGE plpgsql;
+ database/v1.10.0/09_notifications.sql view
@@ -0,0 +1,103 @@+------------------------------------------------------------+-- Notification functions+------------------------------------------------------------++CREATE FUNCTION pgmq._should_notify(queue_name TEXT)+RETURNS BOOLEAN AS $$+DECLARE+    should_notify BOOLEAN;+    throttle_ms INTEGER;+    last_notified TIMESTAMP WITH TIME ZONE;+BEGIN+    -- Try to get throttle settings+    SELECT throttle_interval_ms, last_notified_at+    INTO throttle_ms, last_notified+    FROM pgmq.notify_insert_throttle+    WHERE pgmq.notify_insert_throttle.queue_name = _should_notify.queue_name+    FOR UPDATE SKIP LOCKED;++    -- If no row found, assume notifications are enabled without throttling+    IF throttle_ms IS NULL THEN+        RETURN TRUE;+    END IF;++    -- Check if enough time has passed since last notification+    IF clock_timestamp() >= last_notified + make_interval(secs => throttle_ms::numeric / 1000) THEN+        -- Update the last notification time+        UPDATE pgmq.notify_insert_throttle+        SET last_notified_at = clock_timestamp()+        WHERE pgmq.notify_insert_throttle.queue_name = _should_notify.queue_name;+        RETURN TRUE;+    END IF;++    RETURN FALSE;+END;+$$ LANGUAGE plpgsql;++CREATE FUNCTION pgmq._notify_on_insert()+RETURNS TRIGGER AS $$+BEGIN+    IF pgmq._should_notify(TG_ARGV[0]) THEN+        PERFORM pg_notify(TG_ARGV[0], NEW.msg_id::text);+    END IF;+    RETURN NEW;+END;+$$ LANGUAGE plpgsql;++-- Add trigger to queue table to notify on insert+CREATE FUNCTION pgmq.enable_queue_notifications(queue_name TEXT)+RETURNS void AS $$+DECLARE+    qtable TEXT := pgmq.format_table_name(queue_name, 'q');+    channel TEXT := queue_name;+    trigger_name TEXT := 'notify_insert_' || qtable;+BEGIN+    -- Create trigger on the queue table+    EXECUTE FORMAT(+        $QUERY$+        CREATE TRIGGER %I+        AFTER INSERT ON pgmq.%I+        FOR EACH ROW+        EXECUTE FUNCTION pgmq._notify_on_insert(%L)+        $QUERY$,+        trigger_name, qtable, channel+    );++    -- Insert throttle tracking row if it doesn't exist+    INSERT INTO pgmq.notify_insert_throttle (queue_name, throttle_interval_ms, last_notified_at)+    VALUES (queue_name, 0, to_timestamp(0))+    ON CONFLICT (queue_name) DO NOTHING;+END;+$$ LANGUAGE plpgsql;++-- Remove trigger from queue table+CREATE FUNCTION pgmq.disable_queue_notifications(queue_name TEXT)+RETURNS void AS $$+DECLARE+    qtable TEXT := pgmq.format_table_name(queue_name, 'q');+    trigger_name TEXT := 'notify_insert_' || qtable;+BEGIN+    -- Drop trigger if it exists+    EXECUTE FORMAT(+        $QUERY$+        DROP TRIGGER IF EXISTS %I ON pgmq.%I+        $QUERY$,+        trigger_name, qtable+    );++    -- Remove throttle tracking row+    DELETE FROM pgmq.notify_insert_throttle+    WHERE pgmq.notify_insert_throttle.queue_name = disable_queue_notifications.queue_name;+END;+$$ LANGUAGE plpgsql;++-- Set throttle interval for queue notifications+CREATE FUNCTION pgmq.set_notification_throttle(queue_name TEXT, throttle_interval_ms INTEGER)+RETURNS void AS $$+BEGIN+    INSERT INTO pgmq.notify_insert_throttle (queue_name, throttle_interval_ms, last_notified_at)+    VALUES (queue_name, throttle_interval_ms, to_timestamp(0))+    ON CONFLICT (queue_name) DO UPDATE+    SET throttle_interval_ms = EXCLUDED.throttle_interval_ms;+END;+$$ LANGUAGE plpgsql;
+ database/v1.10.0/10_fifo.sql view
@@ -0,0 +1,571 @@+------------------------------------------------------------+-- FIFO read functions (for pgmq 1.8.0+ and 1.9.0+)+-- Note: last_read_at handling added in pgmq 1.10.0+------------------------------------------------------------++-- read_fifo: reads messages in strict FIFO order+-- This function ensures messages are read in strict FIFO order by:+-- 1. Ordering by msg_id ASC (not just using SKIP LOCKED)+-- 2. Only returning messages that don't skip any earlier pending messages+CREATE FUNCTION pgmq.read_fifo(+    queue_name TEXT,+    vt INTEGER,+    qty INTEGER+)+RETURNS SETOF pgmq.message_record AS $$+DECLARE+    sql TEXT;+    qtable TEXT := pgmq.format_table_name(queue_name, 'q');+BEGIN+    sql := FORMAT(+        $QUERY$+        WITH cte AS+        (+            SELECT msg_id, vt AS current_vt+            FROM pgmq.%I+            ORDER BY msg_id ASC+            LIMIT $1+            FOR UPDATE SKIP LOCKED+        ),+        eligible AS (+            SELECT c.msg_id+            FROM cte c+            WHERE c.msg_id = (+                SELECT MIN(msg_id)+                FROM pgmq.%I+                WHERE vt <= clock_timestamp()+            )+            OR NOT EXISTS (+                SELECT 1+                FROM pgmq.%I q+                WHERE q.msg_id < c.msg_id+                AND q.vt <= clock_timestamp()+                AND q.msg_id NOT IN (SELECT msg_id FROM cte)+            )+        )+        UPDATE pgmq.%I m+        SET+            vt = clock_timestamp() + %L,+            read_ct = read_ct + 1,+            last_read_at = clock_timestamp()+        FROM eligible e+        WHERE m.msg_id = e.msg_id+        AND m.vt <= clock_timestamp()+        RETURNING m.msg_id, m.read_ct, m.enqueued_at, m.last_read_at, m.vt, m.message, m.headers;+        $QUERY$,+        qtable, qtable, qtable, qtable, make_interval(secs => vt)+    );+    RETURN QUERY EXECUTE sql USING qty;+END;+$$ LANGUAGE plpgsql;++-- read_fifo with conditional filter+CREATE FUNCTION pgmq.read_fifo(+    queue_name TEXT,+    vt INTEGER,+    qty INTEGER,+    conditional JSONB+)+RETURNS SETOF pgmq.message_record AS $$+DECLARE+    sql TEXT;+    qtable TEXT := pgmq.format_table_name(queue_name, 'q');+BEGIN+    sql := FORMAT(+        $QUERY$+        WITH cte AS+        (+            SELECT msg_id, vt AS current_vt+            FROM pgmq.%I+            WHERE message @> $2+            ORDER BY msg_id ASC+            LIMIT $1+            FOR UPDATE SKIP LOCKED+        ),+        eligible AS (+            SELECT c.msg_id+            FROM cte c+            WHERE c.msg_id = (+                SELECT MIN(msg_id)+                FROM pgmq.%I+                WHERE vt <= clock_timestamp()+                AND message @> $2+            )+            OR NOT EXISTS (+                SELECT 1+                FROM pgmq.%I q+                WHERE q.msg_id < c.msg_id+                AND q.vt <= clock_timestamp()+                AND q.message @> $2+                AND q.msg_id NOT IN (SELECT msg_id FROM cte)+            )+        )+        UPDATE pgmq.%I m+        SET+            vt = clock_timestamp() + %L,+            read_ct = read_ct + 1,+            last_read_at = clock_timestamp()+        FROM eligible e+        WHERE m.msg_id = e.msg_id+        AND m.vt <= clock_timestamp()+        RETURNING m.msg_id, m.read_ct, m.enqueued_at, m.last_read_at, m.vt, m.message, m.headers;+        $QUERY$,+        qtable, qtable, qtable, qtable, make_interval(secs => vt)+    );+    RETURN QUERY EXECUTE sql USING qty, conditional;+END;+$$ LANGUAGE plpgsql;++-- read_fifo_with_poll: reads messages in strict FIFO order with polling+CREATE FUNCTION pgmq.read_fifo_with_poll(+    queue_name TEXT,+    vt INTEGER,+    qty INTEGER,+    max_poll_seconds INTEGER DEFAULT 5,+    poll_interval_ms INTEGER DEFAULT 100+)+RETURNS SETOF pgmq.message_record AS $$+DECLARE+    r pgmq.message_record;+    stop_at TIMESTAMP;+    sql TEXT;+    qtable TEXT := pgmq.format_table_name(queue_name, 'q');+BEGIN+    stop_at := clock_timestamp() + make_interval(secs => max_poll_seconds);+    LOOP+      IF clock_timestamp() >= stop_at THEN+        RETURN;+      END IF;++      sql := FORMAT(+          $QUERY$+          WITH cte AS+          (+              SELECT msg_id, vt AS current_vt+              FROM pgmq.%I+              ORDER BY msg_id ASC+              LIMIT $1+              FOR UPDATE SKIP LOCKED+          ),+          eligible AS (+              SELECT c.msg_id+              FROM cte c+              WHERE c.msg_id = (+                  SELECT MIN(msg_id)+                  FROM pgmq.%I+                  WHERE vt <= clock_timestamp()+              )+              OR NOT EXISTS (+                  SELECT 1+                  FROM pgmq.%I q+                  WHERE q.msg_id < c.msg_id+                  AND q.vt <= clock_timestamp()+                  AND q.msg_id NOT IN (SELECT msg_id FROM cte)+              )+          )+          UPDATE pgmq.%I m+          SET+              vt = clock_timestamp() + %L,+              read_ct = read_ct + 1,+              last_read_at = clock_timestamp()+          FROM eligible e+          WHERE m.msg_id = e.msg_id+          AND m.vt <= clock_timestamp()+          RETURNING m.msg_id, m.read_ct, m.enqueued_at, m.last_read_at, m.vt, m.message, m.headers;+          $QUERY$,+          qtable, qtable, qtable, qtable, make_interval(secs => vt)+      );++      FOR r IN+        EXECUTE sql USING qty+      LOOP+        RETURN NEXT r;+      END LOOP;+      IF FOUND THEN+        RETURN;+      ELSE+        PERFORM pg_sleep(poll_interval_ms::numeric / 1000);+      END IF;+    END LOOP;+END;+$$ LANGUAGE plpgsql;++-- read_fifo_with_poll with conditional filter+CREATE FUNCTION pgmq.read_fifo_with_poll(+    queue_name TEXT,+    vt INTEGER,+    qty INTEGER,+    conditional JSONB,+    max_poll_seconds INTEGER DEFAULT 5,+    poll_interval_ms INTEGER DEFAULT 100+)+RETURNS SETOF pgmq.message_record AS $$+DECLARE+    r pgmq.message_record;+    stop_at TIMESTAMP;+    sql TEXT;+    qtable TEXT := pgmq.format_table_name(queue_name, 'q');+BEGIN+    stop_at := clock_timestamp() + make_interval(secs => max_poll_seconds);+    LOOP+      IF clock_timestamp() >= stop_at THEN+        RETURN;+      END IF;++      sql := FORMAT(+          $QUERY$+          WITH cte AS+          (+              SELECT msg_id, vt AS current_vt+              FROM pgmq.%I+              WHERE message @> $2+              ORDER BY msg_id ASC+              LIMIT $1+              FOR UPDATE SKIP LOCKED+          ),+          eligible AS (+              SELECT c.msg_id+              FROM cte c+              WHERE c.msg_id = (+                  SELECT MIN(msg_id)+                  FROM pgmq.%I+                  WHERE vt <= clock_timestamp()+                  AND message @> $2+              )+              OR NOT EXISTS (+                  SELECT 1+                  FROM pgmq.%I q+                  WHERE q.msg_id < c.msg_id+                  AND q.vt <= clock_timestamp()+                  AND q.message @> $2+                  AND q.msg_id NOT IN (SELECT msg_id FROM cte)+              )+          )+          UPDATE pgmq.%I m+          SET+              vt = clock_timestamp() + %L,+              read_ct = read_ct + 1,+              last_read_at = clock_timestamp()+          FROM eligible e+          WHERE m.msg_id = e.msg_id+          AND m.vt <= clock_timestamp()+          RETURNING m.msg_id, m.read_ct, m.enqueued_at, m.last_read_at, m.vt, m.message, m.headers;+          $QUERY$,+          qtable, qtable, qtable, qtable, make_interval(secs => vt)+      );++      FOR r IN+        EXECUTE sql USING qty, conditional+      LOOP+        RETURN NEXT r;+      END LOOP;+      IF FOUND THEN+        RETURN;+      ELSE+        PERFORM pg_sleep(poll_interval_ms::numeric / 1000);+      END IF;+    END LOOP;+END;+$$ LANGUAGE plpgsql;++------------------------------------------------------------+-- FIFO index functions (pgmq 1.8.0+)+------------------------------------------------------------++-- _create_fifo_index_if_not_exists+-- internal function to create GIN index on headers for better FIFO performance+CREATE OR REPLACE FUNCTION pgmq._create_fifo_index_if_not_exists(queue_name TEXT)+RETURNS void AS $$+DECLARE+    qtable TEXT := pgmq.format_table_name(queue_name, 'q');+    index_name TEXT := qtable || '_fifo_idx';+BEGIN+    -- Create GIN index on headers for efficient FIFO key lookups+    EXECUTE FORMAT(+        $QUERY$+        CREATE INDEX IF NOT EXISTS %I ON pgmq.%I USING GIN (headers);+        $QUERY$,+        index_name, qtable+    );+END;+$$ LANGUAGE plpgsql;++-- create_fifo_index+-- creates a GIN index on the headers column to improve FIFO read performance+CREATE FUNCTION pgmq.create_fifo_index(queue_name TEXT)+RETURNS void AS $$+BEGIN+    PERFORM pgmq._create_fifo_index_if_not_exists(queue_name);+END;+$$ LANGUAGE plpgsql;++-- create_fifo_indexes_all+-- creates FIFO indexes on all existing queues+CREATE FUNCTION pgmq.create_fifo_indexes_all()+RETURNS void AS $$+DECLARE+    queue_record RECORD;+BEGIN+    FOR queue_record IN SELECT queue_name FROM pgmq.meta LOOP+        PERFORM pgmq.create_fifo_index(queue_record.queue_name);+    END LOOP;+END;+$$ LANGUAGE plpgsql;++------------------------------------------------------------+-- read_grouped functions (pgmq 1.8.0+)+-- AWS SQS FIFO-style batch retrieval behavior+------------------------------------------------------------++-- read_grouped+-- reads messages with AWS SQS FIFO-style batch retrieval behavior+-- attempts to return as many messages as possible from the same message group+CREATE FUNCTION pgmq.read_grouped(+    queue_name TEXT,+    vt INTEGER,+    qty INTEGER+)+RETURNS SETOF pgmq.message_record AS $$+DECLARE+    sql TEXT;+    qtable TEXT := pgmq.format_table_name(queue_name, 'q');+BEGIN+    sql := FORMAT(+        $QUERY$+        WITH fifo_groups AS (+            -- Find the minimum msg_id for each FIFO group that's ready to be processed+            SELECT+                COALESCE(headers->>'x-pgmq-group', '_default_fifo_group') as fifo_key,+                MIN(msg_id) as min_msg_id+            FROM pgmq.%I+            WHERE vt <= clock_timestamp()+            GROUP BY COALESCE(headers->>'x-pgmq-group', '_default_fifo_group')+        ),+        locked_groups AS (+            -- Lock the first available message in each FIFO group+            SELECT+                m.msg_id,+                fg.fifo_key+            FROM pgmq.%I m+            INNER JOIN fifo_groups fg ON+                COALESCE(m.headers->>'x-pgmq-group', '_default_fifo_group') = fg.fifo_key+                AND m.msg_id = fg.min_msg_id+            WHERE m.vt <= clock_timestamp()+            ORDER BY m.msg_id ASC+            FOR UPDATE SKIP LOCKED+        ),+        group_priorities AS (+            -- Assign priority to groups based on their oldest message+            SELECT+                fifo_key,+                msg_id as min_msg_id,+                ROW_NUMBER() OVER (ORDER BY msg_id) as group_priority+            FROM locked_groups+        ),+        available_messages AS (+            -- Get messages prioritizing filling batch from earliest group first+            SELECT+                m.msg_id,+                gp.group_priority,+                ROW_NUMBER() OVER (PARTITION BY gp.fifo_key ORDER BY m.msg_id) as msg_rank_in_group+            FROM pgmq.%I m+            INNER JOIN group_priorities gp ON+                COALESCE(m.headers->>'x-pgmq-group', '_default_fifo_group') = gp.fifo_key+            WHERE m.vt <= clock_timestamp()+            AND m.msg_id >= gp.min_msg_id  -- Only messages from min_msg_id onwards in each group+            AND NOT EXISTS (+                -- Ensure no earlier message in this group is currently being processed+                SELECT 1+                FROM pgmq.%I m2+                WHERE COALESCE(m2.headers->>'x-pgmq-group', '_default_fifo_group') =+                      COALESCE(m.headers->>'x-pgmq-group', '_default_fifo_group')+                AND m2.vt > clock_timestamp()+                AND m2.msg_id < m.msg_id+            )+        ),+        batch_selection AS (+            -- Select messages to fill batch, prioritizing earliest group+            SELECT+                msg_id,+                ROW_NUMBER() OVER (ORDER BY group_priority, msg_rank_in_group) as overall_rank+            FROM available_messages+        ),+        selected_messages AS (+            -- Limit to requested quantity+            SELECT msg_id+            FROM batch_selection+            WHERE overall_rank <= $1+            ORDER BY msg_id+            FOR UPDATE SKIP LOCKED+        )+        UPDATE pgmq.%I m+        SET+            vt = clock_timestamp() + %L,+            read_ct = read_ct + 1,+            last_read_at = clock_timestamp()+        FROM selected_messages sm+        WHERE m.msg_id = sm.msg_id+        RETURNING m.msg_id, m.read_ct, m.enqueued_at, m.last_read_at, m.vt, m.message, m.headers;+        $QUERY$,+        qtable, qtable, qtable, qtable, qtable, make_interval(secs => vt)+    );+    RETURN QUERY EXECUTE sql USING qty;+END;+$$ LANGUAGE plpgsql;++-- read_grouped_with_poll+-- reads messages with AWS SQS FIFO-style batch retrieval behavior, with polling support+CREATE FUNCTION pgmq.read_grouped_with_poll(+    queue_name TEXT,+    vt INTEGER,+    qty INTEGER,+    max_poll_seconds INTEGER DEFAULT 5,+    poll_interval_ms INTEGER DEFAULT 100+)+RETURNS SETOF pgmq.message_record AS $$+DECLARE+    r pgmq.message_record;+    stop_at TIMESTAMP;+BEGIN+    stop_at := clock_timestamp() + make_interval(secs => max_poll_seconds);+    LOOP+      IF (SELECT clock_timestamp() >= stop_at) THEN+        RETURN;+      END IF;++      FOR r IN+        SELECT * FROM pgmq.read_grouped(queue_name, vt, qty)+      LOOP+        RETURN NEXT r;+      END LOOP;+      IF FOUND THEN+        RETURN;+      ELSE+        PERFORM pg_sleep(poll_interval_ms::numeric / 1000);+      END IF;+    END LOOP;+END;+$$ LANGUAGE plpgsql;++------------------------------------------------------------+-- Round-robin FIFO functions (pgmq 1.9.0+)+------------------------------------------------------------++-- read_grouped_rr (round-robin)+-- reads messages while preserving FIFO within groups and interleaving across groups (layered round-robin)+CREATE FUNCTION pgmq.read_grouped_rr(+    queue_name TEXT,+    vt INTEGER,+    qty INTEGER+)+RETURNS SETOF pgmq.message_record AS $$+DECLARE+    sql TEXT;+    qtable TEXT := pgmq.format_table_name(queue_name, 'q');+BEGIN+    sql := FORMAT(+        $QUERY$+        WITH fifo_groups AS (+            -- Determine the absolute head (oldest) message id per FIFO group, regardless of visibility+            SELECT+                COALESCE(headers->>'x-pgmq-group', '_default_fifo_group') AS fifo_key,+                MIN(msg_id) AS head_msg_id+            FROM pgmq.%1$I+            GROUP BY COALESCE(headers->>'x-pgmq-group', '_default_fifo_group')+        ),+        eligible_groups AS (+            -- Only groups whose head message is currently visible+            -- Acquire a transaction-level advisory lock per group to prevent concurrent selection+            SELECT+                g.fifo_key,+                g.head_msg_id,+                ROW_NUMBER() OVER (ORDER BY g.head_msg_id) AS group_priority+            FROM fifo_groups g+            JOIN pgmq.%2$I h ON h.msg_id = g.head_msg_id+            WHERE h.vt <= clock_timestamp()+              AND pg_try_advisory_xact_lock(pg_catalog.hashtextextended(g.fifo_key, 0))+        ),+        available_messages AS (+            -- All currently visible messages starting at the head for each eligible group+            SELECT+                m.msg_id,+                eg.group_priority,+                ROW_NUMBER() OVER (+                    PARTITION BY eg.fifo_key+                    ORDER BY m.msg_id+                ) AS msg_rank_in_group+            FROM pgmq.%3$I m+            JOIN eligible_groups eg+              ON COALESCE(m.headers->>'x-pgmq-group', '_default_fifo_group') = eg.fifo_key+            WHERE m.vt <= clock_timestamp()+              AND m.msg_id >= eg.head_msg_id+        ),+        ordered_messages AS (+            -- Layered round-robin: take rank 1 of all groups by group_priority, then rank 2, etc.+            -- Assign selection order before locking+            SELECT msg_id, ROW_NUMBER() OVER (ORDER BY msg_rank_in_group, group_priority) as selection_order+            FROM available_messages+        ),+        selected_messages AS (+            -- Lock the messages in the correct order, preserving selection_order+            SELECT om.msg_id, om.selection_order+            FROM ordered_messages om+            JOIN pgmq.%4$I m ON m.msg_id = om.msg_id+            WHERE om.selection_order <= $1+            ORDER BY om.selection_order+            FOR UPDATE OF m SKIP LOCKED+        ),+        updated_messages AS (+            UPDATE pgmq.%5$I m+            SET+                vt = clock_timestamp() + %6$L,+                read_ct = read_ct + 1,+                last_read_at = clock_timestamp()+            FROM selected_messages sm+            WHERE m.msg_id = sm.msg_id+              AND m.vt <= clock_timestamp() -- final guard to avoid duplicate reads under races+            RETURNING m.msg_id, m.read_ct, m.enqueued_at, m.last_read_at, m.vt, m.message, m.headers, sm.selection_order+        )+        SELECT msg_id, read_ct, enqueued_at, last_read_at, vt, message, headers+        FROM updated_messages+        ORDER BY selection_order;+        $QUERY$,+        qtable, qtable, qtable, qtable, qtable, make_interval(secs => vt)+    );+    RETURN QUERY EXECUTE sql USING qty;+END;+$$ LANGUAGE plpgsql;++-- read_grouped_rr_with_poll+-- reads messages using round-robin layering across groups, with polling support+CREATE FUNCTION pgmq.read_grouped_rr_with_poll(+    queue_name TEXT,+    vt INTEGER,+    qty INTEGER,+    max_poll_seconds INTEGER DEFAULT 5,+    poll_interval_ms INTEGER DEFAULT 100+)+RETURNS SETOF pgmq.message_record AS $$+DECLARE+    r pgmq.message_record;+    stop_at TIMESTAMP;+BEGIN+    stop_at := clock_timestamp() + make_interval(secs => max_poll_seconds);+    LOOP+      IF (SELECT clock_timestamp() >= stop_at) THEN+        RETURN;+      END IF;++      FOR r IN+        SELECT * FROM pgmq.read_grouped_rr(queue_name, vt, qty)+      LOOP+        RETURN NEXT r;+      END LOOP;+      IF FOUND THEN+        RETURN;+      ELSE+        PERFORM pg_sleep(poll_interval_ms::numeric / 1000);+      END IF;+    END LOOP;+END;+$$ LANGUAGE plpgsql;
+ database/v1.9.0/01_schema.sql view
@@ -0,0 +1,11 @@+------------------------------------------------------------+-- Schema creation and grants (without extension)+------------------------------------------------------------+CREATE SCHEMA IF NOT EXISTS pgmq;++-- Grant permission to pg_monitor to all tables and sequences+GRANT USAGE ON SCHEMA pgmq TO pg_monitor;+GRANT SELECT ON ALL TABLES IN SCHEMA pgmq TO pg_monitor;+GRANT SELECT ON ALL SEQUENCES IN SCHEMA pgmq TO pg_monitor;+ALTER DEFAULT PRIVILEGES IN SCHEMA pgmq GRANT SELECT ON TABLES TO pg_monitor;+ALTER DEFAULT PRIVILEGES IN SCHEMA pgmq GRANT SELECT ON SEQUENCES TO pg_monitor;
+ database/v1.9.0/02_tables.sql view
@@ -0,0 +1,25 @@+------------------------------------------------------------+-- Tables: meta and notify_insert_throttle+------------------------------------------------------------++-- Table where queues and metadata about them is stored+CREATE TABLE IF NOT EXISTS pgmq.meta (+    queue_name VARCHAR UNIQUE NOT NULL,+    is_partitioned BOOLEAN NOT NULL,+    is_unlogged BOOLEAN NOT NULL,+    created_at TIMESTAMP WITH TIME ZONE DEFAULT now() NOT NULL+);++-- Table to track notification throttling for queues+CREATE UNLOGGED TABLE IF NOT EXISTS pgmq.notify_insert_throttle (+    queue_name           VARCHAR UNIQUE NOT NULL -- Queue name (without 'q_' prefix)+       CONSTRAINT notify_insert_throttle_meta_queue_name_fk+            REFERENCES pgmq.meta (queue_name)+            ON DELETE CASCADE,+    throttle_interval_ms INTEGER NOT NULL DEFAULT 0, -- Min milliseconds between notifications (0 = no throttling)+    last_notified_at     TIMESTAMP WITH TIME ZONE NOT NULL DEFAULT to_timestamp(0) -- Timestamp of last sent notification+);++CREATE INDEX IF NOT EXISTS idx_notify_throttle_active+    ON pgmq.notify_insert_throttle (queue_name, last_notified_at)+    WHERE throttle_interval_ms > 0;
+ database/v1.9.0/03_types.sql view
@@ -0,0 +1,32 @@+------------------------------------------------------------+-- Custom types: message_record, queue_record, metrics_result+------------------------------------------------------------++-- This type has the shape of a message in a queue, and is often returned by+-- pgmq functions that return messages+CREATE TYPE pgmq.message_record AS (+    msg_id BIGINT,+    read_ct INTEGER,+    enqueued_at TIMESTAMP WITH TIME ZONE,+    vt TIMESTAMP WITH TIME ZONE,+    message JSONB,+    headers JSONB+);++CREATE TYPE pgmq.queue_record AS (+    queue_name VARCHAR,+    is_partitioned BOOLEAN,+    is_unlogged BOOLEAN,+    created_at TIMESTAMP WITH TIME ZONE+);++-- returned by pgmq.metrics() and pgmq.metrics_all+CREATE TYPE pgmq.metrics_result AS (+    queue_name text,+    queue_length bigint,+    newest_msg_age_sec int,+    oldest_msg_age_sec int,+    total_messages bigint,+    scrape_time timestamp with time zone,+    queue_visible_length bigint+);
+ database/v1.9.0/04_core_functions.sql view
@@ -0,0 +1,69 @@+------------------------------------------------------------+-- Core utility functions+------------------------------------------------------------++-- prevents race conditions during queue creation by acquiring a transaction-level advisory lock+-- uses a transaction advisory lock maintain the lock until transaction commit+-- a race condition would still exist if lock was released before commit+CREATE FUNCTION pgmq.acquire_queue_lock(queue_name TEXT)+RETURNS void AS $$+BEGIN+  PERFORM pg_advisory_xact_lock(hashtext('pgmq.queue_' || queue_name));+END;+$$ LANGUAGE plpgsql;++-- a helper to format table names and check for invalid characters+CREATE FUNCTION pgmq.format_table_name(queue_name text, prefix text)+RETURNS TEXT AS $$+BEGIN+    IF queue_name ~ '\$|;|--|'''+    THEN+        RAISE EXCEPTION 'queue name contains invalid characters: $, ;, --, or \''';+    END IF;+    RETURN lower(prefix || '_' || queue_name);+END;+$$ LANGUAGE plpgsql;++CREATE FUNCTION pgmq.validate_queue_name(queue_name TEXT)+RETURNS void AS $$+BEGIN+  IF length(queue_name) > 47 THEN+    -- complete table identifier must be <= 63+    -- https://www.postgresql.org/docs/17/sql-syntax-lexical.html#SQL-SYNTAX-IDENTIFIERS+    -- e.g. template_pgmq_q_my_queue is an identifier for my_queue when partitioned+    -- template_pgmq_q_ (16) + <a max length queue name> (47) = 63+    RAISE EXCEPTION 'queue name is too long, maximum length is 47 characters';+  END IF;+END;+$$ LANGUAGE plpgsql;++CREATE FUNCTION pgmq._belongs_to_pgmq(table_name TEXT)+RETURNS BOOLEAN AS $$+DECLARE+    sql TEXT;+    result BOOLEAN;+BEGIN+  SELECT EXISTS (+    SELECT 1+    FROM pg_depend+    WHERE refobjid = (SELECT oid FROM pg_extension WHERE extname = 'pgmq')+    AND objid = (+        SELECT oid+        FROM pg_class+        WHERE relname = table_name+    )+  ) INTO result;+  RETURN result;+END;+$$ LANGUAGE plpgsql;++CREATE FUNCTION pgmq._extension_exists(extension_name TEXT)+    RETURNS BOOLEAN+    LANGUAGE SQL+AS $$+SELECT EXISTS (+    SELECT 1+    FROM pg_extension+    WHERE extname = extension_name+)+$$;
+ database/v1.9.0/05_queue_management.sql view
@@ -0,0 +1,261 @@+------------------------------------------------------------+-- Queue management functions+------------------------------------------------------------++CREATE FUNCTION pgmq.create_non_partitioned(queue_name TEXT)+RETURNS void AS $$+DECLARE+  qtable TEXT := pgmq.format_table_name(queue_name, 'q');+  qtable_seq TEXT := qtable || '_msg_id_seq';+  atable TEXT := pgmq.format_table_name(queue_name, 'a');+BEGIN+  PERFORM pgmq.validate_queue_name(queue_name);+  PERFORM pgmq.acquire_queue_lock(queue_name);++  EXECUTE FORMAT(+    $QUERY$+    CREATE TABLE IF NOT EXISTS pgmq.%I (+        msg_id BIGINT PRIMARY KEY GENERATED ALWAYS AS IDENTITY,+        read_ct INT DEFAULT 0 NOT NULL,+        enqueued_at TIMESTAMP WITH TIME ZONE DEFAULT now() NOT NULL,+        vt TIMESTAMP WITH TIME ZONE NOT NULL,+        message JSONB,+        headers JSONB+    )+    $QUERY$,+    qtable+  );++  EXECUTE FORMAT(+    $QUERY$+    CREATE TABLE IF NOT EXISTS pgmq.%I (+      msg_id BIGINT PRIMARY KEY,+      read_ct INT DEFAULT 0 NOT NULL,+      enqueued_at TIMESTAMP WITH TIME ZONE DEFAULT now() NOT NULL,+      archived_at TIMESTAMP WITH TIME ZONE DEFAULT now() NOT NULL,+      vt TIMESTAMP WITH TIME ZONE NOT NULL,+      message JSONB,+      headers JSONB+    );+    $QUERY$,+    atable+  );++  EXECUTE FORMAT(+    $QUERY$+    CREATE INDEX IF NOT EXISTS %I ON pgmq.%I (vt ASC);+    $QUERY$,+    qtable || '_vt_idx', qtable+  );++  EXECUTE FORMAT(+    $QUERY$+    CREATE INDEX IF NOT EXISTS %I ON pgmq.%I (archived_at);+    $QUERY$,+    'archived_at_idx_' || queue_name, atable+  );++  EXECUTE FORMAT(+    $QUERY$+    INSERT INTO pgmq.meta (queue_name, is_partitioned, is_unlogged)+    VALUES (%L, false, false)+    ON CONFLICT+    DO NOTHING;+    $QUERY$,+    queue_name+  );++END;+$$ LANGUAGE plpgsql;++CREATE FUNCTION pgmq.create_unlogged(queue_name TEXT)+RETURNS void AS $$+DECLARE+  qtable TEXT := pgmq.format_table_name(queue_name, 'q');+  qtable_seq TEXT := qtable || '_msg_id_seq';+  atable TEXT := pgmq.format_table_name(queue_name, 'a');+BEGIN+  PERFORM pgmq.validate_queue_name(queue_name);+  PERFORM pgmq.acquire_queue_lock(queue_name);++  EXECUTE FORMAT(+    $QUERY$+    CREATE UNLOGGED TABLE IF NOT EXISTS pgmq.%I (+        msg_id BIGINT PRIMARY KEY GENERATED ALWAYS AS IDENTITY,+        read_ct INT DEFAULT 0 NOT NULL,+        enqueued_at TIMESTAMP WITH TIME ZONE DEFAULT now() NOT NULL,+        vt TIMESTAMP WITH TIME ZONE NOT NULL,+        message JSONB,+        headers JSONB+    )+    $QUERY$,+    qtable+  );++  EXECUTE FORMAT(+    $QUERY$+    CREATE TABLE IF NOT EXISTS pgmq.%I (+      msg_id BIGINT PRIMARY KEY,+      read_ct INT DEFAULT 0 NOT NULL,+      enqueued_at TIMESTAMP WITH TIME ZONE DEFAULT now() NOT NULL,+      archived_at TIMESTAMP WITH TIME ZONE DEFAULT now() NOT NULL,+      vt TIMESTAMP WITH TIME ZONE NOT NULL,+      message JSONB,+      headers JSONB+    );+    $QUERY$,+    atable+  );++  EXECUTE FORMAT(+    $QUERY$+    CREATE INDEX IF NOT EXISTS %I ON pgmq.%I (vt ASC);+    $QUERY$,+    qtable || '_vt_idx', qtable+  );++  EXECUTE FORMAT(+    $QUERY$+    CREATE INDEX IF NOT EXISTS %I ON pgmq.%I (archived_at);+    $QUERY$,+    'archived_at_idx_' || queue_name, atable+  );++  EXECUTE FORMAT(+    $QUERY$+    INSERT INTO pgmq.meta (queue_name, is_partitioned, is_unlogged)+    VALUES (%L, false, true)+    ON CONFLICT+    DO NOTHING;+    $QUERY$,+    queue_name+  );+END;+$$ LANGUAGE plpgsql;++CREATE FUNCTION pgmq.create(queue_name TEXT)+RETURNS void AS $$+BEGIN+    PERFORM pgmq.create_non_partitioned(queue_name);+END;+$$ LANGUAGE plpgsql;++CREATE FUNCTION pgmq.drop_queue(queue_name TEXT, partitioned BOOLEAN)+RETURNS BOOLEAN AS $$+DECLARE+    qtable TEXT := pgmq.format_table_name(queue_name, 'q');+    fq_qtable TEXT := 'pgmq.' || qtable;+    atable TEXT := pgmq.format_table_name(queue_name, 'a');+    fq_atable TEXT := 'pgmq.' || atable;+BEGIN+    RAISE WARNING 'drop_queue(queue_name, partitioned) is deprecated and will be removed in PGMQ v2.0. Use drop_queue(queue_name) instead';++    PERFORM pgmq.drop_queue(queue_name);++    RETURN TRUE;+END;+$$ LANGUAGE plpgsql;++CREATE FUNCTION pgmq.drop_queue(queue_name TEXT)+RETURNS BOOLEAN AS $$+DECLARE+    qtable TEXT := pgmq.format_table_name(queue_name, 'q');+    qtable_seq TEXT := qtable || '_msg_id_seq';+    fq_qtable TEXT := 'pgmq.' || qtable;+    atable TEXT := pgmq.format_table_name(queue_name, 'a');+    fq_atable TEXT := 'pgmq.' || atable;+    partitioned BOOLEAN;+BEGIN+    PERFORM pgmq.acquire_queue_lock(queue_name);+    EXECUTE FORMAT(+        $QUERY$+        SELECT is_partitioned FROM pgmq.meta WHERE queue_name = %L+        $QUERY$,+        queue_name+    ) INTO partitioned;++    -- check if the queue exists+    IF NOT EXISTS (+        SELECT 1+        FROM information_schema.tables+        WHERE table_name = qtable and table_schema = 'pgmq'+    ) THEN+        RAISE NOTICE 'pgmq queue `%` does not exist', queue_name;+        RETURN FALSE;+    END IF;++    EXECUTE FORMAT(+        $QUERY$+        DROP TABLE IF EXISTS pgmq.%I+        $QUERY$,+        qtable+    );++    EXECUTE FORMAT(+        $QUERY$+        DROP TABLE IF EXISTS pgmq.%I+        $QUERY$,+        atable+    );++     IF EXISTS (+          SELECT 1+          FROM information_schema.tables+          WHERE table_name = 'meta' and table_schema = 'pgmq'+     ) THEN+        EXECUTE FORMAT(+            $QUERY$+            DELETE FROM pgmq.meta WHERE queue_name = %L+            $QUERY$,+            queue_name+        );+     END IF;++     IF partitioned THEN+        EXECUTE FORMAT(+          $QUERY$+          DELETE FROM %I.part_config where parent_table in (%L, %L)+          $QUERY$,+          pgmq._get_pg_partman_schema(), fq_qtable, fq_atable+        );+     END IF;++    RETURN TRUE;+END;+$$ LANGUAGE plpgsql;++-- list queues+CREATE FUNCTION pgmq."list_queues"()+RETURNS SETOF pgmq.queue_record AS $$+BEGIN+  RETURN QUERY SELECT * FROM pgmq.meta;+END+$$ LANGUAGE plpgsql;++-- purge queue, deleting all entries in it.+CREATE OR REPLACE FUNCTION pgmq."purge_queue"(queue_name TEXT)+RETURNS BIGINT AS $$+DECLARE+  deleted_count INTEGER;+  qtable TEXT := pgmq.format_table_name(queue_name, 'q');+BEGIN+  -- Get the row count before truncating+  EXECUTE format('SELECT count(*) FROM pgmq.%I', qtable) INTO deleted_count;++  -- Use TRUNCATE for better performance on large tables+  EXECUTE format('TRUNCATE TABLE pgmq.%I', qtable);++  -- Return the number of purged rows+  RETURN deleted_count;+END+$$ LANGUAGE plpgsql;++-- unassign archive, so it can be kept when a queue is deleted+CREATE FUNCTION pgmq."detach_archive"(queue_name TEXT)+RETURNS VOID AS $$+DECLARE+  atable TEXT := pgmq.format_table_name(queue_name, 'a');+BEGIN+  RAISE WARNING 'detach_archive(queue_name) is deprecated and is a no-op. It will be removed in PGMQ v2.0. Archive tables are no longer member objects.';+END+$$ LANGUAGE plpgsql;
+ database/v1.9.0/06_message_ops.sql view
@@ -0,0 +1,626 @@+------------------------------------------------------------+-- Message operations: send, read, delete, archive, pop, set_vt+------------------------------------------------------------++-- send: 2 args, no delay+CREATE FUNCTION pgmq.send(+    queue_name TEXT,+    msg JSONB+)+RETURNS SETOF BIGINT AS $$+BEGIN+    RETURN QUERY SELECT * FROM pgmq.send(queue_name, msg, NULL, clock_timestamp());+END;+$$ LANGUAGE plpgsql;++-- send: 3 args with headers+CREATE FUNCTION pgmq.send(+    queue_name TEXT,+    msg JSONB,+    headers JSONB+)+RETURNS SETOF BIGINT AS $$+BEGIN+    RETURN QUERY SELECT * FROM pgmq.send(queue_name, msg, headers, clock_timestamp());+END;+$$ LANGUAGE plpgsql;++-- send: 3 args with delay+CREATE FUNCTION pgmq.send(+    queue_name TEXT,+    msg JSONB,+    delay INTEGER+)+RETURNS SETOF BIGINT AS $$+BEGIN+    RETURN QUERY SELECT * FROM pgmq.send(queue_name, msg, NULL, clock_timestamp() + make_interval(secs => delay));+END;+$$ LANGUAGE plpgsql;++-- send: 4 args with headers and delay+CREATE FUNCTION pgmq.send(+    queue_name TEXT,+    msg JSONB,+    headers JSONB,+    delay INTEGER+)+RETURNS SETOF BIGINT AS $$+BEGIN+    RETURN QUERY SELECT * FROM pgmq.send(queue_name, msg, headers, clock_timestamp() + make_interval(secs => delay));+END;+$$ LANGUAGE plpgsql;++-- send: actual implementation+CREATE FUNCTION pgmq.send(+    queue_name TEXT,+    msg JSONB,+    headers JSONB,+    vt TIMESTAMP WITH TIME ZONE+)+RETURNS SETOF BIGINT AS $$+DECLARE+    sql TEXT;+    qtable TEXT := pgmq.format_table_name(queue_name, 'q');+BEGIN+    sql := FORMAT(+        $QUERY$+        INSERT INTO pgmq.%I (vt, message, headers)+        VALUES ($1, $2, $3)+        RETURNING msg_id;+        $QUERY$,+        qtable+    );+    RETURN QUERY EXECUTE sql USING vt, msg, headers;+END;+$$ LANGUAGE plpgsql;++-- send_batch: no delay+CREATE FUNCTION pgmq.send_batch(+    queue_name TEXT,+    msgs JSONB[]+)+RETURNS SETOF BIGINT AS $$+BEGIN+    RETURN QUERY SELECT * FROM pgmq.send_batch(queue_name, msgs, NULL, clock_timestamp());+END;+$$ LANGUAGE plpgsql;++-- send_batch: with delay+CREATE FUNCTION pgmq.send_batch(+    queue_name TEXT,+    msgs JSONB[],+    delay INTEGER+)+RETURNS SETOF BIGINT AS $$+BEGIN+    RETURN QUERY SELECT * FROM pgmq.send_batch(queue_name, msgs, NULL, clock_timestamp() + make_interval(secs => delay));+END;+$$ LANGUAGE plpgsql;++-- send_batch: with headers+CREATE FUNCTION pgmq.send_batch(+    queue_name TEXT,+    msgs JSONB[],+    headers JSONB[]+)+RETURNS SETOF BIGINT AS $$+BEGIN+    RETURN QUERY SELECT * FROM pgmq.send_batch(queue_name, msgs, headers, clock_timestamp());+END;+$$ LANGUAGE plpgsql;++-- send_batch: with headers and delay+CREATE FUNCTION pgmq.send_batch(+    queue_name TEXT,+    msgs JSONB[],+    headers JSONB[],+    delay INTEGER+)+RETURNS SETOF BIGINT AS $$+BEGIN+    RETURN QUERY SELECT * FROM pgmq.send_batch(queue_name, msgs, headers, clock_timestamp() + make_interval(secs => delay));+END;+$$ LANGUAGE plpgsql;++-- send_batch: actual implementation+CREATE FUNCTION pgmq.send_batch(+    queue_name TEXT,+    msgs JSONB[],+    headers JSONB[],+    vt TIMESTAMP WITH TIME ZONE+)+RETURNS SETOF BIGINT AS $$+DECLARE+    sql TEXT;+    qtable TEXT := pgmq.format_table_name(queue_name, 'q');+BEGIN+    sql := FORMAT(+        $QUERY$+        INSERT INTO pgmq.%I (vt, message, headers)+        SELECT $1, unnest($2), unnest(coalesce($3, array_fill(NULL::jsonb, array[cardinality($2)])))+        RETURNING msg_id;+        $QUERY$,+        qtable+    );+    RETURN QUERY EXECUTE sql USING vt, msgs, headers;+END;+$$ LANGUAGE plpgsql;++---- read+---- reads a number of messages from a queue, setting a visibility timeout on them+CREATE FUNCTION pgmq.read(+    queue_name TEXT,+    vt INTEGER,+    qty INTEGER+)+RETURNS SETOF pgmq.message_record AS $$+DECLARE+    sql TEXT;+    qtable TEXT := pgmq.format_table_name(queue_name, 'q');+BEGIN+    sql := FORMAT(+        $QUERY$+        WITH cte AS+        (+            SELECT msg_id+            FROM pgmq.%I+            WHERE vt <= clock_timestamp()+            ORDER BY msg_id ASC+            LIMIT $1+            FOR UPDATE SKIP LOCKED+        )+        UPDATE pgmq.%I m+        SET+            vt = clock_timestamp() + %L,+            read_ct = read_ct + 1+        FROM cte+        WHERE m.msg_id = cte.msg_id+        RETURNING m.msg_id, m.read_ct, m.enqueued_at, m.vt, m.message, m.headers;+        $QUERY$,+        qtable, qtable, make_interval(secs => vt)+    );+    RETURN QUERY EXECUTE sql USING qty;+END;+$$ LANGUAGE plpgsql;++---- read+---- reads a number of messages from a queue, setting a visibility timeout on them+CREATE FUNCTION pgmq.read(+    queue_name TEXT,+    vt INTEGER,+    qty INTEGER,+    conditional JSONB+)+RETURNS SETOF pgmq.message_record AS $$+DECLARE+    sql TEXT;+    qtable TEXT := pgmq.format_table_name(queue_name, 'q');+BEGIN+    sql := FORMAT(+        $QUERY$+        WITH cte AS+        (+            SELECT msg_id+            FROM pgmq.%I+            WHERE vt <= clock_timestamp() AND message @> $2+            ORDER BY msg_id ASC+            LIMIT $1+            FOR UPDATE SKIP LOCKED+        )+        UPDATE pgmq.%I m+        SET+            vt = clock_timestamp() + %L,+            read_ct = read_ct + 1+        FROM cte+        WHERE m.msg_id = cte.msg_id+        RETURNING m.msg_id, m.read_ct, m.enqueued_at, m.vt, m.message, m.headers;+        $QUERY$,+        qtable, qtable, make_interval(secs => vt)+    );+    RETURN QUERY EXECUTE sql USING qty, conditional;+END;+$$ LANGUAGE plpgsql;++---- read_with_poll+---- reads a number of messages from a queue, setting a visibility timeout on them+---- 5-arg version without conditional filter+CREATE FUNCTION pgmq.read_with_poll(+    queue_name TEXT,+    vt INTEGER,+    qty INTEGER,+    max_poll_seconds INTEGER DEFAULT 5,+    poll_interval_ms INTEGER DEFAULT 100+)+RETURNS SETOF pgmq.message_record AS $$+DECLARE+    r pgmq.message_record;+    stop_at TIMESTAMP;+    sql TEXT;+    qtable TEXT := pgmq.format_table_name(queue_name, 'q');+BEGIN+    stop_at := clock_timestamp() + make_interval(secs => max_poll_seconds);+    LOOP+      IF clock_timestamp() >= stop_at THEN+        RETURN;+      END IF;++      sql := FORMAT(+          $QUERY$+          WITH cte AS+          (+              SELECT msg_id+              FROM pgmq.%I+              WHERE vt <= clock_timestamp()+              ORDER BY msg_id ASC+              LIMIT $1+              FOR UPDATE SKIP LOCKED+          )+          UPDATE pgmq.%I m+          SET+              vt = clock_timestamp() + %L,+              read_ct = read_ct + 1+          FROM cte+          WHERE m.msg_id = cte.msg_id+          RETURNING m.msg_id, m.read_ct, m.enqueued_at, m.vt, m.message, m.headers;+          $QUERY$,+          qtable, qtable, make_interval(secs => vt)+      );++      FOR r IN+        EXECUTE sql USING qty+      LOOP+        RETURN NEXT r;+      END LOOP;+      IF FOUND THEN+        RETURN;+      ELSE+        PERFORM pg_sleep(poll_interval_ms::numeric / 1000);+      END IF;+    END LOOP;+END;+$$ LANGUAGE plpgsql;++---- read_with_poll+---- 6-arg version with conditional as 6th parameter (matches pgmq-hasql expected order)+---- This is the overload used by pgmq-hasql: read_with_poll($1,$2,$3,$4,$5,$6)+CREATE FUNCTION pgmq.read_with_poll(+    queue_name TEXT,+    vt INTEGER,+    qty INTEGER,+    max_poll_seconds INTEGER,+    poll_interval_ms INTEGER,+    conditional JSONB+)+RETURNS SETOF pgmq.message_record AS $$+DECLARE+    r pgmq.message_record;+    stop_at TIMESTAMP;+    sql TEXT;+    qtable TEXT := pgmq.format_table_name(queue_name, 'q');+BEGIN+    stop_at := clock_timestamp() + make_interval(secs => max_poll_seconds);+    LOOP+      IF clock_timestamp() >= stop_at THEN+        RETURN;+      END IF;++      -- If conditional is NULL or empty, don't filter+      IF conditional IS NULL OR conditional = '{}'::jsonb THEN+          sql := FORMAT(+              $QUERY$+              WITH cte AS+              (+                  SELECT msg_id+                  FROM pgmq.%I+                  WHERE vt <= clock_timestamp()+                  ORDER BY msg_id ASC+                  LIMIT $1+                  FOR UPDATE SKIP LOCKED+              )+              UPDATE pgmq.%I m+              SET+                  vt = clock_timestamp() + %L,+                  read_ct = read_ct + 1+              FROM cte+              WHERE m.msg_id = cte.msg_id+              RETURNING m.msg_id, m.read_ct, m.enqueued_at, m.vt, m.message, m.headers;+              $QUERY$,+              qtable, qtable, make_interval(secs => vt)+          );++          FOR r IN+            EXECUTE sql USING qty+          LOOP+            RETURN NEXT r;+          END LOOP;+      ELSE+          sql := FORMAT(+              $QUERY$+              WITH cte AS+              (+                  SELECT msg_id+                  FROM pgmq.%I+                  WHERE vt <= clock_timestamp() AND message @> $2+                  ORDER BY msg_id ASC+                  LIMIT $1+                  FOR UPDATE SKIP LOCKED+              )+              UPDATE pgmq.%I m+              SET+                  vt = clock_timestamp() + %L,+                  read_ct = read_ct + 1+              FROM cte+              WHERE m.msg_id = cte.msg_id+              RETURNING m.msg_id, m.read_ct, m.enqueued_at, m.vt, m.message, m.headers;+              $QUERY$,+              qtable, qtable, make_interval(secs => vt)+          );++          FOR r IN+            EXECUTE sql USING qty, conditional+          LOOP+            RETURN NEXT r;+          END LOOP;+      END IF;++      IF FOUND THEN+        RETURN;+      ELSE+        PERFORM pg_sleep(poll_interval_ms::numeric / 1000);+      END IF;+    END LOOP;+END;+$$ LANGUAGE plpgsql;++---- read_with_poll+---- 6-arg version with conditional as 4th parameter (legacy order)+CREATE FUNCTION pgmq.read_with_poll(+    queue_name TEXT,+    vt INTEGER,+    qty INTEGER,+    conditional JSONB,+    max_poll_seconds INTEGER DEFAULT 5,+    poll_interval_ms INTEGER DEFAULT 100+)+RETURNS SETOF pgmq.message_record AS $$+DECLARE+    r pgmq.message_record;+    stop_at TIMESTAMP;+    sql TEXT;+    qtable TEXT := pgmq.format_table_name(queue_name, 'q');+BEGIN+    stop_at := clock_timestamp() + make_interval(secs => max_poll_seconds);+    LOOP+      IF clock_timestamp() >= stop_at THEN+        RETURN;+      END IF;++      sql := FORMAT(+          $QUERY$+          WITH cte AS+          (+              SELECT msg_id+              FROM pgmq.%I+              WHERE vt <= clock_timestamp() AND message @> $2+              ORDER BY msg_id ASC+              LIMIT $1+              FOR UPDATE SKIP LOCKED+          )+          UPDATE pgmq.%I m+          SET+              vt = clock_timestamp() + %L,+              read_ct = read_ct + 1+          FROM cte+          WHERE m.msg_id = cte.msg_id+          RETURNING m.msg_id, m.read_ct, m.enqueued_at, m.vt, m.message, m.headers;+          $QUERY$,+          qtable, qtable, make_interval(secs => vt)+      );++      FOR r IN+        EXECUTE sql USING qty, conditional+      LOOP+        RETURN NEXT r;+      END LOOP;+      IF FOUND THEN+        RETURN;+      ELSE+        PERFORM pg_sleep(poll_interval_ms::numeric / 1000);+      END IF;+    END LOOP;+END;+$$ LANGUAGE plpgsql;++-- delete+-- deletes a message id from the queue permanently+CREATE FUNCTION pgmq.delete(+    queue_name TEXT,+    msg_id BIGINT+)+RETURNS BOOLEAN AS $$+DECLARE+    sql TEXT;+    result BIGINT;+    qtable TEXT := pgmq.format_table_name(queue_name, 'q');+BEGIN+    sql := FORMAT(+        $QUERY$+        DELETE FROM pgmq.%I+        WHERE msg_id = $1+        RETURNING msg_id+        $QUERY$,+        qtable+    );+    EXECUTE sql USING msg_id INTO result;+    RETURN NOT (result IS NULL);+END;+$$ LANGUAGE plpgsql;++-- delete+-- deletes an array of message ids from the queue permanently+CREATE FUNCTION pgmq.delete(+    queue_name TEXT,+    msg_ids BIGINT[]+)+RETURNS SETOF BIGINT AS $$+DECLARE+    sql TEXT;+    qtable TEXT := pgmq.format_table_name(queue_name, 'q');+BEGIN+    sql := FORMAT(+        $QUERY$+        DELETE FROM pgmq.%I+        WHERE msg_id = ANY($1)+        RETURNING msg_id+        $QUERY$,+        qtable+    );+    RETURN QUERY EXECUTE sql USING msg_ids;+END;+$$ LANGUAGE plpgsql;++-- archive+-- removes a message from the queue, and sends it to the archive, where its+-- temporary durability is lost and will be deleted on the next archive expiration+-- job.+CREATE FUNCTION pgmq.archive(+    queue_name TEXT,+    msg_id BIGINT+)+RETURNS BOOLEAN AS $$+DECLARE+    sql TEXT;+    result BIGINT;+    qtable TEXT := pgmq.format_table_name(queue_name, 'q');+    atable TEXT := pgmq.format_table_name(queue_name, 'a');+BEGIN+    sql := FORMAT(+        $QUERY$+        WITH archived AS (+            DELETE FROM pgmq.%I+            WHERE msg_id = $1+            RETURNING msg_id, vt, read_ct, enqueued_at, message, headers+        )+        INSERT INTO pgmq.%I (msg_id, vt, read_ct, enqueued_at, message, headers)+        SELECT msg_id, vt, read_ct, enqueued_at, message, headers+        FROM archived+        RETURNING msg_id;+        $QUERY$,+        qtable, atable+    );+    EXECUTE sql USING msg_id INTO result;+    RETURN NOT (result IS NULL);+END;+$$ LANGUAGE plpgsql;++-- archive+-- removes an array of message ids from the queue, and sends it to the archive,+-- where its temporary durability is lost and will be deleted on the next archive+-- expiration job.+CREATE FUNCTION pgmq.archive(+    queue_name TEXT,+    msg_ids BIGINT[]+)+RETURNS SETOF BIGINT AS $$+DECLARE+    sql TEXT;+    qtable TEXT := pgmq.format_table_name(queue_name, 'q');+    atable TEXT := pgmq.format_table_name(queue_name, 'a');+BEGIN+    sql := FORMAT(+        $QUERY$+        WITH archived AS (+            DELETE FROM pgmq.%I+            WHERE msg_id = ANY($1)+            RETURNING msg_id, vt, read_ct, enqueued_at, message, headers+        )+        INSERT INTO pgmq.%I (msg_id, vt, read_ct, enqueued_at, message, headers)+        SELECT msg_id, vt, read_ct, enqueued_at, message, headers+        FROM archived+        RETURNING msg_id;+        $QUERY$,+        qtable, atable+    );+    RETURN QUERY EXECUTE sql USING msg_ids;+END;+$$ LANGUAGE plpgsql;++-- pop messages from queue (atomic read + delete)+-- qty parameter added in pgmq 1.7.0+CREATE FUNCTION pgmq.pop(+    queue_name TEXT,+    qty INTEGER DEFAULT 1+)+RETURNS SETOF pgmq.message_record AS $$+DECLARE+    sql TEXT;+    qtable TEXT := pgmq.format_table_name(queue_name, 'q');+BEGIN+    sql := FORMAT(+        $QUERY$+        WITH cte AS+            (+                SELECT msg_id+                FROM pgmq.%I+                WHERE vt <= clock_timestamp()+                ORDER BY msg_id ASC+                LIMIT $1+                FOR UPDATE SKIP LOCKED+            )+        DELETE from pgmq.%I+        WHERE msg_id IN (select msg_id from cte)+        RETURNING *;+        $QUERY$,+        qtable, qtable+    );+    RETURN QUERY EXECUTE sql USING qty;+END;+$$ LANGUAGE plpgsql;++-- Sets vt of a single message, returns it+CREATE FUNCTION pgmq.set_vt(+    queue_name TEXT,+    msg_id BIGINT,+    vt_offset INTEGER+)+RETURNS SETOF pgmq.message_record AS $$+DECLARE+    sql TEXT;+    qtable TEXT := pgmq.format_table_name(queue_name, 'q');+BEGIN+    sql := FORMAT(+        $QUERY$+        UPDATE pgmq.%I+        SET vt = (clock_timestamp() + %L)+        WHERE msg_id = %L+        RETURNING *;+        $QUERY$,+        qtable, make_interval(secs => vt_offset), msg_id+    );+    RETURN QUERY EXECUTE sql;+END;+$$ LANGUAGE plpgsql;++-- Batch set_vt for multiple messages (pgmq 1.8.0+)+CREATE FUNCTION pgmq.set_vt(+    queue_name TEXT,+    msg_ids BIGINT[],+    vt_offset INTEGER+)+RETURNS SETOF pgmq.message_record AS $$+DECLARE+    sql TEXT;+    qtable TEXT := pgmq.format_table_name(queue_name, 'q');+BEGIN+    sql := FORMAT(+        $QUERY$+        UPDATE pgmq.%I+        SET vt = (clock_timestamp() + %L)+        WHERE msg_id = ANY($1)+        RETURNING *;+        $QUERY$,+        qtable, make_interval(secs => vt_offset)+    );+    RETURN QUERY EXECUTE sql USING msg_ids;+END;+$$ LANGUAGE plpgsql;
+ database/v1.9.0/07_metrics.sql view
@@ -0,0 +1,83 @@+------------------------------------------------------------+-- Metrics functions+------------------------------------------------------------++CREATE FUNCTION pgmq.metrics(queue_name TEXT)+RETURNS pgmq.metrics_result AS $$+DECLARE+    result_row pgmq.metrics_result;+    query TEXT;+    qtable TEXT := pgmq.format_table_name(queue_name, 'q');+BEGIN+    query := FORMAT(+        $QUERY$+        WITH q_summary AS (+            SELECT+                count(*) as queue_length,+                EXTRACT(epoch FROM (NOW() - max(enqueued_at)))::int as newest_msg_age_sec,+                EXTRACT(epoch FROM (NOW() - min(enqueued_at)))::int as oldest_msg_age_sec,+                NOW() as scrape_time,+                count(*) FILTER (WHERE vt <= NOW()) AS queue_visible_length+            FROM pgmq.%I+        ),+        all_metrics AS (+            SELECT CASE+                WHEN is_partitioned THEN+                    %L || '_part_' || pgmq._get_partition_col(%L)+                ELSE+                    %L+            END as queue_name+            FROM pgmq.meta+            WHERE meta.queue_name = %L+        )+        SELECT+            m.queue_name,+            q_summary.queue_length,+            q_summary.newest_msg_age_sec,+            q_summary.oldest_msg_age_sec,+            pgmq._get_pg_stat_get_xact_tuples_inserted(%L)::bigint,+            q_summary.scrape_time,+            q_summary.queue_visible_length+        FROM q_summary+        CROSS JOIN all_metrics m+        $QUERY$,+        qtable, queue_name, queue_name, queue_name, queue_name, 'pgmq.' || qtable+    );+    EXECUTE query INTO result_row;+    RETURN result_row;+END;+$$ LANGUAGE plpgsql;++CREATE FUNCTION pgmq.metrics_all()+RETURNS SETOF pgmq.metrics_result AS $$+DECLARE+    row_name RECORD;+    result_row pgmq.metrics_result;+BEGIN+    FOR row_name IN SELECT queue_name FROM pgmq.meta LOOP+        result_row := pgmq.metrics(row_name.queue_name);+        RETURN NEXT result_row;+    END LOOP;+END;+$$ LANGUAGE plpgsql;++CREATE FUNCTION pgmq._get_pg_stat_get_xact_tuples_inserted(table_name TEXT)+RETURNS bigint AS $$+DECLARE+    result bigint;+BEGIN+    SELECT pg_stat_get_xact_tuples_inserted(oid)+    INTO result+    FROM pg_class+    WHERE relname = table_name;++    IF result IS NULL THEN+        SELECT pg_stat_get_xact_tuples_inserted(oid)+        INTO result+        FROM pg_class+        WHERE relname = SUBSTRING(table_name FROM 6);  -- Remove 'pgmq.' prefix (5 chars + 1 for the dot)+    END IF;++    RETURN COALESCE(result, 0);+END;+$$ LANGUAGE plpgsql;
+ database/v1.9.0/08_partitioning.sql view
@@ -0,0 +1,263 @@+------------------------------------------------------------+-- Partitioning support functions+------------------------------------------------------------++CREATE FUNCTION pgmq._get_pg_partman_schema()+RETURNS TEXT AS $$+    SELECT extnamespace::regnamespace::text+    FROM pg_extension+    WHERE extname = 'pg_partman';+$$ LANGUAGE sql STABLE;++CREATE FUNCTION pgmq._get_partition_col(queue_name TEXT)+RETURNS TEXT AS $$+DECLARE+    _partition_col TEXT;+    _table_name TEXT := pgmq.format_table_name(queue_name, 'q');+BEGIN+    SELECT part_config.partition_interval INTO _partition_col+    FROM pg_class+    JOIN pg_namespace ON pg_class.relnamespace = pg_namespace.oid+    JOIN (SELECT parent_table, partition_interval FROM pgmq._get_partman_config()) AS part_config+      ON part_config.parent_table = pg_namespace.nspname || '.' || pg_class.relname+    WHERE pg_class.relname = _table_name+      AND pg_namespace.nspname = 'pgmq';++    RETURN _partition_col;+END;+$$ LANGUAGE plpgsql;++CREATE FUNCTION pgmq._get_partman_config()+RETURNS TABLE (parent_table text, partition_interval text) AS $$+DECLARE+    schema_name TEXT;+BEGIN+    SELECT pgmq._get_pg_partman_schema() INTO schema_name;+    IF schema_name IS NULL THEN+        RETURN QUERY SELECT NULL::text, NULL::text WHERE FALSE;+        RETURN;+    END IF;+    RETURN QUERY EXECUTE FORMAT(+        $QUERY$+        SELECT parent_table, partition_interval+        FROM %I.part_config+        $QUERY$,+        schema_name+    );+END;+$$ LANGUAGE plpgsql;++-- partitioned, with partition interval+CREATE FUNCTION pgmq.create_partitioned(+  queue_name TEXT,+  partition_interval TEXT DEFAULT '10000',+  retention_interval TEXT DEFAULT '100000'+)+RETURNS void AS $$+DECLARE+  qtable TEXT := pgmq.format_table_name(queue_name, 'q');+  qtable_seq TEXT := qtable || '_msg_id_seq';+  atable TEXT := pgmq.format_table_name(queue_name, 'a');+  atable_seq TEXT := atable || '_msg_id_seq';+BEGIN+  PERFORM pgmq.validate_queue_name(queue_name);+  PERFORM pgmq.acquire_queue_lock(queue_name);++  EXECUTE FORMAT(+    $QUERY$+    CREATE TABLE IF NOT EXISTS pgmq.%I (+        msg_id BIGINT GENERATED ALWAYS AS IDENTITY,+        read_ct INT DEFAULT 0 NOT NULL,+        enqueued_at TIMESTAMP WITH TIME ZONE DEFAULT now() NOT NULL,+        vt TIMESTAMP WITH TIME ZONE NOT NULL,+        message JSONB,+        headers JSONB+    ) PARTITION BY RANGE (msg_id)+    $QUERY$,+    qtable+  );++  EXECUTE FORMAT(+    $QUERY$+    CREATE TABLE IF NOT EXISTS pgmq.%I (+      msg_id BIGINT NOT NULL,+      read_ct INT DEFAULT 0 NOT NULL,+      enqueued_at TIMESTAMP WITH TIME ZONE DEFAULT now() NOT NULL,+      archived_at TIMESTAMP WITH TIME ZONE DEFAULT now() NOT NULL,+      vt TIMESTAMP WITH TIME ZONE NOT NULL,+      message JSONB,+      headers JSONB+    ) PARTITION BY RANGE (msg_id);+    $QUERY$,+    atable+  );++  IF NOT pgmq._extension_exists('pg_partman') THEN+    RAISE EXCEPTION 'pg_partman is required for partitioned queues';+  END IF;++  EXECUTE FORMAT(+    $QUERY$+    SELECT %I.create_parent(+        p_parent_table := 'pgmq.%I',+        p_control := 'msg_id',+        p_type := 'range',+        p_interval := %L,+        p_premake := 1+    )+    $QUERY$,+    pgmq._get_pg_partman_schema(), qtable, partition_interval+  );++  EXECUTE FORMAT(+    $QUERY$+    SELECT %I.create_parent(+        p_parent_table := 'pgmq.%I',+        p_control := 'msg_id',+        p_type := 'range',+        p_interval := %L,+        p_premake := 1+    )+    $QUERY$,+    pgmq._get_pg_partman_schema(), atable, partition_interval+  );++  EXECUTE FORMAT(+    $QUERY$+    UPDATE %I.part_config+    SET+        retention = %L,+        retention_keep_table = false,+        retention_keep_index = true,+        automatic_maintenance = 'on'+    WHERE parent_table = 'pgmq.%I';+    $QUERY$,+    pgmq._get_pg_partman_schema(), retention_interval, qtable+  );++  EXECUTE FORMAT(+    $QUERY$+    UPDATE %I.part_config+    SET+        retention = %L,+        retention_keep_table = false,+        retention_keep_index = true,+        automatic_maintenance = 'on'+    WHERE parent_table = 'pgmq.%I';+    $QUERY$,+    pgmq._get_pg_partman_schema(), retention_interval, atable+  );++  EXECUTE FORMAT(+    $QUERY$+    CREATE INDEX IF NOT EXISTS %I ON pgmq.%I (vt ASC);+    $QUERY$,+    qtable || '_vt_idx', qtable+  );++  EXECUTE FORMAT(+    $QUERY$+    CREATE INDEX IF NOT EXISTS %I ON pgmq.%I (archived_at);+    $QUERY$,+    'archived_at_idx_' || queue_name, atable+  );++  EXECUTE FORMAT(+    $QUERY$+    INSERT INTO pgmq.meta (queue_name, is_partitioned, is_unlogged)+    VALUES (%L, true, false)+    ON CONFLICT+    DO NOTHING;+    $QUERY$,+    queue_name+  );++END;+$$ LANGUAGE plpgsql;++-- convert_archive_partitioned+CREATE FUNCTION pgmq.convert_archive_partitioned(+  queue_name TEXT,+  partition_interval TEXT DEFAULT '10000',+  retention_interval TEXT DEFAULT '100000',+  leading_partition INT DEFAULT 10+)+RETURNS void AS $$+DECLARE+  atable TEXT := pgmq.format_table_name(queue_name, 'a');+BEGIN+  EXECUTE FORMAT(+    $QUERY$+    ALTER TABLE IF EXISTS pgmq.%I RENAME TO %I+    $QUERY$,+    atable, atable || '_old'+  );++  EXECUTE FORMAT(+    $QUERY$+    CREATE TABLE IF NOT EXISTS pgmq.%I (+      msg_id BIGINT NOT NULL,+      read_ct INT DEFAULT 0 NOT NULL,+      enqueued_at TIMESTAMP WITH TIME ZONE DEFAULT now() NOT NULL,+      archived_at TIMESTAMP WITH TIME ZONE DEFAULT now() NOT NULL,+      vt TIMESTAMP WITH TIME ZONE NOT NULL,+      message JSONB,+      headers JSONB+    ) PARTITION BY RANGE (msg_id);+    $QUERY$,+    atable+  );++  IF NOT pgmq._extension_exists('pg_partman') THEN+    RAISE EXCEPTION 'pg_partman is required for partitioned queues';+  END IF;++  EXECUTE FORMAT(+    $QUERY$+    SELECT %I.create_parent(+        p_parent_table := 'pgmq.%I',+        p_control := 'msg_id',+        p_type := 'range',+        p_interval := %L,+        p_premake := %L+    )+    $QUERY$,+    pgmq._get_pg_partman_schema(), atable, partition_interval, leading_partition+  );++  EXECUTE FORMAT(+    $QUERY$+    UPDATE %I.part_config+    SET+        retention = %L,+        retention_keep_table = false,+        retention_keep_index = true,+        automatic_maintenance = 'on'+    WHERE parent_table = 'pgmq.%I';+    $QUERY$,+    pgmq._get_pg_partman_schema(), retention_interval, atable+  );++  EXECUTE FORMAT(+    $QUERY$+    INSERT INTO pgmq.%I (msg_id, read_ct, enqueued_at, archived_at, vt, message, headers)+    SELECT msg_id, read_ct, enqueued_at, archived_at, vt, message, headers FROM pgmq.%I+    $QUERY$,+    atable, atable || '_old'+  );++  EXECUTE FORMAT(+    $QUERY$+    CREATE INDEX IF NOT EXISTS %I ON pgmq.%I (archived_at);+    $QUERY$,+    'archived_at_idx_' || queue_name, atable+  );++  EXECUTE FORMAT(+    $QUERY$+    DROP TABLE IF EXISTS pgmq.%I+    $QUERY$,+    atable || '_old'+  );+END;+$$ LANGUAGE plpgsql;
+ database/v1.9.0/09_notifications.sql view
@@ -0,0 +1,103 @@+------------------------------------------------------------+-- Notification functions+------------------------------------------------------------++CREATE FUNCTION pgmq._should_notify(queue_name TEXT)+RETURNS BOOLEAN AS $$+DECLARE+    should_notify BOOLEAN;+    throttle_ms INTEGER;+    last_notified TIMESTAMP WITH TIME ZONE;+BEGIN+    -- Try to get throttle settings+    SELECT throttle_interval_ms, last_notified_at+    INTO throttle_ms, last_notified+    FROM pgmq.notify_insert_throttle+    WHERE pgmq.notify_insert_throttle.queue_name = _should_notify.queue_name+    FOR UPDATE SKIP LOCKED;++    -- If no row found, assume notifications are enabled without throttling+    IF throttle_ms IS NULL THEN+        RETURN TRUE;+    END IF;++    -- Check if enough time has passed since last notification+    IF clock_timestamp() >= last_notified + make_interval(secs => throttle_ms::numeric / 1000) THEN+        -- Update the last notification time+        UPDATE pgmq.notify_insert_throttle+        SET last_notified_at = clock_timestamp()+        WHERE pgmq.notify_insert_throttle.queue_name = _should_notify.queue_name;+        RETURN TRUE;+    END IF;++    RETURN FALSE;+END;+$$ LANGUAGE plpgsql;++CREATE FUNCTION pgmq._notify_on_insert()+RETURNS TRIGGER AS $$+BEGIN+    IF pgmq._should_notify(TG_ARGV[0]) THEN+        PERFORM pg_notify(TG_ARGV[0], NEW.msg_id::text);+    END IF;+    RETURN NEW;+END;+$$ LANGUAGE plpgsql;++-- Add trigger to queue table to notify on insert+CREATE FUNCTION pgmq.enable_queue_notifications(queue_name TEXT)+RETURNS void AS $$+DECLARE+    qtable TEXT := pgmq.format_table_name(queue_name, 'q');+    channel TEXT := queue_name;+    trigger_name TEXT := 'notify_insert_' || qtable;+BEGIN+    -- Create trigger on the queue table+    EXECUTE FORMAT(+        $QUERY$+        CREATE TRIGGER %I+        AFTER INSERT ON pgmq.%I+        FOR EACH ROW+        EXECUTE FUNCTION pgmq._notify_on_insert(%L)+        $QUERY$,+        trigger_name, qtable, channel+    );++    -- Insert throttle tracking row if it doesn't exist+    INSERT INTO pgmq.notify_insert_throttle (queue_name, throttle_interval_ms, last_notified_at)+    VALUES (queue_name, 0, to_timestamp(0))+    ON CONFLICT (queue_name) DO NOTHING;+END;+$$ LANGUAGE plpgsql;++-- Remove trigger from queue table+CREATE FUNCTION pgmq.disable_queue_notifications(queue_name TEXT)+RETURNS void AS $$+DECLARE+    qtable TEXT := pgmq.format_table_name(queue_name, 'q');+    trigger_name TEXT := 'notify_insert_' || qtable;+BEGIN+    -- Drop trigger if it exists+    EXECUTE FORMAT(+        $QUERY$+        DROP TRIGGER IF EXISTS %I ON pgmq.%I+        $QUERY$,+        trigger_name, qtable+    );++    -- Remove throttle tracking row+    DELETE FROM pgmq.notify_insert_throttle+    WHERE pgmq.notify_insert_throttle.queue_name = disable_queue_notifications.queue_name;+END;+$$ LANGUAGE plpgsql;++-- Set throttle interval for queue notifications+CREATE FUNCTION pgmq.set_notification_throttle(queue_name TEXT, throttle_interval_ms INTEGER)+RETURNS void AS $$+BEGIN+    INSERT INTO pgmq.notify_insert_throttle (queue_name, throttle_interval_ms, last_notified_at)+    VALUES (queue_name, throttle_interval_ms, to_timestamp(0))+    ON CONFLICT (queue_name) DO UPDATE+    SET throttle_interval_ms = EXCLUDED.throttle_interval_ms;+END;+$$ LANGUAGE plpgsql;
+ database/v1.9.0/10_fifo.sql view
@@ -0,0 +1,564 @@+------------------------------------------------------------+-- FIFO read functions (for pgmq 1.8.0+ and 1.9.0+)+------------------------------------------------------------++-- read_fifo: reads messages in strict FIFO order+-- This function ensures messages are read in strict FIFO order by:+-- 1. Ordering by msg_id ASC (not just using SKIP LOCKED)+-- 2. Only returning messages that don't skip any earlier pending messages+CREATE FUNCTION pgmq.read_fifo(+    queue_name TEXT,+    vt INTEGER,+    qty INTEGER+)+RETURNS SETOF pgmq.message_record AS $$+DECLARE+    sql TEXT;+    qtable TEXT := pgmq.format_table_name(queue_name, 'q');+BEGIN+    sql := FORMAT(+        $QUERY$+        WITH cte AS+        (+            SELECT msg_id, vt AS current_vt+            FROM pgmq.%I+            ORDER BY msg_id ASC+            LIMIT $1+            FOR UPDATE SKIP LOCKED+        ),+        eligible AS (+            SELECT c.msg_id+            FROM cte c+            WHERE c.msg_id = (+                SELECT MIN(msg_id)+                FROM pgmq.%I+                WHERE vt <= clock_timestamp()+            )+            OR NOT EXISTS (+                SELECT 1+                FROM pgmq.%I q+                WHERE q.msg_id < c.msg_id+                AND q.vt <= clock_timestamp()+                AND q.msg_id NOT IN (SELECT msg_id FROM cte)+            )+        )+        UPDATE pgmq.%I m+        SET+            vt = clock_timestamp() + %L,+            read_ct = read_ct + 1+        FROM eligible e+        WHERE m.msg_id = e.msg_id+        AND m.vt <= clock_timestamp()+        RETURNING m.msg_id, m.read_ct, m.enqueued_at, m.vt, m.message, m.headers;+        $QUERY$,+        qtable, qtable, qtable, qtable, make_interval(secs => vt)+    );+    RETURN QUERY EXECUTE sql USING qty;+END;+$$ LANGUAGE plpgsql;++-- read_fifo with conditional filter+CREATE FUNCTION pgmq.read_fifo(+    queue_name TEXT,+    vt INTEGER,+    qty INTEGER,+    conditional JSONB+)+RETURNS SETOF pgmq.message_record AS $$+DECLARE+    sql TEXT;+    qtable TEXT := pgmq.format_table_name(queue_name, 'q');+BEGIN+    sql := FORMAT(+        $QUERY$+        WITH cte AS+        (+            SELECT msg_id, vt AS current_vt+            FROM pgmq.%I+            WHERE message @> $2+            ORDER BY msg_id ASC+            LIMIT $1+            FOR UPDATE SKIP LOCKED+        ),+        eligible AS (+            SELECT c.msg_id+            FROM cte c+            WHERE c.msg_id = (+                SELECT MIN(msg_id)+                FROM pgmq.%I+                WHERE vt <= clock_timestamp()+                AND message @> $2+            )+            OR NOT EXISTS (+                SELECT 1+                FROM pgmq.%I q+                WHERE q.msg_id < c.msg_id+                AND q.vt <= clock_timestamp()+                AND q.message @> $2+                AND q.msg_id NOT IN (SELECT msg_id FROM cte)+            )+        )+        UPDATE pgmq.%I m+        SET+            vt = clock_timestamp() + %L,+            read_ct = read_ct + 1+        FROM eligible e+        WHERE m.msg_id = e.msg_id+        AND m.vt <= clock_timestamp()+        RETURNING m.msg_id, m.read_ct, m.enqueued_at, m.vt, m.message, m.headers;+        $QUERY$,+        qtable, qtable, qtable, qtable, make_interval(secs => vt)+    );+    RETURN QUERY EXECUTE sql USING qty, conditional;+END;+$$ LANGUAGE plpgsql;++-- read_fifo_with_poll: reads messages in strict FIFO order with polling+CREATE FUNCTION pgmq.read_fifo_with_poll(+    queue_name TEXT,+    vt INTEGER,+    qty INTEGER,+    max_poll_seconds INTEGER DEFAULT 5,+    poll_interval_ms INTEGER DEFAULT 100+)+RETURNS SETOF pgmq.message_record AS $$+DECLARE+    r pgmq.message_record;+    stop_at TIMESTAMP;+    sql TEXT;+    qtable TEXT := pgmq.format_table_name(queue_name, 'q');+BEGIN+    stop_at := clock_timestamp() + make_interval(secs => max_poll_seconds);+    LOOP+      IF clock_timestamp() >= stop_at THEN+        RETURN;+      END IF;++      sql := FORMAT(+          $QUERY$+          WITH cte AS+          (+              SELECT msg_id, vt AS current_vt+              FROM pgmq.%I+              ORDER BY msg_id ASC+              LIMIT $1+              FOR UPDATE SKIP LOCKED+          ),+          eligible AS (+              SELECT c.msg_id+              FROM cte c+              WHERE c.msg_id = (+                  SELECT MIN(msg_id)+                  FROM pgmq.%I+                  WHERE vt <= clock_timestamp()+              )+              OR NOT EXISTS (+                  SELECT 1+                  FROM pgmq.%I q+                  WHERE q.msg_id < c.msg_id+                  AND q.vt <= clock_timestamp()+                  AND q.msg_id NOT IN (SELECT msg_id FROM cte)+              )+          )+          UPDATE pgmq.%I m+          SET+              vt = clock_timestamp() + %L,+              read_ct = read_ct + 1+          FROM eligible e+          WHERE m.msg_id = e.msg_id+          AND m.vt <= clock_timestamp()+          RETURNING m.msg_id, m.read_ct, m.enqueued_at, m.vt, m.message, m.headers;+          $QUERY$,+          qtable, qtable, qtable, qtable, make_interval(secs => vt)+      );++      FOR r IN+        EXECUTE sql USING qty+      LOOP+        RETURN NEXT r;+      END LOOP;+      IF FOUND THEN+        RETURN;+      ELSE+        PERFORM pg_sleep(poll_interval_ms::numeric / 1000);+      END IF;+    END LOOP;+END;+$$ LANGUAGE plpgsql;++-- read_fifo_with_poll with conditional filter+CREATE FUNCTION pgmq.read_fifo_with_poll(+    queue_name TEXT,+    vt INTEGER,+    qty INTEGER,+    conditional JSONB,+    max_poll_seconds INTEGER DEFAULT 5,+    poll_interval_ms INTEGER DEFAULT 100+)+RETURNS SETOF pgmq.message_record AS $$+DECLARE+    r pgmq.message_record;+    stop_at TIMESTAMP;+    sql TEXT;+    qtable TEXT := pgmq.format_table_name(queue_name, 'q');+BEGIN+    stop_at := clock_timestamp() + make_interval(secs => max_poll_seconds);+    LOOP+      IF clock_timestamp() >= stop_at THEN+        RETURN;+      END IF;++      sql := FORMAT(+          $QUERY$+          WITH cte AS+          (+              SELECT msg_id, vt AS current_vt+              FROM pgmq.%I+              WHERE message @> $2+              ORDER BY msg_id ASC+              LIMIT $1+              FOR UPDATE SKIP LOCKED+          ),+          eligible AS (+              SELECT c.msg_id+              FROM cte c+              WHERE c.msg_id = (+                  SELECT MIN(msg_id)+                  FROM pgmq.%I+                  WHERE vt <= clock_timestamp()+                  AND message @> $2+              )+              OR NOT EXISTS (+                  SELECT 1+                  FROM pgmq.%I q+                  WHERE q.msg_id < c.msg_id+                  AND q.vt <= clock_timestamp()+                  AND q.message @> $2+                  AND q.msg_id NOT IN (SELECT msg_id FROM cte)+              )+          )+          UPDATE pgmq.%I m+          SET+              vt = clock_timestamp() + %L,+              read_ct = read_ct + 1+          FROM eligible e+          WHERE m.msg_id = e.msg_id+          AND m.vt <= clock_timestamp()+          RETURNING m.msg_id, m.read_ct, m.enqueued_at, m.vt, m.message, m.headers;+          $QUERY$,+          qtable, qtable, qtable, qtable, make_interval(secs => vt)+      );++      FOR r IN+        EXECUTE sql USING qty, conditional+      LOOP+        RETURN NEXT r;+      END LOOP;+      IF FOUND THEN+        RETURN;+      ELSE+        PERFORM pg_sleep(poll_interval_ms::numeric / 1000);+      END IF;+    END LOOP;+END;+$$ LANGUAGE plpgsql;++------------------------------------------------------------+-- FIFO index functions (pgmq 1.8.0+)+------------------------------------------------------------++-- _create_fifo_index_if_not_exists+-- internal function to create GIN index on headers for better FIFO performance+CREATE OR REPLACE FUNCTION pgmq._create_fifo_index_if_not_exists(queue_name TEXT)+RETURNS void AS $$+DECLARE+    qtable TEXT := pgmq.format_table_name(queue_name, 'q');+    index_name TEXT := qtable || '_fifo_idx';+BEGIN+    -- Create GIN index on headers for efficient FIFO key lookups+    EXECUTE FORMAT(+        $QUERY$+        CREATE INDEX IF NOT EXISTS %I ON pgmq.%I USING GIN (headers);+        $QUERY$,+        index_name, qtable+    );+END;+$$ LANGUAGE plpgsql;++-- create_fifo_index+-- creates a GIN index on the headers column to improve FIFO read performance+CREATE FUNCTION pgmq.create_fifo_index(queue_name TEXT)+RETURNS void AS $$+BEGIN+    PERFORM pgmq._create_fifo_index_if_not_exists(queue_name);+END;+$$ LANGUAGE plpgsql;++-- create_fifo_indexes_all+-- creates FIFO indexes on all existing queues+CREATE FUNCTION pgmq.create_fifo_indexes_all()+RETURNS void AS $$+DECLARE+    queue_record RECORD;+BEGIN+    FOR queue_record IN SELECT queue_name FROM pgmq.meta LOOP+        PERFORM pgmq.create_fifo_index(queue_record.queue_name);+    END LOOP;+END;+$$ LANGUAGE plpgsql;++------------------------------------------------------------+-- read_grouped functions (pgmq 1.8.0+)+-- AWS SQS FIFO-style batch retrieval behavior+------------------------------------------------------------++-- read_grouped+-- reads messages with AWS SQS FIFO-style batch retrieval behavior+-- attempts to return as many messages as possible from the same message group+CREATE FUNCTION pgmq.read_grouped(+    queue_name TEXT,+    vt INTEGER,+    qty INTEGER+)+RETURNS SETOF pgmq.message_record AS $$+DECLARE+    sql TEXT;+    qtable TEXT := pgmq.format_table_name(queue_name, 'q');+BEGIN+    sql := FORMAT(+        $QUERY$+        WITH fifo_groups AS (+            -- Find the minimum msg_id for each FIFO group that's ready to be processed+            SELECT+                COALESCE(headers->>'x-pgmq-group', '_default_fifo_group') as fifo_key,+                MIN(msg_id) as min_msg_id+            FROM pgmq.%I+            WHERE vt <= clock_timestamp()+            GROUP BY COALESCE(headers->>'x-pgmq-group', '_default_fifo_group')+        ),+        locked_groups AS (+            -- Lock the first available message in each FIFO group+            SELECT+                m.msg_id,+                fg.fifo_key+            FROM pgmq.%I m+            INNER JOIN fifo_groups fg ON+                COALESCE(m.headers->>'x-pgmq-group', '_default_fifo_group') = fg.fifo_key+                AND m.msg_id = fg.min_msg_id+            WHERE m.vt <= clock_timestamp()+            ORDER BY m.msg_id ASC+            FOR UPDATE SKIP LOCKED+        ),+        group_priorities AS (+            -- Assign priority to groups based on their oldest message+            SELECT+                fifo_key,+                msg_id as min_msg_id,+                ROW_NUMBER() OVER (ORDER BY msg_id) as group_priority+            FROM locked_groups+        ),+        available_messages AS (+            -- Get messages prioritizing filling batch from earliest group first+            SELECT+                m.msg_id,+                gp.group_priority,+                ROW_NUMBER() OVER (PARTITION BY gp.fifo_key ORDER BY m.msg_id) as msg_rank_in_group+            FROM pgmq.%I m+            INNER JOIN group_priorities gp ON+                COALESCE(m.headers->>'x-pgmq-group', '_default_fifo_group') = gp.fifo_key+            WHERE m.vt <= clock_timestamp()+            AND m.msg_id >= gp.min_msg_id  -- Only messages from min_msg_id onwards in each group+            AND NOT EXISTS (+                -- Ensure no earlier message in this group is currently being processed+                SELECT 1+                FROM pgmq.%I m2+                WHERE COALESCE(m2.headers->>'x-pgmq-group', '_default_fifo_group') =+                      COALESCE(m.headers->>'x-pgmq-group', '_default_fifo_group')+                AND m2.vt > clock_timestamp()+                AND m2.msg_id < m.msg_id+            )+        ),+        batch_selection AS (+            -- Select messages to fill batch, prioritizing earliest group+            SELECT+                msg_id,+                ROW_NUMBER() OVER (ORDER BY group_priority, msg_rank_in_group) as overall_rank+            FROM available_messages+        ),+        selected_messages AS (+            -- Limit to requested quantity+            SELECT msg_id+            FROM batch_selection+            WHERE overall_rank <= $1+            ORDER BY msg_id+            FOR UPDATE SKIP LOCKED+        )+        UPDATE pgmq.%I m+        SET+            vt = clock_timestamp() + %L,+            read_ct = read_ct + 1+        FROM selected_messages sm+        WHERE m.msg_id = sm.msg_id+        RETURNING m.msg_id, m.read_ct, m.enqueued_at, m.vt, m.message, m.headers;+        $QUERY$,+        qtable, qtable, qtable, qtable, qtable, make_interval(secs => vt)+    );+    RETURN QUERY EXECUTE sql USING qty;+END;+$$ LANGUAGE plpgsql;++-- read_grouped_with_poll+-- reads messages with AWS SQS FIFO-style batch retrieval behavior, with polling support+CREATE FUNCTION pgmq.read_grouped_with_poll(+    queue_name TEXT,+    vt INTEGER,+    qty INTEGER,+    max_poll_seconds INTEGER DEFAULT 5,+    poll_interval_ms INTEGER DEFAULT 100+)+RETURNS SETOF pgmq.message_record AS $$+DECLARE+    r pgmq.message_record;+    stop_at TIMESTAMP;+BEGIN+    stop_at := clock_timestamp() + make_interval(secs => max_poll_seconds);+    LOOP+      IF (SELECT clock_timestamp() >= stop_at) THEN+        RETURN;+      END IF;++      FOR r IN+        SELECT * FROM pgmq.read_grouped(queue_name, vt, qty)+      LOOP+        RETURN NEXT r;+      END LOOP;+      IF FOUND THEN+        RETURN;+      ELSE+        PERFORM pg_sleep(poll_interval_ms::numeric / 1000);+      END IF;+    END LOOP;+END;+$$ LANGUAGE plpgsql;++------------------------------------------------------------+-- Round-robin FIFO functions (pgmq 1.9.0+)+------------------------------------------------------------++-- read_grouped_rr (round-robin)+-- reads messages while preserving FIFO within groups and interleaving across groups (layered round-robin)+CREATE FUNCTION pgmq.read_grouped_rr(+    queue_name TEXT,+    vt INTEGER,+    qty INTEGER+)+RETURNS SETOF pgmq.message_record AS $$+DECLARE+    sql TEXT;+    qtable TEXT := pgmq.format_table_name(queue_name, 'q');+BEGIN+    sql := FORMAT(+        $QUERY$+        WITH fifo_groups AS (+            -- Determine the absolute head (oldest) message id per FIFO group, regardless of visibility+            SELECT+                COALESCE(headers->>'x-pgmq-group', '_default_fifo_group') AS fifo_key,+                MIN(msg_id) AS head_msg_id+            FROM pgmq.%1$I+            GROUP BY COALESCE(headers->>'x-pgmq-group', '_default_fifo_group')+        ),+        eligible_groups AS (+            -- Only groups whose head message is currently visible+            -- Acquire a transaction-level advisory lock per group to prevent concurrent selection+            SELECT+                g.fifo_key,+                g.head_msg_id,+                ROW_NUMBER() OVER (ORDER BY g.head_msg_id) AS group_priority+            FROM fifo_groups g+            JOIN pgmq.%2$I h ON h.msg_id = g.head_msg_id+            WHERE h.vt <= clock_timestamp()+              AND pg_try_advisory_xact_lock(pg_catalog.hashtextextended(g.fifo_key, 0))+        ),+        available_messages AS (+            -- All currently visible messages starting at the head for each eligible group+            SELECT+                m.msg_id,+                eg.group_priority,+                ROW_NUMBER() OVER (+                    PARTITION BY eg.fifo_key+                    ORDER BY m.msg_id+                ) AS msg_rank_in_group+            FROM pgmq.%3$I m+            JOIN eligible_groups eg+              ON COALESCE(m.headers->>'x-pgmq-group', '_default_fifo_group') = eg.fifo_key+            WHERE m.vt <= clock_timestamp()+              AND m.msg_id >= eg.head_msg_id+        ),+        ordered_messages AS (+            -- Layered round-robin: take rank 1 of all groups by group_priority, then rank 2, etc.+            -- Assign selection order before locking+            SELECT msg_id, ROW_NUMBER() OVER (ORDER BY msg_rank_in_group, group_priority) as selection_order+            FROM available_messages+        ),+        selected_messages AS (+            -- Lock the messages in the correct order, preserving selection_order+            SELECT om.msg_id, om.selection_order+            FROM ordered_messages om+            JOIN pgmq.%4$I m ON m.msg_id = om.msg_id+            WHERE om.selection_order <= $1+            ORDER BY om.selection_order+            FOR UPDATE OF m SKIP LOCKED+        ),+        updated_messages AS (+            UPDATE pgmq.%5$I m+            SET+                vt = clock_timestamp() + %6$L,+                read_ct = read_ct + 1+            FROM selected_messages sm+            WHERE m.msg_id = sm.msg_id+              AND m.vt <= clock_timestamp() -- final guard to avoid duplicate reads under races+            RETURNING m.msg_id, m.read_ct, m.enqueued_at, m.vt, m.message, m.headers, sm.selection_order+        )+        SELECT msg_id, read_ct, enqueued_at, vt, message, headers+        FROM updated_messages+        ORDER BY selection_order;+        $QUERY$,+        qtable, qtable, qtable, qtable, qtable, make_interval(secs => vt)+    );+    RETURN QUERY EXECUTE sql USING qty;+END;+$$ LANGUAGE plpgsql;++-- read_grouped_rr_with_poll+-- reads messages using round-robin layering across groups, with polling support+CREATE FUNCTION pgmq.read_grouped_rr_with_poll(+    queue_name TEXT,+    vt INTEGER,+    qty INTEGER,+    max_poll_seconds INTEGER DEFAULT 5,+    poll_interval_ms INTEGER DEFAULT 100+)+RETURNS SETOF pgmq.message_record AS $$+DECLARE+    r pgmq.message_record;+    stop_at TIMESTAMP;+BEGIN+    stop_at := clock_timestamp() + make_interval(secs => max_poll_seconds);+    LOOP+      IF (SELECT clock_timestamp() >= stop_at) THEN+        RETURN;+      END IF;++      FOR r IN+        SELECT * FROM pgmq.read_grouped_rr(queue_name, vt, qty)+      LOOP+        RETURN NEXT r;+      END LOOP;+      IF FOUND THEN+        RETURN;+      ELSE+        PERFORM pg_sleep(poll_interval_ms::numeric / 1000);+      END IF;+    END LOOP;+END;+$$ LANGUAGE plpgsql;
+ pgmq-migration.cabal view
@@ -0,0 +1,69 @@+cabal-version:      3.4+name:               pgmq-migration+version:            0.1.0.0+synopsis:           PGMQ schema migrations without PostgreSQL extension+description:+  Installs the PGMQ schema into PostgreSQL without requiring the pgmq extension.+  Uses hasql-migration for tracking applied migrations.++homepage:           https://github.com/topagentnetwork/pgmq-hs+license:            MIT+license-file:       LICENSE+author:             Nadeem Bitar+maintainer:         nadeem@topagentnetwork.com+category:           Database+build-type:         Simple+extra-doc-files:    CHANGELOG.md+extra-source-files: database/**/*.sql++common warnings+  ghc-options:+    -Wall -Wcompat -Widentities -Wincomplete-uni-patterns+    -Wincomplete-record-updates -Wredundant-constraints+    -fhide-source-paths -Wmissing-export-lists -Wpartial-fields+    -Wmissing-deriving-strategies++library+  import:             warnings+  exposed-modules:+    Pgmq.Migration+    Pgmq.Migration.Migrations+    Pgmq.Migration.Migrations.V1_10_0+    Pgmq.Migration.Migrations.V1_9_0+    Pgmq.Migration.Migrations.V1_9_0_to_V1_10_0+    Pgmq.Migration.Sessions+    Pgmq.Migration.Statements+    Pgmq.Migration.Transactions++  default-extensions:+    ImportQualifiedPost+    OverloadedStrings++  build-depends:+    , base               >=4.18    && <5+    , bytestring         ^>=0.12+    , file-embed         ^>=0.0.16+    , hasql              ^>=1.10+    , hasql-migration    ^>=0.3.1+    , hasql-transaction  ^>=1.2+    , text               ^>=2.1+    , transformers       ^>=0.6++  hs-source-dirs:     src+  default-language:   GHC2024++test-suite pgmq-migration-test+  import:           warnings+  default-language: GHC2024+  type:             exitcode-stdio-1.0+  hs-source-dirs:   test+  main-is:          Main.hs+  build-depends:+    , base               >=4.18  && <5+    , ephemeral-pg       >=0.2.1+    , hasql+    , hasql-migration+    , hasql-transaction+    , pgmq-migration+    , tasty              ^>=1.5+    , tasty-hunit        ^>=0.10
+ src/Pgmq/Migration.hs view
@@ -0,0 +1,147 @@+-- | PGMQ schema migration support+--+-- This module provides functions to install the PGMQ schema into PostgreSQL+-- without requiring the pgmq extension. It uses hasql-migration for tracking+-- applied migrations.+--+-- == Fresh Installation+--+-- For new projects, use 'migrate' to install the complete PGMQ schema:+--+-- @+-- import Hasql.Connection (acquire)+-- import Hasql.Session (run)+-- import Pgmq.Migration (migrate)+--+-- main :: IO ()+-- main = do+--   Right conn <- acquire "host=localhost dbname=mydb"+--   result <- run migrate conn+--   case result of+--     Right (Right ()) -> putStrLn "Migration successful"+--     Right (Left err) -> print err+--     Left sessionErr  -> print sessionErr+-- @+--+-- == Upgrading Existing Installations+--+-- For projects that previously installed PGMQ via this package (e.g., at v1.9.0),+-- use 'upgrade' to apply only the incremental changes:+--+-- @+-- import Hasql.Connection (acquire)+-- import Hasql.Session (run)+-- import Pgmq.Migration (upgrade)+--+-- main :: IO ()+-- main = do+--   Right conn <- acquire "host=localhost dbname=mydb"+--   result <- run upgrade conn+--   case result of+--     Right (Right ()) -> putStrLn "Upgrade successful"+--     Right (Left err) -> print err+--     Left sessionErr  -> print sessionErr+-- @+--+-- The hasql-migration library tracks which migrations have been applied,+-- so 'upgrade' will only apply migrations that haven't run yet.+--+-- == Which Function Should I Use?+--+-- * __New project, fresh database__: Use 'migrate'+-- * __Existing project using pgmq-migration__: Use 'upgrade'+-- * __Not sure__: Use 'upgrade' - it's safe on fresh databases too+--   (it will apply all needed migrations)+module Pgmq.Migration+  ( -- * Migration Operations+    migrate,+    upgrade,+    validate,+    getMigrations,++    -- * Migration Info+    version,+    migrations,+    upgradeMigrations,++    -- * Re-exports+    MigrationCommand,+    MigrationError (..),+    SchemaMigration (..),+  )+where++import Control.Monad (foldM)+import Hasql.Migration (MigrationCommand, MigrationError (..), SchemaMigration (..))+import Hasql.Migration qualified as Migration+import Hasql.Session (Session)+import Pgmq.Migration.Migrations qualified as Migrations+import Pgmq.Migration.Sessions qualified as Sessions++-- | Current version of the PGMQ schema+version :: String+version = Migrations.version++-- | Full migration commands for the current version.+--+-- Use this for fresh installations. Applies the complete schema.+migrations :: [MigrationCommand]+migrations = Migrations.migrations++-- | Incremental upgrade migrations.+--+-- Use this to upgrade existing installations that were set up via this package.+-- Contains only the delta migrations (e.g., v1.9.0 -> v1.10.0).+upgradeMigrations :: [MigrationCommand]+upgradeMigrations = Migrations.upgradeMigrations++-- | Run all PGMQ migrations for a fresh installation.+--+-- This function is idempotent - it will only run migrations that haven't+-- been applied yet. Returns 'Left' 'MigrationError' if a migration fails.+--+-- Use this for new projects with a fresh database.+migrate :: Session (Either MigrationError ())+migrate = runMigrations migrations++-- | Run upgrade migrations for existing installations.+--+-- This function applies only the incremental migrations needed to upgrade+-- from a previous version installed via this package (e.g., v1.9.0 -> v1.10.0).+--+-- It is idempotent and safe to run on any database - migrations that have+-- already been applied will be skipped.+--+-- Use this for existing projects that need to upgrade their PGMQ schema.+upgrade :: Session (Either MigrationError ())+upgrade = runMigrations upgradeMigrations++-- | Run a list of migrations+runMigrations :: [MigrationCommand] -> Session (Either MigrationError ())+runMigrations cmds = foldM runIfOk (Right ()) cmds+  where+    runIfOk :: Either MigrationError () -> MigrationCommand -> Session (Either MigrationError ())+    runIfOk (Left err) _ = pure (Left err)+    runIfOk (Right ()) cmd = Sessions.runMigrationSession cmd++-- | Validate all migrations without applying them+--+-- Checks that previously applied migrations haven't changed.+-- Returns a list of validation errors, or an empty list if validation passes.+validate :: Session [MigrationError]+validate = do+  results <- mapM validateOne validationCommands+  pure $ concat results+  where+    validationCommands = map Migration.MigrationValidation migrations++    validateOne :: MigrationCommand -> Session [MigrationError]+    validateOne cmd = do+      result <- Sessions.runMigrationSession cmd+      pure $ case result of+        Left err -> [err]+        Right () -> []++-- | Get all applied migrations+getMigrations :: Session [SchemaMigration]+getMigrations = Sessions.getMigrationsSession
+ src/Pgmq/Migration/Migrations.hs view
@@ -0,0 +1,63 @@+-- | Aggregation of all PGMQ migrations+--+-- This module provides two migration paths:+--+-- == Fresh Installation+--+-- For new projects, use 'migrations' which installs the latest PGMQ schema directly:+--+-- @+-- import Pgmq.Migration.Migrations (migrations)+--+-- -- Apply migrations to a fresh database+-- runMigrations conn migrations+-- @+--+-- == Upgrading Existing Installations+--+-- For projects that previously installed PGMQ via this package, use 'upgradeMigrations'+-- to apply only the incremental changes:+--+-- @+-- import Pgmq.Migration.Migrations (upgradeMigrations)+--+-- -- Apply only upgrade migrations (v1.9.0 -> v1.10.0, etc.)+-- runMigrations conn upgradeMigrations+-- @+--+-- The hasql-migration library tracks which migrations have been applied,+-- so running 'upgradeMigrations' on a database with v1.9.0 installed will+-- only apply the delta migrations needed to reach the current version.+module Pgmq.Migration.Migrations+  ( version,+    migrations,+    upgradeMigrations,+  )+where++import Hasql.Migration (MigrationCommand (..))+import Pgmq.Migration.Migrations.V1_10_0 qualified as V1_10_0+import Pgmq.Migration.Migrations.V1_9_0_to_V1_10_0 qualified as V1_9_0_to_V1_10_0++-- | Current version of the PGMQ schema+version :: String+version = V1_10_0.version++-- | Full migration commands for the current version.+--+-- Use this for fresh installations. Applies the complete v1.10.0 schema.+migrations :: [MigrationCommand]+migrations = V1_10_0.migrations++-- | Incremental upgrade migrations.+--+-- Use this to upgrade existing installations that were set up via this package.+-- Contains only the delta migrations (e.g., v1.9.0 -> v1.10.0).+--+-- The hasql-migration library will skip migrations that have already been applied,+-- so this is safe to run on any database regardless of its current version.+upgradeMigrations :: [MigrationCommand]+upgradeMigrations =+  [ MigrationInitialization+  ]+    ++ V1_9_0_to_V1_10_0.migrations
+ src/Pgmq/Migration/Migrations/V1_10_0.hs view
@@ -0,0 +1,63 @@+{-# LANGUAGE TemplateHaskell #-}++-- | PGMQ v1.10.0 migrations+-- Main change: Added last_read_at column to track when messages were last read+module Pgmq.Migration.Migrations.V1_10_0+  ( version,+    migrations,+  )+where++import Data.ByteString (ByteString)+import Data.FileEmbed (embedFile)+import Hasql.Migration (MigrationCommand (..))++-- | Version string for this migration set+version :: String+version = "v1.10.0"++-- | All migration commands for v1.10.0+migrations :: [MigrationCommand]+migrations =+  [ MigrationInitialization,+    MigrationScript "pgmq_v1.10.0_01_schema" schema,+    MigrationScript "pgmq_v1.10.0_02_tables" tables,+    MigrationScript "pgmq_v1.10.0_03_types" types,+    MigrationScript "pgmq_v1.10.0_04_core_functions" coreFunctions,+    MigrationScript "pgmq_v1.10.0_05_queue_management" queueManagement,+    MigrationScript "pgmq_v1.10.0_06_message_ops" messageOps,+    MigrationScript "pgmq_v1.10.0_07_metrics" metrics,+    MigrationScript "pgmq_v1.10.0_08_partitioning" partitioning,+    MigrationScript "pgmq_v1.10.0_09_notifications" notifications,+    MigrationScript "pgmq_v1.10.0_10_fifo" fifo+  ]++schema :: ByteString+schema = $(embedFile "database/v1.10.0/01_schema.sql")++tables :: ByteString+tables = $(embedFile "database/v1.10.0/02_tables.sql")++types :: ByteString+types = $(embedFile "database/v1.10.0/03_types.sql")++coreFunctions :: ByteString+coreFunctions = $(embedFile "database/v1.10.0/04_core_functions.sql")++queueManagement :: ByteString+queueManagement = $(embedFile "database/v1.10.0/05_queue_management.sql")++messageOps :: ByteString+messageOps = $(embedFile "database/v1.10.0/06_message_ops.sql")++metrics :: ByteString+metrics = $(embedFile "database/v1.10.0/07_metrics.sql")++partitioning :: ByteString+partitioning = $(embedFile "database/v1.10.0/08_partitioning.sql")++notifications :: ByteString+notifications = $(embedFile "database/v1.10.0/09_notifications.sql")++fifo :: ByteString+fifo = $(embedFile "database/v1.10.0/10_fifo.sql")
+ src/Pgmq/Migration/Migrations/V1_9_0.hs view
@@ -0,0 +1,62 @@+{-# LANGUAGE TemplateHaskell #-}++-- | PGMQ v1.9.0 migrations+module Pgmq.Migration.Migrations.V1_9_0+  ( version,+    migrations,+  )+where++import Data.ByteString (ByteString)+import Data.FileEmbed (embedFile)+import Hasql.Migration (MigrationCommand (..))++-- | Version string for this migration set+version :: String+version = "v1.9.0"++-- | All migration commands for v1.9.0+migrations :: [MigrationCommand]+migrations =+  [ MigrationInitialization,+    MigrationScript "pgmq_v1.9.0_01_schema" schema,+    MigrationScript "pgmq_v1.9.0_02_tables" tables,+    MigrationScript "pgmq_v1.9.0_03_types" types,+    MigrationScript "pgmq_v1.9.0_04_core_functions" coreFunctions,+    MigrationScript "pgmq_v1.9.0_05_queue_management" queueManagement,+    MigrationScript "pgmq_v1.9.0_06_message_ops" messageOps,+    MigrationScript "pgmq_v1.9.0_07_metrics" metrics,+    MigrationScript "pgmq_v1.9.0_08_partitioning" partitioning,+    MigrationScript "pgmq_v1.9.0_09_notifications" notifications,+    MigrationScript "pgmq_v1.9.0_10_fifo" fifo+  ]++schema :: ByteString+schema = $(embedFile "database/v1.9.0/01_schema.sql")++tables :: ByteString+tables = $(embedFile "database/v1.9.0/02_tables.sql")++types :: ByteString+types = $(embedFile "database/v1.9.0/03_types.sql")++coreFunctions :: ByteString+coreFunctions = $(embedFile "database/v1.9.0/04_core_functions.sql")++queueManagement :: ByteString+queueManagement = $(embedFile "database/v1.9.0/05_queue_management.sql")++messageOps :: ByteString+messageOps = $(embedFile "database/v1.9.0/06_message_ops.sql")++metrics :: ByteString+metrics = $(embedFile "database/v1.9.0/07_metrics.sql")++partitioning :: ByteString+partitioning = $(embedFile "database/v1.9.0/08_partitioning.sql")++notifications :: ByteString+notifications = $(embedFile "database/v1.9.0/09_notifications.sql")++fifo :: ByteString+fifo = $(embedFile "database/v1.9.0/10_fifo.sql")
+ src/Pgmq/Migration/Migrations/V1_9_0_to_V1_10_0.hs view
@@ -0,0 +1,21 @@+{-# LANGUAGE TemplateHaskell #-}++-- | PGMQ v1.9.0 to v1.10.0 migration+-- Changes: Added last_read_at column to track when messages were last read+module Pgmq.Migration.Migrations.V1_9_0_to_V1_10_0+  ( migrations,+  )+where++import Data.ByteString (ByteString)+import Data.FileEmbed (embedFile)+import Hasql.Migration (MigrationCommand (..))++-- | Migration commands to upgrade from v1.9.0 to v1.10.0+migrations :: [MigrationCommand]+migrations =+  [ MigrationScript "pgmq_v1.9.0_to_v1.10.0_upgrade" upgrade+  ]++upgrade :: ByteString+upgrade = $(embedFile "database/migrations/v1.9.0_to_v1.10.0.sql")
+ src/Pgmq/Migration/Sessions.hs view
@@ -0,0 +1,44 @@+-- | Session-level migration operations+module Pgmq.Migration.Sessions+  ( runMigrationSession,+    getMigrationsSession,+  )+where++import Control.Monad.Trans.Except (ExceptT (..), runExceptT)+import Data.Text (Text)+import Hasql.Migration (MigrationCommand, MigrationError, SchemaMigration)+import Hasql.Session (Session)+import Hasql.Session qualified as Session+import Hasql.Transaction.Sessions qualified as Transaction.Sessions+import Pgmq.Migration.Statements qualified as Statements+import Pgmq.Migration.Transactions qualified as Transactions++-- | Run a single migration command in a session+-- Returns Left MigrationError on failure, Right () on success+runMigrationSession :: MigrationCommand -> Session (Either MigrationError ())+runMigrationSession cmd = runExceptT $ do+  -- Get the current search path+  searchPath <- ExceptT $ fmap Right $ Session.statement () Statements.getSearchPath+  -- Run the migration in a transaction with proper isolation+  result <- ExceptT $ fmap Right $ runMigrationTransaction searchPath cmd+  -- Convert Maybe MigrationError to Either+  case result of+    Nothing -> pure ()+    Just err -> ExceptT $ pure $ Left err++-- | Run migration in a transaction with serializable isolation+runMigrationTransaction :: Text -> MigrationCommand -> Session (Maybe MigrationError)+runMigrationTransaction searchPath cmd =+  Transaction.Sessions.transaction+    Transaction.Sessions.Serializable+    Transaction.Sessions.Write+    (Transactions.runMigrationWithSearchPath searchPath cmd)++-- | Get all applied migrations+getMigrationsSession :: Session [SchemaMigration]+getMigrationsSession =+  Transaction.Sessions.transaction+    Transaction.Sessions.ReadCommitted+    Transaction.Sessions.Read+    Transactions.getMigrationsTransaction
+ src/Pgmq/Migration/Statements.hs view
@@ -0,0 +1,29 @@+{-# LANGUAGE QuasiQuotes #-}++-- | Search path management statements for migrations+module Pgmq.Migration.Statements+  ( getSearchPath,+    setSearchPath,+  )+where++import Data.Text (Text)+import Hasql.Decoders qualified as Decoders+import Hasql.Encoders qualified as Encoders+import Hasql.Statement (Statement, unpreparable)++-- | Get the current search path+getSearchPath :: Statement () Text+getSearchPath =+  unpreparable+    "SHOW search_path"+    Encoders.noParams+    (Decoders.singleRow (Decoders.column (Decoders.nonNullable Decoders.text)))++-- | Set the search path+setSearchPath :: Statement Text ()+setSearchPath =+  unpreparable+    "SELECT set_config('search_path', $1, true)"+    (Encoders.param (Encoders.nonNullable Encoders.text))+    Decoders.noResult
+ src/Pgmq/Migration/Transactions.hs view
@@ -0,0 +1,26 @@+-- | Transaction-level migration operations+module Pgmq.Migration.Transactions+  ( runMigrationWithSearchPath,+    getMigrationsTransaction,+  )+where++import Control.Monad (void)+import Data.Text (Text)+import Hasql.Migration (MigrationCommand, MigrationError, SchemaMigration)+import Hasql.Migration qualified as Migration+import Hasql.Transaction (Transaction)+import Hasql.Transaction qualified as Transaction+import Pgmq.Migration.Statements qualified as Statements++-- | Run a migration command, preserving and restoring the search path+runMigrationWithSearchPath :: Text -> MigrationCommand -> Transaction (Maybe MigrationError)+runMigrationWithSearchPath searchPath cmd = do+  -- Set the search path for this transaction+  void $ Transaction.statement searchPath Statements.setSearchPath+  -- Run the migration+  Migration.runMigration cmd++-- | Get all applied migrations+getMigrationsTransaction :: Transaction [SchemaMigration]+getMigrationsTransaction = Migration.getMigrations
+ test/Main.hs view
@@ -0,0 +1,250 @@+{-# LANGUAGE OverloadedStrings #-}++module Main (main) where++import Control.Monad (foldM)+import Data.List (isInfixOf)+import EphemeralPg+  ( StartError,+    connectionSettings,+    withCached,+  )+import Hasql.Connection qualified as Connection+import Hasql.Decoders qualified as Decoders+import Hasql.Encoders qualified as Encoders+import Hasql.Migration (MigrationCommand, MigrationError)+import Hasql.Session (Session)+import Hasql.Session qualified as Session+import Hasql.Statement (preparable)+import Hasql.Statement qualified as Statement+import Pgmq.Migration qualified as Migration+import Pgmq.Migration.Migrations.V1_9_0 qualified as V1_9_0+import Pgmq.Migration.Sessions qualified as Sessions+import Test.Tasty (TestTree, defaultMain, testGroup)+import Test.Tasty.HUnit (assertBool, assertFailure, testCase, (@?=))++main :: IO ()+main = do+  result <- withCached $ \db -> do+    let connSettings = connectionSettings db+    connResult <- Connection.acquire connSettings+    case connResult of+      Left err -> error $ "Failed to connect: " <> show err+      Right conn ->+        defaultMain (tests conn)+  case result of+    Left startErr -> error $ "Failed to start temp database: " <> show startErr+    Right () -> pure ()++tests :: Connection.Connection -> TestTree+tests conn =+  testGroup+    "pgmq-migration"+    [ testGroup+        "fresh install"+        [ testCase "migrate on fresh database succeeds" (testMigrateFresh conn),+          testCase "migrate is idempotent" (testMigrateIdempotent conn),+          testCase "getMigrations returns applied migrations" (testGetMigrations conn),+          testCase "version is v1.10.0" testVersion+        ],+      testGroup+        "upgrade"+        [ testCase "upgrade from v1.9.0 succeeds" (testUpgradeFromV1_9_0 conn),+          testCase "upgrade is idempotent" (testUpgradeIdempotent conn),+          testCase "upgrade adds last_read_at column" (testUpgradeAddsLastReadAt conn)+        ]+    ]++-- | Reset the database to a clean state by dropping the pgmq schema+-- and migration tracking table+resetDb :: Connection.Connection -> IO ()+resetDb conn = do+  resetResult <- Connection.use conn resetSession+  case resetResult of+    Left err -> error $ "Failed to reset database: " <> show err+    Right () -> pure ()+  where+    resetSession :: Session ()+    resetSession = do+      Session.statement () dropPgmqSchema+      Session.statement () dropMigrationTable++    dropPgmqSchema :: Statement.Statement () ()+    dropPgmqSchema =+      preparable+        "DROP SCHEMA IF EXISTS pgmq CASCADE"+        Encoders.noParams+        Decoders.noResult++    dropMigrationTable :: Statement.Statement () ()+    dropMigrationTable =+      preparable+        "DROP TABLE IF EXISTS public.schema_migrations"+        Encoders.noParams+        Decoders.noResult++-- | Run a test with a clean database+withCleanDb :: Connection.Connection -> (Connection.Connection -> IO ()) -> IO ()+withCleanDb conn action = do+  resetDb conn+  action conn++testMigrateFresh :: Connection.Connection -> IO ()+testMigrateFresh conn = withCleanDb conn $ \c -> do+  migResult <- Connection.use c Migration.migrate+  case migResult of+    Left sessionErr -> assertFailure $ "Session error: " <> show sessionErr+    Right (Left migrationErr) -> assertFailure $ "Migration error: " <> show migrationErr+    Right (Right ()) -> pure ()++testMigrateIdempotent :: Connection.Connection -> IO ()+testMigrateIdempotent conn = withCleanDb conn $ \c -> do+  -- Run migration first time+  result1 <- Connection.use c Migration.migrate+  case result1 of+    Left sessionErr -> assertFailure $ "First migration session error: " <> show sessionErr+    Right (Left migrationErr) -> assertFailure $ "First migration error: " <> show migrationErr+    Right (Right ()) -> pure ()++  -- Run migration second time - should succeed without error+  result2 <- Connection.use c Migration.migrate+  case result2 of+    Left sessionErr -> assertFailure $ "Second migration session error: " <> show sessionErr+    Right (Left migrationErr) -> assertFailure $ "Second migration error: " <> show migrationErr+    Right (Right ()) -> pure ()++testGetMigrations :: Connection.Connection -> IO ()+testGetMigrations conn = withCleanDb conn $ \c -> do+  -- Run migrations first+  _ <- Connection.use c Migration.migrate++  -- Get applied migrations+  migrationsResult <- Connection.use c Migration.getMigrations+  case migrationsResult of+    Left sessionErr -> assertFailure $ "Session error: " <> show sessionErr+    Right appliedMigrations -> do+      -- Should have applied the initialization + 10 SQL migrations+      assertBool "Should have applied migrations" (length appliedMigrations >= 10)++testVersion :: IO ()+testVersion =+  Migration.version @?= "v1.10.0"++-- | Helper to run a list of migrations+runMigrations :: [MigrationCommand] -> Session (Either MigrationError ())+runMigrations cmds = foldM runIfOk (Right ()) cmds+  where+    runIfOk :: Either MigrationError () -> MigrationCommand -> Session (Either MigrationError ())+    runIfOk (Left err) _ = pure (Left err)+    runIfOk (Right ()) cmd = Sessions.runMigrationSession cmd++-- | Install v1.9.0 schema+installV1_9_0 :: Session (Either MigrationError ())+installV1_9_0 = runMigrations V1_9_0.migrations++testUpgradeFromV1_9_0 :: Connection.Connection -> IO ()+testUpgradeFromV1_9_0 conn = withCleanDb conn $ \c -> do+  -- First install v1.9.0 schema+  v190Result <- Connection.use c installV1_9_0+  case v190Result of+    Left sessionErr -> assertFailure $ "v1.9.0 install session error: " <> show sessionErr+    Right (Left migrationErr) -> assertFailure $ "v1.9.0 install migration error: " <> show migrationErr+    Right (Right ()) -> pure ()++  -- Now run upgrade+  upgradeResult <- Connection.use c Migration.upgrade+  case upgradeResult of+    Left sessionErr -> assertFailure $ "Upgrade session error: " <> show sessionErr+    Right (Left migrationErr) -> assertFailure $ "Upgrade migration error: " <> show migrationErr+    Right (Right ()) -> pure ()++  -- Verify upgrade migration was applied+  migrationsResult <- Connection.use c Migration.getMigrations+  case migrationsResult of+    Left sessionErr -> assertFailure $ "getMigrations session error: " <> show sessionErr+    Right appliedMigrations -> do+      let migrationNames = map show appliedMigrations+          hasUpgrade = any ("v1.9.0_to_v1.10.0" `isInfixOf`) migrationNames+      assertBool "Should have applied upgrade migration" hasUpgrade++testUpgradeIdempotent :: Connection.Connection -> IO ()+testUpgradeIdempotent conn = withCleanDb conn $ \c -> do+  -- First install v1.9.0 schema+  v190Result <- Connection.use c installV1_9_0+  case v190Result of+    Left sessionErr -> assertFailure $ "v1.9.0 install session error: " <> show sessionErr+    Right (Left migrationErr) -> assertFailure $ "v1.9.0 install migration error: " <> show migrationErr+    Right (Right ()) -> pure ()++  -- Run upgrade first time+  upgrade1Result <- Connection.use c Migration.upgrade+  case upgrade1Result of+    Left sessionErr -> assertFailure $ "First upgrade session error: " <> show sessionErr+    Right (Left migrationErr) -> assertFailure $ "First upgrade migration error: " <> show migrationErr+    Right (Right ()) -> pure ()++  -- Run upgrade second time - should succeed+  upgrade2Result <- Connection.use c Migration.upgrade+  case upgrade2Result of+    Left sessionErr -> assertFailure $ "Second upgrade session error: " <> show sessionErr+    Right (Left migrationErr) -> assertFailure $ "Second upgrade migration error: " <> show migrationErr+    Right (Right ()) -> pure ()++testUpgradeAddsLastReadAt :: Connection.Connection -> IO ()+testUpgradeAddsLastReadAt conn = withCleanDb conn $ \c -> do+  -- First install v1.9.0 schema+  v190Result <- Connection.use c installV1_9_0+  case v190Result of+    Left sessionErr -> assertFailure $ "v1.9.0 install session error: " <> show sessionErr+    Right (Left migrationErr) -> assertFailure $ "v1.9.0 install migration error: " <> show migrationErr+    Right (Right ()) -> pure ()++  -- Create a queue to test the schema+  createQueueResult <- Connection.use c createTestQueue+  case createQueueResult of+    Left sessionErr -> assertFailure $ "Create queue session error: " <> show sessionErr+    Right () -> pure ()++  -- Run upgrade+  upgradeResult <- Connection.use c Migration.upgrade+  case upgradeResult of+    Left sessionErr -> assertFailure $ "Upgrade session error: " <> show sessionErr+    Right (Left migrationErr) -> assertFailure $ "Upgrade migration error: " <> show migrationErr+    Right (Right ()) -> pure ()++  -- Verify last_read_at column exists in the queue table+  checkResult <- Connection.use c checkLastReadAtColumn+  case checkResult of+    Left sessionErr -> assertFailure $ "Check column session error: " <> show sessionErr+    Right hasColumn ->+      assertBool "Queue table should have last_read_at column after upgrade" hasColumn++-- | Create a test queue using pgmq.create+createTestQueue :: Session ()+createTestQueue = Session.statement () createQueueStmt+  where+    createQueueStmt :: Statement.Statement () ()+    createQueueStmt =+      preparable sql encoder decoder+      where+        sql = "SELECT pgmq.create('test_upgrade_queue')"+        encoder = Encoders.noParams+        decoder = Decoders.noResult++-- | Check if last_read_at column exists in the test queue table+checkLastReadAtColumn :: Session Bool+checkLastReadAtColumn = Session.statement () checkColumnStmt+  where+    checkColumnStmt :: Statement.Statement () Bool+    checkColumnStmt =+      preparable sql encoder decoder+      where+        sql =+          "SELECT EXISTS ( \+          \  SELECT 1 FROM information_schema.columns \+          \  WHERE table_schema = 'pgmq' \+          \  AND table_name = 'q_test_upgrade_queue' \+          \  AND column_name = 'last_read_at' \+          \)"+        encoder = Encoders.noParams+        decoder = Decoders.singleRow (Decoders.column (Decoders.nonNullable Decoders.bool))