From 8590a9ddbc4a365cc4d44f6b42cda56cab30b8d4 Mon Sep 17 00:00:00 2001 From: samo-agent <280144521+samo-agent@users.noreply.github.com> Date: Thu, 1 Oct 2026 14:18:14 +0000 Subject: [PATCH 1/5] fix: backport receive safety to 0.2.1 Refs #365. Reject oversized plain and cooperative batches, correct documentation, and retain the 0.2 stable feature set. --- .github/workflows/ci.yml | 5 +- build/transform.sh | 2 +- clients/go/README.md | 2 +- clients/python/README.md | 14 +- docs/examples.md | 2 +- docs/reference.md | 10 +- docs/three-latencies.md | 2 +- docs/upgrading.md | 16 +- sql/pgque-additions/lifecycle.sql | 2 +- sql/pgque-api/cooperative_consumers.sql | 5 +- sql/pgque-api/receive.sql | 5 +- sql/pgque-tle.sql | 24 +- sql/pgque.sql | 14 +- tests/run_all.sql | 1 + tests/test_api_receive.sql | 24 +- tests/test_pgque_config.sql | 4 +- tests/test_receive_overflow.sql | 367 ++++++++++++++++++++++++ 17 files changed, 458 insertions(+), 41 deletions(-) create mode 100644 tests/test_receive_overflow.sql diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index d212c076..c22bc795 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -1,10 +1,11 @@ name: CI on: + workflow_dispatch: push: - branches: [main] + branches: [main, maintenance-0-2, release-0-2-1] pull_request: - branches: [main] + branches: [main, maintenance-0-2, release-0-2-1] jobs: test: diff --git a/build/transform.sh b/build/transform.sh index 01b38003..b667c749 100755 --- a/build/transform.sh +++ b/build/transform.sh @@ -646,7 +646,7 @@ apply_idempotency_guards() { # Start with header cat > "${INSTALL_FILE}" << 'HEADER' -- pgque.sql -- PgQ Universal Edition --- Version: 0.2.0 +-- Version: 0.2.1 -- Copyright 2026 Nikolay Samokhvalov. Apache-2.0 license. -- Includes code derived from PgQ (ISC license, Marko Kreen / Skype Technologies OU). -- diff --git a/clients/go/README.md b/clients/go/README.md index dfca23c2..dd6ea0e5 100644 --- a/clients/go/README.md +++ b/clients/go/README.md @@ -93,7 +93,7 @@ func main() { | Option | Default | Notes | | --------------------------------------- | -------------- | --------------------------------------------------------------------- | | `WithPollInterval(d time.Duration)` | `30s` | Idle backoff between polls when the queue is empty. | -| `WithMaxMessages(n int)` | `math.MaxInt32` | Per-Receive limit. The default requests the whole PgQ batch before `Ack`. If you lower it below the real batch size, `Ack` still finishes the batch and unreturned rows are skipped. | +| `WithMaxMessages(n int)` | `math.MaxInt32` | Complete-batch safety ceiling. Servers with the overflow guard raise if the batch is larger; retry with a sufficient ceiling within your resource budget. Older servers can truncate, so upgrade before relying on the guard. Process the whole batch before `Ack`. | | `WithUnknownHandlerPolicy(p)` | `NackUnknown` | `AckUnknown` logs and skips messages with no registered handler. | | `WithRetryAfter(d time.Duration)` | `60s` | Retry delay for Consumer-issued `Nack` calls on handler failure or unknown type. | diff --git a/clients/python/README.md b/clients/python/README.md index 0d59ea1c..6f24b0ac 100644 --- a/clients/python/README.md +++ b/clients/python/README.md @@ -71,12 +71,14 @@ consumer.start() # blocks until SIGTERM / SIGINT ### Consumer options `Consumer(..., max_messages=...)` controls the per-`receive` limit. -The default is PostgreSQL's `int` maximum, so the consumer requests -the whole PgQ batch before acknowledging it. `ack()` finishes the -entire underlying PgQ batch, including rows beyond `max_messages`; -only lower this value when it is at least as large as the queue's -worst-case batch size, otherwise rows past the limit are silently -skipped by the batch ack. +The default is the Postgres `int` maximum, so the consumer requests +an entire PgQ batch before acknowledging it. On servers with the overflow +guard, a batch larger than this ceiling raises an error instead of returning +partial results. Roll back failed transactions and retry with a sufficient +ceiling within your resource budget; monitor errors and consumer lag. +The ticker threshold does not cap batch size. Older servers can truncate +results, so upgrade the server before relying on this protection. +`ack()` always finishes the entire batch: process every message first. ### Handling unknown event types diff --git a/docs/examples.md b/docs/examples.md index 01ea8b63..559c84b5 100644 --- a/docs/examples.md +++ b/docs/examples.md @@ -67,7 +67,7 @@ begin; commit; ``` -The `inserted` CTE runs to completion even though the main query does not reference it (data-modifying CTEs always execute). Every row in `msgs` shares the same `batch_id`, so the scalar subquery picks any one of them and `pgque.ack` runs exactly once. **Batch-ownership caveat:** `pgque.ack(batch_id)` advances the consumer past the entire underlying batch, even if `receive()` returned fewer rows than the batch contains (due to `max_return`). Either consume the full batch before acking, or use `max_return >= ticker_max_count` (default 500) to ensure all rows are returned. +The `inserted` CTE runs to completion even though the main query does not reference it (data-modifying CTEs always execute). Every row in `msgs` shares the same `batch_id`, so the scalar subquery picks any one of them and `pgque.ack` runs exactly once. **Complete batch or error:** `max_return` is a safety ceiling, not pagination. An oversized batch raises an error rather than returning partial results. Roll back and retry with a sufficient ceiling within your resource budget; never acknowledge a failed receive. Process every event before calling `ack()`, which finishes the whole batch. `ticker_max_count` is a tick-trigger threshold, not a hard batch-size cap. > **Anti-pattern: send + receive in one transaction.** Above merges `receive` + writes + `ack` into one tx — correct. Do **not** also merge `send` / `force_next_tick` / `ticker` into the same tx; the ticker's snapshot must be taken *after* `send` commits. > diff --git a/docs/reference.md b/docs/reference.md index c696a212..497f64ba 100644 --- a/docs/reference.md +++ b/docs/reference.md @@ -119,14 +119,16 @@ All consume-side functions (`receive`, `ack`, `nack`, `subscribe`, `unsubscribe` #### `pgque.receive(queue text, consumer text, max_return int default 100) → setof pgque.message` -Pulls the next batch for `consumer` on `queue` and streams up to `max_return` messages. `max_return` must be >= 1; passing 0 or a negative value raises an error. Returns an empty set if no batch is available. Each row is a `pgque.message` composite (see [§Message type](#message-type)). +Pulls the complete next batch for `consumer` on `queue`, or raises an error if it exceeds `max_return`. `max_return` must be >= 1; passing 0 or a negative value raises an error. Returns an empty set if no batch is available. Each row is a `pgque.message` composite (see [§Message type](#message-type)). Grant: `pgque_reader`. Source: `sql/pgque-api/receive.sql`. ```sql select * from pgque.receive('orders', 'processor', 100); ``` -**Batch-ownership caveat.** `max_return` limits the number of rows returned to the caller, but `ack(batch_id)` advances the consumer cursor past the entire underlying batch. If `max_return < ticker_max_count`, calling `ack()` after a partial receive will drop the unreturned rows from the consumer's perspective. Either consume the full batch before acking, or use `max_return >= ticker_max_count` for safe pagination. +**Complete batch or error.** `max_return` is a safety ceiling, not a page size. Exactly that many events is valid; encountering one more raises an error, with no successful result or consumer advancement from that call. Roll back a failed transaction and retry with a larger ceiling within your resource budget, or use the low-level full-batch interface. Process every event before calling `ack(batch_id)`, which finishes the entire batch. Do not acknowledge after an overflow error. Repeating `receive()` with the same ceiling does not paginate or resolve overflow. + +`ticker_max_count` is a tick-trigger threshold, not a hard batch-size cap. Setting `max_return` to that value cannot guarantee success during bursts. Monitor receive errors and consumer lag so an undersized ceiling does not silently stall processing. #### `pgque.ack(batch_id bigint) → integer` @@ -170,9 +172,9 @@ Receives messages for one subconsumer. `max_return` must be >= 1. `dead_interval **Auto-registration.** If the logical `consumer` or `subconsumer` is not yet registered, `receive_coop()` registers them on the fly (creates the `coop_main` row on first call, then the `coop_member` row), so a worker can call `receive_coop()` cold without a prior `register_subconsumer`. Use the explicit `register_subconsumer(..., convert_normal => true)` call only when you need to convert an existing normal consumer into a cooperative main. -**Empty tick windows are auto-finished.** When the current batch's tick window holds no events, `receive_coop()` calls `finish_batch` internally and returns the empty set. Callers polling a quiet queue do not see (and do not need to ack) a `batch_id`; this differs from `receive()`, which still returns an active batch token even when the result set is empty. +**Empty tick windows are auto-finished.** When the current batch's tick window holds no events, `receive_coop()` calls `finish_batch` internally and returns the empty set. Callers polling a quiet queue do not see (and do not need to ack) a `batch_id`; `receive()` also auto-finishes empty batches. -**Batch-ownership caveat.** As with `receive()`, `max_return` limits only returned rows; `ack(batch_id)` advances the cooperative cursor past the whole underlying batch. Use `max_return >= ticker_max_count` or consume the full batch before acking. +**Complete batch or error.** As with `receive()`, `max_return` is a safety ceiling, not pagination. An oversized batch raises an error and rolls back allocation or takeover performed by that call. Retry with a sufficient ceiling within your resource budget, process the complete batch, then acknowledge. The ticker threshold does not cap batch size. **Throughput note.** Cooperative allocation serializes on a `FOR UPDATE` of the `coop_main` subscription row, so many workers polling tiny batches contend on a single hot row. If you scale workers, also tune `ticker_max_count` and tick cadence so each batch is large enough to amortize the lock. diff --git a/docs/three-latencies.md b/docs/three-latencies.md index 8ea6846b..8767d320 100644 --- a/docs/three-latencies.md +++ b/docs/three-latencies.md @@ -44,6 +44,6 @@ Per-queue thresholds (`queue_ticker_max_lag` default `3 seconds`, `queue_ticker_ ## Load behavior: PgQue vs. UPDATE/DELETE designs -The key property of the tick model: **e2e does not grow with load.** The ticker fires at its configured rate regardless of backlog, so under pressure batch size grows (up to `queue_ticker_max_count`) — not e2e. +Tick cadence is independent of batch size, provided the ticker keeps up. Under pressure, batch size can grow beyond `queue_ticker_max_count`: that setting is a tick-trigger threshold, not a hard cap. Overloaded tickers or consumers can also increase end-to-end latency. Size receive ceilings for bursts and monitor consumer lag. UPDATE/DELETE-based systems use a different model: a consumer call returns messages immediately, marking them consumed via UPDATE (claim) and DELETE (ack) rather than advancing a snapshot cursor. So e2e ≈ consumer poll interval — sub-ms when the consumer is actively polling, up to the poll interval otherwise. Drain rate is `batch_size / poll_interval`; if producers outrun that, queue depth grows and e2e grows with it until consumers scale out. Separately, those UPDATEs and DELETEs produce dead tuples that autovacuum cannot reclaim under MVCC pressure (long-running tx, idle-in-transaction, lagging logical replication slot, physical standby with `hot_standby_feedback=on`) — the bloat failure mode [PgQue avoids by construction](../README.md#why-pgque). diff --git a/docs/upgrading.md b/docs/upgrading.md index abefe4bb..b4fc5900 100644 --- a/docs/upgrading.md +++ b/docs/upgrading.md @@ -12,6 +12,20 @@ The installer is idempotent: it preserves queues, consumers, subscriptions, retry rows, DLQ rows, and existing event tables while adding new functions, columns, grants, and constraints required by the target release. +## v0.2.0 to v0.2.1 + +Re-run the installer using the command above. This maintenance release makes +`pgque.receive()` and `pgque.receive_coop()` reject batches larger than +`max_return`, rather than returning a truncated result that could be +acknowledged as a complete batch. The upgrade replaces functions only; it does not +change tables or queue state. + +Applications do not need client-library updates for this server-side fix. +An undersized receive ceiling now causes an error: roll back the failed +transaction and retry with a sufficient ceiling within your resource budget. +Never acknowledge after a receive error. Ticker thresholds do not cap batch +size, and repeated receive calls are not pagination. + ## v0.1.0 to v0.2.0 The supported v0.1.0 → v0.2.0 path is the same re-install procedure: @@ -32,7 +46,7 @@ After upgrading, verify the installed version: ```sql select pgque.version(); --- 0.2.0-rc.1, or the exact release you installed +-- 0.2.1, or the exact release you installed ``` You can also run the idempotency smoke test from the repository: diff --git a/sql/pgque-additions/lifecycle.sql b/sql/pgque-additions/lifecycle.sql index de8ba674..486e815f 100644 --- a/sql/pgque-additions/lifecycle.sql +++ b/sql/pgque-additions/lifecycle.sql @@ -438,7 +438,7 @@ $$ language plpgsql security definer set search_path = pgque, pg_catalog; create or replace function pgque.version() returns text as $$ begin - return '0.2.0'; + return '0.2.1'; end; $$ language plpgsql security definer set search_path = pgque, pg_catalog; diff --git a/sql/pgque-api/cooperative_consumers.sql b/sql/pgque-api/cooperative_consumers.sql index d110f1bd..14e3f8ff 100644 --- a/sql/pgque-api/cooperative_consumers.sql +++ b/sql/pgque-api/cooperative_consumers.sql @@ -1143,6 +1143,10 @@ begin ev_extra4 from pgque.get_batch_events(v_batch_id) loop + if cnt = i_max_return then + raise exception 'pgque.receive_coop: batch exceeds max_return of %', i_max_return + using hint = 'Retry with a larger resource-safe max_return to receive the complete batch. Do not acknowledge after this error.'; + end if; return next row( ev.ev_id, v_batch_id, @@ -1156,7 +1160,6 @@ begin ev.ev_extra4 )::pgque.message; cnt := cnt + 1; - exit when cnt >= i_max_return; end loop; -- Empty batch: release the member token so the subconsumer is not wedged diff --git a/sql/pgque-api/receive.sql b/sql/pgque-api/receive.sql index 4484e7e4..419968fc 100644 --- a/sql/pgque-api/receive.sql +++ b/sql/pgque-api/receive.sql @@ -50,13 +50,16 @@ begin ev_extra1, ev_extra2, ev_extra3, ev_extra4 from pgque.get_batch_events(v_batch_id) loop + if cnt = i_max_return then + raise exception 'pgque.receive: batch exceeds max_return of %', i_max_return + using hint = 'Retry with a larger resource-safe max_return to receive the complete batch. Do not acknowledge after this error.'; + end if; return next row( ev.ev_id, v_batch_id, ev.ev_type, ev.ev_data, ev.ev_retry, ev.ev_time, ev.ev_extra1, ev.ev_extra2, ev.ev_extra3, ev.ev_extra4 )::pgque.message; cnt := cnt + 1; - exit when cnt >= i_max_return; end loop; -- Empty batch: finish immediately to advance the consumer cursor. diff --git a/sql/pgque-tle.sql b/sql/pgque-tle.sql index 2601a24f..33aaea03 100644 --- a/sql/pgque-tle.sql +++ b/sql/pgque-tle.sql @@ -68,8 +68,8 @@ begin from pgtle.available_extensions() where name = 'pgque'; - if existing_version = '0.2.0' then - raise notice 'pgque 0.2.0 already registered with pg_tle; skipping install_extension().'; + if existing_version = '0.2.1' then + raise notice 'pgque 0.2.1 already registered with pg_tle; skipping install_extension().'; return; end if; @@ -78,16 +78,16 @@ begin 'but this script registers version %. Run ' 'sql/pgque-tle-uninstall.sql first to remove the existing ' 'registration, then re-run this script.', - existing_version, '0.2.0'; + existing_version, '0.2.1'; end if; perform pgtle.install_extension( 'pgque', - '0.2.0', + '0.2.1', 'PgQue — PgQ Universal Edition (zero-bloat Postgres queue)', $pgque_extension_body$ -- pgque.sql -- PgQ Universal Edition --- Version: 0.2.0 +-- Version: 0.2.1 -- Copyright 2026 Nikolay Samokhvalov. Apache-2.0 license. -- Includes code derived from PgQ (ISC license, Marko Kreen / Skype Technologies OU). -- @@ -4658,7 +4658,7 @@ $$ language plpgsql security definer set search_path = pgque, pg_catalog; create or replace function pgque.version() returns text as $$ begin - return '0.2.0'; + return '0.2.1'; end; $$ language plpgsql security definer set search_path = pgque, pg_catalog; @@ -5451,13 +5451,16 @@ begin ev_extra1, ev_extra2, ev_extra3, ev_extra4 from pgque.get_batch_events(v_batch_id) loop + if cnt = i_max_return then + raise exception 'pgque.receive: batch exceeds max_return of %', i_max_return + using hint = 'Retry with a larger resource-safe max_return to receive the complete batch. Do not acknowledge after this error.'; + end if; return next row( ev.ev_id, v_batch_id, ev.ev_type, ev.ev_data, ev.ev_retry, ev.ev_time, ev.ev_extra1, ev.ev_extra2, ev.ev_extra3, ev.ev_extra4 )::pgque.message; cnt := cnt + 1; - exit when cnt >= i_max_return; end loop; -- Empty batch: finish immediately to advance the consumer cursor. @@ -6705,6 +6708,10 @@ begin ev_extra4 from pgque.get_batch_events(v_batch_id) loop + if cnt = i_max_return then + raise exception 'pgque.receive_coop: batch exceeds max_return of %', i_max_return + using hint = 'Retry with a larger resource-safe max_return to receive the complete batch. Do not acknowledge after this error.'; + end if; return next row( ev.ev_id, v_batch_id, @@ -6718,7 +6725,6 @@ begin ev.ev_extra4 )::pgque.message; cnt := cnt + 1; - exit when cnt >= i_max_return; end loop; -- Empty batch: release the member token so the subconsumer is not wedged @@ -7135,5 +7141,5 @@ $pgque_extension_body$ end $wrapper$; \echo '' -\echo 'PgQue 0.2.0 registered with pg_tle.' +\echo 'PgQue 0.2.1 registered with pg_tle.' \echo 'Run create extension pgque; to materialise the schema in this database.' diff --git a/sql/pgque.sql b/sql/pgque.sql index e3ac01a7..60a9885c 100644 --- a/sql/pgque.sql +++ b/sql/pgque.sql @@ -1,5 +1,5 @@ -- pgque.sql -- PgQ Universal Edition --- Version: 0.2.0 +-- Version: 0.2.1 -- Copyright 2026 Nikolay Samokhvalov. Apache-2.0 license. -- Includes code derived from PgQ (ISC license, Marko Kreen / Skype Technologies OU). -- @@ -4570,7 +4570,7 @@ $$ language plpgsql security definer set search_path = pgque, pg_catalog; create or replace function pgque.version() returns text as $$ begin - return '0.2.0'; + return '0.2.1'; end; $$ language plpgsql security definer set search_path = pgque, pg_catalog; @@ -5363,13 +5363,16 @@ begin ev_extra1, ev_extra2, ev_extra3, ev_extra4 from pgque.get_batch_events(v_batch_id) loop + if cnt = i_max_return then + raise exception 'pgque.receive: batch exceeds max_return of %', i_max_return + using hint = 'Retry with a larger resource-safe max_return to receive the complete batch. Do not acknowledge after this error.'; + end if; return next row( ev.ev_id, v_batch_id, ev.ev_type, ev.ev_data, ev.ev_retry, ev.ev_time, ev.ev_extra1, ev.ev_extra2, ev.ev_extra3, ev.ev_extra4 )::pgque.message; cnt := cnt + 1; - exit when cnt >= i_max_return; end loop; -- Empty batch: finish immediately to advance the consumer cursor. @@ -6617,6 +6620,10 @@ begin ev_extra4 from pgque.get_batch_events(v_batch_id) loop + if cnt = i_max_return then + raise exception 'pgque.receive_coop: batch exceeds max_return of %', i_max_return + using hint = 'Retry with a larger resource-safe max_return to receive the complete batch. Do not acknowledge after this error.'; + end if; return next row( ev.ev_id, v_batch_id, @@ -6630,7 +6637,6 @@ begin ev.ev_extra4 )::pgque.message; cnt := cnt + 1; - exit when cnt >= i_max_return; end loop; -- Empty batch: release the member token so the subconsumer is not wedged diff --git a/tests/run_all.sql b/tests/run_all.sql index 7ed6a1bf..8c5235b0 100644 --- a/tests/run_all.sql +++ b/tests/run_all.sql @@ -84,6 +84,7 @@ \echo 'Running: test_api_receive' \i tests/test_api_receive.sql +\i tests/test_receive_overflow.sql \echo 'Running: test_api_dlq' \i tests/test_api_dlq.sql diff --git a/tests/test_api_receive.sql b/tests/test_api_receive.sql index 3beb0598..2742b8b3 100644 --- a/tests/test_api_receive.sql +++ b/tests/test_api_receive.sql @@ -58,7 +58,7 @@ begin assert v_count = 0, 'should have no more messages after ack'; end $$; --- Step 6: partial receive still acks the whole underlying batch +-- Step 6: a receive ceiling smaller than the batch fails closed do $$ begin perform pgque.create_queue('test_recv_partial'); @@ -80,18 +80,30 @@ end $$; do $$ declare - v_msg pgque.message; + v_raised boolean := false; v_count int := 0; v_batch_id bigint; + v_msg pgque.message; begin - for v_msg in select * from pgque.receive('test_recv_partial', 'c1', 1) + begin + perform * from pgque.receive('test_recv_partial', 'c1', 1); + exception + when others then + v_raised := true; + assert sqlerrm like '%batch exceeds max_return of 1%', + 'unexpected overflow error: ' || sqlerrm; + end; + assert v_raised, 'receive(..., 1) must reject a three-event batch'; + + for v_msg in select * from pgque.receive('test_recv_partial', 'c1', 3) loop v_count := v_count + 1; v_batch_id := v_msg.batch_id; end loop; - assert v_count = 1, 'receive(..., 1) should return exactly 1 row'; - assert v_batch_id is not null, 'batch_id should be set for partial receive'; + assert v_count = 3, + 'receive retry must return the complete three-event batch'; + assert v_batch_id is not null, 'batch_id should be set for complete receive'; perform pgque.ack(v_batch_id); end $$; @@ -107,7 +119,7 @@ begin end loop; assert v_count = 0, - 'ack(batch_id) should finish the whole batch, even if receive(..., 1) returned one row'; + 'ack(batch_id) should finish the complete batch'; end $$; -- Step 7: send(text) fast path must store payload byte-for-byte diff --git a/tests/test_pgque_config.sql b/tests/test_pgque_config.sql index 7740dbf9..3e9e28fd 100644 --- a/tests/test_pgque_config.sql +++ b/tests/test_pgque_config.sql @@ -15,8 +15,8 @@ begin end; -- Version function works - assert pgque.version() = '0.2.0', - 'version should be 0.2.0, got ' || pgque.version(); + assert pgque.version() = '0.2.1', + 'version should be 0.2.1, got ' || pgque.version(); raise notice 'PASS: pgque_config'; end $$; diff --git a/tests/test_receive_overflow.sql b/tests/test_receive_overflow.sql new file mode 100644 index 00000000..b5672e57 --- /dev/null +++ b/tests/test_receive_overflow.sql @@ -0,0 +1,367 @@ +\set ON_ERROR_STOP on + +-- Regression: receive ceilings must fail closed instead of truncating a batch. +-- Copyright 2026 Nikolay Samokhvalov. Apache-2.0 license. + +-- Plain receive: N+1 raises, rolls back batch allocation, and all events retry. +do $$ +begin + perform pgque.create_queue('recv_overflow'); + perform pgque.register_consumer('recv_overflow', 'c1'); + perform pgque.send('recv_overflow', 'ev', '{"n":1}'::text); + perform pgque.send('recv_overflow', 'ev', '{"n":2}'::text); + perform pgque.send('recv_overflow', 'ev', '{"n":3}'::text); +end $$; + +select pgque.force_next_tick('recv_overflow'); +select pgque.ticker(); + +do $$ +declare + v_raised boolean := false; + v_count int := 0; + v_batch_id bigint; + v_payloads text[] := array[]::text[]; + v_msg pgque.message; + v_before record; + v_after record; + v_hint text; +begin + select s.* into v_before + from pgque.subscription as s + join pgque.queue as q on q.queue_id = s.sub_queue + join pgque.consumer as c on c.co_id = s.sub_consumer + where q.queue_name = 'recv_overflow' and c.co_name = 'c1'; + + begin + perform * from pgque.receive('recv_overflow', 'c1', 2); + exception + when others then + v_raised := true; + get stacked diagnostics v_hint = pg_exception_hint; + assert sqlerrm like '%batch exceeds max_return of 2%', + 'unexpected receive overflow error: ' || sqlerrm; + assert v_hint like '%Do not acknowledge after this error%', + 'receive overflow must provide an actionable no-ack hint'; + end; + assert v_raised, 'receive must raise when a batch contains N+1 events'; + + select s.* into v_after + from pgque.subscription as s + join pgque.queue as q on q.queue_id = s.sub_queue + join pgque.consumer as c on c.co_id = s.sub_consumer + where q.queue_name = 'recv_overflow' and c.co_name = 'c1'; + assert v_after.sub_batch is not distinct from v_before.sub_batch, + 'overflow must roll back the active batch assignment'; + assert v_after.sub_last_tick is not distinct from v_before.sub_last_tick, + 'overflow must roll back the consumer cursor'; + assert v_after.sub_next_tick is not distinct from v_before.sub_next_tick, + 'overflow must roll back the next-tick boundary'; + + for v_msg in select * from pgque.receive('recv_overflow', 'c1', 3) + loop + v_count := v_count + 1; + v_batch_id := v_msg.batch_id; + v_payloads := array_append(v_payloads, v_msg.payload); + end loop; + + assert v_count = 3, + format('receive retry must return all 3 events, got %s', v_count); + assert v_payloads @> array['{"n":1}', '{"n":2}', '{"n":3}'], + format('receive retry lost events: %s', v_payloads); + perform pgque.ack(v_batch_id); +end $$; + +-- An already allocated batch must also fail closed and remain retryable. +do $$ +begin + perform pgque.create_queue('recv_active'); + perform pgque.register_consumer('recv_active', 'c1'); + perform pgque.send('recv_active', 'ev', 'one'); + perform pgque.send('recv_active', 'ev', 'two'); + perform pgque.send('recv_active', 'ev', 'three'); +end $$; + +select pgque.force_next_tick('recv_active'); +select pgque.ticker(); + +do $$ +declare + v_active_batch bigint; + v_after_batch bigint; + v_count int := 0; + v_raised boolean := false; + v_msg pgque.message; +begin + select batch_id into v_active_batch + from pgque.receive('recv_active', 'c1', 3) + limit 1; + assert v_active_batch is not null, 'expected an allocated batch'; + + begin + perform * from pgque.receive('recv_active', 'c1', 2); + exception + when others then + v_raised := true; + end; + assert v_raised, 'overflow must also reject an already allocated batch'; + + select s.sub_batch into v_after_batch + from pgque.subscription as s + join pgque.queue as q on q.queue_id = s.sub_queue + join pgque.consumer as c on c.co_id = s.sub_consumer + where q.queue_name = 'recv_active' and c.co_name = 'c1'; + assert v_after_batch = v_active_batch, + 'overflow must preserve the already allocated batch token'; + + for v_msg in select * from pgque.receive('recv_active', 'c1', 3) + loop + v_count := v_count + 1; + assert v_msg.batch_id = v_active_batch, + 'overflow retry must retain the existing batch token'; + end loop; + assert v_count = 3, + format('existing-batch retry must return all 3 events, got %s', v_count); + perform pgque.ack(v_active_batch); +end $$; + +-- Exactly N remains valid, including the largest int ceiling on an empty batch. +do $$ +begin + perform pgque.create_queue('recv_exact'); + perform pgque.register_consumer('recv_exact', 'c1'); + perform pgque.send('recv_exact', 'ev', 'one'); + perform pgque.send('recv_exact', 'ev', 'two'); +end $$; + +-- A batch containing N-1 events also succeeds. +do $$ +begin + perform pgque.create_queue('recv_nminus1'); + perform pgque.register_consumer('recv_nminus1', 'c1'); + perform pgque.send('recv_nminus1', 'ev', 'one'); + perform pgque.send('recv_nminus1', 'ev', 'two'); +end $$; + +select pgque.force_next_tick('recv_nminus1'); +select pgque.ticker(); + +do $$ +declare + v_count int := 0; + v_batch_id bigint; + v_msg pgque.message; +begin + for v_msg in select * from pgque.receive('recv_nminus1', 'c1', 3) + loop + v_count := v_count + 1; + v_batch_id := v_msg.batch_id; + end loop; + assert v_count = 2, + format('N-1 events must succeed, got %s', v_count); + perform pgque.ack(v_batch_id); +end $$; + +select pgque.force_next_tick('recv_exact'); +select pgque.ticker(); + +do $$ +declare + v_count int := 0; + v_batch_id bigint; + v_msg pgque.message; +begin + for v_msg in select * from pgque.receive('recv_exact', 'c1', 2) + loop + v_count := v_count + 1; + v_batch_id := v_msg.batch_id; + end loop; + assert v_count = 2, format('exactly N events must succeed, got %s', v_count); + perform pgque.ack(v_batch_id); + + perform * from pgque.receive('recv_exact', 'c1', 2147483647); +end $$; + +-- Preserve the SQL NULL behavior: it means no explicit ceiling. +do $$ +begin + perform pgque.create_queue('recv_null'); + perform pgque.register_consumer('recv_null', 'c1'); + perform pgque.send('recv_null', 'ev', 'one'); + perform pgque.send('recv_null', 'ev', 'two'); + perform pgque.send('recv_null', 'ev', 'three'); +end $$; + +select pgque.force_next_tick('recv_null'); +select pgque.ticker(); + +do $$ +declare + v_count int := 0; + v_batch_id bigint; + v_msg pgque.message; +begin + for v_msg in select * from pgque.receive('recv_null', 'c1', null) + loop + v_count := v_count + 1; + v_batch_id := v_msg.batch_id; + end loop; + assert v_count = 3, + format('NULL max_return must preserve unbounded receive, got %s', v_count); + perform pgque.ack(v_batch_id); +end $$; + +-- Cooperative receive has the same fail-closed and retry contract. +do $$ +begin + perform pgque.create_queue('recv_coop_overflow'); + perform pgque.register_subconsumer('recv_coop_overflow', 'main_c', 'w1'); + perform pgque.send('recv_coop_overflow', 'ev', 'one'); + perform pgque.send('recv_coop_overflow', 'ev', 'two'); + perform pgque.send('recv_coop_overflow', 'ev', 'three'); +end $$; + +select pgque.force_next_tick('recv_coop_overflow'); +select pgque.ticker(); + +do $$ +declare + v_raised boolean := false; + v_count int := 0; + v_batch_id bigint; + v_payloads text[] := array[]::text[]; + v_msg pgque.message; + v_hint text; +begin + begin + perform * from pgque.receive_coop( + 'recv_coop_overflow', 'main_c', 'w1', 2 + ); + exception + when others then + v_raised := true; + get stacked diagnostics v_hint = pg_exception_hint; + assert sqlerrm like '%batch exceeds max_return of 2%', + 'unexpected receive_coop overflow error: ' || sqlerrm; + assert v_hint like '%Do not acknowledge after this error%', + 'receive_coop overflow must provide an actionable no-ack hint'; + end; + assert v_raised, 'receive_coop must raise when a batch contains N+1 events'; + + for v_msg in + select * from pgque.receive_coop( + 'recv_coop_overflow', 'main_c', 'w1', 3 + ) + loop + v_count := v_count + 1; + v_batch_id := v_msg.batch_id; + v_payloads := array_append(v_payloads, v_msg.payload); + end loop; + assert v_count = 3, + format('receive_coop retry must return all 3 events, got %s', v_count); + assert v_payloads @> array['one', 'two', 'three'], + format('receive_coop retry lost events: %s', v_payloads); + perform pgque.ack(v_batch_id); +end $$; + +-- A failed stale takeover must roll back ownership to the original member. +do $$ +begin + perform pgque.create_queue('recv_coop_takeover'); + perform pgque.register_subconsumer('recv_coop_takeover', 'main_c', 'w1'); + perform pgque.register_subconsumer('recv_coop_takeover', 'main_c', 'w2'); + perform pgque.send('recv_coop_takeover', 'ev', 'one'); + perform pgque.send('recv_coop_takeover', 'ev', 'two'); + perform pgque.send('recv_coop_takeover', 'ev', 'three'); +end $$; + +select pgque.force_next_tick('recv_coop_takeover'); +select pgque.ticker(); + +do $$ +declare + v_old_batch bigint; + v_w1_batch bigint; + v_w2_batch bigint; + v_new_batch bigint; + v_count int := 0; + v_raised boolean := false; + v_payloads text[] := array[]::text[]; + v_msg pgque.message; +begin + select batch_id into v_old_batch + from pgque.receive_coop('recv_coop_takeover', 'main_c', 'w1', 3) + limit 1; + + update pgque.subscription as s + set sub_active = now() - interval '10 minutes' + from pgque.queue as q + cross join pgque.consumer as c + where q.queue_name = 'recv_coop_takeover' + and c.co_name = 'main_c.w1' + and s.sub_queue = q.queue_id + and s.sub_consumer = c.co_id; + + begin + perform * from pgque.receive_coop( + 'recv_coop_takeover', 'main_c', 'w2', 2, interval '1 minute' + ); + exception + when others then + v_raised := true; + end; + assert v_raised, 'overflow must abort a stale takeover'; + + select s.sub_batch into v_w1_batch + from pgque.subscription as s + join pgque.queue as q on q.queue_id = s.sub_queue + join pgque.consumer as c on c.co_id = s.sub_consumer + where q.queue_name = 'recv_coop_takeover' and c.co_name = 'main_c.w1'; + select s.sub_batch into v_w2_batch + from pgque.subscription as s + join pgque.queue as q on q.queue_id = s.sub_queue + join pgque.consumer as c on c.co_id = s.sub_consumer + where q.queue_name = 'recv_coop_takeover' and c.co_name = 'main_c.w2'; + assert v_w1_batch = v_old_batch, + 'failed takeover must preserve the stale owner batch'; + assert v_w2_batch is null, + 'failed takeover must not leave the new member owning a batch'; + + for v_msg in + select * from pgque.receive_coop( + 'recv_coop_takeover', 'main_c', 'w2', 3, interval '1 minute' + ) + loop + v_count := v_count + 1; + v_new_batch := v_msg.batch_id; + v_payloads := array_append(v_payloads, v_msg.payload); + end loop; + assert v_count = 3 and v_payloads @> array['one', 'two', 'three'], + format('takeover retry lost events: count=%s payloads=%s', v_count, v_payloads); + assert v_new_batch is not null and v_new_batch <> v_old_batch, + 'successful takeover retry must issue a fresh batch token'; + perform pgque.ack(v_new_batch); +end $$; + + +do $$ +begin + perform pgque.unregister_consumer('recv_overflow', 'c1'); + perform pgque.drop_queue('recv_overflow'); + perform pgque.unregister_consumer('recv_active', 'c1'); + perform pgque.drop_queue('recv_active'); + perform pgque.unregister_consumer('recv_exact', 'c1'); + perform pgque.drop_queue('recv_exact'); + perform pgque.unregister_consumer('recv_nminus1', 'c1'); + perform pgque.drop_queue('recv_nminus1'); + perform pgque.unregister_consumer('recv_null', 'c1'); + perform pgque.drop_queue('recv_null'); + perform pgque.unregister_subconsumer('recv_coop_overflow', 'main_c', 'w1'); + perform pgque.unregister_consumer('recv_coop_overflow', 'main_c'); + perform pgque.drop_queue('recv_coop_overflow'); + perform pgque.unregister_subconsumer('recv_coop_takeover', 'main_c', 'w1'); + perform pgque.unregister_subconsumer('recv_coop_takeover', 'main_c', 'w2'); + perform pgque.unregister_consumer('recv_coop_takeover', 'main_c'); + perform pgque.drop_queue('recv_coop_takeover'); + raise notice 'PASS: receive ceilings fail closed and preserve complete batches'; +end $$; From c167cfd93bc8353ac1333822db07591ddb0fe4f1 Mon Sep 17 00:00:00 2001 From: samo-agent <280144521+samo-agent@users.noreply.github.com> Date: Thu, 1 Oct 2026 14:24:31 +0000 Subject: [PATCH 2/5] test: assert SDK overflow recovery Refs #365. Backport complete-batch retry assertions. --- clients/go/integration_test.go | 44 +++++++++++++++++++++++----- clients/python/tests/test_receive.py | 20 +++++++++++-- 2 files changed, 54 insertions(+), 10 deletions(-) diff --git a/clients/go/integration_test.go b/clients/go/integration_test.go index 67c9c900..15c483bc 100644 --- a/clients/go/integration_test.go +++ b/clients/go/integration_test.go @@ -4,6 +4,7 @@ package pgque_test import ( "context" + "errors" "strings" "testing" "time" @@ -76,8 +77,9 @@ func TestSend_MultipleEventsOneBatch(t *testing.T) { } } -// TestReceive_RespectsMaxBatch ensures Receive returns at most maxMessages. -func TestReceive_RespectsMaxBatch(t *testing.T) { +// TestReceive_RejectsBatchOverMax ensures Receive fails closed and leaves the +// complete batch available for a retry with a sufficient ceiling. +func TestReceive_RejectsBatchOverMax(t *testing.T) { client := connectOrSkip(t) defer client.Close() queue, consumer := setupFreshQueue(t, client) @@ -94,14 +96,42 @@ func TestReceive_RespectsMaxBatch(t *testing.T) { tick(t, client, queue) msgs, err := client.Receive(ctx, queue, consumer, 10) + if err == nil { + t.Fatal("expected oversized batch to raise") + } + if len(msgs) != 0 { + t.Fatalf("oversized receive returned %d partial messages", len(msgs)) + } + var sqlErr *pgque.SQLError + if !errors.As(err, &sqlErr) || sqlErr.SQLSTATE != "P0001" { + t.Fatalf("expected propagated P0001 SQLError, got %T: %v", err, err) + } + if !strings.Contains(err.Error(), "batch exceeds max_return of 10") { + t.Fatalf("unexpected overflow error: %v", err) + } + + msgs, err = client.Receive(ctx, queue, consumer, total) if err != nil { - t.Fatal(err) + t.Fatalf("complete retry failed: %v", err) } - if len(msgs) > 10 { - t.Fatalf("Receive returned %d messages, expected ≤ 10", len(msgs)) + if len(msgs) != total { + t.Fatalf("complete retry returned %d messages, expected %d", len(msgs), total) } - if len(msgs) > 0 { - _, _ = client.Ack(ctx, msgs[0].BatchID) + batchID := msgs[0].BatchID + for _, msg := range msgs { + if msg.BatchID != batchID { + t.Fatalf("retry returned mixed batch ids %d and %d", batchID, msg.BatchID) + } + } + if _, err := client.Ack(ctx, batchID); err != nil { + t.Fatalf("ack after complete retry failed: %v", err) + } + msgs, err = client.Receive(ctx, queue, consumer, total) + if err != nil { + t.Fatalf("receive after ack failed: %v", err) + } + if len(msgs) != 0 { + t.Fatalf("expected no messages after ack, got %d", len(msgs)) } } diff --git a/clients/python/tests/test_receive.py b/clients/python/tests/test_receive.py index da1d9b08..77bea147 100644 --- a/clients/python/tests/test_receive.py +++ b/clients/python/tests/test_receive.py @@ -2,7 +2,10 @@ """Consumer-side tests: ``Client.receive`` / ``Client.ack``.""" +import json + import pgque +import pytest def test_receive_empty_when_no_tick(conn, setup_queue): @@ -52,7 +55,7 @@ def test_ack_advances_position(conn, setup_queue): assert msgs2 == [] -def test_receive_returns_at_most_max_messages(conn, setup_queue): +def test_receive_rejects_batch_over_max_and_retries_complete(conn, setup_queue): queue, consumer = setup_queue client = pgque.PgqueClient(conn) for i in range(5): @@ -61,10 +64,21 @@ def test_receive_returns_at_most_max_messages(conn, setup_queue): conn.execute("select pgque.force_next_tick(%s)", (queue,)) conn.execute("select pgque.ticker(%s)", (queue,)) conn.commit() - msgs = client.receive(queue, consumer, max_messages=3) - assert len(msgs) == 3 + with pytest.raises(pgque.PgqueError, match="batch exceeds max_return of 3"): + client.receive(queue, consumer, max_messages=3) + + # The failed statement returns no partial list and rolls back its batch + # allocation. Reset the failed transaction before retrying. + conn.rollback() + msgs = client.receive(queue, consumer, max_messages=5) + assert len(msgs) == 5 + assert { + (m.payload if isinstance(m.payload, dict) else json.loads(m.payload))["i"] + for m in msgs + } == set(range(5)) client.ack(msgs[0].batch_id) conn.commit() + assert client.receive(queue, consumer, max_messages=5) == [] def test_receive_preserves_event_type(conn, setup_queue): From 22e0b240e18c60eb46f7f525fff97bb4460a96f3 Mon Sep 17 00:00:00 2001 From: samo-agent <280144521+samo-agent@users.noreply.github.com> Date: Thu, 1 Oct 2026 14:39:26 +0000 Subject: [PATCH 3/5] fix: preserve pg_tle state on 0.2.1 update Refs #365. Register a function-only update edge, classify overflow as SQLSTATE54000, and tighten regression assertions. --- .github/workflows/ci.yml | 10 + build/transform.sh | 94 ++++++-- clients/go/integration_test.go | 4 +- clients/python/tests/test_receive.py | 5 +- docs/reference.md | 8 +- docs/upgrading.md | 13 ++ sql/pgque-api/cooperative_consumers.sql | 4 +- sql/pgque-api/receive.sql | 4 +- sql/pgque-tle-updates/pgque--0.2.0--0.2.1.sql | 125 +++++++++++ sql/pgque-tle.sql | 211 +++++++++++++++--- sql/pgque.sql | 8 +- tests/test_receive_overflow.sql | 64 ++++-- tests/test_tle_upgrade_v0_2.sql | 170 ++++++++++++++ 13 files changed, 643 insertions(+), 77 deletions(-) create mode 100644 sql/pgque-tle-updates/pgque--0.2.0--0.2.1.sql create mode 100644 tests/test_tle_upgrade_v0_2.sql diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index c22bc795..c1c9b497 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -173,6 +173,7 @@ jobs: - uses: actions/checkout@v4 with: submodules: recursive + fetch-depth: 0 - name: Build pgque run: bash build/transform.sh @@ -217,6 +218,15 @@ jobs: PGPASSWORD=pgque_test psql -h localhost -U postgres -d pgque_test \ -v ON_ERROR_STOP=1 -f tests/test_tle_install.sql + - name: Run pg_tle 0.2.0 upgrade test + run: | + set -Eeuo pipefail + git show v0.2.0:sql/pgque-tle.sql > /tmp/pgque-tle-v0.2.0.sql + PGPASSWORD=pgque_test psql -h localhost -U postgres -d postgres \ + -v ON_ERROR_STOP=1 -c 'create database pgque_tle_upgrade' + PGPASSWORD=pgque_test psql -h localhost -U postgres -d pgque_tle_upgrade \ + -v ON_ERROR_STOP=1 -f tests/test_tle_upgrade_v0_2.sql + - name: Run full regression suite via pg_tle install run: | set -Eeuo pipefail diff --git a/build/transform.sh b/build/transform.sh index b667c749..ab6bc05f 100755 --- a/build/transform.sh +++ b/build/transform.sh @@ -992,12 +992,15 @@ echo "=== Packaging pg_tle install script ===" PGTLE_FILE="${SQL_DIR}/pgque-tle.sql" PGTLE_DOLLAR_TAG="pgque_extension_body" +PGTLE_UPDATE_DOLLAR_TAG="pgque_update_body" +PGTLE_UPDATE_DIR="${SQL_DIR}/pgque-tle-updates" +PGTLE_UPDATE_FILE="${PGTLE_UPDATE_DIR}/pgque--0.2.0--0.2.1.sql" # Sanity check: neither dollar-quote tag (the body tag and the outer wrapper # tag) may appear in the body content, otherwise the literal terminates early # and pg_tle.install_extension receives a truncated body. Use grep -F so the # `$` characters are matched literally instead of as end-of-line anchors. -for tag in "${PGTLE_DOLLAR_TAG}" wrapper; do +for tag in "${PGTLE_DOLLAR_TAG}" "${PGTLE_UPDATE_DOLLAR_TAG}" wrapper; do if grep -qF "\$${tag}\$" "${INSTALL_FILE}"; then echo "FAIL: dollar-quote tag \$${tag}\$ collides with body content" >&2 exit 1 @@ -1027,6 +1030,35 @@ if ! [[ "${PGTLE_VERSION}" =~ ^[0-9]+\.[0-9]+\.[0-9]+(-[A-Za-z0-9.]+)?$ ]]; then exit 1 fi +# Generate the minimal 0.2.0 -> 0.2.1 update body from the same source +# definitions used by the default installer. This function-only update keeps +# all queue, subscription, active-batch, retry, and DLQ table state intact. +mkdir -p "${PGTLE_UPDATE_DIR}" +cat > "${PGTLE_UPDATE_FILE}" << 'HEADER' +-- pgque--0.2.0--0.2.1.sql -- pg_tle update body. +-- Auto-generated by build/transform.sh. Do not edit or execute directly. +-- Copyright 2026 Nikolay Samokhvalov. Apache-2.0 license. +-- Includes code derived from PgQ (ISC license, Marko Kreen / Skype Technologies OU). + +HEADER + +sed -n \ + '/^create or replace function pgque\.receive(/,/^\$\$ language plpgsql security definer set search_path = pgque, pg_catalog;$/p' \ + "${API_DIR}/receive.sql" >> "${PGTLE_UPDATE_FILE}" +echo "" >> "${PGTLE_UPDATE_FILE}" +sed -n \ + '/^create or replace function pgque\.receive_coop(/,/^\$\$ language plpgsql security definer set search_path = pgque, pg_catalog;$/p' \ + "${API_DIR}/cooperative_consumers.sql" >> "${PGTLE_UPDATE_FILE}" +echo "" >> "${PGTLE_UPDATE_FILE}" +sed -n \ + '/^create or replace function pgque\.version()/,/^\$\$ language plpgsql security definer set search_path = pgque, pg_catalog;$/p' \ + "${ADDITIONS_DIR}/lifecycle.sql" >> "${PGTLE_UPDATE_FILE}" + +if [[ $(grep -c '^create or replace function pgque\.' "${PGTLE_UPDATE_FILE}") -ne 3 ]]; then + echo "FAIL: expected exactly 3 functions in ${PGTLE_UPDATE_FILE}" >&2 + exit 1 +fi + cat > "${PGTLE_FILE}" << HEADER -- pgque-tle.sql -- Install PgQue as a pg_tle (Trusted Language Extension). -- Auto-generated by build/transform.sh from sql/pgque.sql. Do not edit by hand. @@ -1085,14 +1117,27 @@ begin end if; end \$\$; --- Step 3: register the extension body with pg_tle. +-- Step 3: register the extension body or the supported update path with pg_tle. -- Same version already registered -> no-op (so deployment scripts can rerun). --- Different version already registered -> raise so the user goes through the --- explicit uninstall + reinstall path; pg_tle has no managed upgrade path --- between unrelated registrations of an extension. +-- Version 0.2.0 -> register the non-destructive 0.2.1 update path. do \$wrapper\$ declare existing_version text; + update_path_exists boolean; + extension_sql text := \$${PGTLE_DOLLAR_TAG}\$ +HEADER + +cat "${INSTALL_FILE}" >> "${PGTLE_FILE}" + +cat >> "${PGTLE_FILE}" << HEADER +\$${PGTLE_DOLLAR_TAG}\$; + update_sql text := \$${PGTLE_UPDATE_DOLLAR_TAG}\$ +HEADER + +cat "${PGTLE_UPDATE_FILE}" >> "${PGTLE_FILE}" + +cat >> "${PGTLE_FILE}" << FOOTER +\$${PGTLE_UPDATE_DOLLAR_TAG}\$; begin select default_version into existing_version from pgtle.available_extensions() @@ -1103,11 +1148,28 @@ begin return; end if; + if existing_version = '0.2.0' and '${PGTLE_VERSION}' = '0.2.1' then + select exists ( + select 1 + from pgtle.extension_update_paths('pgque') + where source = '0.2.0' + and target = '0.2.1' + and path is not null + ) into update_path_exists; + + if not update_path_exists then + perform pgtle.install_update_path( + 'pgque', '0.2.0', '0.2.1', update_sql + ); + end if; + perform pgtle.set_default_version('pgque', '0.2.1'); + raise notice 'pgque 0.2.0 -> 0.2.1 update path registered with pg_tle.'; + return; + end if; + if existing_version is not null then raise exception 'pgque is already registered with pg_tle at version % ' - 'but this script registers version %. Run ' - 'sql/pgque-tle-uninstall.sql first to remove the existing ' - 'registration, then re-run this script.', + 'but this script only supports a managed update from 0.2.0 to %.', existing_version, '${PGTLE_VERSION}'; end if; @@ -1115,19 +1177,14 @@ begin 'pgque', '${PGTLE_VERSION}', 'PgQue — PgQ Universal Edition (zero-bloat Postgres queue)', -\$${PGTLE_DOLLAR_TAG}\$ -HEADER - -cat "${INSTALL_FILE}" >> "${PGTLE_FILE}" - -cat >> "${PGTLE_FILE}" << FOOTER -\$${PGTLE_DOLLAR_TAG}\$ + extension_sql ); end \$wrapper\$; \\echo '' \\echo 'PgQue ${PGTLE_VERSION} registered with pg_tle.' -\\echo 'Run create extension pgque; to materialise the schema in this database.' +\\echo 'For a new install, run: create extension pgque;' +\\echo 'For an installed 0.2.0 extension, apply ALTER EXTENSION UPDATE as described in the upgrade documentation.' FOOTER pgtle_lines=$(wc -l < "${PGTLE_FILE}") @@ -1146,6 +1203,11 @@ if ! grep -q 'create schema if not exists pgque' "${PGTLE_FILE}"; then pgtle_errors=$((pgtle_errors + 1)) fi +if ! grep -q 'pgtle.install_update_path' "${PGTLE_FILE}"; then + echo "FAIL: pg_tle install script missing managed update path" + pgtle_errors=$((pgtle_errors + 1)) +fi + if [[ ${pgtle_errors} -eq 0 ]]; then echo "=== pg_tle PACKAGING COMPLETE ===" else diff --git a/clients/go/integration_test.go b/clients/go/integration_test.go index 15c483bc..dbbeeb2b 100644 --- a/clients/go/integration_test.go +++ b/clients/go/integration_test.go @@ -103,8 +103,8 @@ func TestReceive_RejectsBatchOverMax(t *testing.T) { t.Fatalf("oversized receive returned %d partial messages", len(msgs)) } var sqlErr *pgque.SQLError - if !errors.As(err, &sqlErr) || sqlErr.SQLSTATE != "P0001" { - t.Fatalf("expected propagated P0001 SQLError, got %T: %v", err, err) + if !errors.As(err, &sqlErr) || sqlErr.SQLSTATE != "54000" { + t.Fatalf("expected propagated 54000 SQLError, got %T: %v", err, err) } if !strings.Contains(err.Error(), "batch exceeds max_return of 10") { t.Fatalf("unexpected overflow error: %v", err) diff --git a/clients/python/tests/test_receive.py b/clients/python/tests/test_receive.py index 77bea147..2f125af0 100644 --- a/clients/python/tests/test_receive.py +++ b/clients/python/tests/test_receive.py @@ -64,8 +64,11 @@ def test_receive_rejects_batch_over_max_and_retries_complete(conn, setup_queue): conn.execute("select pgque.force_next_tick(%s)", (queue,)) conn.execute("select pgque.ticker(%s)", (queue,)) conn.commit() - with pytest.raises(pgque.PgqueError, match="batch exceeds max_return of 3"): + with pytest.raises( + pgque.PgqueError, match="batch exceeds max_return of 3" + ) as exc_info: client.receive(queue, consumer, max_messages=3) + assert exc_info.value.__cause__.sqlstate == "54000" # The failed statement returns no partial list and rolls back its batch # allocation. Reset the failed transaction before retrying. diff --git a/docs/reference.md b/docs/reference.md index 497f64ba..e00e1e8b 100644 --- a/docs/reference.md +++ b/docs/reference.md @@ -119,14 +119,16 @@ All consume-side functions (`receive`, `ack`, `nack`, `subscribe`, `unsubscribe` #### `pgque.receive(queue text, consumer text, max_return int default 100) → setof pgque.message` -Pulls the complete next batch for `consumer` on `queue`, or raises an error if it exceeds `max_return`. `max_return` must be >= 1; passing 0 or a negative value raises an error. Returns an empty set if no batch is available. Each row is a `pgque.message` composite (see [§Message type](#message-type)). +Pulls the complete next batch for `consumer` on `queue`, or raises an error if it exceeds `max_return`. A non-null `max_return` must be >= 1; passing 0 or a negative value raises an error. SQL `null` disables the explicit ceiling for compatibility; prefer an explicit resource-safe ceiling. Returns an empty set if no batch is available. Each row is a `pgque.message` composite (see [§Message type](#message-type)). Grant: `pgque_reader`. Source: `sql/pgque-api/receive.sql`. ```sql select * from pgque.receive('orders', 'processor', 100); ``` -**Complete batch or error.** `max_return` is a safety ceiling, not a page size. Exactly that many events is valid; encountering one more raises an error, with no successful result or consumer advancement from that call. Roll back a failed transaction and retry with a larger ceiling within your resource budget, or use the low-level full-batch interface. Process every event before calling `ack(batch_id)`, which finishes the entire batch. Do not acknowledge after an overflow error. Repeating `receive()` with the same ceiling does not paginate or resolve overflow. +Overflow reports SQLSTATE `54000` (`program_limit_exceeded`). Like other SQL errors, it aborts the enclosing transaction unless recovered through a savepoint. + +**Complete batch or error.** `max_return` is a safety ceiling, not a page size. Exactly that many events is valid; encountering one more raises an error, with no successful result or consumer advancement from that call. Roll back a failed transaction and retry with a larger ceiling within your resource budget, or explicitly manage the complete batch with `pgque.next_batch()`, `pgque.get_batch_events()` and `pgque.finish_batch()` (the low-level PgQ API). Process every event before calling `ack(batch_id)`, which finishes the entire batch. Do not acknowledge after an overflow error. Repeating `receive()` with the same ceiling does not paginate or resolve overflow. `ticker_max_count` is a tick-trigger threshold, not a hard batch-size cap. Setting `max_return` to that value cannot guarantee success during bursts. Monitor receive errors and consumer lag so an undersized ceiling does not silently stall processing. @@ -168,7 +170,7 @@ Grant: `pgque_reader`. Source: `sql/pgque-api/cooperative_consumers.sql`. #### `pgque.receive_coop(queue text, consumer text, subconsumer text, max_return int default 100, dead_interval interval default null) → setof pgque.message` -Receives messages for one subconsumer. `max_return` must be >= 1. `dead_interval` enables stale-batch takeover from another inactive member; takeover allocates a fresh `batch_id`, so old tokens cannot ack/nack the new owner's state. The cooperative group is a trust boundary: callers allowed to use the same `(queue, consumer)` can steal stale batches from each other by design, so do not share one cooperative group across mutually untrusted workers. +Receives messages for one subconsumer. A non-null `max_return` must be >= 1; SQL `null` disables the explicit ceiling as in `receive()`. `dead_interval` enables stale-batch takeover from another inactive member; takeover allocates a fresh `batch_id`, so old tokens cannot ack/nack the new owner's state. The cooperative group is a trust boundary: callers allowed to use the same `(queue, consumer)` can steal stale batches from each other by design, so do not share one cooperative group across mutually untrusted workers. **Auto-registration.** If the logical `consumer` or `subconsumer` is not yet registered, `receive_coop()` registers them on the fly (creates the `coop_main` row on first call, then the `coop_member` row), so a worker can call `receive_coop()` cold without a prior `register_subconsumer`. Use the explicit `register_subconsumer(..., convert_normal => true)` call only when you need to convert an existing normal consumer into a cooperative main. diff --git a/docs/upgrading.md b/docs/upgrading.md index b4fc5900..024a4bf4 100644 --- a/docs/upgrading.md +++ b/docs/upgrading.md @@ -26,6 +26,19 @@ transaction and retry with a sufficient ceiling within your resource budget. Never acknowledge after a receive error. Ticker thresholds do not cap batch size, and repeated receive calls are not pagination. +For a pg_tle-managed v0.2.0 installation, register the update path and apply +it through Postgres extension management: + +```sql +\i sql/pgque-tle.sql +alter extension pgque update to '0.2.1'; +select extversion from pg_extension where extname = 'pgque'; +select pgque.version(); +``` + +Both version queries must return `0.2.1`. The update replaces functions only; +do not drop or unregister the existing extension before upgrading. + ## v0.1.0 to v0.2.0 The supported v0.1.0 → v0.2.0 path is the same re-install procedure: diff --git a/sql/pgque-api/cooperative_consumers.sql b/sql/pgque-api/cooperative_consumers.sql index 14e3f8ff..f14a3341 100644 --- a/sql/pgque-api/cooperative_consumers.sql +++ b/sql/pgque-api/cooperative_consumers.sql @@ -1145,7 +1145,9 @@ begin loop if cnt = i_max_return then raise exception 'pgque.receive_coop: batch exceeds max_return of %', i_max_return - using hint = 'Retry with a larger resource-safe max_return to receive the complete batch. Do not acknowledge after this error.'; + using + errcode = '54000', + hint = 'Retry with a larger resource-safe max_return to receive the complete batch. Do not acknowledge after this error.'; end if; return next row( ev.ev_id, diff --git a/sql/pgque-api/receive.sql b/sql/pgque-api/receive.sql index 419968fc..721170da 100644 --- a/sql/pgque-api/receive.sql +++ b/sql/pgque-api/receive.sql @@ -52,7 +52,9 @@ begin loop if cnt = i_max_return then raise exception 'pgque.receive: batch exceeds max_return of %', i_max_return - using hint = 'Retry with a larger resource-safe max_return to receive the complete batch. Do not acknowledge after this error.'; + using + errcode = '54000', + hint = 'Retry with a larger resource-safe max_return to receive the complete batch. Do not acknowledge after this error.'; end if; return next row( ev.ev_id, v_batch_id, ev.ev_type, ev.ev_data, diff --git a/sql/pgque-tle-updates/pgque--0.2.0--0.2.1.sql b/sql/pgque-tle-updates/pgque--0.2.0--0.2.1.sql new file mode 100644 index 00000000..37f8f7e0 --- /dev/null +++ b/sql/pgque-tle-updates/pgque--0.2.0--0.2.1.sql @@ -0,0 +1,125 @@ +-- pgque--0.2.0--0.2.1.sql -- pg_tle update body. +-- Auto-generated by build/transform.sh. Do not edit or execute directly. +-- Copyright 2026 Nikolay Samokhvalov. Apache-2.0 license. +-- Includes code derived from PgQ (ISC license, Marko Kreen / Skype Technologies OU). + +create or replace function pgque.receive( + i_queue text, i_consumer text, i_max_return int default 100) +returns setof pgque.message as $$ +declare + v_batch_id bigint; + ev record; + cnt int := 0; +begin + if i_max_return < 1 then + raise exception 'pgque.receive: max_return must be >= 1, got %', i_max_return; + end if; + + -- Get next batch (may return NULL if no tick window is ready) + v_batch_id := pgque.next_batch(i_queue, i_consumer); + if v_batch_id is null then + return; + end if; + + -- Yield messages from the batch + for ev in + select ev_id, ev_type, ev_data, ev_retry, ev_time, + ev_extra1, ev_extra2, ev_extra3, ev_extra4 + from pgque.get_batch_events(v_batch_id) + loop + if cnt = i_max_return then + raise exception 'pgque.receive: batch exceeds max_return of %', i_max_return + using + errcode = '54000', + hint = 'Retry with a larger resource-safe max_return to receive the complete batch. Do not acknowledge after this error.'; + end if; + return next row( + ev.ev_id, v_batch_id, ev.ev_type, ev.ev_data, + ev.ev_retry, ev.ev_time, + ev.ev_extra1, ev.ev_extra2, ev.ev_extra3, ev.ev_extra4 + )::pgque.message; + cnt := cnt + 1; + end loop; + + -- Empty batch: finish immediately to advance the consumer cursor. + if cnt = 0 then + perform pgque.finish_batch(v_batch_id); + end if; + + return; +end; +$$ language plpgsql security definer set search_path = pgque, pg_catalog; + +create or replace function pgque.receive_coop( + i_queue text, + i_consumer text, + i_subconsumer text, + i_max_return int default 100, + i_dead_interval interval default null) +returns setof pgque.message as $$ +declare + v_batch_id bigint; + ev record; + cnt int := 0; +begin + if i_max_return < 1 then + raise exception 'pgque.receive_coop: max_return must be >= 1, got %', i_max_return; + end if; + + v_batch_id := pgque.next_batch(i_queue, i_consumer, i_subconsumer, i_dead_interval); + if v_batch_id is null then + return; + end if; + + for ev in + select + ev_id, + ev_type, + ev_data, + ev_retry, + ev_time, + ev_extra1, + ev_extra2, + ev_extra3, + ev_extra4 + from pgque.get_batch_events(v_batch_id) + loop + if cnt = i_max_return then + raise exception 'pgque.receive_coop: batch exceeds max_return of %', i_max_return + using + errcode = '54000', + hint = 'Retry with a larger resource-safe max_return to receive the complete batch. Do not acknowledge after this error.'; + end if; + return next row( + ev.ev_id, + v_batch_id, + ev.ev_type, + ev.ev_data, + ev.ev_retry, + ev.ev_time, + ev.ev_extra1, + ev.ev_extra2, + ev.ev_extra3, + ev.ev_extra4 + )::pgque.message; + cnt := cnt + 1; + end loop; + + -- Empty batch: release the member token so the subconsumer is not wedged + -- on a tick window with no visible events. finish_batch on a coop_member + -- clears sub_batch + sub_last_tick + sub_next_tick (it does not advance + -- the main cursor, which already moved when the batch was allocated). + if cnt = 0 then + perform pgque.finish_batch(v_batch_id); + end if; + + return; +end; +$$ language plpgsql security definer set search_path = pgque, pg_catalog; + +create or replace function pgque.version() +returns text as $$ +begin + return '0.2.1'; +end; +$$ language plpgsql security definer set search_path = pgque, pg_catalog; diff --git a/sql/pgque-tle.sql b/sql/pgque-tle.sql index 33aaea03..f04f1be7 100644 --- a/sql/pgque-tle.sql +++ b/sql/pgque-tle.sql @@ -55,37 +55,14 @@ begin end if; end $$; --- Step 3: register the extension body with pg_tle. +-- Step 3: register the extension body or the supported update path with pg_tle. -- Same version already registered -> no-op (so deployment scripts can rerun). --- Different version already registered -> raise so the user goes through the --- explicit uninstall + reinstall path; pg_tle has no managed upgrade path --- between unrelated registrations of an extension. +-- Version 0.2.0 -> register the non-destructive 0.2.1 update path. do $wrapper$ declare existing_version text; -begin - select default_version into existing_version - from pgtle.available_extensions() - where name = 'pgque'; - - if existing_version = '0.2.1' then - raise notice 'pgque 0.2.1 already registered with pg_tle; skipping install_extension().'; - return; - end if; - - if existing_version is not null then - raise exception 'pgque is already registered with pg_tle at version % ' - 'but this script registers version %. Run ' - 'sql/pgque-tle-uninstall.sql first to remove the existing ' - 'registration, then re-run this script.', - existing_version, '0.2.1'; - end if; - - perform pgtle.install_extension( - 'pgque', - '0.2.1', - 'PgQue — PgQ Universal Edition (zero-bloat Postgres queue)', -$pgque_extension_body$ + update_path_exists boolean; + extension_sql text := $pgque_extension_body$ -- pgque.sql -- PgQ Universal Edition -- Version: 0.2.1 -- Copyright 2026 Nikolay Samokhvalov. Apache-2.0 license. @@ -5453,7 +5430,9 @@ begin loop if cnt = i_max_return then raise exception 'pgque.receive: batch exceeds max_return of %', i_max_return - using hint = 'Retry with a larger resource-safe max_return to receive the complete batch. Do not acknowledge after this error.'; + using + errcode = '54000', + hint = 'Retry with a larger resource-safe max_return to receive the complete batch. Do not acknowledge after this error.'; end if; return next row( ev.ev_id, v_batch_id, ev.ev_type, ev.ev_data, @@ -6710,7 +6689,9 @@ begin loop if cnt = i_max_return then raise exception 'pgque.receive_coop: batch exceeds max_return of %', i_max_return - using hint = 'Retry with a larger resource-safe max_return to receive the complete batch. Do not acknowledge after this error.'; + using + errcode = '54000', + hint = 'Retry with a larger resource-safe max_return to receive the complete batch. Do not acknowledge after this error.'; end if; return next row( ev.ev_id, @@ -7136,10 +7117,178 @@ revoke execute on function pgque.insert_event_bulk(text, text, text[]) revoke execute on all functions in schema pgque from public; -$pgque_extension_body$ +$pgque_extension_body$; + update_sql text := $pgque_update_body$ +-- pgque--0.2.0--0.2.1.sql -- pg_tle update body. +-- Auto-generated by build/transform.sh. Do not edit or execute directly. +-- Copyright 2026 Nikolay Samokhvalov. Apache-2.0 license. +-- Includes code derived from PgQ (ISC license, Marko Kreen / Skype Technologies OU). + +create or replace function pgque.receive( + i_queue text, i_consumer text, i_max_return int default 100) +returns setof pgque.message as $$ +declare + v_batch_id bigint; + ev record; + cnt int := 0; +begin + if i_max_return < 1 then + raise exception 'pgque.receive: max_return must be >= 1, got %', i_max_return; + end if; + + -- Get next batch (may return NULL if no tick window is ready) + v_batch_id := pgque.next_batch(i_queue, i_consumer); + if v_batch_id is null then + return; + end if; + + -- Yield messages from the batch + for ev in + select ev_id, ev_type, ev_data, ev_retry, ev_time, + ev_extra1, ev_extra2, ev_extra3, ev_extra4 + from pgque.get_batch_events(v_batch_id) + loop + if cnt = i_max_return then + raise exception 'pgque.receive: batch exceeds max_return of %', i_max_return + using + errcode = '54000', + hint = 'Retry with a larger resource-safe max_return to receive the complete batch. Do not acknowledge after this error.'; + end if; + return next row( + ev.ev_id, v_batch_id, ev.ev_type, ev.ev_data, + ev.ev_retry, ev.ev_time, + ev.ev_extra1, ev.ev_extra2, ev.ev_extra3, ev.ev_extra4 + )::pgque.message; + cnt := cnt + 1; + end loop; + + -- Empty batch: finish immediately to advance the consumer cursor. + if cnt = 0 then + perform pgque.finish_batch(v_batch_id); + end if; + + return; +end; +$$ language plpgsql security definer set search_path = pgque, pg_catalog; + +create or replace function pgque.receive_coop( + i_queue text, + i_consumer text, + i_subconsumer text, + i_max_return int default 100, + i_dead_interval interval default null) +returns setof pgque.message as $$ +declare + v_batch_id bigint; + ev record; + cnt int := 0; +begin + if i_max_return < 1 then + raise exception 'pgque.receive_coop: max_return must be >= 1, got %', i_max_return; + end if; + + v_batch_id := pgque.next_batch(i_queue, i_consumer, i_subconsumer, i_dead_interval); + if v_batch_id is null then + return; + end if; + + for ev in + select + ev_id, + ev_type, + ev_data, + ev_retry, + ev_time, + ev_extra1, + ev_extra2, + ev_extra3, + ev_extra4 + from pgque.get_batch_events(v_batch_id) + loop + if cnt = i_max_return then + raise exception 'pgque.receive_coop: batch exceeds max_return of %', i_max_return + using + errcode = '54000', + hint = 'Retry with a larger resource-safe max_return to receive the complete batch. Do not acknowledge after this error.'; + end if; + return next row( + ev.ev_id, + v_batch_id, + ev.ev_type, + ev.ev_data, + ev.ev_retry, + ev.ev_time, + ev.ev_extra1, + ev.ev_extra2, + ev.ev_extra3, + ev.ev_extra4 + )::pgque.message; + cnt := cnt + 1; + end loop; + + -- Empty batch: release the member token so the subconsumer is not wedged + -- on a tick window with no visible events. finish_batch on a coop_member + -- clears sub_batch + sub_last_tick + sub_next_tick (it does not advance + -- the main cursor, which already moved when the batch was allocated). + if cnt = 0 then + perform pgque.finish_batch(v_batch_id); + end if; + + return; +end; +$$ language plpgsql security definer set search_path = pgque, pg_catalog; + +create or replace function pgque.version() +returns text as $$ +begin + return '0.2.1'; +end; +$$ language plpgsql security definer set search_path = pgque, pg_catalog; +$pgque_update_body$; +begin + select default_version into existing_version + from pgtle.available_extensions() + where name = 'pgque'; + + if existing_version = '0.2.1' then + raise notice 'pgque 0.2.1 already registered with pg_tle; skipping install_extension().'; + return; + end if; + + if existing_version = '0.2.0' and '0.2.1' = '0.2.1' then + select exists ( + select 1 + from pgtle.extension_update_paths('pgque') + where source = '0.2.0' + and target = '0.2.1' + and path is not null + ) into update_path_exists; + + if not update_path_exists then + perform pgtle.install_update_path( + 'pgque', '0.2.0', '0.2.1', update_sql + ); + end if; + perform pgtle.set_default_version('pgque', '0.2.1'); + raise notice 'pgque 0.2.0 -> 0.2.1 update path registered with pg_tle.'; + return; + end if; + + if existing_version is not null then + raise exception 'pgque is already registered with pg_tle at version % ' + 'but this script only supports a managed update from 0.2.0 to %.', + existing_version, '0.2.1'; + end if; + + perform pgtle.install_extension( + 'pgque', + '0.2.1', + 'PgQue — PgQ Universal Edition (zero-bloat Postgres queue)', + extension_sql ); end $wrapper$; \echo '' \echo 'PgQue 0.2.1 registered with pg_tle.' -\echo 'Run create extension pgque; to materialise the schema in this database.' +\echo 'For a new install, run: create extension pgque;' +\echo 'For an installed 0.2.0 extension, apply ALTER EXTENSION UPDATE as described in the upgrade documentation.' diff --git a/sql/pgque.sql b/sql/pgque.sql index 60a9885c..65035ab5 100644 --- a/sql/pgque.sql +++ b/sql/pgque.sql @@ -5365,7 +5365,9 @@ begin loop if cnt = i_max_return then raise exception 'pgque.receive: batch exceeds max_return of %', i_max_return - using hint = 'Retry with a larger resource-safe max_return to receive the complete batch. Do not acknowledge after this error.'; + using + errcode = '54000', + hint = 'Retry with a larger resource-safe max_return to receive the complete batch. Do not acknowledge after this error.'; end if; return next row( ev.ev_id, v_batch_id, ev.ev_type, ev.ev_data, @@ -6622,7 +6624,9 @@ begin loop if cnt = i_max_return then raise exception 'pgque.receive_coop: batch exceeds max_return of %', i_max_return - using hint = 'Retry with a larger resource-safe max_return to receive the complete batch. Do not acknowledge after this error.'; + using + errcode = '54000', + hint = 'Retry with a larger resource-safe max_return to receive the complete batch. Do not acknowledge after this error.'; end if; return next row( ev.ev_id, diff --git a/tests/test_receive_overflow.sql b/tests/test_receive_overflow.sql index b5672e57..8833fa5c 100644 --- a/tests/test_receive_overflow.sql +++ b/tests/test_receive_overflow.sql @@ -26,6 +26,7 @@ declare v_before record; v_after record; v_hint text; + v_state text; begin select s.* into v_before from pgque.subscription as s @@ -38,7 +39,11 @@ begin exception when others then v_raised := true; - get stacked diagnostics v_hint = pg_exception_hint; + get stacked diagnostics + v_hint = pg_exception_hint, + v_state = returned_sqlstate; + assert v_state = '54000', + 'receive overflow SQLSTATE must be 54000, got ' || v_state; assert sqlerrm like '%batch exceeds max_return of 2%', 'unexpected receive overflow error: ' || sqlerrm; assert v_hint like '%Do not acknowledge after this error%', @@ -134,16 +139,7 @@ begin perform pgque.send('recv_exact', 'ev', 'two'); end $$; --- A batch containing N-1 events also succeeds. -do $$ -begin - perform pgque.create_queue('recv_nminus1'); - perform pgque.register_consumer('recv_nminus1', 'c1'); - perform pgque.send('recv_nminus1', 'ev', 'one'); - perform pgque.send('recv_nminus1', 'ev', 'two'); -end $$; - -select pgque.force_next_tick('recv_nminus1'); +select pgque.force_next_tick('recv_exact'); select pgque.ticker(); do $$ @@ -151,18 +147,42 @@ declare v_count int := 0; v_batch_id bigint; v_msg pgque.message; + v_active_batch bigint; begin - for v_msg in select * from pgque.receive('recv_nminus1', 'c1', 3) + for v_msg in select * from pgque.receive('recv_exact', 'c1', 2) loop v_count := v_count + 1; v_batch_id := v_msg.batch_id; end loop; - assert v_count = 2, - format('N-1 events must succeed, got %s', v_count); + assert v_count = 2, format('exactly N events must succeed, got %s', v_count); perform pgque.ack(v_batch_id); + + v_count := 0; + for v_msg in select * from pgque.receive('recv_exact', 'c1', 2147483647) + loop + v_count := v_count + 1; + end loop; + assert v_count = 0, + format('INT_MAX receive after ack must be empty, got %s', v_count); + select s.sub_batch into v_active_batch + from pgque.subscription as s + join pgque.queue as q on q.queue_id = s.sub_queue + join pgque.consumer as c on c.co_id = s.sub_consumer + where q.queue_name = 'recv_exact' and c.co_name = 'c1'; + assert v_active_batch is null, + 'empty INT_MAX receive must not leave an active batch'; end $$; -select pgque.force_next_tick('recv_exact'); +-- A batch containing N-1 events also succeeds. +do $$ +begin + perform pgque.create_queue('recv_nminus1'); + perform pgque.register_consumer('recv_nminus1', 'c1'); + perform pgque.send('recv_nminus1', 'ev', 'one'); + perform pgque.send('recv_nminus1', 'ev', 'two'); +end $$; + +select pgque.force_next_tick('recv_nminus1'); select pgque.ticker(); do $$ @@ -171,15 +191,14 @@ declare v_batch_id bigint; v_msg pgque.message; begin - for v_msg in select * from pgque.receive('recv_exact', 'c1', 2) + for v_msg in select * from pgque.receive('recv_nminus1', 'c1', 3) loop v_count := v_count + 1; v_batch_id := v_msg.batch_id; end loop; - assert v_count = 2, format('exactly N events must succeed, got %s', v_count); + assert v_count = 2, + format('N-1 events must succeed, got %s', v_count); perform pgque.ack(v_batch_id); - - perform * from pgque.receive('recv_exact', 'c1', 2147483647); end $$; -- Preserve the SQL NULL behavior: it means no explicit ceiling. @@ -232,6 +251,7 @@ declare v_payloads text[] := array[]::text[]; v_msg pgque.message; v_hint text; + v_state text; begin begin perform * from pgque.receive_coop( @@ -240,7 +260,11 @@ begin exception when others then v_raised := true; - get stacked diagnostics v_hint = pg_exception_hint; + get stacked diagnostics + v_hint = pg_exception_hint, + v_state = returned_sqlstate; + assert v_state = '54000', + 'receive_coop overflow SQLSTATE must be 54000, got ' || v_state; assert sqlerrm like '%batch exceeds max_return of 2%', 'unexpected receive_coop overflow error: ' || sqlerrm; assert v_hint like '%Do not acknowledge after this error%', diff --git a/tests/test_tle_upgrade_v0_2.sql b/tests/test_tle_upgrade_v0_2.sql new file mode 100644 index 00000000..94bf6082 --- /dev/null +++ b/tests/test_tle_upgrade_v0_2.sql @@ -0,0 +1,170 @@ +-- test_tle_upgrade_v0_2.sql -- Non-destructive pg_tle 0.2.0 to 0.2.1 upgrade. +-- Copyright 2026 Nikolay Samokhvalov. Apache-2.0 license. +-- +-- The CI caller writes the tagged v0.2.0 pg_tle installer to +-- /tmp/pgque-tle-v0.2.0.sql before running this test. + +\set ON_ERROR_STOP on + +create extension pg_tle; +\i /tmp/pgque-tle-v0.2.0.sql +create extension pgque; + +select pgque.create_queue('tle_upgrade_overflow'); +select pgque.register_consumer('tle_upgrade_overflow', 'processor'); +select pgque.send('tle_upgrade_overflow', 'event', '{"n":1}'::jsonb); +select pgque.send('tle_upgrade_overflow', 'event', '{"n":2}'::jsonb); +select pgque.send('tle_upgrade_overflow', 'event', '{"n":3}'::jsonb); + +select pgque.create_queue('tle_upgrade_active'); +select pgque.register_consumer('tle_upgrade_active', 'processor'); +select pgque.send('tle_upgrade_active', 'event', '{"active":true}'::jsonb); + +select pgque.create_queue('tle_upgrade_coop'); +select pgque.register_subconsumer('tle_upgrade_coop', 'workers', 'worker_1'); +select pgque.send('tle_upgrade_coop', 'event', '{"n":1}'::jsonb); +select pgque.send('tle_upgrade_coop', 'event', '{"n":2}'::jsonb); +select pgque.send('tle_upgrade_coop', 'event', '{"n":3}'::jsonb); + +select pgque.force_next_tick('tle_upgrade_overflow'); +select pgque.force_next_tick('tle_upgrade_active'); +select pgque.force_next_tick('tle_upgrade_coop'); +select pgque.ticker(); + +do $$ +declare + v_msg pgque.message; +begin + select * into strict v_msg + from pgque.receive('tle_upgrade_active', 'processor', 10); + assert v_msg.batch_id is not null, 'expected active batch before upgrade'; +end $$; + +create temporary table tle_upgrade_state as +select q.queue_name, s.sub_batch, s.sub_last_tick, s.sub_next_tick +from pgque.subscription as s +join pgque.queue as q on q.queue_id = s.sub_queue +join pgque.consumer as c on c.co_id = s.sub_consumer +where (q.queue_name in ('tle_upgrade_overflow', 'tle_upgrade_active') + and c.co_name = 'processor') + or (q.queue_name = 'tle_upgrade_coop' and c.co_name = 'workers.worker_1'); + +\i sql/pgque-tle.sql +\i sql/pgque-tle.sql + +do $$ +begin + assert (select extversion = '0.2.0' + from pg_catalog.pg_extension where extname = 'pgque'), + 'registering the update must not mutate the installed extension'; + assert exists ( + select 1 + from pgtle.extension_update_paths('pgque') + where source = '0.2.0' and target = '0.2.1' and path is not null + ), 'pg_tle must expose the 0.2.0 to 0.2.1 update path'; + assert (select default_version = '0.2.1' + from pgtle.available_extensions() where name = 'pgque'), + 'pg_tle default version must be 0.2.1'; +end $$; + +alter extension pgque update to '0.2.1'; + +do $$ +declare + v_before record; + v_after record; + v_msg pgque.message; + v_count integer := 0; + v_batch_id bigint; + v_raised boolean := false; + v_sqlstate text; +begin + assert (select extversion = '0.2.1' + from pg_catalog.pg_extension where extname = 'pgque'), + 'pg_extension.extversion must be 0.2.1'; + assert pgque.version() = '0.2.1', + format('pgque.version() must be 0.2.1, got %s', pgque.version()); + + for v_before in select * from tle_upgrade_state + loop + select s.sub_batch, s.sub_last_tick, s.sub_next_tick + into strict v_after + from pgque.subscription as s + join pgque.queue as q on q.queue_id = s.sub_queue + join pgque.consumer as c on c.co_id = s.sub_consumer + where q.queue_name = v_before.queue_name + and c.co_name = case + when v_before.queue_name = 'tle_upgrade_coop' + then 'workers.worker_1' + else 'processor' + end; + + assert v_after.sub_batch is not distinct from v_before.sub_batch + and v_after.sub_last_tick is not distinct from v_before.sub_last_tick + and v_after.sub_next_tick is not distinct from v_before.sub_next_tick, + format('subscription state changed for queue %s', v_before.queue_name); + end loop; + + begin + perform * from pgque.receive('tle_upgrade_overflow', 'processor', 2); + exception + when sqlstate '54000' then + get stacked diagnostics v_sqlstate = returned_sqlstate; + v_raised := true; + end; + assert v_raised and v_sqlstate = '54000', + 'upgraded receive must reject overflow with SQLSTATE 54000'; + + for v_msg in + select * from pgque.receive('tle_upgrade_overflow', 'processor', 3) + loop + v_count := v_count + 1; + v_batch_id := v_msg.batch_id; + end loop; + assert v_count = 3, + format('expected all 3 overflow events after retry, got %s', v_count); + perform pgque.ack(v_batch_id); + + v_raised := false; + v_sqlstate := null; + begin + perform * from pgque.receive_coop( + 'tle_upgrade_coop', 'workers', 'worker_1', 2 + ); + exception + when sqlstate '54000' then + get stacked diagnostics v_sqlstate = returned_sqlstate; + v_raised := true; + end; + assert v_raised and v_sqlstate = '54000', + 'upgraded receive_coop must reject overflow with SQLSTATE 54000'; + + v_count := 0; + v_batch_id := null; + for v_msg in + select * from pgque.receive_coop( + 'tle_upgrade_coop', 'workers', 'worker_1', 3 + ) + loop + v_count := v_count + 1; + v_batch_id := v_msg.batch_id; + end loop; + assert v_count = 3, + format('expected all 3 cooperative events after retry, got %s', v_count); + perform pgque.ack(v_batch_id); +end $$; + +-- The update edge must also make the new default usable for a fresh create. +drop extension pgque cascade; +create extension pgque; + +do $$ +begin + assert (select extversion = '0.2.1' + from pg_catalog.pg_extension where extname = 'pgque'), + 'fresh create after update registration must install 0.2.1'; + assert pgque.version() = '0.2.1', + 'fresh create after update registration must expose version 0.2.1'; +end $$; + +\echo '=== test_tle_upgrade_v0_2: ALL PASSED ===' From dc0b698afce03bd4d9043d9a75b6d30584e78253 Mon Sep 17 00:00:00 2001 From: samo-agent <280144521+samo-agent@users.noreply.github.com> Date: Thu, 1 Oct 2026 14:57:12 +0000 Subject: [PATCH 4/5] docs: clarify complete-batch receive contract Resolve review findings on SDK documentation and default ceilings. Refs #365. --- clients/go/README.md | 4 ++-- clients/go/concurrency_test.go | 6 ++---- clients/go/options.go | 21 ++++++++++----------- clients/go/pgque.go | 17 ++++++++++------- clients/python/README.md | 13 +++++-------- clients/python/pgque/client.py | 22 +++++++++++----------- clients/typescript/README.md | 4 ++-- clients/typescript/src/client.ts | 18 +++++++++--------- clients/typescript/src/types.ts | 10 ++++------ docs/reference.md | 6 +++++- docs/upgrading.md | 14 ++++++++++++++ tests/run_all.sql | 1 + 12 files changed, 75 insertions(+), 61 deletions(-) diff --git a/clients/go/README.md b/clients/go/README.md index dd6ea0e5..a705640f 100644 --- a/clients/go/README.md +++ b/clients/go/README.md @@ -93,7 +93,7 @@ func main() { | Option | Default | Notes | | --------------------------------------- | -------------- | --------------------------------------------------------------------- | | `WithPollInterval(d time.Duration)` | `30s` | Idle backoff between polls when the queue is empty. | -| `WithMaxMessages(n int)` | `math.MaxInt32` | Complete-batch safety ceiling. Servers with the overflow guard raise if the batch is larger; retry with a sufficient ceiling within your resource budget. Older servers can truncate, so upgrade before relying on the guard. Process the whole batch before `Ack`. | +| `WithMaxMessages(n int)` | `math.MaxInt32` | Complete-batch safety ceiling. A larger batch fails with SQLSTATE `54000` and no partial result. Roll back and retry with a resource-safe larger ceiling; never `Ack` the failed receive. Ticker thresholds do not cap batch size. | | `WithUnknownHandlerPolicy(p)` | `NackUnknown` | `AckUnknown` logs and skips messages with no registered handler. | | `WithRetryAfter(d time.Duration)` | `60s` | Retry delay for Consumer-issued `Nack` calls on handler failure or unknown type. | @@ -239,7 +239,7 @@ Options: | --- | --- | | `WithSubconsumer(name)` | High-level `Consumer`: enables coop mode. | | `WithDeadInterval(d)` | High-level `Consumer`: passes a takeover window to `ReceiveCoop`. | -| `WithCoopMaxMessages(n)` | `ReceiveCoop` per-call row cap (default 100). | +| `WithCoopMaxMessages(n)` | `ReceiveCoop` complete-batch safety ceiling (default 100). Overflow returns no partial result and SQLSTATE `54000`; roll back and retry with a resource-safe larger ceiling, and never `Ack` the failed receive. | | `WithCoopDeadInterval(d)` | `ReceiveCoop` takeover window for one call. | | `WithBatchHandlingRetry()` | `UnsubscribeSubconsumer`: route active batch through retry/DLQ instead of erroring. | diff --git a/clients/go/concurrency_test.go b/clients/go/concurrency_test.go index d346f1f2..e7ca4770 100644 --- a/clients/go/concurrency_test.go +++ b/clients/go/concurrency_test.go @@ -45,10 +45,8 @@ func TestRace_ConcurrentSend(t *testing.T) { tick(t, client, queue) expected := goroutines * perGoroutine - // pgque.receive truncates yielded rows at i_max_return, but Ack - // finishes the entire batch — events past the cap are not yielded - // again without a fresh tick. Size the cap above the total so the - // whole tick window flows out in one call. + // The complete batch must fit within maxMessages. This test knows its + // maximum producer count, so use a ceiling above that count. total := 0 for { msgs, err := client.Receive(ctx, queue, consumer, 2*expected) diff --git a/clients/go/options.go b/clients/go/options.go index 43f66d79..cfb2cd93 100644 --- a/clients/go/options.go +++ b/clients/go/options.go @@ -15,15 +15,14 @@ func WithPollInterval(d time.Duration) ConsumerOption { return func(c *Consumer) { c.pollInterval = d } } -// WithMaxMessages sets the per-Receive limit. By default the Consumer +// WithMaxMessages sets the complete-batch safety ceiling. By default the Consumer // requests PostgreSQL's int maximum so it drains the whole PgQ batch before // acknowledging it. Panics if n <= 0. // -// WARNING: pgque.ack(batch_id) finishes the entire underlying PgQ batch, -// including rows the consumer never received because of this limit. If you -// set maxMessages below the real batch size, unreturned rows are skipped -// after ack. Only lower this value when it is at least as large as the -// queue's possible batch size for your workload. +// A larger batch makes Receive fail with SQLSTATE 54000 and returns no partial +// result. Roll back, retry with a resource-safe larger ceiling, process every +// message, then Ack. Never Ack a failed Receive. Ticker thresholds do not cap +// batch size. func WithMaxMessages(n int) ConsumerOption { if n <= 0 { panic("pgque: WithMaxMessages requires n > 0") @@ -123,14 +122,14 @@ func newReceiveCoopConfig() *receiveCoopConfig { // ReceiveCoopOption tunes a single Client.ReceiveCoop call. type ReceiveCoopOption func(*receiveCoopConfig) -// WithCoopMaxMessages sets the per-call message limit (maps to the +// WithCoopMaxMessages sets the per-call complete-batch safety ceiling (maps to the // i_max_return argument of pgque.receive_coop). Default is 100. Panics // if n <= 0. // -// As with Receive, ack(batch_id) finishes the entire underlying batch -// regardless of how many rows are returned; a low limit can therefore -// drop rows. Match this to ticker_max_count (or larger) when you care -// about per-message dispatch. +// A larger batch makes ReceiveCoop fail with SQLSTATE 54000 and returns no +// partial result. Roll back, retry with a resource-safe larger ceiling, +// process every message, then Ack. Never Ack a failed ReceiveCoop. Ticker +// thresholds do not cap batch size. func WithCoopMaxMessages(n int) ReceiveCoopOption { if n <= 0 { panic("pgque: WithCoopMaxMessages requires n > 0") diff --git a/clients/go/pgque.go b/clients/go/pgque.go index b1aebfe1..7faeba41 100644 --- a/clients/go/pgque.go +++ b/clients/go/pgque.go @@ -119,11 +119,12 @@ func (c *Client) Unsubscribe(ctx context.Context, queue, consumer string) (int64 return n, nil } -// Receive fetches up to maxMessages from the next batch for the named -// consumer. Returns an empty slice when no batch is available; in that -// case the caller should sleep before polling again. Each returned -// Message carries a BatchID that must be passed to Ack once all -// messages in the batch have been processed. +// Receive fetches the complete next batch for the named consumer when it fits +// within maxMessages. A larger batch returns no partial result and fails with +// SQLSTATE 54000; roll back and retry with a resource-safe larger ceiling, and +// never Ack the failed call. Ticker thresholds do not cap batch size. Returns +// an empty slice when no batch is available. Ack only after processing every +// returned Message. func (c *Client) Receive(ctx context.Context, queue, consumer string, maxMessages int) ([]Message, error) { rows, err := c.pool.Query(ctx, "SELECT * FROM pgque.receive($1, $2, $3)", queue, consumer, maxMessages) @@ -299,8 +300,10 @@ func (c *Client) UnsubscribeSubconsumer(ctx context.Context, queue, consumer, su // // receive_coop auto-registers the cooperative main row and the // subconsumer on first call, so an explicit SubscribeSubconsumer is -// not required. WithCoopMaxMessages tunes the per-call row cap -// (default 100); WithCoopDeadInterval enables stale-worker takeover. +// not required. WithCoopMaxMessages sets the complete-batch safety ceiling +// (default 100): overflow returns no partial result and SQLSTATE 54000, so roll +// back and retry with a resource-safe larger ceiling without acknowledging the +// failed call. WithCoopDeadInterval enables stale-worker takeover. // // Experimental in PgQue 0.2. func (c *Client) ReceiveCoop(ctx context.Context, queue, consumer, subconsumer string, opts ...ReceiveCoopOption) ([]Message, error) { diff --git a/clients/python/README.md b/clients/python/README.md index 6f24b0ac..fde02083 100644 --- a/clients/python/README.md +++ b/clients/python/README.md @@ -71,14 +71,11 @@ consumer.start() # blocks until SIGTERM / SIGINT ### Consumer options `Consumer(..., max_messages=...)` controls the per-`receive` limit. -The default is the Postgres `int` maximum, so the consumer requests -an entire PgQ batch before acknowledging it. On servers with the overflow -guard, a batch larger than this ceiling raises an error instead of returning -partial results. Roll back failed transactions and retry with a sufficient -ceiling within your resource budget; monitor errors and consumer lag. -The ticker threshold does not cap batch size. Older servers can truncate -results, so upgrade the server before relying on this protection. -`ack()` always finishes the entire batch: process every message first. +The default is the Postgres `int` maximum. A batch larger than the configured +complete-batch safety ceiling returns no partial result and raises SQLSTATE +`54000`. Roll back and retry with a resource-safe larger ceiling; never +acknowledge the failed receive. Ticker thresholds do not cap batch size. +Process every message before `ack()`. ### Handling unknown event types diff --git a/clients/python/pgque/client.py b/clients/python/pgque/client.py index 08379dcd..e9e00c41 100644 --- a/clients/python/pgque/client.py +++ b/clients/python/pgque/client.py @@ -226,16 +226,16 @@ def receive( Maps to ``pgque.receive(queue, consumer, max_messages)``, which opens a batch via ``next_batch`` internally. The caller must ``ack()`` the batch (with the ``batch_id`` from any returned - message) to advance the consumer past it. ``ack()`` finishes the - whole underlying PgQ batch, including rows beyond ``max_messages``; - direct callers should pass a value large enough for the queue's - possible batch size before acknowledging. + message) to advance the consumer past it. ``max_messages`` is a + complete-batch safety ceiling. A larger batch returns no partial + result and raises SQLSTATE 54000. Roll back and retry with a + resource-safe larger ceiling; never acknowledge the failed receive. + Ticker thresholds do not cap batch size. Args: queue: Queue name. consumer: Consumer name (must be registered on the queue). - max_messages: Maximum number of messages to return from the - current batch. + max_messages: Complete-batch safety ceiling. Returns: List of ``Message`` objects, possibly empty if no batch is @@ -408,11 +408,11 @@ def receive_coop( queue: Queue name. consumer: Logical consumer (the ``coop_main`` row). subconsumer: Per-worker member name. - max_messages: Maximum rows to return from the current batch. - ``ack(batch_id)`` advances the cooperative cursor past - the entire underlying batch, so set this >= the queue's - worst-case batch size or consume the full batch before - acking. + max_messages: Complete-batch safety ceiling. A larger batch + returns no partial result and raises SQLSTATE 54000. Roll + back and retry with a resource-safe larger ceiling; never + acknowledge the failed receive. Ticker thresholds do not + cap batch size. dead_interval: Optional PostgreSQL interval syntax (e.g. ``"5 minutes"``). When set, allows takeover of a stale sibling's batch under a fresh ``batch_id``; the old diff --git a/clients/typescript/README.md b/clients/typescript/README.md index b8b950fb..644e45b3 100644 --- a/clients/typescript/README.md +++ b/clients/typescript/README.md @@ -74,7 +74,7 @@ try { | `connect(dsn, poolOptions?)` | Connect via `pg.Pool`. Eagerly probes the connection. | | `client.send(queue, event)` | Publish; returns event id (`bigint`). | | `client.sendBatch(queue, type, payloads)` | Publish a same-type batch atomically; returns event ids (`bigint[]`). | -| `client.receive(queue, consumer, max?)` | Fetch up to `max` (default 100) messages from the next batch. If you later call `ack(batchId)`, PgQue finishes the whole underlying batch, including rows beyond `max`; size `max` for your queue or use the high-level consumer default. | +| `client.receive(queue, consumer, max?)` | Fetch the complete next batch when it fits within the safety ceiling (default 100). Overflow returns no partial result and SQLSTATE `54000`; roll back and retry with a resource-safe larger ceiling, and never acknowledge the failed receive. Ticker thresholds do not cap batch size. | | `client.ack(batchId)` | Finish the batch. Returns `1` on success, `0` if the batch was already finished or not found (stale/double ack — log at warn level, not an error). | | `client.nack(batchId, msg, opts?)` | Single-message retry/DLQ. | | `client.subscribe(queue, consumer)` | Wraps `pgque.register_consumer`. | @@ -121,7 +121,7 @@ identically to the non-cooperative form. |---|---| | `client.subscribeSubconsumer(queue, consumer, subconsumer)` | Register a subconsumer. Returns `1` first call, `0` if already registered. | | `client.unsubscribeSubconsumer(queue, consumer, subconsumer, { batchHandling? })` | Remove a subconsumer. Default raises if a batch is active; `batchHandling: 1` routes the active batch through retry/DLQ first. | -| `client.receiveCoop(queue, consumer, subconsumer, { maxMessages?, deadInterval? })` | Cooperative receive. Auto-registers main + subconsumer rows on first call. `deadInterval` enables stale-batch takeover; the new owner gets a fresh `batchId`. | +| `client.receiveCoop(queue, consumer, subconsumer, { maxMessages?, deadInterval? })` | Cooperative receive with a complete-batch safety ceiling. Overflow returns no partial result and SQLSTATE `54000`; roll back and retry with a resource-safe larger ceiling, and never acknowledge the failed receive. Auto-registers main + subconsumer rows; `deadInterval` enables stale-batch takeover under a fresh `batchId`. | | `client.touchSubconsumer(queue, consumer, subconsumer)` | Refresh the subconsumer heartbeat so a long handler is not stolen. The high-level consumer does not call this automatically. | Throughput note: cooperative allocation serializes on a `FOR UPDATE` of the diff --git a/clients/typescript/src/client.ts b/clients/typescript/src/client.ts index 663ff965..e1864cd9 100644 --- a/clients/typescript/src/client.ts +++ b/clients/typescript/src/client.ts @@ -149,12 +149,11 @@ export class Client { } /** - * Fetch up to `maxMessages` from the next batch for `consumer` on `queue`. - * Returns an empty array when no batch is currently available. - * - * WARNING: `ack(batchId)` finishes the whole underlying PgQ batch, including - * rows beyond `maxMessages`. Direct receive callers should pass a value large - * enough for the queue's possible batch size before acknowledging the batch. + * Fetch the complete next batch for `consumer` when it fits within + * `maxMessages`. A larger batch returns no partial result and fails with + * SQLSTATE 54000. Roll back and retry with a resource-safe larger ceiling; + * never acknowledge the failed receive. Ticker thresholds do not cap batch + * size. Returns an empty array when no batch is currently available. */ async receive(queue: string, consumer: string, maxMessages = 100): Promise { if (!queue) { @@ -373,9 +372,10 @@ export class Client { * Wraps `pgque.receive_coop`. The cooperative main and subconsumer rows * are auto-registered on first call. * - * `options.maxMessages` defaults to `100` (the SQL default). `ack(batchId)` - * still finishes the entire underlying batch, so size `maxMessages` - * appropriately or use the high-level `Consumer` default. + * `options.maxMessages` is a complete-batch safety ceiling and defaults to + * `100`. A larger batch returns no partial result and fails with SQLSTATE + * 54000. Roll back and retry with a resource-safe larger ceiling; never + * acknowledge the failed receive. Ticker thresholds do not cap batch size. * * `options.deadInterval` is a PostgreSQL `interval` text (e.g. * `"5 minutes"`); when set, `receive_coop` may steal a stale sibling's diff --git a/clients/typescript/src/types.ts b/clients/typescript/src/types.ts index 1aff5fec..7313d1a7 100644 --- a/clients/typescript/src/types.ts +++ b/clients/typescript/src/types.ts @@ -56,15 +56,13 @@ export interface ConsumerOptions { */ pollInterval?: number; /** - * Maximum messages returned per `receive()` call. By default the + * Complete-batch safety ceiling for each `receive()` call. By default the * high-level consumer requests the PostgreSQL `int` maximum so it drains * the whole PgQ batch before acknowledging it. * - * WARNING: `pgque.ack(batch_id)` finishes the entire underlying batch, - * including rows the client never returned. If you set `maxMessages` - * below the real batch size, unreturned rows are skipped after ack. - * Only lower this value when it is at least as large as the queue's - * possible batch size for your workload. + * A larger batch returns no partial result and fails with SQLSTATE 54000. + * Roll back and retry with a resource-safe larger ceiling; never acknowledge + * the failed receive. Ticker thresholds do not cap batch size. */ maxMessages?: number; /** diff --git a/docs/reference.md b/docs/reference.md index e00e1e8b..07a03c6c 100644 --- a/docs/reference.md +++ b/docs/reference.md @@ -122,8 +122,10 @@ All consume-side functions (`receive`, `ack`, `nack`, `subscribe`, `unsubscribe` Pulls the complete next batch for `consumer` on `queue`, or raises an error if it exceeds `max_return`. A non-null `max_return` must be >= 1; passing 0 or a negative value raises an error. SQL `null` disables the explicit ceiling for compatibility; prefer an explicit resource-safe ceiling. Returns an empty set if no batch is available. Each row is a `pgque.message` composite (see [§Message type](#message-type)). Grant: `pgque_reader`. Source: `sql/pgque-api/receive.sql`. +**Choose a ceiling explicitly.** The SQL default is 100, below the default ticker event-count threshold of 500. Ordinary bursts can therefore exceed the default receive ceiling. Pass an explicit ceiling sized for your expected batches and resource budget; alert on SQLSTATE `54000` and consumer lag. The example value below is illustrative, not a batch-size guarantee. + ```sql -select * from pgque.receive('orders', 'processor', 100); +select * from pgque.receive('orders', 'processor', 1000); ``` Overflow reports SQLSTATE `54000` (`program_limit_exceeded`). Like other SQL errors, it aborts the enclosing transaction unless recovered through a savepoint. @@ -176,6 +178,8 @@ Receives messages for one subconsumer. A non-null `max_return` must be >= 1; SQL **Empty tick windows are auto-finished.** When the current batch's tick window holds no events, `receive_coop()` calls `finish_batch` internally and returns the empty set. Callers polling a quiet queue do not see (and do not need to ack) a `batch_id`; `receive()` also auto-finishes empty batches. +The default cooperative ceiling is also 100; pass an explicit burst-sized ceiling rather than assuming the ticker threshold caps the batch. + **Complete batch or error.** As with `receive()`, `max_return` is a safety ceiling, not pagination. An oversized batch raises an error and rolls back allocation or takeover performed by that call. Retry with a sufficient ceiling within your resource budget, process the complete batch, then acknowledge. The ticker threshold does not cap batch size. **Throughput note.** Cooperative allocation serializes on a `FOR UPDATE` of the `coop_main` subscription row, so many workers polling tiny batches contend on a single hot row. If you scale workers, also tune `ticker_max_count` and tick cadence so each batch is large enough to amortize the lock. diff --git a/docs/upgrading.md b/docs/upgrading.md index 024a4bf4..9f5708e3 100644 --- a/docs/upgrading.md +++ b/docs/upgrading.md @@ -20,6 +20,15 @@ Re-run the installer using the command above. This maintenance release makes acknowledged as a complete batch. The upgrade replaces functions only; it does not change tables or queue state. +**Before upgrading, size receive ceilings explicitly.** The SQL default +`max_return` is 100, below the default ticker event-count threshold of 500; +ordinary batches can exceed 100. Direct Python and TypeScript receive calls +and Go `ReceiveCoop` also default to 100. Callers relying on those defaults +must pass a sufficient resource-safe ceiling before upgrading. Repeating an +oversized receive with the same ceiling cannot make progress. Alert on +SQLSTATE `54000` and consumer lag. The high-level consumer loops default to +the Postgres integer maximum and are not constrained by the SQL default. + Applications do not need client-library updates for this server-side fix. An undersized receive ceiling now causes an error: roll back the failed transaction and retry with a sufficient ceiling within your resource budget. @@ -36,6 +45,11 @@ select extversion from pg_extension where extname = 'pgque'; select pgque.version(); ``` +The managed pg_tle update path supports exactly a registered `0.2.0` origin. +Other origins are rejected without changing the installation; do not drop a +populated extension as a workaround. Plan and test a separate migration for +those origins. + Both version queries must return `0.2.1`. The update replaces functions only; do not drop or unregister the existing extension before upgrading. diff --git a/tests/run_all.sql b/tests/run_all.sql index 8c5235b0..57c42bf0 100644 --- a/tests/run_all.sql +++ b/tests/run_all.sql @@ -84,6 +84,7 @@ \echo 'Running: test_api_receive' \i tests/test_api_receive.sql +\echo 'Running: test_receive_overflow' \i tests/test_receive_overflow.sql \echo 'Running: test_api_dlq' From 034f3f6b353caeedd00d3f8cb1041cde58520692 Mon Sep 17 00:00:00 2001 From: samo-agent <280144521+samo-agent@users.noreply.github.com> Date: Thu, 1 Oct 2026 15:07:56 +0000 Subject: [PATCH 5/5] docs: remove remaining batch-cap wording Clarify ticker thresholds and consumer comments. Refs #365. --- clients/go/concurrency_test.go | 2 +- clients/python/README.md | 2 +- clients/typescript/src/consumer.ts | 4 ++-- docs/pgq-concepts.md | 2 +- 4 files changed, 5 insertions(+), 5 deletions(-) diff --git a/clients/go/concurrency_test.go b/clients/go/concurrency_test.go index e7ca4770..d3231940 100644 --- a/clients/go/concurrency_test.go +++ b/clients/go/concurrency_test.go @@ -46,7 +46,7 @@ func TestRace_ConcurrentSend(t *testing.T) { expected := goroutines * perGoroutine // The complete batch must fit within maxMessages. This test knows its - // maximum producer count, so use a ceiling above that count. + // maximum message count, so use a ceiling above that count. total := 0 for { msgs, err := client.Receive(ctx, queue, consumer, 2*expected) diff --git a/clients/python/README.md b/clients/python/README.md index fde02083..84db45f7 100644 --- a/clients/python/README.md +++ b/clients/python/README.md @@ -70,7 +70,7 @@ consumer.start() # blocks until SIGTERM / SIGINT ### Consumer options -`Consumer(..., max_messages=...)` controls the per-`receive` limit. +`Consumer(..., max_messages=...)` sets the complete-batch safety ceiling for each `receive`. The default is the Postgres `int` maximum. A batch larger than the configured complete-batch safety ceiling returns no partial result and raises SQLSTATE `54000`. Roll back and retry with a resource-safe larger ceiling; never diff --git a/clients/typescript/src/consumer.ts b/clients/typescript/src/consumer.ts index fee27aaa..9d1177e0 100644 --- a/clients/typescript/src/consumer.ts +++ b/clients/typescript/src/consumer.ts @@ -6,8 +6,8 @@ import type { ConsumerOptions, HandlerFunc, Message } from './types.js'; /** * Default `maxMessages` for the high-level Consumer. PostgreSQL `int4` max - * (`2^31 - 1`); request the whole PgQ batch by default so a subsequent - * `pgque.ack(batch_id)` does not strand events the client never saw. + * (`2^31 - 1`); request the whole PgQ batch by default so ordinary bursts + * do not exceed a smaller receive ceiling (SQLSTATE 54000). */ export const DEFAULT_MAX_MESSAGES = 2_147_483_647; diff --git a/docs/pgq-concepts.md b/docs/pgq-concepts.md index 051b7423..2ffb2b8b 100644 --- a/docs/pgq-concepts.md +++ b/docs/pgq-concepts.md @@ -64,7 +64,7 @@ short name below; the function auto-prefixes `queue_` internally. - `ticker_max_lag` — max wall time between ticks. - `ticker_idle_period` — tick interval when idle. -- `ticker_max_count` — force tick at N events (batch-size cap). +- `ticker_max_count` — event-count threshold for creating a tick; not a batch-size cap. - `rotation_period` — table rotation period (disk vs. history). - `max_retries` — retry ceiling before a message goes to `pgque.dead_letter`.