Skip to content

feat: durable paged batch consumption for 0.3 - #368

Draft
NikolayS wants to merge 10 commits into
mainfrom
feat-batch-pages
Draft

NikolayS wants to merge 10 commits into
mainfrom
feat-batch-pages

Conversation

@NikolayS

@NikolayS NikolayS commented Oct 1, 2026 •

Copy link
Copy Markdown
Owner

Summary

Implements durable bounded consumption for the 0.3 development line. Closes #364 after acceptance. Frozen sql/ / 0.2.1 artifacts remain untouched; no release or client-package publication in this PR.

  • Add ordinary, cooperative and partitioned receive_page, ack_page, and renew_page APIs.
  • Receiving never advances progress. Atomic page acknowledgments checkpoint exact event keys; only terminal acknowledgment finishes the underlying batch.
  • Preserve pending boundaries/checkpoints through takeover, fence stale tokens and partition epochs, and retain the most recent committed ack receipt per subscription.
  • Atomically route explicit failures to retry/DLQ; reject changed receipt requests and ambiguous duplicate IDs.
  • Guard legacy whole-batch finish/retry/reset/unregister bypasses, including cooperative-main cursor reset while a member has active pages.
  • Add independent page types and one-shot helpers in Python, Go, TypeScript and Ruby without changing existing Message types or consumer interfaces.
  • Add SQL reference, executable documentation example, and CI paging/concurrency coverage.

Contract and limitations

See blueprints/PAGED_BATCHES.md and docs/paged-batches.md.

  • External effects remain at-least-once and require idempotency.
  • One outstanding page per subscription; latest-ack receipt only (not an unlimited receipt store).
  • Page buffers/returned rows are bounded; underlying batch scan/sort I/O is not promised O(N).
  • Engine-managed batch membership must remain immutable. Privileged history mutation/rewind behind an acknowledged key is unsupported.
  • SDK one-shot helpers do not run background renewal; caller must choose sufficient lease or explicitly renew.
  • Python synchronous helpers reject async/generator handlers and lazy results instead of acknowledging unexecuted work.

Verification so far

  • PG14 and PG18 full regression passed, including fencing, three-slot filtering, legacy guards and administrative force deletion.
  • Fresh pg_tle full regression passed.
  • PG18 acceptance and multi-backend crash/concurrency tests passed, including force-drop fail-fast lock handling.
  • All four SDKs passed live fresh-PG18 tests: handler failure/redelivery, failure routing, ack replay, and maximum int8 identity decoding.
  • Python, Ruby and TypeScript page unit tests passed; lazy handlers cannot be acknowledged.
  • Original review at d1c56c1 reported seven blockers. Review fixes are pushed; updated CI and exact-head re-review are pending.
  • Frozen stable installers remain untouched.

Reproduction:

bash build/transform.sh
psql "$DATABASE_URL" -v ON_ERROR_STOP=1 -f devel/sql/pgque.sql
psql "$DATABASE_URL" -v ON_ERROR_STOP=1 -f tests/run_all.sql
psql "$DATABASE_URL" -v ON_ERROR_STOP=1 -f tests/acceptance/run_acceptance.sql
PGQUE_TEST_IMAGE=postgres:18 bash tests/test_paged_concurrency.sh

Remaining gates

Draft until full CI is green, exact-head SamoRev/REV is posted with all blocking findings resolved, and post-review executable SQL/SDK evidence is posted. No feature-completion claim yet.

Specify ownership, retries and legacy guards for 0.3. Refs #364.
Separate committed receipts from pending ownership checks. Refs #364.
Implement fenced page checkpoints and atomic failure routing.\nAdd additive SDK helpers and paging acceptance coverage.\n\nRefs #364; targets the 0.3 development line.
Keep plain SQL reapply separate from extension-owned installs.
Exercise real backend termination and producer snapshot boundaries.
Add scalar SDK regressions from live integration testing.

Refs #364.
@NikolayS

NikolayS commented Oct 1, 2026

Copy link
Copy Markdown
Owner Author

REV Code Review Report

  • PR: feat: durable paged batch consumption for 0.3 #368 - feat: durable paged batch consumption for 0.3
  • Author: @NikolayS
  • Reviewed head: d1c56c1ee9af9c81fc38f1fee627f8995be0128a (local checkout verified equal to PR head)
  • Scope: all source/test/SDK/doc changes from merge-base b8933a8 (4 commits, 39 files, +6580/-67), not only the last commit
  • AI-Assisted: Yes
  • Mode: --blocking
Pipeline Coverage
✅ 18/18 checks success at this head (PG14–19beta1 matrix, pg_tle, pg_cron, pg_timetable, upgrade v0.1.0→HEAD, frozen-sql smoke, Python/Go/TS/Ruby jobs) Not reported

Result: FAILED: 7 blocking findings (1 correctness, 1 SDK, 1 docs, 4 test-coverage). The core protocol held up under source reading and live two-session reproduction: crash/ack atomicity, lease/epoch fencing incl. same-worker ABA, cooperative transfer with independent receipts, lock order, membership/retry identity, and role permissions. No CRITICAL/HIGH correctness defect was found in the SQL protocol.

Specialized reviewers (all 6 completed):

  1. SQL concurrency/bug hunter (opus): reproduced on throwaway PG18 incl. two-session races
  2. Security/permissions (opus): live ACL diff base vs HEAD, fresh install + upgrade re-run, role-based repro
  3. SDK bug hunter (opus): live PG18 repro for Python/TypeScript/Ruby; Go by reading only
  4. Test analyzer (sonnet): contract→test coverage matrix; ran TS (5/5) and Ruby (4/4) page unit tests
  5. Guidelines (sonnet): PgQue CLAUDE.md + postgres-ai rules
  6. Docs (sonnet): checked each docs claim against the source

Generated-SQL assembly: build/transform.sh was rerun on a clean git archive of HEAD in a temp dir. devel/sql/pgque.sql, pgque-tle.sql, and both uninstall scripts rebuild byte-identical to the committed files. Duplicated generated text was not counted as extra logic. Final definitions win as intended: the guarded finish_batch, event_retry x2, batch_retry, register_consumer_at and unregister_consumer are the last definitions. Frozen sql/ is untouched.


BLOCKING ISSUES (7)

MEDIUM devel/sql/pgque-api/cooperative_consumers.sql:78 - Forced drop_queue(q, true) is blocked by any active paged batch, which contradicts the contract

The new _assert_unpaged in unregister_consumer also fires inside pgque.drop_queue(text, bool) (force branch → perform pgque.unregister_consumer(...), generated pgque.sql:1513). Reproduced on PG18: receive_page('dq','c','w1',1) → page, then drop_queue('dq', true) → ERROR: batch 1 is managed by paged delivery (55000). blueprints/PAGED_BATCHES.md:159 and docs/paged-batches.md:152 both promise that "administrative queue destruction remains destructive". page_state is revoked from pgque_admin, so the only way out is to receive and ack every page with a fresh worker, which falsely records unprocessed events as processed. No test force-drops a queue while a page is active; every test drop happens after draining.
Fix: In the drop_queue force path, delete subscriptions directly under the queue lock (FK cascade removes page_state), or call a private unguarded unregister core. Keep the guard on public unregister_consumer/unsubscribe. Add tests that force-drop with a page outstanding and between pages.

MEDIUM clients/typescript/src/client.ts:278 (also clients/ruby/lib/pgque/client.rb:149) - The one-shot processPage/process_page helpers ack pages whose handlers never ran (lazy handlers)

TS: for (const message of page.messages) await handler(message);. Ruby: page.messages.each { |message| handler.call(message) }. Neither checks the handler's return value. Reproduced on PG18: a TS function* handler gave {status:'page', processedCount:2} with 0 side effects, and the next receive returned page 2. A Ruby handler returning Enumerator.new {...} behaved the same way. blueprints/PAGED_BATCHES.md:230 requires: "Unknown handler types must be explicit errors, not silently skipped acknowledged messages." Python already guards this case, so the SDKs are inconsistent.
Fix: TS: reject GeneratorFunction/AsyncGeneratorFunction before receiving, and reject any non-undefined result (including an awaited value) before ackPage. Ruby: raise TypeError unless the result is nil (reject Enumerator/Enumerator::Lazy) before ack_page. Add a unit test for each.

MEDIUM clients/python/README.md:179 - The Python paging example holds receive locks through processing, which contradicts docs/paged-batches.md:54

The example runs receive_page → handle() → ack_page → conn.commit() in one transaction on the default non-autocommit connection. _receive_page takes FOR UPDATE locks on the subscription and page_state rows (paged_batches.sql:76-83). Those locks are held for the whole handler run, so other workers block instead of getting busy, and the issued page is never durably recorded during processing. docs/paged-batches.md:54 says the opposite: "For external work, commit receive first, process…, then ack in another transaction." (This validates the CLI report's LOW docs finding and upgrades it, because the two documents contradict each other.)
Fix: Show receive_page → commit() → process → ack_page (+ DB side effects) → commit(). State that process_page() is only appropriate with autocommit=True or DB-only handlers. Also note that msg_id in failures must be a decimal string.

MEDIUM tests/test_paged_batches.sql:116, tests/test_paged_modes.sql:197 - Fencing negatives are untested: wrong worker, forged token, and the pending-token partition epoch fence

Every ack_page/renew_page call in the suite uses the correct worker. The only PQP01 cases are tokens already superseded (batches:334, modes:124, modes:306). The modes:306 ABA test passes on the token check, not the epoch check: _validate_pending_page's lease_owner/epoch branch (paged_batches.sql:293-300) is never reached. Removing pending_worker is distinct from i_worker, the last_ack_worker check (paged_batches.sql:384), or the epoch predicate would leave the suite green. Blueprint lines 239-242 require stale-token and same-worker ABA coverage.
Fix: Add wrong-worker, random-UUID, null-worker, and wrong-worker receipt replay → PQP01 for normal/coop/partition. Add the sequence "pending token; slot lease expires; worker-b claims (new epoch); worker-a calls ack_page/renew_page with no re-receive" → page partition epoch fenced.

MEDIUM tests/test_paged_modes.sql:201 - Partition hash filtering before pagination is untested (n=1 only)

The only partition fixtures use subscribe_slot(..., 0, 1) (modes:201, concurrency.sh:203). With n=1 the _page_messages hash predicate (paged_batches.sql:31-35) is trivially true, so a wrong modulus, wrong null-key rule, or wrong lookahead/is_last on the filtered set would all pass. The concurrency reviewer confirmed by reading that the expression matches receive_partitioned. Only the tests are missing.
Fix: Use n≥3 with keys in different slots plus a null ev_extra1. Page with size 1–2 and assert per-slot membership, order and is_last. Add a filtered-empty window → advanced.

MEDIUM tests/test_paged_legacy.sql:6, tests/test_paged_modes.sql:10 - Legacy-bypass guards and coop takeover gating are only partly tested

No test hits the guards in pgque.receive (receive.sql:58), pgque.nack (receive.sql:201), receive_partitioned/nack_partitioned (partition_keys.sql:523,619), direct legacy next_batch_custom, or unregister_subconsumer (cooperative_consumers.sql:1077). event_retry/batch_retry/register_consumer_at are covered only with a hand-inserted page_state fixture. Coop gating covers only the case where both the dead interval and the lease have expired. Untested: a live lease refusing takeover, takeover between pages, a legacy allocator skipping a paged victim (cooperative_consumers.sql:~800), the 40001 "victim renewed" branch, and coop busy. Deleting any of these guards or predicates leaves CI green. (The concurrency reviewer confirmed by live repro that the guards and gating currently behave correctly.)
Fix: Using real receive_page* state (outstanding and between pages), assert 55000 for each guarded function and that the cursor and checkpoint are unchanged afterwards. Add the coop gating scenarios listed above.

MEDIUM clients/go/page_test.go:1 (all SDKs) - CI never runs the SDK paging code against a real database

Every SDK CI job already installs devel/sql/pgque.sql into a live container (ci.yml:565, 622, 666, 724), but every page test uses mocks or fake connections. Go's page tests cover only JSON names and durationInterval, with no ReceivePage/AckPage/ProcessPage. No test checks the SQL text, parameter order, ::uuid casts, or int8 round-trips through pg/pgx/psycopg/ruby-pg. The PR body itself says Ruby live verification is pending. (Validates CLI report finding #2.)
Fix: Add one live test per SDK covering receive_page → a handler failure that leaves the page for redelivery → ack_page with a msg_id of 9223372036854775807 → receipt replay → batch_finished.


NON-BLOCKING (3)

LOW devel/sql/pgque-api/paged_batches.sql:339-364,381 - The pre-auth failures array is O(n²) and unbounded

Suggestion: Two reviewers reproduced this independently: any pgque_reader with a random UUID can make one backend run for 3.3–3.5 s at 5k items and 53–58 s at 20k items before getting stale page token. Lock and validate the token first, reject arrays longer than pending_page_size, use a set-based duplicate check, and build the message→row map once instead of re-running unnest per failure.

INFO devel/sql/pgque-api/paged_batches.sql:1 - The header cites an issue number (#364) and does not follow the sibling header form

Suggestion: CLAUDE.md forbids issue numbers in code. Use -- pgque-api/paged_batches.sql -- Durable bounded batch pages (see blueprints/PAGED_BATCHES.md).

INFO docs/paged-batches.md:2 (+ docs/reference.md:21,25, docs/examples.md:16,20, nav label) - Version tags and change-history framing in docs/

Suggestion: "(0.3 development)" and "not available in the frozen 0.2.1 installer" break the CLAUDE.md docs rule. Use "development installer only (devel/sql/pgque.sql)" and move the release state to release notes.


POTENTIAL ISSUES (9)

MEDIUM clients/python/pgque/client.py:349-354 - process_page rejects every handler result other than None and calls close() on arbitrary returned objects (confidence: 6/10)

The common lambda m: conn.execute(...) returns a Cursor, so every attempt runs the side effect and then raises. Reproduced: the same message was inserted 3 times with no ack (a poison loop). This matches what the README documents, but it is a sharp edge.
Suggestion: Reject only lazy or awaitable results (coroutine/generator/iterator/awaitable), and close only those.

LOW clients/python/pgque/client.py:322, clients/ruby/lib/pgque/client.rb:115 (TS/Go types) - Failure msg_id must be a string, but each SDK's own Message.msg_id is an integer (confidence: 7/10)

ack_page(t, w, [{"msg_id": m.msg_id}]) fails with 22023 in Python and Ruby; in TS a bigint throws at JSON.stringify. It fails loudly, never silently. Ruby failures: nil → nil.map NoMethodError (found by reading only).
Suggestion: Have the SDKs convert int/bigint IDs to decimal strings; treat Ruby nil failures as [].

LOW SDK error mapping (Python _wrap_sql_error, TS mapPgError, Ruby) - No typed or exposed SQLSTATE for PQP01/PQP02 (confidence: 5/10)

Suggestion: Expose sqlstate or add typed stale-token and receipt-mismatch errors so callers can classify an ambiguous ack without inspecting __cause__.

LOW Tests, additional coverage gaps (confidence: 7/10)

Untested:

  • failure-descriptor validation rules (only duplicate IDs and PQP02 are covered);
  • expired owner ack/renew before any successor;
  • lease overflow/infinity, dead_interval <= 0, # consumer name;
  • coop/partition advanced, and the 21000 "pending page membership changed" branches;
  • writer denied and reader ack/renew/coop/partitioned (has_function_privilege);
  • coop takeover-vs-renew and double concurrent ack_page races;
  • upgrade from 0.1/0.2 with an in-flight batch then receive_page; the reinstall test skips under pg_tle;
  • wall-clock-sleep takeover tests (batches:317) may flake.

Suggestion: Add these as table-driven SQL cases; use DML-forced expiry.

LOW docs/README.md:20, docs/installation.md:340, docs/reference.md:21 - The new page is not linked from the docs index, and the grants are not listed (confidence: 7/10)

Suggestion: Link paged-batches.md from Guides, and add the page functions to the pgque_reader row. State that pgque_writer alone cannot call them.

LOW docs/paged-batches.md:117,140-147 - DLQ routing condition (queue_max_retries, default 5) not stated; the 55000/21000/40001 rows are narrower than the code (confidence: 6/10)

Suggestion: Add one sentence on max-retries → dead_letter, and widen the SQLSTATE meanings: 55000 also for incompatible active batch, 21000 also for membership changed, 40001 also for victim renewed.

LOW docs/paged-batches.md:61 - Cooperative paged receive does not auto-register members (it calls _validate_coop_names, not register_subconsumer) and has no setup snippet (confidence: 5/10)

The doc does say "subscribe explicitly … including cooperative members". What's missing is the function name, the fact that dead_interval := null disables takeover, and the name restrictions.

INFO devel/sql/pgque-api/paged_batches.sql:21-155,243-244 - SQL style (confidence: 6/10)

Several multi-argument clauses and assignments are packed onto one line. The multi-line comment at 243-244 uses -- where CLAUDE.md requires /* */ for 2+ lines. There is no consolidated design-notes block stating the slot → subscription → page_state lock order.

INFO devel/sql/pgque-api/paged_batches.sql:408 - ack_page calls renew_page (which extends the partition slot lease) with no comment (confidence: 4/10)


Validation of CLI report (samorev-cli-report.md) against source

CLI finding Verdict
HIGH: review target is a draft TRUE, isDraft=true at d1c56c1. This is a merge gate, not a code defect; the PR body keeps it in draft until this review.
MEDIUM: CI never runs the SDK paging code against a real DB TRUE, confirmed in ci.yml and the test files; carried as blocking above.
LOW: Python README non-autocommit guidance TRUE, upgraded to MEDIUM blocking: it contradicts docs/paged-batches.md:54.
Defender-dropped: stale type/table on reinstall Agree dropped (new objects; the reinstall-with-active-page test exists, though it is skipped under TLE).
Defender-dropped: TS camelCase failure names Agree dropped: page.test.ts asserts the snake_case JSON.
Defender-dropped: Python __version__ Agree dropped: release-time work.
Defender-dropped: principle-1 re-review reminder Agree dropped as a finding, but it still applies: any new commit needs a delta re-review.

Contested by this review and dropped: the docs reviewer's claim that takeover requirements are "stated too strongly". With no outstanding page, the lease condition is vacuously satisfied, so the doc is accurate.

Focus-area verdicts

Area Verdict
Crash/ack atomicity OK. Checkpoint + retry/DLQ + receipt + guard clear + finish_batch happen in one transaction. Receive never advances the cursor. Lookahead-based is_last. Crash before receive/ack commit is covered by real backend kills.
Lease / epoch fencing / ABA OK in code. Expired owner may ack until replaced, as documented. Partition epoch bumps on every owner change incl. release_slot. Tests missing (blocking above).
Cooperative transfer + independent receipts OK, reproduced: progress and boundary move, a new token is issued, both members' receipts replay, the victim ack during takeover gets stale page token, and legacy allocators skip paged victims.
Lock order OK: slot → subscription → page_state; coop main → member → victim (skip locked) → dest page → victim page; register_consumer_at consumer → main → members. No inversion found in reading or races.
Immutable membership / retry IDs OK: fixed tick snapshots + ev_id keyset; retry reinserts get a new txid; rotation pinned by min(sub_last_tick); adjacent-duplicate detection is sound at the N+1 boundary.
Legacy mutation bypasses All application-role mutators are guarded (finish/ack/ack_partitioned, receive/receive_coop/receive_partitioned, nack/nack_partitioned, event_retry x2, batch_retry, register_consumer_at incl. coop members, unregister_consumer/unsubscribe/unsubscribe_slot, unregister_subconsumer). Over-guarded: drop_queue(force) (blocking above).
Row/role permissions PASS (live ACL diff): page_state is owner-only; all 13 internals are revoked; public page API is reader-only; no PUBLIC execute; replaced legacy functions keep their base ACLs. Page tokens are not a security boundary between readers: a reader who knows a worker string can obtain that worker's token. This is consistent with the existing reader trust model, but worth documenting.
Managed Postgres / pg_tle PASS by reading plus the green CI pg_tle job: no superuser or C code; create type in DO, if not exists, custom PQP SQLSTATEs, and core gen_random_uuid are all fine.
SDK scalar/int8 decoding PASS (live) for Python/TS/Ruby; Go correct by reading (typed pgx int64/*time.Time), not executed.
SDK handler failures Raise/reject propagate with no ack in all SDKs. TS/Ruby lazy handlers get acked (blocking above); Python's strict-None rule is a potential issue.

NOT REVIEWED / limits

  • The Go SDK was not executed (no toolchain, and the host disk filled while pulling an image); reviewed by reading only.
  • The SDK receive_page_coop/receive_page_partitioned paths were not exercised live through the SDKs.
  • An actual pg_tle install was not run by reviewers; this relies on the green CI pg_tle install path job plus reading.
  • The repo SQL suites and test_paged_concurrency.sh were not re-run by the review; this relies on green CI at this head. Reviewer repros were separate.
  • Not checked: per-language style of client code, web/astro.config.mjs beyond the nav label, or the prose of docs/paged-batches.md beyond the claims above.
  • Not reported: a consumer switching mid-batch from legacy receive to receive_page silently adopts the in-flight batch (judged not a regression, not proven harmful).
  • SOC2 checks were skipped, as PgQue instructs.
  • Real-environment testing (CLAUDE.md: CI green → review → real testing) is still pending after these fixes.

Summary

Area Findings Potential Filtered
CI/Pipeline 0 0 0
Security 1 (LOW) 0 0
Bugs (SQL) 1 0 0
SDK 1 3 0
Tests 4 1 0
Guidelines 2 (INFO) 1 0
Docs 1 4 1
Metadata Draft PR (gate) 0 0

Verdict: FAILED (blocking). Not ready for merge. Each later commit needs a delta re-review at the new exact head (CLAUDE.md principle 1).


REV-assisted review (AI analysis by postgres-ai/rev), SamoRev specialized-agent run at d1c56c1ee9af9c81fc38f1fee627f8995be0128a

@NikolayS

NikolayS commented Oct 1, 2026

Copy link
Copy Markdown
Owner Author

REV Code Review Report (re-review)

  • PR: feat: durable paged batch consumption for 0.3 #368 - feat: durable paged batch consumption for 0.3
  • Author: @NikolayS
  • Reviewed head: e53b974d76940a4c61d90146280adccc29b27fa0 (local checkout verified equal to the PR head)
  • Previous review: d1c56c1 (comment), 7 blockers
  • Scope:
    • All source, test, SDK and doc changes from merge-base b8933a8 (6 commits, 48 files, +8052/-80).
    • Explicit delta review of d1c56c1..e53b974 (41ac2b6, e53b974; 27 files, +1534/-75).
    • Full-source coverage from the previous round is carried forward, after checking that the delta did not touch those paths.
  • AI-Assisted: Yes
  • Mode: --blocking
Pipeline Coverage
✅ 18/18 checks success at e53b974 (PG14–19beta1 matrix, pg_tle, pg_cron, pg_timetable, upgrade v0.1.0→HEAD, frozen-sql smoke, Python/Go/TS/Ruby jobs) Not reported

Result: FAILED, 4 blocking findings (1 SQL behaviour regression, 1 SDK, 2 test-coverage).

5 of the 7 previous blockers are resolved; the other 2 are only partly resolved. The new SQL (private allocator, authenticated and bounded failure validation, force drop) was reproduced live and is correct. However, the force-drop fix introduces a new admin regression for ordinary (non-paged) consumers, and the Python lazy-result fix still rejects the most common database handler after its side effect has run.

No CRITICAL or HIGH defect was found in the paging protocol.

Specialized reviewers (all 6 completed):

  1. SQL concurrency / bug hunter (opus). Ran a live PG18 two-session repro. Ran the full tests/run_all.sql, test_paged_destroy, test_paged_review and test_paged_validation, and test_paged_concurrency.sh through a local docker shim.
  2. Security / permissions / managed-PG (opus). Compared live ACLs for base, HEAD, and base→HEAD upgrade. Made 76 legacy-mutator calls as reader and admin, with a hash of the state tables taken before and after each. Timed DoS attempts.
  3. SDK bug hunter (opus). Ran all four SDK integration tests against a local PG18, including Go 1.25.6. Checked the CI logs at e53b974 to confirm the live tests actually ran rather than being skipped. Ran adversarial handler repros.
  4. Test analyzer (sonnet). Built a coverage matrix and mutation-tested a git archive copy on PG18.
  5. Guidelines (sonnet). PgQue CLAUDE.md plus postgres-ai rules.
  6. Docs (sonnet). Checked each docs claim against the source.

Generated-SQL assembly:

  • build/transform.sh was rerun on a clean git archive of HEAD, with the pgq submodule at v3.5.1. devel/sql/pgque.sql and pgque-tle.sql rebuild byte-identical. The uninstall scripts are hand-maintained and not changed by this PR; drop schema … cascade covers the new objects.

  • Final definitions win in the intended order:

    Function pgque.sql pgque-tle.sql
    _next_batch_custom (6-arg) :5920 :6010
    next_batch_custom SQL wrapper :6152 :6242
    7-arg coop next_batch_custom :6729 —
    drop_queue(text,bool) override :8619 (over :1513) :8709
  • Frozen sql/ is untouched. Duplicated generated text was not counted as extra reviewed logic.


Previous blockers: resolution at e53b974

# Previous blocker Status Evidence
1 Forced drop_queue blocked by an active page RESOLVED (but see new blocker B1) Override at paged_legacy.sql:267-333 deletes the subscriptions directly; page_state goes with them through the FK cascade. Live: force drop works with a page outstanding, between pages, and in partition and coop modes (test_paged_destroy.sql). No cleanup is lost compared with the original. Coop queues now drop cleanly; at base the coop unregister_consumer raised on them. ACLs are unchanged (pgque_admin and owner).
2 TS/Ruby lazy handlers acked RESOLVED Live: TS rejects sync and async generators, Promise→iterator and {next} before ack, with 0 side effects. Ruby rejects Enumerator and Enumerator::Lazy. Residual deferred-callable gap is LOW (N2).
3 Python README holds receive locks RESOLVED clients/python/README.md:179-201 commits right after receive_page, then processes, acks and commits, which matches docs/paged-batches.md:53-56. But the advice to use autocommit with process_page() interacts badly with B2.
4 Fencing negatives untested PARTIAL (B3) Wrong worker, null worker, forged token, wrong-worker receipt replay and a pending-token epoch fence are now tested (test_paged_review.sql:89-192); mutations are killed. However, the epoch predicate alone is not covered.
5 Partition hash filtering untested (n=1) RESOLVED test_paged_review.sql:195-312: n=3, null key, page size 1, per-slot membership, order, is_last, and filtered-empty → advanced. Four hash mutations (null rule, no filter, modulus, wrong slot) all go red.
6 Legacy guards and coop gating partly tested PARTIAL (B4) Now covered with real paged state:
- receive, nack, receive_partitioned, nack_partitioned, direct next_batch_custom, unregister_subconsumer;
- coop live-lease refusal, coop busy, takeover between pages, legacy allocator skipping a paged victim.
Still missing: the 40001 victim-renewed branch, and event_retry/batch_retry/register_consumer_at against real receive_page state.
7 CI never runs SDK paging live RESOLVED One live test per SDK. CI logs at e53b974 show them executed, not skipped:
- Python test_page_live_round_trip PASSED
- Go --- PASS: TestPageLiveRoundTrip
- TS page-integration.test.ts (1 test)
- Ruby TestPageIntegration#test_live_page_round_trip, 0 skips
Reviewer reran all four locally: pass. The int8 failure-path gap is LOW (N3).

The previous non-blocking items are also resolved:

  • The pre-auth O(n²) failures array: the token is now locked and authenticated first and the array is bounded by pending_page_size. A forged token with 20k items now fails in 0.5–22 ms; it took 53–58 s before.
  • The (#364) source header.
  • Docs version tags.
  • The docs index link and grants.
  • DLQ and max_retries wording.
  • The coop setup snippet.

BLOCKING ISSUES (4)

MEDIUM devel/sql/pgque-api/paged_legacy.sql:282-289,325-327 - Forced drop_queue now fails immediately with 40001 whenever any subscription row is locked, including ordinary legacy consumers. This regresses from base, and the new error is undocumented.

The override does perform 1 from pgque.subscription where sub_queue = … for update nowait, and maps lock_not_available to 40001 queue is in use; retry administrative force drop. A legacy consumer using the documented single-transaction pattern (receive → writes → ack, docs/concepts.md:42, docs/examples.md:132) holds its subscription row lock for the whole transaction.

Reproduced on PG18 with one consumer in begin; receive('lq','c1',10); pg_sleep(3):

  • base b8933a8: drop_queue('lq', true) waited 2.02 s, then returned 1;
  • HEAD: 40001 after 13 ms.

With two consumers looping receive; pg_sleep(0.3); commit, the force drop succeeded 0 of 40 attempts made 100 ms apart. An operator cannot destroy a busy queue without first stopping every consumer.

docs/reference.md:239-241 still says force "unregisters all attached consumers first". Neither the reference nor the SQLSTATE table in docs/paged-batches.md mentions this 40001, and the blueprint omits it too.

Fix: Preferred: take the partition_slot → subscription locks with waiting, in the same order the ack path uses, before taking the queue row. Then lock the queue and re-check for subscriptions created in the meantime. This makes the drop wait like base did, with no deadlock cycle against ack's subscription → queue key-share. Alternatively, keep NOWAIT but document it: the 40001, the retry, and that live consumers block the drop. Either way, add a test where a legacy consumer has an open transaction during a force drop.

MEDIUM clients/python/pgque/client.py:351-358 - process_page still rejects a psycopg Cursor returned by a handler, after its side effect has run, and closes the caller's cursor

The new lazy check includes isinstance(result, collections.abc.Iterator). Verified: issubclass(psycopg.Cursor, Iterator) is True (psycopg 3.3.3 and 3.3.6). So the most common database handler, lambda m: conn.execute("insert …"), runs its insert and then gets TypeError("handler must complete…"). The page is never acked, and result.close() closes the caller's cursor.

The README (clients/python/README.md:197-201) tells users to run process_page() on an autocommit connection. There, every retry commits the side effect again. Live repro: 3 attempts produced 3 TypeErrors and 3 inserted rows, and the page was never acked; a caller-owned cursor came back closed=True.

This is the previous potential issue, still present after a fix that targeted it, and the README now steers users into it.

Fix: Drop the Iterator clause. The existing pre-receive iscoroutinefunction/isgeneratorfunction/isasyncgenfunction check, plus a post-call isawaitable/isgenerator/isasyncgen check, covers lazy results. Call close() only on generators and coroutines. Add a unit test whose handler returns a psycopg Cursor.

MEDIUM tests/test_paged_review.sql:140-192 - The partition epoch predicate itself is not covered by any test (previous blocker #4, partial)

The pending-token fence test fences with a different worker (worker-b). lease_owner = i_worker and epoch = partition_epoch in _validate_pending_page (paged_batches.sql:296-303) therefore both fail together, and each masks the other. Mutation-confirmed: deleting only and epoch = i_state.partition_epoch leaves every SQL suite green.

The case where only the epoch fences is same-worker ABA. Same-owner claim_slot keeps the epoch (partition_keys.sql:409-417), so the sequence is:

  1. worker-a has a pending page;
  2. the lease expires and worker-b claims (epoch+1);
  3. b releases or expires and worker-a re-claims (epoch+2).

Now a's old token passes the token, worker and lease_owner checks, and only the epoch rejects it. The SQL reviewer reproduced this by hand and the code currently fences it correctly ("page partition epoch fenced"), but nothing in CI protects it. The blueprint (lines 239-242) explicitly requires a same-worker ABA test.

Fix: Add exactly that A→B→A sequence without a re-receive, and assert PQP01 page partition epoch fenced for both ack_page and renew_page with the old token.

MEDIUM devel/sql/pgque-api/cooperative_consumers.sql:~855-860, tests/test_paged_legacy.sql:30-88 - Coop 40001 "victim renewed" and legacy retry/registration guards on real paged state are still untested (previous blocker #6, partial)

  • Replacing the takeover recheck if found then with if false then leaves all SQL tests green. No two-session test holds the victim's renewal while a takeover runs.
  • event_retry (both forms), batch_retry and register_consumer_at are still exercised only against a hand-inserted page_state row (test_paged_legacy.sql:30-34). There is no outstanding-page or between-pages case from a real receive_page, and no check that subscription.sub_last_tick stays unchanged. The security reviewer confirmed live that these guards currently hold (76-call state-hash matrix); what is missing is regression protection.
  • The "legacy allocator skips paged victim" check (test_paged_review.sql:503-510) is only killed incidentally, through the 40001 recheck aborting the block. It does not assert the expected outcome.

Fix:

  • Add a two-session case to test_paged_concurrency.sh: A locks or renews the victim, B's takeover gets 40001.
  • Reuse the test_paged_review.sql before/after pattern for event_retry, batch_retry and register_consumer_at, outstanding and between pages, comparing page_state and sub_last_tick.
  • Assert explicitly that the legacy steal returns null and the victim keeps its sub_batch.

NON-BLOCKING (8)

LOW devel/sql/pgque-api/paged_legacy.sql:325-327 - exception when lock_not_available wraps the whole drop_queue body

Suggestion: A lock_timeout (also 55P03) on the queue row or on drop table is rewritten into 40001 "retry administrative force drop", even for the non-force drop_queue(text). Reproduced with lock_timeout='200ms'. Wrap only the two NOWAIT statements in an inner block. (This may become moot after the B1 fix.)

LOW clients/typescript/src/client.ts:278-284, clients/ruby/lib/pgque/client.rb:149-155 - Handlers returning a deferred callable are still acked with their body never run

Suggestion: Live examples: TS () => () => {ran++} gives processedCount:1, ran=0; Ruby -> { ran += 1 } and Fiber.new {} give ACKED, body 0. The blueprint (PAGED_BATCHES.md:230-231) says unknown handler types must be explicit errors. Reject Function, Proc and Fiber results as well.

LOW clients/*/…page_integration… - The max-int8 ID is delivered but never used as a failure msg_id

Suggestion: Put the 9223372036854775807 event on the failed page, send {"msg_id":"9223372036854775807"}, and assert the retry_queue row in all four SDKs. Only Python checks routing today.

LOW docs/paged-batches.md:175-179 - The worker name, not the token, is the effective credential

Suggestion: receive_page with the same worker string returns the same pending token (paged_batches.sql:132-137,167). This is consistent with the existing reader trust model, since readers can already finish_batch any batch. State that worker IDs should be random per process, not secret-bearing, and not logged.

LOW docs/paged-batches.md:149-156, blueprints/PAGED_BATCHES.md:3,18,28-29,74-76,99-102 - The SQLSTATE table and the blueprint are not updated to the as-built state

Suggestion:

  • SQLSTATE table: receive-time coop and slot setup errors are P0001, not 22023 or PQP01; busy never occurs in partition mode.
  • Blueprint: it still says "Proposed" and omits 40001 and the private _next_batch_custom/_next_batch_coop(..., i_paged) allocator.

LOW SDK READMEs (TS, Ruby, Go) - The page helpers and their handler contracts are documented only for Python

Suggestion: Add a short per-SDK section covering:

  • the method names;
  • the handler contract (Python lazy rejection, TS iterator rejection, Ruby Enumerator, Go error);
  • the failures key spelling: TS msgId/retryAfterSeconds, Ruby both spellings, Python and Go snake_case;
  • that msg_id must be a decimal string.

LOW clients/python/tests/conftest.py:54-67 and tests/test_paged_review.sql (end) - Cleanup leaks after a paged failure

Suggestion: Python teardown calls unregister_consumer before drop_queue, so a mid-page failure raises 55000 and leaks the queue; call drop_queue(..., true) first. test_paged_review.sql leaves its review_* queues with open pages in the shared CI database; force-drop them at the end, which also exercises B1.

INFO Style (paged_batches.sql:246-247, cooperative_consumers.sql new coop recheck comment, several tests/test_paged_*.sql comments; packed assignments in _receive_page; tests/test_paged_batches.sql:4 still cites (#364); commit b58cb4b body has literal \n)

Suggestion:

  • Use /* */ for 2+ line comments, and one assignment or argument per line (CLAUDE.md).
  • Add a consolidated lock-order design-notes block covering slot → subscription → page, coop, and drop.
  • Remove the issue number from the test header.
  • Fix the squash-merge message.

POTENTIAL ISSUES (3)

LOW Python/Ruby/TS ack_page - SDKs pass failure msg_id through as-is, while their own Message.msg_id is an integer (confidence: 7/10)

Carried from the previous review, unchanged. It always fails loudly:

  • Python: 22023.
  • TS bigint: "Do not know how to serialize a BigInt", mis-wrapped as a SQL error.
  • Ruby failures: nil: NoMethodError.

Go forces a string type.

Suggestion: Convert int and bigint to a decimal string; treat Ruby nil as [].

LOW SDK error mapping - No typed or exposed SQLSTATE for PQP01/PQP02 (confidence: 5/10)

Carried from the previous review, unchanged.

INFO docs/paged-batches.md:12 - "Receiving does not advance progress" is not true for advanced (an empty window is finished inside receive) (confidence: 4/10)


Validation of CLI report (samorev-cli-review2.md) against source

CLI finding Verdict
HIGH: review target is a draft TRUE as a state, not a code defect. isDraft=true at e53b974. It is intentionally held in draft until the gates pass.
LOW: Python integration test sets autocommit=True on the shared conn fixture and never restores it FALSE. clients/python/tests/conftest.py:30-35: conn is function-scoped, with a fresh psycopg.connect per test that is closed at teardown, so nothing leaks to other tests. The real nearby issue is the teardown order (N7).
Defender-dropped: Closes #364 wording Agree dropped (PR metadata, not code).
Defender-dropped: blockers can't be confirmed from a truncated diff Agree dropped. This review read the full sources and verified each blocker above.
Defender-dropped: renamed next_batch_custom leaves an unguarded public function Agree dropped. Public 5-arg is redefined (create or replace, cooperative_consumers.sql:414-437) as a SQL wrapper calling _next_batch_custom(..., false). ACLs are preserved; the private 6-arg is executable by the owner only. Live: 55000 for next_batch/next_batch_info/next_batch_custom on a paged batch.
Defender-dropped: Ruby test file missing copyright header Agree dropped (cosmetic; other test files have none either).

The CLI missed all four blockers above. The B2 trigger (psycopg Cursor is an Iterator) and the B1 regression only show up by running the code.

Focus-area verdicts

Area Verdict
Crash/ack atomicity OK, unchanged by the delta. Checkpoint, retry/DLQ, receipt, guard clear and finish_batch happen in one transaction. Failure validation now runs after auth but still inside the same ack transaction.
Lease / epoch fencing / ABA OK in code. A same-worker A→B→A repro is fenced by the epoch. The epoch predicate is not regression-tested (B3).
Cooperative transfer + independent receipts OK, unchanged. Coop live-lease refusal, busy and takeover between pages are now tested. The 40001 victim-renewed branch is untested (B4).
Lock order OK for the paging paths. drop_queue takes queue → slot NOWAIT → subscription NOWAIT, so it cannot deadlock, but it gives up instead of waiting (B1).
Immutable membership / retry IDs OK. Unchanged by the delta; set-based failure normalization keeps the previous semantics (ordering, canonical text ID, "1"/"01" duplicates, default 60, null reason).
Legacy mutation bypasses PASS (live, 76 calls, reader and admin, outstanding and between pages). The private allocator closes next_batch/next_batch_info/next_batch_custom(5/7). No remaining bypass found. Some outer guards (receive, nack, *_partitioned) are shadowed by inner guards; that is defence-in-depth.
Row/role permissions PASS (live ACL diff for fresh install and upgrade). _next_batch_custom is owner-only. The wrapper keeps base ACLs and search_path. drop_queue ACLs are unchanged. page_state is owner-only. There is no PUBLIC execute. The reinstall grant drift to pgque_admin predates this PR.
Managed Postgres / pg_tle PASS. No superuser-only constructs or C. The TLE wrapper tag is $pgque_extension_body$, which the new $$ bodies do not clash with. Final-definition order is correct. CI pg_tle install path runs the full run_all.sql, including all test_paged_*, green.
SDK scalar/int8 decoding PASS (live, all four SDKs incl. Go). Failure-path max-int8 is not exercised (N3).
SDK handler failures Raise/reject → no ack, page redelivered, in all four (live + CI). TS/Ruby lazy handlers are fixed. Python rejects psycopg Cursor after the side effect (B2). Deferred-callable gap (N2).

NOT REVIEWED / limits

  • pg_tle: no reviewer ran a real install; this relies on the green CI pg_tle install path job plus reading pgque-tle.sql.
  • PG versions: reviewer repros ran on PG18 only. PG14–17 and 19beta1 rely on the green CI matrix.
  • Concurrency suite: test_paged_concurrency.sh was run through a docker shim against a local cluster, not under real docker.
  • SDK coop/partition paths: receive_page_coop and receive_page_partitioned were not exercised live through the SDKs.
  • CI runners: Ruby was run directly (no bundler on the host), and TS under node 24 rather than Bun as CI uses.
  • Mutation testing: not done for coop between-pages takeover, coop busy, filtered-empty advanced, direct next_batch_custom guard, or the SDK integration tests.
  • Client code style: not reviewed in depth for Python/TS/Ruby/Go.
  • Pre-existing, out of scope: drop_queue builds DDL from table names stored in pgque.queue, which pgque_admin can edit. The same code exists on base. It is worth a separate ticket.
  • SOC2: skipped, as PgQue instructs.
  • Real-environment testing: still pending. PgQue CLAUDE.md sets the order as CI green → review → real testing.

Summary

Area Findings Potential Filtered
CI/Pipeline 0 0 0
Security 0 0 0
Bugs (SQL) 1 blocking + 1 LOW 0 1 (orphan consumer row under SKIP LOCKED, harmless)
SDK 1 blocking + 2 LOW 2 0
Tests 2 blocking + 1 LOW 0 0
Guidelines 1 (INFO, grouped) 0 0
Docs 3 LOW 1 0
Metadata Draft PR (intentional gate) 0 0

Verdict: FAILED (blocking). Not ready for merge. Each later commit needs a delta re-review at the new exact head (CLAUDE.md principle 1).


REV-assisted review (AI analysis by postgres-ai/rev), SamoRev specialized-agent re-review at e53b974d76940a4c61d90146280adccc29b27fa0

Refs #364. Preserve completed SDK handler results, cover epoch and legacy guards, and document administrative contention.
@NikolayS

NikolayS commented Oct 1, 2026

Copy link
Copy Markdown
Owner Author

Re-review fixes and clarification

This revision addresses the second review's SQL/SDK defects and regression gaps:

  • B1: uses the review's explicit-documentation alternative: administrative force drop retains deadlock-safe NOWAIT for all slot/subscription rows, including ordinary consumers. Documents 40001, whole-transaction retry, and pausing consumers for reliable destruction. A real legacy receive transaction now verifies rejection is atomic and retry succeeds after unlock. Only those NOWAIT locks translate to 40001; unrelated force/non-force lock timeouts retain 55P03 (red/green test).
  • B2: Python accepts an executed psycopg cursor and never closes caller-owned cursors. Only actual coroutine/generator results are closed. The live regression inserts effects, checks their IDs exactly once, verifies returned cursors stay open, and verifies acknowledgment replay. TS/Ruby also reject deferred callables.
  • B3: same-worker A→B→A, without re-receiving, leaves token/worker valid and makes only the epoch reject both ack and renewal. Removing the epoch predicate makes the new test fail.
  • B4: real outstanding/between-page event_retry (both overloads), batch_retry, and register_consumer_at must preserve page state and sub_last_tick; legacy stealing explicitly returns NULL and preserves the victim batch. Removing the event_retry guard makes the regression fail.

B4 concurrency clarification for the next reviewer

The previous request for "A holds a real renewal, B waits and gets 40001" does not match the implementation: the victim scan uses FOR UPDATE OF candidate SKIP LOCKED. renew_page locks that subscription and refreshes sub_active together with the lease. An in-flight renewal is skipped; a committed one is excluded by the refreshed predicate. Once takeover owns the subscription lock, public renewal cannot pass it before the defensive recheck.

The deterministic real-API test now checks the actual contract: takeover returns idle promptly while renewal is open, still cannot steal after renewal commits, and the original token, live lease, and assignment survive. The existing defensive 40001 recheck is retained. Please adjudicate the prior branch-coverage demand against this locking protocol; no test-only source hook was added to manufacture an impossible public-API interleaving.

Non-blocking review items

  • Added max-int8 failure routing assertions to all four live SDK tests without removing successful helper/replay coverage.
  • Added SDK helper/handler contracts; corrected Python and SQL test cleanup.
  • Updated setup SQLSTATEs, partition status semantics, worker-identity trust, the advanced exception, as-built private allocators, and consolidated lock-order notes.
  • Failure descriptors intentionally continue to require decimal-string IDs, with documented examples and live int8 coverage; implicit ID coercion and a cross-SDK typed-error redesign are outside this correctness patch. Existing driver error behavior is retained.
  • Existing historical commit messages are not rewritten. Any future squash message should use real newlines.

Fresh exact-head CI and SamoRev review are still required. No merge or release in this step.

Refs #364. Exercise JavaScript runtime rejection without weakening the public PageHandler type.
@NikolayS

NikolayS commented Oct 1, 2026

Copy link
Copy Markdown
Owner Author

REV Code Review Report (re-review 3)

  • PR: feat: durable paged batch consumption for 0.3 #368 - feat: durable paged batch consumption for 0.3
  • Author: @NikolayS · AI-Assisted: Yes · Mode: --blocking
  • Reviewed head: 6bb17158cda5228b4ed98b07721910593a3741f6 (local == PR head, verified)
  • Previous review: e53b974 (comment), 4 blockers
  • Scope: delta e53b974..6bb1715 (4f4ab4e, 6bb1715; 26 files, +628/-123), source read directly, with its interaction with the full protocol. Full-source coverage from rounds 1–2 is carried forward for unchanged paths. The only SQL source change is paged_legacy.sql drop_queue(text,bool).
Pipeline Coverage
✅ 18/18 checks success at 6bb1715 (PG14–19beta1, pg_tle, pg_cron, pg_timetable, upgrade v0.1.0→HEAD, frozen-sql smoke, Python/Go/TS/Ruby) Not reported

Result: FAILED, 1 blocking finding (Python SDK). B1, B3 and B4 are resolved. B2 is partial: executed psycopg cursors are now accepted, but the fix also lets real lazy and deferred results through to ack. No defect was found in the SQL protocol.

Specialized reviewers (all 6 completed):

  1. SQL concurrency (opus): live PG18/PG14, mutants, assembly rebuild.
  2. Test analyzer (opus): 16 mutants on an archive copy, local PG18.3, CI log audit.
  3. SDK (opus): all four SDK suites live on PG18 plus adversarial handler probes.
  4. Security / permissions / managed-PG / pg_tle (opus): live ACL diff HEAD vs e53b974 vs base, 156-call bypass matrix, local pg_tle 1.5.2 install.
  5. Docs (sonnet).
  6. Guidelines (sonnet). Partial: its shell was unavailable (disk full), so the orchestrator checked commit messages and file modes.

Environment caveat: the host root filesystem was at 100% during the review. Some runs used /dev/shm, a local cluster, or a docker shim instead of real containers. Every result below was re-run green after a transient ENOSPC failure.

Generated-SQL assembly: build/transform.sh on a clean git archive HEAD (pgq v3.5.1) produces byte-identical pgque.sql and pgque-tle.sql. The drop_queue(text,bool) override is the final definition in both: pgque.sql:8619 (base :1513) and pgque-tle.sql:8709. Duplicated generated text was not counted as extra reviewed logic.


Previous blockers: resolution at 6bb1715

# Blocker Status Evidence
B1 Force drop_queue NOWAIT 40001 regression for legacy consumers RESOLVED (review-approved alternative) See below.
B2 Python rejected a psycopg Cursor after its side effect PARTIAL Cursor accepted, acked, and not closed (live, psycopg 3.3.3). New blocker N-B1 below.
B3 Epoch predicate untested RESOLVED See below.
B4 Coop 40001 victim-renewed test and legacy real-state guards RESOLVED (the coop demand is corrected below) See below.

B1 evidence:

  • The 40001 mapping now wraps only the two NOWAIT statements (paged_legacy.sql:282-294).
  • Live PG18, all 55P03:
    • queue-row lock_timeout, force and non-force;
    • drop table blocked, non-force;
    • drop table blocked, force (with no partial state).
  • Live PG18, legacy begin; receive():
    • force drop → 40001 in 85 ms, everything intact;
    • retry after commit → 1.
  • Lock order: queue → slot NOWAIT → subscription NOWAIT. This has no cycle with ack (slot → sub → page → queue key-share).
  • Mutants:
    • The e53b974 whole-body mapping is killed by test_paged_drop_timeout.sh.
    • Removing NOWAIT is killed by test_paged_concurrency.sh (legacy_drop_busy).
  • Docs: reference.md:239-248, paged-batches.md and the blueprint document the 40001, retry, and pausing consumers.
  • CI: the script runs in all 6 matrix jobs.

B3 evidence:

  • test_paged_review.sql runs an isolated same-worker A→B→A sequence with no re-receive: old token, pending_worker = a, old epoch. Both ack_page and renew_page raise PQP01.
  • Deleting only and epoch = i_state.partition_epoch is now killed (test_paged_review.sql:201).

B4 legacy evidence:

  • Real receive_page state, outstanding and between pages.
  • Removing the guard from event_retry (both forms), batch_retry or register_consumer_at is killed by test_paged_review.sql alone.
  • An explicit "legacy allocator must return null / victim keeps sub_batch" assert kills the victim-join mutant even with the recheck neutralised.

B4 coop evidence:

  • The new two-session case is deterministic: the renewer holds its locks behind a pg_stat_activity PgSleep barrier.
  • It asserts:
    • takeover → idle, with no wait, in under 3 s while the renewal is open;
    • idle again after commit;
    • the victim keeps its sub_batch, token and a live lease.
  • SKIP LOCKED → plain FOR UPDATE and SKIP LOCKED → NOWAIT mutants are both killed.

Adjudication of review2's coop 40001 demand

Review2 asked for "A renews, B's takeover gets 40001". That demand was based on wrong lock semantics and is withdrawn.

  • renew_page locks the victim's subscription row (paged_batches.sql:264-266) and refreshes sub_active and the lease (:325).
  • The takeover selects candidates with for update of candidate skip locked (cooperative_consumers.sql:825-848).
  • So a held renewal makes the takeover skip the victim and return idle. After commit, the live lease excludes the victim. Both are the intended outcomes, not defects.

The defensive 40001 recheck (:849-860) is unchanged and correct.

  • It is reachable only in a narrow READ COMMITTED race: the renewal commits between the takeover's snapshot and its row lock, and the renewal transaction was held open longer than the dead interval. The sub_active re-read passes, but the unlocked page_state join is still stale.
  • Reproduced: 1 × 40001 in 454,313 attempts over 28 s, with 0 steals.
  • It cannot be tested deterministically without instrumentation, and none is required.

BLOCKING ISSUES (1)

MEDIUM clients/python/pgque/client.py:349-357 - The B2 fix acks handlers that return unexecuted lazy or deferred results. Their side effect never runs, so this loses data. map/filter results regressed from e53b974.

The isinstance(result, Iterator) rejection was removed outright, and no returned-callable rejection was added.

Live (psycopg 3.3.3, PG18):

Handler returns Outcome Side effects Redelivered
map(...) ACCEPTED, processed_count=1 0 no
filter(...) ACCEPTED 0 no
returned lambda ACCEPTED 0 no
functools.partial ACCEPTED 0 no

At e53b974, map and filter were rejected.

This contradicts:

  • the blueprint: "Unknown handler types must be explicit errors" (blueprints/PAGED_BATCHES.md:259);
  • the Python README: "lazy/async handlers leave [the page outstanding]" (clients/python/README.md:198);
  • TS and Ruby, which now reject callables.

Review2 itself suggested "drop the Iterator clause", so part of this regression comes from that advice. The requirement stated for this round is to reject actual lazy and deferred results.

Fix:

  • Reject callable(result).
  • Reject isinstance(result, collections.abc.Iterator) unless it is a psycopg cursor (isinstance(result, psycopg.cursor.BaseCursor), which covers client and server cursors). Close nothing that the caller owns.
  • Add unit tests for map(...) and a returned lambda: both must raise TypeError, with no ack.

NON-BLOCKING (9)

LOW clients/ruby/lib/pgque/client.rb:155 - A returned Method or UnboundMethod is still acked with 0 side effects (live).

Suggestion: Reject any outcome that respond_to?(:call), or add Method and UnboundMethod.

LOW paged_batches.sql:299 / tests/test_paged_review.sql - The lease_owner = i_worker predicate alone is not mutation-killed.

Its only independent case is "worker releases its own slot, then acks or renews the old token". release_slot keeps the epoch and sets the owner to null (partition_keys.sql:467-470). The code fences this correctly today.

Suggestion: Add claim → receive → release_slot → ack and renew with the old token → expect PQP01.

LOW tests/test_paged_review.sql:6-33 - expect_page_error checks only the SQLSTATE.

Both the generic stale-token check and the epoch check raise PQP01.

Suggestion: Also assert the message page partition epoch fenced for the A→B→A cases.

LOW tests/test_paged_legacy.sql:139-195 - The coop-main register_consumer_at member-page guard is killed only by hand-inserted state.

Suggestion: Repeat it in the review_coop_gate block with a real receive_page_coop page.

LOW tests/test_paged_concurrency.sh:312-385 - The renewal-race case has no positive control.

Suggestion: Assert the victim is stealable (old sub_active, expired lease) before the renewer starts.

Also, the renew_page sub_active refresh is not mutation-detected: the lease predicate fences after commit. Document it as a liveness signal.

LOW docs/paged-batches.md, docs/reference.md - Docs omit that a queue-row lock_timeout during drop surfaces as 55P03, not 40001.

The P0001 row should also mention receive-time slot-lease and subscription preconditions. PQP01 should be scoped to ack and renew.

LOW SDK READMEs - The page sections don't say the API requires the development installer.

The Python README lacks the failure-descriptor example (str(m.msg_id)) and the accepted-cursor contract.

INFO Housekeeping:

  • tests/test_paged_drop_timeout.sh is mode 100644. CI invokes it via bash, so this is cosmetic.
  • docs/tutorial.md:367 and docs/monitoring.md:136 show force drop with no 40001 caveat.
  • The blueprints/PAGED_BATCHES.md:213 "Proposed safe initial rule" is now implemented.

INFO Carried from review2 and still present: tests/test_paged_batches.sql:4 cites (#364).

The other style items (multi-line -- comments, packed assignments, lock-order notes block) were not re-verified this round.

POTENTIAL ISSUES (1)

LOW tests/test_paged_review.sql:358-447 - The sub_last_tick before/after asserts cannot fail (confidence: 6/10).

Each guarded call runs inside an exception subtransaction, so any write is rolled back anyway. The guards are still covered by the 55000 expectations.


Validation of the CLI report (samorev-cli-review3.md)

CLI finding Verdict
HIGH: target is a draft TRUE as a state, not a code defect. It is held in draft until the gates pass.
LOW: test_paged_drop_timeout.sh mode 100644 vs 100755 TRUE but cosmetic (INFO above). CI runs it via bash tests/... (ci.yml:53-54), and no doc says ./.
Defender-dropped (4: batch_page type guard, msgId type ergonomics, teardown unregister removal, principle-1 restatement) Agree all dropped. The teardown change is correct: unregister_consumer raises 55000 mid-page.

The CLI missed the blocking Python regression, which only shows up when the code is run.

Focus-area verdicts

Area Verdict
Crash/ack atomicity OK. Unchanged; carried forward.
Lease / epoch fencing / ABA OK. Same-worker ABA now regression-tested. lease_owner-alone is untested (LOW).
Coop transfer + independent receipts OK. SKIP LOCKED renewal contention tested; 40001 recheck correct (rare race, reproduced).
Lock order OK. No cycle; drop is NOWAIT by documented policy.
Immutable membership / retry IDs OK. Unchanged. Max-int8 failure routing now tested in all four SDKs (N3 resolved).
Legacy mutation bypasses PASS (156 calls, reader and admin, outstanding and between pages, standard, partition and coop; 0 state changes; identical to e53b974).
Row/role permissions PASS. Full ACL dump identical HEAD vs e53b974; _next_batch_custom owner-only; no PUBLIC execute.
Managed PG / pg_tle PASS. Local pg_tle 1.5.2 install with run_all.sql green; CI pg_tle job runs all 10 test_paged_*.
SDK scalar/int8 decoding PASS (live, all four).
SDK handler failures Raise/reject → no ack + redelivery in all four. Python lazy/deferred results acked (blocker). Ruby Method (LOW).
CI executes SDK integrations PASS. Run 36926331767: Python test_page_live_round_trip PASSED (85/0 skipped), Go TestPageLiveRoundTrip PASS, TS page-integration.test.ts ✓ (93/93), Ruby test_live_page_round_trip (84 runs, 0 skips).

NOT REVIEWED / limits

  • PG versions: reviewer repros ran on PG18 and PG14 only (shell suites). PG15–17 and 19beta1 rely on the CI matrix.
  • pg_tle: no managed-service run (RDS/Aurora). Local pg_tle was checked only with the outstanding-page bypass matrix.
  • Mutation testing: local PG18.3 through a docker shim (disk full), not real containers. SDK code was not mutated.
  • Toolchains: Ruby 3.2.3 without bundler; no PG19beta1 SDK run.
  • Guidelines: the guidelines agent was partial. Org .mdc rules were not opened (CLAUDE.md takes precedence), and review2 style INFO items were not re-verified.
  • SOC2: skipped, as PgQue instructs.
  • Real-environment testing: pending. PgQue CLAUDE.md order is CI green → review → real testing.

Summary

Area Findings Potential Filtered
CI/Pipeline 0 0 0
Security 0 0 0
Bugs (SQL) 0 0 0
SDK 1 blocking + 1 LOW 0 0
Tests 4 LOW 1 0
Guidelines 2 INFO 0 0
Docs 2 LOW (+ INFO) 0 0
Metadata Draft (intentional gate) 0 0

Verdict: FAILED (blocking, 1 MEDIUM). Not ready for merge. Fix the Python lazy/deferred acceptance, then delta re-review at the new exact head (CLAUDE.md principle 1).


REV-assisted review (AI analysis by postgres-ai/rev), SamoRev specialized-agent re-review at 6bb17158cda5228b4ed98b07721910593a3741f6

Refs #364. Preserve completed psycopg cursors while rejecting lazy iterators and returned callables before acknowledgment.
Refs #364. Remove accidentally staged local test artifacts.
@NikolayS

NikolayS commented Oct 1, 2026

Copy link
Copy Markdown
Owner Author

Review 3 follow-up

Head: 220f3f0 (Python/Ruby source fix plus removal of accidentally staged local bytecode).

  • Blocking Python deferred-result finding fixed: callables and unknown iterators raise before acknowledgment. The exception is the public synchronous psycopg.Cursor family (including server/client cursor subclasses), not internal BaseCursor or arbitrary cursor-shaped objects. Caller-owned cursors stay open.
  • Ruby LOW finding fixed: returned bound/unbound methods and custom callable objects are rejected as deferred work.
  • Python unit tests: 5 regressions fail on the reviewed code, then all 12 pass on the fix. Ruby deferred-callable regression fails before the fix; all 7 tests / 29 assertions pass after it.
  • Fresh PG18 integration passed for all four SDKs. Python additionally proves map/filter/partial rejection leaves the same page token available, with zero application effects, and that an executed psycopg cursor is accepted without closing it.
  • Python README now documents accepted cursor ownership.

Remaining non-blocking review items are deferred beyond this handler correction: independent lease-owner mutation coverage, more specific fencing error assertions, real-state coop-main registration coverage, renewal positive-control coverage, and documentation refinements. Existing SQL source is unchanged from the reviewed head; review 3 found its guards correct and its CI/mutation coverage passing. The exception-subtransaction state assertions are supplementary: the SQLSTATE assertions remain the guard-discrimination test. Script mode is intentional in CI (bash tests/test_paged_drop_timeout.sh). No managed-service validation is claimed.

Updated CI and exact-head review must pass before final fresh-database verification. This is development work targeting 0.3, not a change to released 0.2.1, and no merge or publication is performed here.

@NikolayS

NikolayS commented Oct 1, 2026

Copy link
Copy Markdown
Owner Author

Exact-head review gate blocked

Head: 220f3f0. All 18 CI checks pass: https://github.com/NikolayS/PgQue/actions/runs/36939170323

The re-review attempt did not complete:

  • SamoRev CLI reports LLM reviewer unavailable and fails closed.
  • Specialized Claude review exits with a weekly usage-limit error (reported reset: October 6, 12:00 UTC).
  • GPT-5.6 Sol delegation attempts report model capacity errors.

No review PASS is claimed. The prior completed review is for 6bb1715; it identified the Python blocker addressed by this head. Local red/green unit and fresh-database SDK evidence is posted above, but is not a replacement for exact-head review. Final post-review smoke remains pending, and the PR remains draft and unmerged.

This branch has not been deployed

No deployments
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

feat: design safe paged batch consumption

2 participants