pgmq-migration 0.4.0.1 → 0.5.0.0
raw patch · 5 files changed
+404/−17 lines, 5 files
Files
- CHANGELOG.md +44/−0
- migrations/0003-notify-crash-safety-and-locking.sql +297/−0
- migrations/manifest +1/−0
- pgmq-migration.cabal +1/−1
- test/Main.hs +61/−16
CHANGELOG.md view
@@ -1,5 +1,49 @@ # Changelog for pgmq-migration +## 0.5.0.0 -- 2026-08-06++### Bug Fixes++Migration `0003-notify-crash-safety-and-locking.sql` re-creates three functions. It is+appended to the ledger; `0001` and `0002` are unchanged.++* Insert notifications no longer stop permanently after a PostgreSQL crash.+ `pgmq.notify_insert_throttle` is `UNLOGGED`, so crash recovery truncates it, and+ `pgmq.notify_queue_listeners` notified only when its throttle `UPDATE` matched a row.+ Zero updated rows had two causes — throttled, and row absent — and the trigger treated+ both as "stay quiet", so after a crash it fired, matched nothing, and silently never+ notified again until an application restart re-enabled notify. Sends succeeded and+ messages accumulated while listeners starved. A `NOT EXISTS` probe on the already-failed+ branch now tells the two cases apart and notifies unthrottled until the next reconcile+ restores the configured interval. The row is deliberately not re-inserted: the+ configured interval is the crash's data loss, and inventing one in the hot path would+ silently change throttling.+* Concurrent `pgmq.enable_notify_insert` calls for the same queue no longer race. Two+ replicas reconciling the same config at startup could both pass the function's internal+ `DROP TRIGGER IF EXISTS` — on a fresh queue it finds nothing and takes no lock — and the+ loser then blocked on the throttle row's unique constraint, resumed, and created a+ trigger that now existed, failing with SQLSTATE 42710 (duplicate_object) and taking down+ that replica's entire startup reconcile. Measured at roughly a 28% collision rate over+ 200 concurrent iterations. The function now takes the per-queue advisory lock, as+ `pgmq.create` and `pgmq.create_partitioned` already do, which makes concurrent callers+ convergent.+* `pgmq.enable_notify_insert` coalesces a NULL `throttle_interval_ms` to the documented+ 250 ms, so non-Haskell callers get the same guarantee the pgmq-hasql statement provides+ client-side.+* `pgmq.create_partitioned` is now re-entrant. The advisory lock serialized concurrent+ creators, but the second one still called `partman.create_parent` on a parent the first+ had just registered: its `CREATE TABLE IF NOT EXISTS` is a no-op while `create_parent`+ rejects an already-managed parent. Both `create_parent` calls are now guarded by a+ `part_config` probe.++### Other Changes++* The test suite derives the migration ledger from the plan rather than enumerating it+ positionally, so appending a migration no longer turns unrelated tests red. The ledger+ is spelled out once, in `testNativeComponent`.+* Added a pg_partman-gated re-entry test that reports its skip rather than passing+ silently when the extension is unavailable.+ ## 0.4.0.1 -- 2026-07-14 ### Bug Fixes
+ migrations/0003-notify-crash-safety-and-locking.sql view
@@ -0,0 +1,297 @@+-- Notification crash safety, reconcile locking, and partitioned re-entry.+--+-- Three deliberate divergences from upstream pgmq 1.11.0, all re-created with+-- their original signatures so this file is a drop-in replacement:+--+-- 1. pgmq.notify_queue_listeners() fails open when its throttle row is+-- missing. pgmq.notify_insert_throttle is UNLOGGED, so crash recovery+-- truncates it; the trigger's UPDATE then matched zero rows and the+-- updated_count > 0 gate suppressed PG_NOTIFY forever, silently starving+-- every listener until an application re-enabled notify. Losing the+-- throttle interval in a crash is acceptable; losing deliveries is not.+--+-- 2. pgmq.enable_notify_insert() takes the per-queue advisory lock and+-- coalesces a NULL throttle to the documented 250 ms. Without the lock,+-- two replicas reconciling the same queue at startup both pass the+-- internal DROP TRIGGER IF EXISTS, and the loser then fails with SQLSTATE+-- 42710 when it creates a trigger the winner just committed, taking down+-- that replica's whole reconcile. Without the COALESCE, a bound SQL NULL+-- raises 23502: a plpgsql parameter DEFAULT applies only to an omitted+-- argument, never to an explicit NULL.+--+-- 3. pgmq.create_partitioned() skips partman.create_parent when the table is+-- already registered in part_config. The advisory lock serializes+-- concurrent creators, but the second replica still calls create_parent on+-- a parent the first just registered — its CREATE TABLE IF NOT EXISTS is a+-- no-op, while create_parent rejects an already-managed parent. Idempotence+-- is the missing property, not locking.+--+-- See docs/plans/14-make-insert-notifications-survive-crashes-and-document-the-channel-contract.md.++CREATE OR REPLACE FUNCTION pgmq.notify_queue_listeners()+RETURNS TRIGGER AS $$+DECLARE+ queue_name_extracted TEXT; -- Queue name extracted from trigger table name+ updated_count INTEGER; -- Number of rows updated (0 or 1)+BEGIN+ queue_name_extracted := substring(TG_TABLE_NAME from 3);++ UPDATE pgmq.notify_insert_throttle+ SET last_notified_at = clock_timestamp()+ WHERE queue_name = queue_name_extracted+ AND (+ throttle_interval_ms = 0 -- No throttling configured+ OR clock_timestamp() - last_notified_at >=+ (throttle_interval_ms * INTERVAL '1 millisecond') -- Throttle interval has elapsed+ );++ -- Check how many rows were updated (will be 0 or 1)+ GET DIAGNOSTICS updated_count = ROW_COUNT;++ IF updated_count > 0 THEN+ PERFORM PG_NOTIFY('pgmq.' || TG_TABLE_NAME || '.' || TG_OP, NULL);+ ELSIF NOT EXISTS (+ SELECT 1 FROM pgmq.notify_insert_throttle nit+ WHERE nit.queue_name = queue_name_extracted+ ) THEN+ -- Fail open: the trigger exists but its throttle row does not. Zero updated+ -- rows has two causes -- throttled (suppress, correct) and row absent -- and+ -- only this probe tells them apart. The row lives in an UNLOGGED table that+ -- crash recovery truncates, so losing it must not silently stop delivery.+ -- Notify unthrottled until pgmq.enable_notify_insert restores the row.+ --+ -- The row is deliberately NOT re-inserted here: the configured interval is+ -- the crash's data loss, and inventing one in the trigger's hot path would+ -- silently change throttling. The reconciler restores the configured value.+ PERFORM PG_NOTIFY('pgmq.' || TG_TABLE_NAME || '.' || TG_OP, NULL);+ END IF;++RETURN NEW;+END;+$$ LANGUAGE plpgsql;++CREATE OR REPLACE FUNCTION pgmq.enable_notify_insert(queue_name TEXT, throttle_interval_ms INTEGER DEFAULT 250)+RETURNS void AS $$+DECLARE+ qtable TEXT := pgmq.format_table_name(queue_name, 'q');+ v_queue_name TEXT := queue_name;+ -- A bound SQL NULL never triggers the parameter DEFAULT above, so the+ -- documented 250 ms has to be applied here for non-Haskell callers.+ v_throttle_interval_ms INTEGER := COALESCE(throttle_interval_ms, 250);+BEGIN+ -- Serialize the whole disable-then-create sequence per queue, exactly as+ -- pgmq.create and pgmq.create_partitioned already do. That makes concurrent+ -- callers convergent: the second one drops and re-creates the first one's+ -- identical trigger, an idempotent no-op in effect, with no error-code+ -- matching anywhere. Accepted side effect: a concurrent second enable resets+ -- last_notified_at to the epoch, exactly as a sequential re-enable does.+ PERFORM pgmq.acquire_queue_lock(queue_name);++ -- Validate that throttle_interval_ms is non-negative+ IF v_throttle_interval_ms < 0 THEN+ RAISE EXCEPTION 'throttle_interval_ms must be non-negative';+ END IF;++ -- Validate that the queue table exists+ IF NOT EXISTS (SELECT 1 FROM information_schema.tables WHERE table_schema = 'pgmq' AND table_name = qtable) THEN+ RAISE EXCEPTION 'Queue "%" does not exist. Create it first using pgmq.create()', v_queue_name;+ END IF;++ PERFORM pgmq.disable_notify_insert(v_queue_name);++ INSERT INTO pgmq.notify_insert_throttle (queue_name, throttle_interval_ms)+ VALUES (v_queue_name, v_throttle_interval_ms)+ ON CONFLICT ON CONSTRAINT notify_insert_throttle_queue_name_key DO UPDATE+ SET throttle_interval_ms = EXCLUDED.throttle_interval_ms,+ last_notified_at = to_timestamp(0);++ EXECUTE FORMAT(+ $QUERY$+ CREATE CONSTRAINT TRIGGER trigger_notify_queue_insert_listeners+ AFTER INSERT ON pgmq.%I+ DEFERRABLE FOR EACH ROW+ EXECUTE PROCEDURE pgmq.notify_queue_listeners()+ $QUERY$,+ qtable+ );+END;+$$ LANGUAGE plpgsql;++CREATE OR REPLACE FUNCTION pgmq.create_partitioned(+ queue_name TEXT,+ partition_interval TEXT DEFAULT '10000',+ retention_interval TEXT DEFAULT '100000'+)+RETURNS void AS $$+DECLARE+ partition_col TEXT;+ a_partition_col TEXT;+ 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');+ fq_qtable TEXT := 'pgmq.' || qtable;+ fq_atable TEXT := 'pgmq.' || atable;+ l_already_managed BOOLEAN; -- Parent already registered in pg_partman's part_config+BEGIN+ PERFORM pgmq.validate_queue_name(queue_name);+ PERFORM pgmq.acquire_queue_lock(queue_name);+ PERFORM pgmq._ensure_pg_partman_installed();+ SELECT pgmq._get_partition_col(partition_interval) INTO partition_col;++ 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 (%I)+ $QUERY$,+ qtable, partition_col+ );++ -- Re-entry guard: the CREATE TABLE above is IF NOT EXISTS, but create_parent+ -- rejects a parent that is already registered, so a second caller would fail+ -- on a queue the first one finished creating.+ EXECUTE FORMAT(+ $QUERY$+ SELECT EXISTS (SELECT 1 FROM %I.part_config WHERE parent_table = %L)+ $QUERY$,+ pgmq._get_pg_partman_schema(),+ fq_qtable+ ) INTO l_already_managed;++ IF NOT l_already_managed THEN+ -- https://github.com/pgpartman/pg_partman/blob/master/doc/pg_partman.md+ -- p_parent_table - the existing parent table. MUST be schema qualified, even if in public schema.+ EXECUTE FORMAT(+ $QUERY$+ SELECT %I.create_parent(+ p_parent_table := %L,+ p_control := %L,+ p_interval := %L,+ p_type := case+ when pgmq._get_pg_partman_major_version() = 5 then 'range'+ else 'native'+ end+ )+ $QUERY$,+ pgmq._get_pg_partman_schema(),+ fq_qtable,+ partition_col,+ partition_interval+ );+ END IF;++ EXECUTE FORMAT(+ $QUERY$+ CREATE INDEX IF NOT EXISTS %I ON pgmq.%I (%I);+ $QUERY$,+ qtable || '_part_idx', qtable, partition_col+ );++ EXECUTE FORMAT(+ $QUERY$+ UPDATE %I.part_config+ SET+ retention = %L,+ retention_keep_table = false,+ retention_keep_index = true,+ automatic_maintenance = 'on'+ WHERE parent_table = %L;+ $QUERY$,+ pgmq._get_pg_partman_schema(),+ retention_interval,+ 'pgmq.' || qtable+ );++ EXECUTE FORMAT(+ $QUERY$+ INSERT INTO pgmq.meta (queue_name, is_partitioned, is_unlogged)+ VALUES (%L, true, false)+ ON CONFLICT+ DO NOTHING;+ $QUERY$,+ queue_name+ );++ IF partition_col = 'enqueued_at' THEN+ a_partition_col := 'archived_at';+ ELSE+ a_partition_col := partition_col;+ END IF;++ 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 (%I);+ $QUERY$,+ atable, a_partition_col+ );++ -- Same re-entry guard for the archive table.+ EXECUTE FORMAT(+ $QUERY$+ SELECT EXISTS (SELECT 1 FROM %I.part_config WHERE parent_table = %L)+ $QUERY$,+ pgmq._get_pg_partman_schema(),+ fq_atable+ ) INTO l_already_managed;++ IF NOT l_already_managed THEN+ -- https://github.com/pgpartman/pg_partman/blob/master/doc/pg_partman.md+ -- p_parent_table - the existing parent table. MUST be schema qualified, even if in public schema.+ EXECUTE FORMAT(+ $QUERY$+ SELECT %I.create_parent(+ p_parent_table := %L,+ p_control := %L,+ p_interval := %L,+ p_type := case+ when pgmq._get_pg_partman_major_version() = 5 then 'range'+ else 'native'+ end+ )+ $QUERY$,+ pgmq._get_pg_partman_schema(),+ fq_atable,+ a_partition_col,+ partition_interval+ );+ END IF;++ EXECUTE FORMAT(+ $QUERY$+ UPDATE %I.part_config+ SET+ retention = %L,+ retention_keep_table = false,+ retention_keep_index = true,+ automatic_maintenance = 'on'+ WHERE parent_table = %L;+ $QUERY$,+ pgmq._get_pg_partman_schema(),+ retention_interval,+ 'pgmq.' || atable+ );++ EXECUTE FORMAT(+ $QUERY$+ CREATE INDEX IF NOT EXISTS %I ON pgmq.%I (archived_at);+ $QUERY$,+ 'archived_at_idx_' || queue_name, atable+ );++END;+$$ LANGUAGE plpgsql;
migrations/manifest view
@@ -1,2 +1,3 @@ 0001-install-v1.11.0.sql 0002-schema-management-comment.sql+0003-notify-crash-safety-and-locking.sql
pgmq-migration.cabal view
@@ -1,6 +1,6 @@ cabal-version: 3.4 name: pgmq-migration-version: 0.4.0.1+version: 0.5.0.0 synopsis: PGMQ schema migrations without PostgreSQL extension description: Installs the PGMQ schema into PostgreSQL without requiring the pgmq extension.
test/Main.hs view
@@ -47,8 +47,13 @@ ( ComponentDescription (..), PlanDescription (..), componentNameText,+ migrationIdName,+ migrationNameText, planDescription, )+-- MigrationDescription's 'migrationId' field would shadow the smart constructor+-- of the same name imported above, so it is reached through this alias only.+import Database.PostgreSQL.Migrate.Internal qualified as MigrateInternal import EphemeralPg ( connectionSettings, withCached,@@ -93,7 +98,7 @@ [ testGroup "native definition" [ testCase "baseline bytes equal vendored pgmq.sql" testNativePayload,- testCase "component pgmq has two migrations and no dependencies" testNativeComponent+ testCase "component pgmq lists the ledger in order and has no dependencies" testNativeComponent ], testGroup "native runner"@@ -118,6 +123,9 @@ vendored <- ByteString.readFile vendorPath native @?= vendored +-- | The one place the ledger is spelled out. Every other expectation in this+-- suite derives from 'nativeMigrationNames', so appending a migration means+-- reviewing this list and nothing else. testNativeComponent :: IO () testNativeComponent = do component <- either (assertFailure . show) pure Migration.pgmqMigrations@@ -127,9 +135,39 @@ [ComponentDescription {name, dependencies, migrations}] -> do componentNameText name @?= "pgmq" dependencies @?= mempty- length migrations @?= 2+ migrationNames migrations+ @?= [ "0001-install-v1.11.0",+ "0002-schema-management-comment",+ "0003-notify-crash-safety-and-locking"+ ] actual -> assertFailure ("unexpected native PGMQ plan: " <> show actual) +migrationNames :: (Foldable f) => f MigrateInternal.MigrationDescription -> [Text]+migrationNames descriptions =+ [ migrationNameText (migrationIdName (MigrateInternal.migrationId description))+ | description <- toList descriptions+ ]++-- | Migration names in the native PGMQ ledger, in manifest order.+nativeMigrationNames :: IO [Text]+nativeMigrationNames = do+ component <- either (assertFailure . show) pure Migration.pgmqMigrations+ plan <- either (assertFailure . show) pure (migrationPlan (component :| []))+ let PlanDescription components = planDescription plan+ case toList components of+ [ComponentDescription {migrations}] -> pure (migrationNames migrations)+ actual -> assertFailure ("unexpected native PGMQ plan: " <> show actual)++-- | Every migration a history import leaves unapplied. An import records the+-- immutable baseline only, so everything after it is still pending.+pendingAfterBaseline :: IO [VerificationIssue]+pendingAfterBaseline = do+ names <- nativeMigrationNames+ traverse pendingIssue (drop 1 names)+ where+ pendingIssue name =+ PendingMigration <$> either (assertFailure . show) pure (migrationId "pgmq" name)+ findFile :: [FilePath] -> IO FilePath findFile candidates = do existing <- filterM doesFileExist candidates@@ -181,14 +219,15 @@ testNativeRunner :: Settings.Settings -> Connection.Connection -> IO () testNativeRunner settings conn = withCleanDb conn $ \c -> do plan <- nativePlan+ ledgerLength <- length <$> nativeMigrationNames first <- runMigrationPlan defaultRunOptions settings plan case first of Left err -> assertFailure ("fresh native migration failed: " <> show err)- Right report -> (outcome <$> toList (results report)) @?= [AppliedNow, AppliedNow]+ Right report -> (outcome <$> toList (results report)) @?= replicate ledgerLength AppliedNow second <- runMigrationPlan defaultRunOptions settings plan case second of Left err -> assertFailure ("repeated native migration failed: " <> show err)- Right report -> (outcome <$> toList (results report)) @?= [AlreadyApplied, AlreadyApplied]+ Right report -> (outcome <$> toList (results report)) @?= replicate ledgerLength AlreadyApplied functionExists c "pgmq.metrics_all()" >>= (@?= True) hasCanaryComment c >>= (@?= True) @@ -236,15 +275,17 @@ historyOutcomes second @?= [AlreadyImported] plan <- nativePlan- canaryId <- either (assertFailure . show) pure (migrationId "pgmq" "0002-schema-management-comment")+ pending <- pendingAfterBaseline+ ledgerLength <- length <$> nativeMigrationNames beforeCanary <- verifyMigrationPlan defaultRunOptions settings plan case beforeCanary of Left err -> assertFailure ("native verify failed after direct import: " <> show err)- Right (VerificationReport issues _ _ _) -> issues @?= [PendingMigration canaryId]+ Right (VerificationReport issues _ _ _) -> issues @?= pending nativeRun <- runMigrationPlan defaultRunOptions settings plan case nativeRun of Left err -> assertFailure ("native runner failed after direct import: " <> show err)- Right report -> (outcome <$> toList (results report)) @?= [AlreadyApplied, AppliedNow]+ Right report ->+ (outcome <$> toList (results report)) @?= AlreadyApplied : replicate (ledgerLength - 1) AppliedNow afterCanary <- verifyMigrationPlan defaultRunOptions settings plan case afterCanary of Left err -> assertFailure ("native verify failed after direct canary: " <> show err)@@ -252,7 +293,7 @@ repeated <- runMigrationPlan defaultRunOptions settings plan case repeated of Left err -> assertFailure ("native rerun failed after direct canary: " <> show err)- Right report -> (outcome <$> toList (results report)) @?= [AlreadyApplied, AlreadyApplied]+ Right report -> (outcome <$> toList (results report)) @?= replicate ledgerLength AlreadyApplied functionExists c "pgmq.metrics_all()" >>= (@?= False) hasCanaryComment c >>= (@?= True) @@ -319,15 +360,17 @@ historyOutcomes second @?= [AlreadyImported] plan <- nativePlan- canaryId <- either (assertFailure . show) pure (migrationId "pgmq" "0002-schema-management-comment")+ pending <- pendingAfterBaseline+ ledgerLength <- length <$> nativeMigrationNames beforeCanary <- verifyMigrationPlan defaultRunOptions settings plan case beforeCanary of Left err -> assertFailure ("native verify failed after equivalent import: " <> show err)- Right (VerificationReport issues _ _ _) -> issues @?= [PendingMigration canaryId]+ Right (VerificationReport issues _ _ _) -> issues @?= pending nativeRun <- runMigrationPlan defaultRunOptions settings plan case nativeRun of Left err -> assertFailure ("native runner failed after equivalent import: " <> show err)- Right report -> (outcome <$> toList (results report)) @?= [AlreadyApplied, AppliedNow]+ Right report ->+ (outcome <$> toList (results report)) @?= AlreadyApplied : replicate (ledgerLength - 1) AppliedNow afterCanary <- verifyMigrationPlan defaultRunOptions settings plan case afterCanary of Left err -> assertFailure ("native verify failed after equivalent canary: " <> show err)@@ -335,7 +378,7 @@ repeated <- runMigrationPlan defaultRunOptions settings plan case repeated of Left err -> assertFailure ("native rerun failed after equivalent canary: " <> show err)- Right report -> (outcome <$> toList (results report)) @?= [AlreadyApplied, AlreadyApplied]+ Right report -> (outcome <$> toList (results report)) @?= replicate ledgerLength AlreadyApplied hasCanaryComment c >>= (@?= True) testEquivalentContractRejections :: Settings.Settings -> Connection.Connection -> IO ()@@ -549,19 +592,21 @@ assertNativeCanaryLifecycle :: Settings.Settings -> Connection.Connection -> IO () assertNativeCanaryLifecycle settings connection = do plan <- nativePlan- canaryId <- either (assertFailure . show) pure (migrationId "pgmq" "0002-schema-management-comment")+ pending <- pendingAfterBaseline+ ledgerLength <- length <$> nativeMigrationNames beforeCanary <- verifyMigrationPlan defaultRunOptions settings plan case beforeCanary of Left err -> assertFailure ("native verify failed after shared-ledger import: " <> show err)- Right (VerificationReport issues _ _ _) -> issues @?= [PendingMigration canaryId]+ Right (VerificationReport issues _ _ _) -> issues @?= pending nativeRun <- runMigrationPlan defaultRunOptions settings plan case nativeRun of Left err -> assertFailure ("native runner failed after shared-ledger import: " <> show err)- Right report -> (outcome <$> toList (results report)) @?= [AlreadyApplied, AppliedNow]+ Right report ->+ (outcome <$> toList (results report)) @?= AlreadyApplied : replicate (ledgerLength - 1) AppliedNow repeated <- runMigrationPlan defaultRunOptions settings plan case repeated of Left err -> assertFailure ("native rerun failed after shared-ledger import: " <> show err)- Right report -> (outcome <$> toList (results report)) @?= [AlreadyApplied, AlreadyApplied]+ Right report -> (outcome <$> toList (results report)) @?= replicate ledgerLength AlreadyApplied hasCanaryComment connection >>= (@?= True) canaryCommentStatement :: Statement.Statement () Bool