From dbbad12012a8b04c04a2a96a576dbd79fbbc04f2 Mon Sep 17 00:00:00 2001 From: Marco Bambini Date: Sat, 19 Sep 2026 11:40:05 +0200 Subject: [PATCH 1/4] Restore PostgreSQL savepoint callers without a fixed depth cap --- docs/internal/deep-savepoints.md | 13 ++ src/postgresql/database_postgresql.c | 100 +++++---- test/postgresql/62_deep_savepoints.sql | 294 +++++++++++++++++++++++++ test/postgresql/full_test.sql | 1 + 4 files changed, 369 insertions(+), 39 deletions(-) create mode 100644 docs/internal/deep-savepoints.md create mode 100644 test/postgresql/62_deep_savepoints.sql diff --git a/docs/internal/deep-savepoints.md b/docs/internal/deep-savepoints.md new file mode 100644 index 00000000..701a8e9b --- /dev/null +++ b/docs/internal/deep-savepoints.md @@ -0,0 +1,13 @@ +# PostgreSQL savepoint depth + +A heap scan feeding `cloudsync_payload_apply(payload)` failed with `buffer pin ... is not owned by resource owner SubTransaction` at 126 user savepoints. The caller's resource owner and memory context were stored in a 128-entry array; Cloudsync's internal savepoints consumed additional levels and silently exceeded the array. + +The backend now keeps one dynamically allocated frame per Cloudsync subtransaction. Each frame records the caller's owner, memory context and subtransaction ID. Commit and rollback restore a local copy after PostgreSQL ends the subtransaction. Subtransaction callbacks discard frames on both normal completion and external abort; transaction callbacks clear the stack before its owning memory context disappears. Snapshot replacement remains restricted to the outermost Cloudsync savepoint. + +This removes Cloudsync's fixed depth cap; PostgreSQL's own resource limits still apply. Allocations happen before opening a new subtransaction. + +## Validation + +`test/postgresql/62_deep_savepoints.sql` is included in `full_test.sql`. It exercises depths 1, 125, 126, 127, 128, 256, 1024 and 2048, reads payloads from a heap table, checks rollback, reapplies in the same transaction and commits. It also catches 100 trigger failures at depth 256 and verifies a subsequent successful apply in the same backend. + +The full suite passes all 546 checks on PostgreSQL 15.19, 17.11 and 18.6. Restoring the old implementation reproduces the buffer-owner error at depth 126. These are local database tests, not tests of a deployed cloud server. diff --git a/src/postgresql/database_postgresql.c b/src/postgresql/database_postgresql.c index 0701f674..da4e55d6 100644 --- a/src/postgresql/database_postgresql.c +++ b/src/postgresql/database_postgresql.c @@ -3062,32 +3062,48 @@ static int database_refresh_snapshot (void) { // pin of the scan feeding cloudsync_payload_apply, say) is then charged to the wrong // owner: "buffer pin ... is not owned by resource owner TopTransaction". So, as the // procedural languages do, remember both at the start of each subtransaction and -// restore them when it ends. Indexed by nesting level so an owner never leaks across -// a subtransaction that was aborted elsewhere. -#define CLOUDSYNC_SAVEPOINT_MAX_DEPTH 128 -static ResourceOwner savepoint_owner[CLOUDSYNC_SAVEPOINT_MAX_DEPTH]; -static MemoryContext savepoint_context[CLOUDSYNC_SAVEPOINT_MAX_DEPTH]; - -// Only the outermost savepoint opened here may swap the active snapshot when it ends: a -// snapshot pushed while an enclosing savepoint is still open belongs to that -// subtransaction, and rolling the enclosing one back pops it — in place of the caller's -// snapshot the swap popped, leaving the caller's portal without one (an assertion -// failure in EnsurePortalSnapshotExists). Nested, advancing the command counter is -// enough to make the changes visible: every SPI statement takes a fresh snapshot. -// True when the subtransaction at level is nested inside another savepoint opened here. -static bool savepoint_is_nested (int level) { - for (int k = 2; k < level && k < CLOUDSYNC_SAVEPOINT_MAX_DEPTH; k++) { - if (savepoint_owner[k]) return true; - } - return false; -} - -static void savepoint_restore_caller (int level) { - if (level <= 0 || level >= CLOUDSYNC_SAVEPOINT_MAX_DEPTH || !savepoint_owner[level]) return; - MemoryContextSwitchTo(savepoint_context[level]); - CurrentResourceOwner = savepoint_owner[level]; - savepoint_owner[level] = NULL; - savepoint_context[level] = NULL; +// restore them when it ends. Frames live in TopTransactionContext and are removed +// by transaction callbacks even when a caller aborts outside these wrappers. +typedef struct cloudsync_savepoint_frame { + ResourceOwner owner; + MemoryContext context; + SubTransactionId subid; + struct cloudsync_savepoint_frame *previous; +} cloudsync_savepoint_frame; + +static cloudsync_savepoint_frame *savepoint_stack; +static bool savepoint_callbacks_registered; + +static void savepoint_xact_callback (XactEvent event, void *arg) { + if (event == XACT_EVENT_COMMIT || event == XACT_EVENT_ABORT || + event == XACT_EVENT_PARALLEL_COMMIT || event == XACT_EVENT_PARALLEL_ABORT || + event == XACT_EVENT_PREPARE) + savepoint_stack = NULL; // TopTransactionContext owns the allocations. +} + +static void savepoint_subxact_callback (SubXactEvent event, SubTransactionId subid, + SubTransactionId parent, void *arg) { + if (event != SUBXACT_EVENT_COMMIT_SUB && event != SUBXACT_EVENT_ABORT_SUB) return; + if (savepoint_stack && savepoint_stack->subid == subid) { + cloudsync_savepoint_frame *frame = savepoint_stack; + savepoint_stack = frame->previous; + pfree(frame); + } +} + +// Only the outermost cloudsync savepoint may replace the caller's snapshot. +// Copy the frame before ending the subtransaction: its callback frees the frame. +static cloudsync_savepoint_frame savepoint_caller (void) { + cloudsync_savepoint_frame caller = {0}; + if (savepoint_stack && savepoint_stack->subid == GetCurrentSubTransactionId()) + caller = *savepoint_stack; + return caller; +} + +static void savepoint_restore_caller (cloudsync_savepoint_frame caller) { + if (!caller.owner) return; + MemoryContextSwitchTo(caller.context); + CurrentResourceOwner = caller.owner; } int database_begin_savepoint (cloudsync_context *data, const char *savepoint_name) { @@ -3098,12 +3114,20 @@ int database_begin_savepoint (cloudsync_context *data, const char *savepoint_nam ResourceOwner oldowner = CurrentResourceOwner; PG_TRY(); { - BeginInternalSubTransaction(NULL); - int level = GetCurrentTransactionNestLevel(); - if (level > 0 && level < CLOUDSYNC_SAVEPOINT_MAX_DEPTH) { - savepoint_owner[level] = oldowner; - savepoint_context[level] = oldcontext; + if (!savepoint_callbacks_registered) { + RegisterXactCallback(savepoint_xact_callback, NULL); + RegisterSubXactCallback(savepoint_subxact_callback, NULL); + savepoint_callbacks_registered = true; } + // Allocate before opening a subtransaction so an allocation error cannot + // leave an untracked rollback boundary behind. + cloudsync_savepoint_frame *frame = MemoryContextAlloc(TopTransactionContext, sizeof(*frame)); + frame->owner = oldowner; + frame->context = oldcontext; + frame->previous = savepoint_stack; + BeginInternalSubTransaction(NULL); + frame->subid = GetCurrentSubTransactionId(); + savepoint_stack = frame; // Keep allocating in the caller's context; the subtransaction's resource owner // stays current so what the savepoint acquires is released with it. MemoryContextSwitchTo(oldcontext); @@ -3128,13 +3152,12 @@ int database_commit_savepoint (cloudsync_context *data, const char *savepoint_na int rc = DBRES_OK; MemoryContext oldcontext = CurrentMemoryContext; - int level = GetCurrentTransactionNestLevel(); + cloudsync_savepoint_frame caller = savepoint_caller(); PG_TRY(); { ReleaseCurrentSubTransaction(); - bool nested = savepoint_is_nested(level); - savepoint_restore_caller(level); - if (nested) CommandCounterIncrement(); + savepoint_restore_caller(caller); + if (caller.previous) CommandCounterIncrement(); else database_refresh_snapshot(); } PG_CATCH(); @@ -3157,13 +3180,12 @@ int database_rollback_savepoint (cloudsync_context *data, const char *savepoint_ int rc = DBRES_OK; MemoryContext oldcontext = CurrentMemoryContext; - int level = GetCurrentTransactionNestLevel(); + cloudsync_savepoint_frame caller = savepoint_caller(); PG_TRY(); { RollbackAndReleaseCurrentSubTransaction(); - bool nested = savepoint_is_nested(level); - savepoint_restore_caller(level); - if (nested) CommandCounterIncrement(); + savepoint_restore_caller(caller); + if (caller.previous) CommandCounterIncrement(); else database_refresh_snapshot(); } PG_CATCH(); diff --git a/test/postgresql/62_deep_savepoints.sql b/test/postgresql/62_deep_savepoints.sql new file mode 100644 index 00000000..144d8107 --- /dev/null +++ b/test/postgresql/62_deep_savepoints.sql @@ -0,0 +1,294 @@ +-- Exercise caller-owned buffers across internal subtransactions at deep nesting. +\set testid '62-deep-savepoints' +\ir helper_test_init.sql +\connect postgres +\ir helper_psql_conn_setup.sql +DROP DATABASE IF EXISTS cloudsync_test_62_source; +DROP DATABASE IF EXISTS cloudsync_test_62_target; +CREATE DATABASE cloudsync_test_62_source; +CREATE DATABASE cloudsync_test_62_target; +\connect cloudsync_test_62_source +\ir helper_psql_conn_setup.sql +CREATE EXTENSION cloudsync; +CREATE TABLE t(id TEXT PRIMARY KEY NOT NULL, value TEXT); +SELECT cloudsync_init('t'); +INSERT INTO t SELECT i::text, repeat('value', 50) FROM generate_series(1, 20) i; +SELECT encode(cloudsync_payload_encode(tbl,pk,col_name,col_value,col_version,db_version,site_id,cl,seq),'hex') AS payload FROM cloudsync_changes \gset +\connect cloudsync_test_62_target +\ir helper_psql_conn_setup.sql +CREATE EXTENSION cloudsync; +CREATE TABLE t(id TEXT PRIMARY KEY NOT NULL, value TEXT); +SELECT cloudsync_init('t'); +CREATE TABLE transport(payload BYTEA); +INSERT INTO transport SELECT decode(:'payload','hex') FROM generate_series(1, 10); + +BEGIN; +SELECT 'SAVEPOINT user_sp' FROM generate_series(1, 1) \gexec +-- A heap scan owns the input buffer; an aggregate consumes every apply result. +SELECT sum(cloudsync_payload_apply(payload)) FROM transport; +SELECT count(*) = 20 AS ok FROM t \gset +\if :ok +\echo [PASS] (:testid) apply at depth 1 +\else +SELECT (:fail::int + 1) AS fail \gset +\echo [FAIL] (:testid) apply at depth 1 +\endif +ROLLBACK TO user_sp; +SELECT count(*) = 0 AS ok FROM t \gset +\if :ok +\echo [PASS] (:testid) rollback at depth 1 +\else +SELECT (:fail::int + 1) AS fail \gset +\echo [FAIL] (:testid) rollback at depth 1 +\endif +SELECT sum(cloudsync_payload_apply(payload)) FROM transport; +COMMIT; +SELECT count(*) = 20 AS ok FROM t \gset +\if :ok +\echo [PASS] (:testid) reuse and commit at depth 1 +\else +SELECT (:fail::int + 1) AS fail \gset +\echo [FAIL] (:testid) reuse and commit at depth 1 +\endif +TRUNCATE t, t_cloudsync; + +BEGIN; +SELECT 'SAVEPOINT user_sp' FROM generate_series(1, 125) \gexec +-- A heap scan owns the input buffer; an aggregate consumes every apply result. +SELECT sum(cloudsync_payload_apply(payload)) FROM transport; +SELECT count(*) = 20 AS ok FROM t \gset +\if :ok +\echo [PASS] (:testid) apply at depth 125 +\else +SELECT (:fail::int + 1) AS fail \gset +\echo [FAIL] (:testid) apply at depth 125 +\endif +ROLLBACK TO user_sp; +SELECT count(*) = 0 AS ok FROM t \gset +\if :ok +\echo [PASS] (:testid) rollback at depth 125 +\else +SELECT (:fail::int + 1) AS fail \gset +\echo [FAIL] (:testid) rollback at depth 125 +\endif +SELECT sum(cloudsync_payload_apply(payload)) FROM transport; +COMMIT; +SELECT count(*) = 20 AS ok FROM t \gset +\if :ok +\echo [PASS] (:testid) reuse and commit at depth 125 +\else +SELECT (:fail::int + 1) AS fail \gset +\echo [FAIL] (:testid) reuse and commit at depth 125 +\endif +TRUNCATE t, t_cloudsync; + +BEGIN; +SELECT 'SAVEPOINT user_sp' FROM generate_series(1, 126) \gexec +-- A heap scan owns the input buffer; an aggregate consumes every apply result. +SELECT sum(cloudsync_payload_apply(payload)) FROM transport; +SELECT count(*) = 20 AS ok FROM t \gset +\if :ok +\echo [PASS] (:testid) apply at depth 126 +\else +SELECT (:fail::int + 1) AS fail \gset +\echo [FAIL] (:testid) apply at depth 126 +\endif +ROLLBACK TO user_sp; +SELECT count(*) = 0 AS ok FROM t \gset +\if :ok +\echo [PASS] (:testid) rollback at depth 126 +\else +SELECT (:fail::int + 1) AS fail \gset +\echo [FAIL] (:testid) rollback at depth 126 +\endif +SELECT sum(cloudsync_payload_apply(payload)) FROM transport; +COMMIT; +SELECT count(*) = 20 AS ok FROM t \gset +\if :ok +\echo [PASS] (:testid) reuse and commit at depth 126 +\else +SELECT (:fail::int + 1) AS fail \gset +\echo [FAIL] (:testid) reuse and commit at depth 126 +\endif +TRUNCATE t, t_cloudsync; + +BEGIN; +SELECT 'SAVEPOINT user_sp' FROM generate_series(1, 127) \gexec +-- A heap scan owns the input buffer; an aggregate consumes every apply result. +SELECT sum(cloudsync_payload_apply(payload)) FROM transport; +SELECT count(*) = 20 AS ok FROM t \gset +\if :ok +\echo [PASS] (:testid) apply at depth 127 +\else +SELECT (:fail::int + 1) AS fail \gset +\echo [FAIL] (:testid) apply at depth 127 +\endif +ROLLBACK TO user_sp; +SELECT count(*) = 0 AS ok FROM t \gset +\if :ok +\echo [PASS] (:testid) rollback at depth 127 +\else +SELECT (:fail::int + 1) AS fail \gset +\echo [FAIL] (:testid) rollback at depth 127 +\endif +SELECT sum(cloudsync_payload_apply(payload)) FROM transport; +COMMIT; +SELECT count(*) = 20 AS ok FROM t \gset +\if :ok +\echo [PASS] (:testid) reuse and commit at depth 127 +\else +SELECT (:fail::int + 1) AS fail \gset +\echo [FAIL] (:testid) reuse and commit at depth 127 +\endif +TRUNCATE t, t_cloudsync; + +BEGIN; +SELECT 'SAVEPOINT user_sp' FROM generate_series(1, 128) \gexec +-- A heap scan owns the input buffer; an aggregate consumes every apply result. +SELECT sum(cloudsync_payload_apply(payload)) FROM transport; +SELECT count(*) = 20 AS ok FROM t \gset +\if :ok +\echo [PASS] (:testid) apply at depth 128 +\else +SELECT (:fail::int + 1) AS fail \gset +\echo [FAIL] (:testid) apply at depth 128 +\endif +ROLLBACK TO user_sp; +SELECT count(*) = 0 AS ok FROM t \gset +\if :ok +\echo [PASS] (:testid) rollback at depth 128 +\else +SELECT (:fail::int + 1) AS fail \gset +\echo [FAIL] (:testid) rollback at depth 128 +\endif +SELECT sum(cloudsync_payload_apply(payload)) FROM transport; +COMMIT; +SELECT count(*) = 20 AS ok FROM t \gset +\if :ok +\echo [PASS] (:testid) reuse and commit at depth 128 +\else +SELECT (:fail::int + 1) AS fail \gset +\echo [FAIL] (:testid) reuse and commit at depth 128 +\endif +TRUNCATE t, t_cloudsync; + +BEGIN; +SELECT 'SAVEPOINT user_sp' FROM generate_series(1, 256) \gexec +-- A heap scan owns the input buffer; an aggregate consumes every apply result. +SELECT sum(cloudsync_payload_apply(payload)) FROM transport; +SELECT count(*) = 20 AS ok FROM t \gset +\if :ok +\echo [PASS] (:testid) apply at depth 256 +\else +SELECT (:fail::int + 1) AS fail \gset +\echo [FAIL] (:testid) apply at depth 256 +\endif +ROLLBACK TO user_sp; +SELECT count(*) = 0 AS ok FROM t \gset +\if :ok +\echo [PASS] (:testid) rollback at depth 256 +\else +SELECT (:fail::int + 1) AS fail \gset +\echo [FAIL] (:testid) rollback at depth 256 +\endif +SELECT sum(cloudsync_payload_apply(payload)) FROM transport; +COMMIT; +SELECT count(*) = 20 AS ok FROM t \gset +\if :ok +\echo [PASS] (:testid) reuse and commit at depth 256 +\else +SELECT (:fail::int + 1) AS fail \gset +\echo [FAIL] (:testid) reuse and commit at depth 256 +\endif +TRUNCATE t, t_cloudsync; + +BEGIN; +SELECT 'SAVEPOINT user_sp' FROM generate_series(1, 1024) \gexec +-- A heap scan owns the input buffer; an aggregate consumes every apply result. +SELECT sum(cloudsync_payload_apply(payload)) FROM transport; +SELECT count(*) = 20 AS ok FROM t \gset +\if :ok +\echo [PASS] (:testid) apply at depth 1024 +\else +SELECT (:fail::int + 1) AS fail \gset +\echo [FAIL] (:testid) apply at depth 1024 +\endif +ROLLBACK TO user_sp; +SELECT count(*) = 0 AS ok FROM t \gset +\if :ok +\echo [PASS] (:testid) rollback at depth 1024 +\else +SELECT (:fail::int + 1) AS fail \gset +\echo [FAIL] (:testid) rollback at depth 1024 +\endif +SELECT sum(cloudsync_payload_apply(payload)) FROM transport; +COMMIT; +SELECT count(*) = 20 AS ok FROM t \gset +\if :ok +\echo [PASS] (:testid) reuse and commit at depth 1024 +\else +SELECT (:fail::int + 1) AS fail \gset +\echo [FAIL] (:testid) reuse and commit at depth 1024 +\endif +TRUNCATE t, t_cloudsync; + +BEGIN; +SELECT 'SAVEPOINT user_sp' FROM generate_series(1, 2048) \gexec +-- A heap scan owns the input buffer; an aggregate consumes every apply result. +SELECT sum(cloudsync_payload_apply(payload)) FROM transport; +SELECT count(*) = 20 AS ok FROM t \gset +\if :ok +\echo [PASS] (:testid) apply at depth 2048 +\else +SELECT (:fail::int + 1) AS fail \gset +\echo [FAIL] (:testid) apply at depth 2048 +\endif +ROLLBACK TO user_sp; +SELECT count(*) = 0 AS ok FROM t \gset +\if :ok +\echo [PASS] (:testid) rollback at depth 2048 +\else +SELECT (:fail::int + 1) AS fail \gset +\echo [FAIL] (:testid) rollback at depth 2048 +\endif +SELECT sum(cloudsync_payload_apply(payload)) FROM transport; +COMMIT; +SELECT count(*) = 20 AS ok FROM t \gset +\if :ok +\echo [PASS] (:testid) reuse and commit at depth 2048 +\else +SELECT (:fail::int + 1) AS fail \gset +\echo [FAIL] (:testid) reuse and commit at depth 2048 +\endif +TRUNCATE t, t_cloudsync; + +-- Catch errors outside the internal wrappers, then reuse the same backend. +CREATE FUNCTION reject_row() RETURNS trigger LANGUAGE plpgsql AS $$ +BEGIN RAISE EXCEPTION 'deep apply rejected'; END $$; +CREATE TRIGGER reject_row BEFORE INSERT ON t FOR EACH ROW EXECUTE FUNCTION reject_row(); +BEGIN; +SELECT 'SAVEPOINT user_sp' FROM generate_series(1, 256) \gexec +DO $$ +BEGIN + FOR i IN 1..100 LOOP + BEGIN + PERFORM cloudsync_payload_apply(payload) FROM transport; + RAISE EXCEPTION 'expected trigger rejection'; + EXCEPTION WHEN OTHERS THEN + IF SQLERRM NOT LIKE '%deep apply rejected%' THEN RAISE; END IF; + END; + END LOOP; +END $$; +DROP TRIGGER reject_row ON t; +SELECT sum(cloudsync_payload_apply(payload)) FROM transport; +COMMIT; +SELECT count(*) = 20 AS ok FROM t \gset +\if :ok +\echo [PASS] (:testid) 100 caught errors followed by successful apply +\else +SELECT (:fail::int + 1) AS fail \gset +\echo [FAIL] (:testid) error recovery +\endif +\connect postgres +DROP DATABASE cloudsync_test_62_source; +DROP DATABASE cloudsync_test_62_target; diff --git a/test/postgresql/full_test.sql b/test/postgresql/full_test.sql index 8916c8df..39ecc0fa 100644 --- a/test/postgresql/full_test.sql +++ b/test/postgresql/full_test.sql @@ -70,6 +70,7 @@ \ir 60_fragment_concurrency.sql \ir 61_fragment_cleanup_backlog.sql \ir 62_deferred_fk_caller_commit.sql +\ir 62_deep_savepoints.sql -- 'Test summary' \echo '\nTest summary:' From 5ac466a910948aa78a77d4d3433792082ac9aae0 Mon Sep 17 00:00:00 2001 From: Andrea Donetti Date: Tue, 22 Sep 2026 08:53:22 -0600 Subject: [PATCH 2/4] test(postgres): renumber the deep savepoint test to 63 Test 62 is now 62_deferred_fk_caller_commit.sql, merged with #65. Co-Authored-By: Claude Opus 5 (1M context) --- docs/internal/deep-savepoints.md | 2 +- ...p_savepoints.sql => 63_deep_savepoints.sql} | 18 +++++++++--------- test/postgresql/full_test.sql | 2 +- 3 files changed, 11 insertions(+), 11 deletions(-) rename test/postgresql/{62_deep_savepoints.sql => 63_deep_savepoints.sql} (96%) diff --git a/docs/internal/deep-savepoints.md b/docs/internal/deep-savepoints.md index 701a8e9b..24237bf4 100644 --- a/docs/internal/deep-savepoints.md +++ b/docs/internal/deep-savepoints.md @@ -8,6 +8,6 @@ This removes Cloudsync's fixed depth cap; PostgreSQL's own resource limits still ## Validation -`test/postgresql/62_deep_savepoints.sql` is included in `full_test.sql`. It exercises depths 1, 125, 126, 127, 128, 256, 1024 and 2048, reads payloads from a heap table, checks rollback, reapplies in the same transaction and commits. It also catches 100 trigger failures at depth 256 and verifies a subsequent successful apply in the same backend. +`test/postgresql/63_deep_savepoints.sql` is included in `full_test.sql`. It exercises depths 1, 125, 126, 127, 128, 256, 1024 and 2048, reads payloads from a heap table, checks rollback, reapplies in the same transaction and commits. It also catches 100 trigger failures at depth 256 and verifies a subsequent successful apply in the same backend. The full suite passes all 546 checks on PostgreSQL 15.19, 17.11 and 18.6. Restoring the old implementation reproduces the buffer-owner error at depth 126. These are local database tests, not tests of a deployed cloud server. diff --git a/test/postgresql/62_deep_savepoints.sql b/test/postgresql/63_deep_savepoints.sql similarity index 96% rename from test/postgresql/62_deep_savepoints.sql rename to test/postgresql/63_deep_savepoints.sql index 144d8107..b6f69859 100644 --- a/test/postgresql/62_deep_savepoints.sql +++ b/test/postgresql/63_deep_savepoints.sql @@ -1,20 +1,20 @@ -- Exercise caller-owned buffers across internal subtransactions at deep nesting. -\set testid '62-deep-savepoints' +\set testid '63-deep-savepoints' \ir helper_test_init.sql \connect postgres \ir helper_psql_conn_setup.sql -DROP DATABASE IF EXISTS cloudsync_test_62_source; -DROP DATABASE IF EXISTS cloudsync_test_62_target; -CREATE DATABASE cloudsync_test_62_source; -CREATE DATABASE cloudsync_test_62_target; -\connect cloudsync_test_62_source +DROP DATABASE IF EXISTS cloudsync_test_63_source; +DROP DATABASE IF EXISTS cloudsync_test_63_target; +CREATE DATABASE cloudsync_test_63_source; +CREATE DATABASE cloudsync_test_63_target; +\connect cloudsync_test_63_source \ir helper_psql_conn_setup.sql CREATE EXTENSION cloudsync; CREATE TABLE t(id TEXT PRIMARY KEY NOT NULL, value TEXT); SELECT cloudsync_init('t'); INSERT INTO t SELECT i::text, repeat('value', 50) FROM generate_series(1, 20) i; SELECT encode(cloudsync_payload_encode(tbl,pk,col_name,col_value,col_version,db_version,site_id,cl,seq),'hex') AS payload FROM cloudsync_changes \gset -\connect cloudsync_test_62_target +\connect cloudsync_test_63_target \ir helper_psql_conn_setup.sql CREATE EXTENSION cloudsync; CREATE TABLE t(id TEXT PRIMARY KEY NOT NULL, value TEXT); @@ -290,5 +290,5 @@ SELECT (:fail::int + 1) AS fail \gset \echo [FAIL] (:testid) error recovery \endif \connect postgres -DROP DATABASE cloudsync_test_62_source; -DROP DATABASE cloudsync_test_62_target; +DROP DATABASE cloudsync_test_63_source; +DROP DATABASE cloudsync_test_63_target; diff --git a/test/postgresql/full_test.sql b/test/postgresql/full_test.sql index 39ecc0fa..66c355be 100644 --- a/test/postgresql/full_test.sql +++ b/test/postgresql/full_test.sql @@ -70,7 +70,7 @@ \ir 60_fragment_concurrency.sql \ir 61_fragment_cleanup_backlog.sql \ir 62_deferred_fk_caller_commit.sql -\ir 62_deep_savepoints.sql +\ir 63_deep_savepoints.sql -- 'Test summary' \echo '\nTest summary:' From 0b37f4583902e783ec69d6e9e5318519f6e582c6 Mon Sep 17 00:00:00 2001 From: Andrea Donetti Date: Tue, 22 Sep 2026 09:28:43 -0600 Subject: [PATCH 3/4] docs: changelog entry for the PostgreSQL savepoint depth fix Co-Authored-By: Claude Opus 5 (1M context) --- CHANGELOG.md | 1 + 1 file changed, 1 insertion(+) diff --git a/CHANGELOG.md b/CHANGELOG.md index 42e8fe5c..703e3ab0 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -9,6 +9,7 @@ The format is based on [Keep a Changelog](https://keepachangelog.com/en/1.0.0/). ### Fixed - **SQLite: a payload whose commit fails no longer leaves its transaction open.** When `cloudsync_payload_apply` started the transaction itself and the commit then failed — a deferred foreign key violated at commit, or `SQLITE_BUSY` because a reader held the database — the transaction stayed open: the uncommitted rows remained visible on the connection and the next `BEGIN` failed. The failed transaction is now rolled back and the original error is returned. Changes from earlier source versions that were already committed are kept, the receive checkpoint does not move, and the rolled-back rows are no longer counted as applied, so delivering the payload again applies it. A transaction or savepoint opened by the caller is still left to the caller. +- **PostgreSQL: applying a payload read from a table works at any savepoint depth.** Inside 126 or more savepoints, `SELECT cloudsync_payload_apply(payload) FROM some_table` still failed with `buffer pin ... is not owned by resource owner SubTransaction` (and a caught error at that depth could abort an assertion-enabled server): cloudsync recorded the caller's resource owner and memory context for at most 128 nesting levels, counting its own internal savepoints, and silently stopped restoring them beyond that. The fixed limit is gone; only PostgreSQL's own resource limits apply. ## [1.1.4] - 2026-09-21 From c46822a42675e29ec045d90e84a9df43d3bf9ea2 Mon Sep 17 00:00:00 2001 From: Andrea Donetti Date: Tue, 22 Sep 2026 09:28:43 -0600 Subject: [PATCH 4/4] test(postgres): report a failed deep-savepoint apply instead of aborting the suite Under ON_ERROR_STOP=on a failing apply stopped full_test.sql with no [FAIL] line. Each depth now records the SQLSTATE, recovers at the user savepoint and reports [FAIL] per check, so the rest of the suite runs. Co-Authored-By: Claude Opus 5 (1M context) --- test/postgresql/63_deep_savepoints.sql | 118 +++++++++++++++++++------ 1 file changed, 93 insertions(+), 25 deletions(-) diff --git a/test/postgresql/63_deep_savepoints.sql b/test/postgresql/63_deep_savepoints.sql index b6f69859..4685ccbc 100644 --- a/test/postgresql/63_deep_savepoints.sql +++ b/test/postgresql/63_deep_savepoints.sql @@ -25,15 +25,20 @@ INSERT INTO transport SELECT decode(:'payload','hex') FROM generate_series(1, 10 BEGIN; SELECT 'SAVEPOINT user_sp' FROM generate_series(1, 1) \gexec -- A heap scan owns the input buffer; an aggregate consumes every apply result. +-- A failed apply aborts the transaction: report it and recover at user_sp. +\set ok false +\set ON_ERROR_STOP off SELECT sum(cloudsync_payload_apply(payload)) FROM transport; +\set apply_state :SQLSTATE SELECT count(*) = 20 AS ok FROM t \gset +ROLLBACK TO user_sp; +\set ON_ERROR_STOP on \if :ok \echo [PASS] (:testid) apply at depth 1 \else SELECT (:fail::int + 1) AS fail \gset -\echo [FAIL] (:testid) apply at depth 1 +\echo [FAIL] (:testid) apply at depth 1: SQLSTATE :apply_state \endif -ROLLBACK TO user_sp; SELECT count(*) = 0 AS ok FROM t \gset \if :ok \echo [PASS] (:testid) rollback at depth 1 @@ -41,29 +46,37 @@ SELECT count(*) = 0 AS ok FROM t \gset SELECT (:fail::int + 1) AS fail \gset \echo [FAIL] (:testid) rollback at depth 1 \endif +\set ON_ERROR_STOP off SELECT sum(cloudsync_payload_apply(payload)) FROM transport; +\set apply_state :SQLSTATE +\set ON_ERROR_STOP on COMMIT; SELECT count(*) = 20 AS ok FROM t \gset \if :ok \echo [PASS] (:testid) reuse and commit at depth 1 \else SELECT (:fail::int + 1) AS fail \gset -\echo [FAIL] (:testid) reuse and commit at depth 1 +\echo [FAIL] (:testid) reuse and commit at depth 1: SQLSTATE :apply_state \endif TRUNCATE t, t_cloudsync; BEGIN; SELECT 'SAVEPOINT user_sp' FROM generate_series(1, 125) \gexec -- A heap scan owns the input buffer; an aggregate consumes every apply result. +-- A failed apply aborts the transaction: report it and recover at user_sp. +\set ok false +\set ON_ERROR_STOP off SELECT sum(cloudsync_payload_apply(payload)) FROM transport; +\set apply_state :SQLSTATE SELECT count(*) = 20 AS ok FROM t \gset +ROLLBACK TO user_sp; +\set ON_ERROR_STOP on \if :ok \echo [PASS] (:testid) apply at depth 125 \else SELECT (:fail::int + 1) AS fail \gset -\echo [FAIL] (:testid) apply at depth 125 +\echo [FAIL] (:testid) apply at depth 125: SQLSTATE :apply_state \endif -ROLLBACK TO user_sp; SELECT count(*) = 0 AS ok FROM t \gset \if :ok \echo [PASS] (:testid) rollback at depth 125 @@ -71,29 +84,37 @@ SELECT count(*) = 0 AS ok FROM t \gset SELECT (:fail::int + 1) AS fail \gset \echo [FAIL] (:testid) rollback at depth 125 \endif +\set ON_ERROR_STOP off SELECT sum(cloudsync_payload_apply(payload)) FROM transport; +\set apply_state :SQLSTATE +\set ON_ERROR_STOP on COMMIT; SELECT count(*) = 20 AS ok FROM t \gset \if :ok \echo [PASS] (:testid) reuse and commit at depth 125 \else SELECT (:fail::int + 1) AS fail \gset -\echo [FAIL] (:testid) reuse and commit at depth 125 +\echo [FAIL] (:testid) reuse and commit at depth 125: SQLSTATE :apply_state \endif TRUNCATE t, t_cloudsync; BEGIN; SELECT 'SAVEPOINT user_sp' FROM generate_series(1, 126) \gexec -- A heap scan owns the input buffer; an aggregate consumes every apply result. +-- A failed apply aborts the transaction: report it and recover at user_sp. +\set ok false +\set ON_ERROR_STOP off SELECT sum(cloudsync_payload_apply(payload)) FROM transport; +\set apply_state :SQLSTATE SELECT count(*) = 20 AS ok FROM t \gset +ROLLBACK TO user_sp; +\set ON_ERROR_STOP on \if :ok \echo [PASS] (:testid) apply at depth 126 \else SELECT (:fail::int + 1) AS fail \gset -\echo [FAIL] (:testid) apply at depth 126 +\echo [FAIL] (:testid) apply at depth 126: SQLSTATE :apply_state \endif -ROLLBACK TO user_sp; SELECT count(*) = 0 AS ok FROM t \gset \if :ok \echo [PASS] (:testid) rollback at depth 126 @@ -101,29 +122,37 @@ SELECT count(*) = 0 AS ok FROM t \gset SELECT (:fail::int + 1) AS fail \gset \echo [FAIL] (:testid) rollback at depth 126 \endif +\set ON_ERROR_STOP off SELECT sum(cloudsync_payload_apply(payload)) FROM transport; +\set apply_state :SQLSTATE +\set ON_ERROR_STOP on COMMIT; SELECT count(*) = 20 AS ok FROM t \gset \if :ok \echo [PASS] (:testid) reuse and commit at depth 126 \else SELECT (:fail::int + 1) AS fail \gset -\echo [FAIL] (:testid) reuse and commit at depth 126 +\echo [FAIL] (:testid) reuse and commit at depth 126: SQLSTATE :apply_state \endif TRUNCATE t, t_cloudsync; BEGIN; SELECT 'SAVEPOINT user_sp' FROM generate_series(1, 127) \gexec -- A heap scan owns the input buffer; an aggregate consumes every apply result. +-- A failed apply aborts the transaction: report it and recover at user_sp. +\set ok false +\set ON_ERROR_STOP off SELECT sum(cloudsync_payload_apply(payload)) FROM transport; +\set apply_state :SQLSTATE SELECT count(*) = 20 AS ok FROM t \gset +ROLLBACK TO user_sp; +\set ON_ERROR_STOP on \if :ok \echo [PASS] (:testid) apply at depth 127 \else SELECT (:fail::int + 1) AS fail \gset -\echo [FAIL] (:testid) apply at depth 127 +\echo [FAIL] (:testid) apply at depth 127: SQLSTATE :apply_state \endif -ROLLBACK TO user_sp; SELECT count(*) = 0 AS ok FROM t \gset \if :ok \echo [PASS] (:testid) rollback at depth 127 @@ -131,29 +160,37 @@ SELECT count(*) = 0 AS ok FROM t \gset SELECT (:fail::int + 1) AS fail \gset \echo [FAIL] (:testid) rollback at depth 127 \endif +\set ON_ERROR_STOP off SELECT sum(cloudsync_payload_apply(payload)) FROM transport; +\set apply_state :SQLSTATE +\set ON_ERROR_STOP on COMMIT; SELECT count(*) = 20 AS ok FROM t \gset \if :ok \echo [PASS] (:testid) reuse and commit at depth 127 \else SELECT (:fail::int + 1) AS fail \gset -\echo [FAIL] (:testid) reuse and commit at depth 127 +\echo [FAIL] (:testid) reuse and commit at depth 127: SQLSTATE :apply_state \endif TRUNCATE t, t_cloudsync; BEGIN; SELECT 'SAVEPOINT user_sp' FROM generate_series(1, 128) \gexec -- A heap scan owns the input buffer; an aggregate consumes every apply result. +-- A failed apply aborts the transaction: report it and recover at user_sp. +\set ok false +\set ON_ERROR_STOP off SELECT sum(cloudsync_payload_apply(payload)) FROM transport; +\set apply_state :SQLSTATE SELECT count(*) = 20 AS ok FROM t \gset +ROLLBACK TO user_sp; +\set ON_ERROR_STOP on \if :ok \echo [PASS] (:testid) apply at depth 128 \else SELECT (:fail::int + 1) AS fail \gset -\echo [FAIL] (:testid) apply at depth 128 +\echo [FAIL] (:testid) apply at depth 128: SQLSTATE :apply_state \endif -ROLLBACK TO user_sp; SELECT count(*) = 0 AS ok FROM t \gset \if :ok \echo [PASS] (:testid) rollback at depth 128 @@ -161,29 +198,37 @@ SELECT count(*) = 0 AS ok FROM t \gset SELECT (:fail::int + 1) AS fail \gset \echo [FAIL] (:testid) rollback at depth 128 \endif +\set ON_ERROR_STOP off SELECT sum(cloudsync_payload_apply(payload)) FROM transport; +\set apply_state :SQLSTATE +\set ON_ERROR_STOP on COMMIT; SELECT count(*) = 20 AS ok FROM t \gset \if :ok \echo [PASS] (:testid) reuse and commit at depth 128 \else SELECT (:fail::int + 1) AS fail \gset -\echo [FAIL] (:testid) reuse and commit at depth 128 +\echo [FAIL] (:testid) reuse and commit at depth 128: SQLSTATE :apply_state \endif TRUNCATE t, t_cloudsync; BEGIN; SELECT 'SAVEPOINT user_sp' FROM generate_series(1, 256) \gexec -- A heap scan owns the input buffer; an aggregate consumes every apply result. +-- A failed apply aborts the transaction: report it and recover at user_sp. +\set ok false +\set ON_ERROR_STOP off SELECT sum(cloudsync_payload_apply(payload)) FROM transport; +\set apply_state :SQLSTATE SELECT count(*) = 20 AS ok FROM t \gset +ROLLBACK TO user_sp; +\set ON_ERROR_STOP on \if :ok \echo [PASS] (:testid) apply at depth 256 \else SELECT (:fail::int + 1) AS fail \gset -\echo [FAIL] (:testid) apply at depth 256 +\echo [FAIL] (:testid) apply at depth 256: SQLSTATE :apply_state \endif -ROLLBACK TO user_sp; SELECT count(*) = 0 AS ok FROM t \gset \if :ok \echo [PASS] (:testid) rollback at depth 256 @@ -191,29 +236,37 @@ SELECT count(*) = 0 AS ok FROM t \gset SELECT (:fail::int + 1) AS fail \gset \echo [FAIL] (:testid) rollback at depth 256 \endif +\set ON_ERROR_STOP off SELECT sum(cloudsync_payload_apply(payload)) FROM transport; +\set apply_state :SQLSTATE +\set ON_ERROR_STOP on COMMIT; SELECT count(*) = 20 AS ok FROM t \gset \if :ok \echo [PASS] (:testid) reuse and commit at depth 256 \else SELECT (:fail::int + 1) AS fail \gset -\echo [FAIL] (:testid) reuse and commit at depth 256 +\echo [FAIL] (:testid) reuse and commit at depth 256: SQLSTATE :apply_state \endif TRUNCATE t, t_cloudsync; BEGIN; SELECT 'SAVEPOINT user_sp' FROM generate_series(1, 1024) \gexec -- A heap scan owns the input buffer; an aggregate consumes every apply result. +-- A failed apply aborts the transaction: report it and recover at user_sp. +\set ok false +\set ON_ERROR_STOP off SELECT sum(cloudsync_payload_apply(payload)) FROM transport; +\set apply_state :SQLSTATE SELECT count(*) = 20 AS ok FROM t \gset +ROLLBACK TO user_sp; +\set ON_ERROR_STOP on \if :ok \echo [PASS] (:testid) apply at depth 1024 \else SELECT (:fail::int + 1) AS fail \gset -\echo [FAIL] (:testid) apply at depth 1024 +\echo [FAIL] (:testid) apply at depth 1024: SQLSTATE :apply_state \endif -ROLLBACK TO user_sp; SELECT count(*) = 0 AS ok FROM t \gset \if :ok \echo [PASS] (:testid) rollback at depth 1024 @@ -221,29 +274,37 @@ SELECT count(*) = 0 AS ok FROM t \gset SELECT (:fail::int + 1) AS fail \gset \echo [FAIL] (:testid) rollback at depth 1024 \endif +\set ON_ERROR_STOP off SELECT sum(cloudsync_payload_apply(payload)) FROM transport; +\set apply_state :SQLSTATE +\set ON_ERROR_STOP on COMMIT; SELECT count(*) = 20 AS ok FROM t \gset \if :ok \echo [PASS] (:testid) reuse and commit at depth 1024 \else SELECT (:fail::int + 1) AS fail \gset -\echo [FAIL] (:testid) reuse and commit at depth 1024 +\echo [FAIL] (:testid) reuse and commit at depth 1024: SQLSTATE :apply_state \endif TRUNCATE t, t_cloudsync; BEGIN; SELECT 'SAVEPOINT user_sp' FROM generate_series(1, 2048) \gexec -- A heap scan owns the input buffer; an aggregate consumes every apply result. +-- A failed apply aborts the transaction: report it and recover at user_sp. +\set ok false +\set ON_ERROR_STOP off SELECT sum(cloudsync_payload_apply(payload)) FROM transport; +\set apply_state :SQLSTATE SELECT count(*) = 20 AS ok FROM t \gset +ROLLBACK TO user_sp; +\set ON_ERROR_STOP on \if :ok \echo [PASS] (:testid) apply at depth 2048 \else SELECT (:fail::int + 1) AS fail \gset -\echo [FAIL] (:testid) apply at depth 2048 +\echo [FAIL] (:testid) apply at depth 2048: SQLSTATE :apply_state \endif -ROLLBACK TO user_sp; SELECT count(*) = 0 AS ok FROM t \gset \if :ok \echo [PASS] (:testid) rollback at depth 2048 @@ -251,14 +312,17 @@ SELECT count(*) = 0 AS ok FROM t \gset SELECT (:fail::int + 1) AS fail \gset \echo [FAIL] (:testid) rollback at depth 2048 \endif +\set ON_ERROR_STOP off SELECT sum(cloudsync_payload_apply(payload)) FROM transport; +\set apply_state :SQLSTATE +\set ON_ERROR_STOP on COMMIT; SELECT count(*) = 20 AS ok FROM t \gset \if :ok \echo [PASS] (:testid) reuse and commit at depth 2048 \else SELECT (:fail::int + 1) AS fail \gset -\echo [FAIL] (:testid) reuse and commit at depth 2048 +\echo [FAIL] (:testid) reuse and commit at depth 2048: SQLSTATE :apply_state \endif TRUNCATE t, t_cloudsync; @@ -268,6 +332,7 @@ BEGIN RAISE EXCEPTION 'deep apply rejected'; END $$; CREATE TRIGGER reject_row BEFORE INSERT ON t FOR EACH ROW EXECUTE FUNCTION reject_row(); BEGIN; SELECT 'SAVEPOINT user_sp' FROM generate_series(1, 256) \gexec +\set ON_ERROR_STOP off DO $$ BEGIN FOR i IN 1..100 LOOP @@ -279,15 +344,18 @@ BEGIN END; END LOOP; END $$; +\set recovery_state :SQLSTATE DROP TRIGGER reject_row ON t; SELECT sum(cloudsync_payload_apply(payload)) FROM transport; +\set ON_ERROR_STOP on COMMIT; SELECT count(*) = 20 AS ok FROM t \gset \if :ok \echo [PASS] (:testid) 100 caught errors followed by successful apply \else SELECT (:fail::int + 1) AS fail \gset -\echo [FAIL] (:testid) error recovery +\echo [FAIL] (:testid) error recovery: caught-errors SQLSTATE :recovery_state +DROP TRIGGER IF EXISTS reject_row ON t; \endif \connect postgres DROP DATABASE cloudsync_test_63_source;