Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
15 changes: 13 additions & 2 deletions .github/workflows/ci.yml
Original file line number Diff line number Diff line change
@@ -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:
Expand Down Expand Up @@ -172,6 +173,7 @@ jobs:
- uses: actions/checkout@v4
with:
submodules: recursive
fetch-depth: 0

- name: Build pgque
run: bash build/transform.sh
Expand Down Expand Up @@ -216,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
Expand Down
96 changes: 79 additions & 17 deletions build/transform.sh
Original file line number Diff line number Diff line change
Expand Up @@ -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).
--
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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.
Expand Down Expand Up @@ -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()
Expand All @@ -1103,31 +1148,43 @@ 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;

perform pgtle.install_extension(
'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}")
Expand All @@ -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
Expand Down
4 changes: 2 additions & 2 deletions clients/go/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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. 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. |

Expand Down Expand Up @@ -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. |

Expand Down
6 changes: 2 additions & 4 deletions clients/go/concurrency_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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 message count, so use a ceiling above that count.
total := 0
for {
msgs, err := client.Receive(ctx, queue, consumer, 2*expected)
Expand Down
44 changes: 37 additions & 7 deletions clients/go/integration_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,7 @@ package pgque_test

import (
"context"
"errors"
"strings"
"testing"
"time"
Expand Down Expand Up @@ -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)
Expand All @@ -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 != "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)
}

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))
}
}

Expand Down
21 changes: 10 additions & 11 deletions clients/go/options.go
Original file line number Diff line number Diff line change
Expand Up @@ -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")
Expand Down Expand Up @@ -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")
Expand Down
17 changes: 10 additions & 7 deletions clients/go/pgque.go
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down Expand Up @@ -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) {
Expand Down
13 changes: 6 additions & 7 deletions clients/python/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -70,13 +70,12 @@ 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.
`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
acknowledge the failed receive. Ticker thresholds do not cap batch size.
Process every message before `ack()`.

### Handling unknown event types

Expand Down
Loading
Loading