From 29379812f4bc50cb78fe752b708c11be8aa5b363 Mon Sep 17 00:00:00 2001 From: Andrea Donetti Date: Thu, 24 Sep 2026 16:14:29 -0600 Subject: [PATCH 1/2] fix(postgres): give every transaction its own db_version and restart seq With more than one synced table, the query that reloads the cached db_version returned one row per meta table and only the first was read, so consecutive transactions writing another table shared one db_version. Wrap it in MAX() and include pre_alter_dbversion, as on SQLite. Close each transaction from a PostgreSQL transaction callback, as the SQLite commit/rollback hooks do: seq restarts at 0 and the next transaction takes a new db_version even while another session pins xmin. Inside one transaction, cloudsync_payload_apply now closes the local db_version at every source db_version boundary, so remote versions no longer collapse into duplicate (db_version, seq) pairs. Fixes #69 Co-Authored-By: Claude Opus 5.5 (1M context) --- CHANGELOG.md | 2 + src/cloudsync.c | 15 +++ src/cloudsync.h | 1 + src/postgresql/cloudsync_postgresql.c | 26 ++++++ src/postgresql/sql_postgresql.c | 6 +- .../66_db_version_per_transaction.sql | 91 +++++++++++++++++++ test/postgresql/full_test.sql | 1 + test/review_regressions.c | 32 +++++++ 8 files changed, 172 insertions(+), 2 deletions(-) create mode 100644 test/postgresql/66_db_version_per_transaction.sql diff --git a/CHANGELOG.md b/CHANGELOG.md index e4c1938b..75a23fcf 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -8,6 +8,8 @@ The format is based on [Keep a Changelog](https://keepachangelog.com/en/1.0.0/). ### Fixed +- **PostgreSQL: separate transactions no longer share a db_version.** With more than one synced table, the current db_version was read from a single table's metadata instead of the maximum across all of them, so consecutive transactions writing to another table all took the same db_version. Holding a transaction open in another session had the same effect on every table. A change that lands on a db_version already served is skipped by the next download, whose window starts after that version, so peers could miss it. Each transaction now takes its own db_version, as on SQLite, and `seq` restarts at 0 in each one instead of growing for the life of the connection, where a long-lived connection could eventually overflow it. +- **PostgreSQL: applying a payload keeps each source db_version apart.** A whole `cloudsync_payload_apply` runs in one transaction, so every db_version in the payload was stored under the same local db_version while each change kept its source `seq`, producing duplicate `(db_version, seq)` pairs that the chunked download relies on being unique. Each source db_version now takes its own local db_version, as on SQLite. - **SQLite: paging a chunked download no longer gets slower with every chunk.** The positional cursor on `cloudsync_payload_chunks` was meant to seek straight to where the previous call stopped, but it stated its resume point only inside `(db_version > ? OR (db_version = ? AND seq >= ?))`, whose two arms carry distinct parameters. SQLite does not derive a range from that, so the scan over `cloudsync_changes` ran with an upper bound only and re-read the window from the beginning on every call, discarding rows until it reached the resume point — making a full drain quadratic in the number of chunks, and long enough on a large tenant to hit a server-side deadline and never complete. An explicit `db_version >= ?` is now stated alongside the disjunction; it selects exactly the same rows and lets the scan seek. Locally, draining a 188-chunk window went from 6110 ms to 210 ms, and the per-chunk cost no longer depends on how large the window is. PostgreSQL was never affected. - **A row rewritten in one statement no longer keeps the old blocks of a shorter value.** Writing a whole row rewrites its block column from the first position, so any block past the end of the new value stayed stored. On SQLite `INSERT OR REPLACE` skips the old row's delete trigger unless `recursive_triggers` is on, so those leftovers kept their metadata and were delivered as content: replacing `AAA\nBBB\nCCC` with `ZZZ` left `ZZZ\nBBB\nCCC` on the peers and on a later local read. On PostgreSQL they carried no metadata, so peers were unaffected, but they stayed in the blocks table for the life of the row. Blocks the new value does not cover are now retired with it. - **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. diff --git a/src/cloudsync.c b/src/cloudsync.c index 41d8dae6..0407b036 100644 --- a/src/cloudsync.c +++ b/src/cloudsync.c @@ -2660,6 +2660,16 @@ void cloudsync_rollback_hook (void *ctx) { data->seq = 0; } +// Runs the commit or rollback hook where the database does not: at the end of every +// PostgreSQL transaction, and between source db_versions of a payload applied inside one. +// It acts only when a db_version was taken, so the cached one survives read-only +// transactions. +void cloudsync_transaction_end (cloudsync_context *data, bool committed) { + if (!data || data->pending_db_version == CLOUDSYNC_VALUE_NOTSET) return; + if (committed) cloudsync_commit_hook(data); + else cloudsync_rollback_hook(data); +} + int cloudsync_begin_alter (cloudsync_context *data, const char *table_name) { // init cloudsync_settings if (cloudsync_context_init(data) == NULL) { @@ -4575,6 +4585,11 @@ int cloudsync_payload_apply (cloudsync_context *data, const char *payload, int b in_savepoint = false; } + // Inside a caller's transaction (always on PostgreSQL) no commit hook runs between + // groups: close the group's db_version here, or every source db_version would get + // the same local one and their seqs would collide. A no-op after the RELEASE above. + if (db_version_changed) cloudsync_transaction_end(data, true); + if (!in_savepoint && db_version_changed && !database_in_transaction(data)) { rc = database_begin_savepoint(data, "cloudsync_payload_apply"); if (rc != DBRES_OK) { diff --git a/src/cloudsync.h b/src/cloudsync.h index fdd5a75d..ffafc4b6 100644 --- a/src/cloudsync.h +++ b/src/cloudsync.h @@ -124,6 +124,7 @@ void cloudsync_apply_stats_reset (cloudsync_context *data); int cloudsync_apply_rows_count (cloudsync_context *data); int cloudsync_commit_hook (void *ctx); void cloudsync_rollback_hook (void *ctx); +void cloudsync_transaction_end (cloudsync_context *data, bool committed); void cloudsync_set_schema (cloudsync_context *data, const char *schema); const char *cloudsync_schema (cloudsync_context *data); const char *cloudsync_table_schema (cloudsync_context *data, const char *table_name); diff --git a/src/postgresql/cloudsync_postgresql.c b/src/postgresql/cloudsync_postgresql.c index 7834561b..9f7145a6 100644 --- a/src/postgresql/cloudsync_postgresql.c +++ b/src/postgresql/cloudsync_postgresql.c @@ -126,6 +126,29 @@ static cloudsync_context *get_cloudsync_context(void) { return pg_cloudsync_context; } +// PostgreSQL has no commit or rollback hook like SQLite's: the transaction callback closes +// the transaction's db_version and restarts seq as the SQLite hooks do, so the next +// transaction takes a new db_version even while another session pins the snapshot xmin. +// It only assigns fields, so it is safe during commit. +static void cloudsync_xact_callback (XactEvent event, void *arg) { + UNUSED_PARAMETER(arg); + cloudsync_context *data = pg_cloudsync_context; + if (!data) return; + switch (event) { + case XACT_EVENT_COMMIT: + case XACT_EVENT_PARALLEL_COMMIT: + case XACT_EVENT_PREPARE: + cloudsync_transaction_end(data, true); + break; + case XACT_EVENT_ABORT: + case XACT_EVENT_PARALLEL_ABORT: + cloudsync_transaction_end(data, false); + break; + default: + break; + } +} + // MARK: - Extension Entry Points - void _PG_init (void) { @@ -138,11 +161,14 @@ void _PG_init (void) { // Set fractional-indexing allocator to use cloudsync memory block_init_allocator(); + + RegisterXactCallback(cloudsync_xact_callback, NULL); } void _PG_fini (void) { // Extension cleanup elog(DEBUG1, "CloudSync extension unloading"); + UnregisterXactCallback(cloudsync_xact_callback, NULL); // Free global context if it exists if (pg_cloudsync_context) { diff --git a/src/postgresql/sql_postgresql.c b/src/postgresql/sql_postgresql.c index c6a630be..2d73efdc 100644 --- a/src/postgresql/sql_postgresql.c +++ b/src/postgresql/sql_postgresql.c @@ -97,10 +97,12 @@ const char * const SQL_DBVERSION_BUILD_QUERY = "), " "query_parts AS (" "SELECT tbl_name, " - "format('SELECT COALESCE(MAX(db_version), 0) FROM %s', tbl_name) as part " + "format('SELECT MAX(db_version) AS version FROM %s', tbl_name) as part " "FROM table_names" ") " - "SELECT string_agg(part, ' UNION ALL ') FROM query_parts;"; + "SELECT 'SELECT COALESCE(MAX(version), 0) FROM (' || string_agg(part, ' UNION ALL ') || " + "' UNION ALL SELECT value::bigint FROM cloudsync_settings WHERE key = ''pre_alter_dbversion'') v' " + "FROM query_parts;"; const char * const SQL_CHANGES_INSERT_ROW = "INSERT INTO cloudsync_changes(tbl, pk, col_name, col_value, col_version, db_version, site_id, cl, seq) " diff --git a/test/postgresql/66_db_version_per_transaction.sql b/test/postgresql/66_db_version_per_transaction.sql new file mode 100644 index 00000000..f6bffe5c --- /dev/null +++ b/test/postgresql/66_db_version_per_transaction.sql @@ -0,0 +1,91 @@ +-- Every transaction, implicit or explicit, takes one db_version; its changes share it +-- with seq restarting at 0, and a rolled back transaction takes none. Same scenario and +-- expectations as test_db_version_per_transaction in test/review_regressions.c, so +-- SQLite and PostgreSQL number changes the same way. PostgreSQL also has to hold while +-- another session keeps an old transaction open. + +\set testid '66-db-version' +\ir helper_test_init.sql + +\connect postgres +\ir helper_psql_conn_setup.sql +DROP DATABASE IF EXISTS cloudsync_test_66; +CREATE DATABASE cloudsync_test_66; +\connect cloudsync_test_66 +\ir helper_psql_conn_setup.sql +CREATE EXTENSION IF NOT EXISTS cloudsync; +CREATE EXTENSION IF NOT EXISTS dblink; +CREATE TABLE t(id TEXT PRIMARY KEY NOT NULL, v TEXT); +SELECT cloudsync_init('t') AS _init \gset +CREATE VIEW changes_t AS +SELECT string_agg(cloudsync_pk_decode(pk,1) || ':' || col_name || '@' || db_version || '/' || seq, ' ' ORDER BY db_version, seq) AS list +FROM cloudsync_changes WHERE tbl = 't'; + +INSERT INTO t VALUES ('a1','x'); +INSERT INTO t VALUES ('a2','x'); +INSERT INTO t VALUES ('b1','x'),('b2','x'); +SELECT list = 'a1:v@1/0 a2:v@2/0 b1:v@3/0 b2:v@3/1' AS autocommit_ok, list AS autocommit_list FROM changes_t \gset +\if :autocommit_ok +\echo [PASS] (:testid) autocommit statements take one db_version each, seq restarts at 0 +\else +\echo [FAIL] (:testid) autocommit statements: :autocommit_list +SELECT (:fail::int + 1) AS fail \gset +\endif + +BEGIN; +INSERT INTO t VALUES ('c1','x'); +SELECT cloudsync_db_version() AS mid_dbv \gset +INSERT INTO t VALUES ('c2','x'); +UPDATE t SET v='y' WHERE id='a1'; +COMMIT; +UPDATE t SET v='z' WHERE id='a2'; +BEGIN; +UPDATE t SET v='y' WHERE id='b1'; +UPDATE t SET v='y' WHERE id='b2'; +COMMIT; +BEGIN; +INSERT INTO t VALUES ('r1','x'); +ROLLBACK; +INSERT INTO t VALUES ('d1','x'); +BEGIN; +UPDATE t SET v='w' WHERE id='a1'; +UPDATE t SET v='q' WHERE id='a1'; +DELETE FROM t WHERE id='b2'; +COMMIT; +SELECT list = 'c1:v@4/0 c2:v@4/1 a2:v@5/0 b1:v@6/0 d1:v@7/0 a1:v@8/1 b2:__[RIP]__@8/2' + AND :mid_dbv = 3 AND cloudsync_db_version() = 8 AS tx_ok, list AS tx_list FROM changes_t \gset +\if :tx_ok +\echo [PASS] (:testid) explicit transactions share one db_version, a rollback takes none +\else +\echo [FAIL] (:testid) explicit transactions: :tx_list (db_version inside tx1: :mid_dbv) +SELECT (:fail::int + 1) AS fail \gset +\endif + +-- Another session holds an old transaction open, pinning the snapshot xmin that the +-- cached db_version is checked against: later transactions must still take new ones. +SELECT dblink_connect('old', format('dbname=%s user=%s', current_database(), current_user)) AS _c \gset +SELECT dblink_exec('old', 'BEGIN') AS _b \gset +SELECT x AS _xid FROM dblink('old', 'SELECT txid_current()') AS r(x BIGINT) \gset +DELETE FROM t; +INSERT INTO t VALUES ('p1','x'); +INSERT INTO t VALUES ('p2','x'); +BEGIN; +INSERT INTO t VALUES ('p3','x'); +INSERT INTO t VALUES ('p4','x'); +COMMIT; +INSERT INTO t VALUES ('p5','x'); +SELECT dblink_exec('old', 'COMMIT') AS _e \gset +SELECT dblink_disconnect('old') AS _d \gset +SELECT string_agg(cloudsync_pk_decode(pk,1) || '@' || db_version || '/' || seq, ' ' ORDER BY db_version, seq) AS pinned_list +FROM cloudsync_changes WHERE tbl = 't' AND cloudsync_pk_decode(pk,1) LIKE 'p%' \gset +SELECT :'pinned_list' = 'p1@10/0 p2@11/0 p3@12/0 p4@12/1 p5@13/0' AS pinned_ok \gset +\if :pinned_ok +\echo [PASS] (:testid) an old open transaction elsewhere does not make transactions share a db_version +\else +\echo [FAIL] (:testid) with an old transaction open elsewhere: :pinned_list +SELECT (:fail::int + 1) AS fail \gset +\endif + +\connect postgres +\ir helper_psql_conn_setup.sql +DROP DATABASE IF EXISTS cloudsync_test_66; diff --git a/test/postgresql/full_test.sql b/test/postgresql/full_test.sql index 26864223..fe21bd0f 100644 --- a/test/postgresql/full_test.sql +++ b/test/postgresql/full_test.sql @@ -72,6 +72,7 @@ \ir 62_deferred_fk_caller_commit.sql \ir 63_deep_savepoints.sql \ir 64_block_rewrite_leftovers.sql +\ir 66_db_version_per_transaction.sql -- 'Test summary' \echo '\nTest summary:' diff --git a/test/review_regressions.c b/test/review_regressions.c index d6d434c5..d80856c9 100644 --- a/test/review_regressions.c +++ b/test/review_regressions.c @@ -364,6 +364,37 @@ static void test_fragment_retention(void) { for (int i = 0; i < frag_count; i++) free(frag_data[i]); } +static bool text_is(sqlite3 *db, const char *query, const char *expected) { + sqlite3_stmt *vm = NULL; + bool ok = sqlite3_prepare_v2(db, query, -1, &vm, NULL) == SQLITE_OK && sqlite3_step(vm) == SQLITE_ROW && + sqlite3_column_text(vm, 0) && strcmp((const char *)sqlite3_column_text(vm, 0), expected) == 0; + if (!ok) fprintf(stderr, " got: %s\n expected: %s\n", vm && sqlite3_column_text(vm, 0) ? (const char *)sqlite3_column_text(vm, 0) : "(null)", expected); + sqlite3_finalize(vm); + return ok; +} +static void test_db_version_per_transaction(void) { + // Every transaction, implicit or explicit, takes one db_version; its changes share it + // with seq restarting at 0, and a rolled back transaction takes none. Same scenario + // and expectations as test/postgresql/66_db_version_per_transaction.sql. + const char *changes = "SELECT string_agg(cloudsync_pk_decode(pk,1) || ':' || col_name || '@' || db_version || '/' || seq, ' ' ORDER BY db_version, seq) FROM cloudsync_changes"; + sqlite3 *db = open_db(); + CHECK(sql(db, "CREATE TABLE t(id TEXT PRIMARY KEY NOT NULL, v TEXT); SELECT cloudsync_init('t');") == SQLITE_OK); + CHECK(sql(db, "INSERT INTO t VALUES ('a1','x')") == SQLITE_OK); + CHECK(sql(db, "INSERT INTO t VALUES ('a2','x')") == SQLITE_OK); + CHECK(sql(db, "INSERT INTO t VALUES ('b1','x'),('b2','x')") == SQLITE_OK); + CHECK(text_is(db, changes, "a1:v@1/0 a2:v@2/0 b1:v@3/0 b2:v@3/1")); + CHECK(sql(db, "BEGIN; INSERT INTO t VALUES ('c1','x');") == SQLITE_OK); + CHECK(scalar(db, "SELECT cloudsync_db_version()") == 3); // the last committed one + CHECK(sql(db, "INSERT INTO t VALUES ('c2','x'); UPDATE t SET v='y' WHERE id='a1'; COMMIT;") == SQLITE_OK); + CHECK(sql(db, "UPDATE t SET v='z' WHERE id='a2'") == SQLITE_OK); + CHECK(sql(db, "BEGIN; UPDATE t SET v='y' WHERE id='b1'; UPDATE t SET v='y' WHERE id='b2'; COMMIT;") == SQLITE_OK); + CHECK(sql(db, "BEGIN; INSERT INTO t VALUES ('r1','x'); ROLLBACK;") == SQLITE_OK); + CHECK(sql(db, "INSERT INTO t VALUES ('d1','x')") == SQLITE_OK); + CHECK(sql(db, "BEGIN; UPDATE t SET v='w' WHERE id='a1'; UPDATE t SET v='q' WHERE id='a1'; DELETE FROM t WHERE id='b2'; COMMIT;") == SQLITE_OK); + CHECK(text_is(db, changes, "c1:v@4/0 c2:v@4/1 a2:v@5/0 b1:v@6/0 d1:v@7/0 a1:v@8/1 b2:__[RIP]__@8/2")); + CHECK(scalar(db, "SELECT cloudsync_db_version()") == 8); + CHECK(close_db(db) == SQLITE_OK); +} static void test_block_write_errors(void) { for (int update = 0; update < 2; update++) { sqlite3 *db = open_db(); @@ -645,6 +676,7 @@ int main(void) { test_resurrected_group_rollback(); test_batched_update_missing_row(); test_fragment_retention(); + test_db_version_per_transaction(); test_block_write_errors(); test_block_materialize_errors(); test_block_migration_orphan(); From e0a551440fb8ef0ac9bb2c09a0edbec662c4ebd2 Mon Sep 17 00:00:00 2001 From: Andrea Donetti Date: Thu, 24 Sep 2026 16:35:12 -0600 Subject: [PATCH 2/2] fix: keep an applied payload's db_versions apart from local writes Closing a payload group inside a transaction reused the commit hook, which copied an uncommitted version into the committed cache: a local write after the apply reused the last merged db_version with seq back at 0, and a rolled back apply left the cache ahead of the tables. A group now only marks the pending version closed, so the next change takes pending + 1 while the committed version stays untouched. The apply also closes it after the last group and around each fragment value. Co-Authored-By: Claude Opus 5.5 (1M context) --- CHANGELOG.md | 2 +- src/cloudsync.c | 28 ++++++-- .../66_db_version_per_transaction.sql | 65 +++++++++++++++++++ test/review_regressions.c | 30 +++++++++ 4 files changed, 119 insertions(+), 6 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 75a23fcf..3af908e6 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -9,7 +9,7 @@ The format is based on [Keep a Changelog](https://keepachangelog.com/en/1.0.0/). ### Fixed - **PostgreSQL: separate transactions no longer share a db_version.** With more than one synced table, the current db_version was read from a single table's metadata instead of the maximum across all of them, so consecutive transactions writing to another table all took the same db_version. Holding a transaction open in another session had the same effect on every table. A change that lands on a db_version already served is skipped by the next download, whose window starts after that version, so peers could miss it. Each transaction now takes its own db_version, as on SQLite, and `seq` restarts at 0 in each one instead of growing for the life of the connection, where a long-lived connection could eventually overflow it. -- **PostgreSQL: applying a payload keeps each source db_version apart.** A whole `cloudsync_payload_apply` runs in one transaction, so every db_version in the payload was stored under the same local db_version while each change kept its source `seq`, producing duplicate `(db_version, seq)` pairs that the chunked download relies on being unique. Each source db_version now takes its own local db_version, as on SQLite. +- **PostgreSQL: applying a payload keeps each source db_version apart.** A whole `cloudsync_payload_apply` runs in one transaction, so every db_version in the payload was stored under the same local db_version while each change kept its source `seq`, producing duplicate `(db_version, seq)` pairs that the chunked download relies on being unique. Each source db_version now takes its own local db_version, as on SQLite, apart from any local writes made in the same transaction. - **SQLite: paging a chunked download no longer gets slower with every chunk.** The positional cursor on `cloudsync_payload_chunks` was meant to seek straight to where the previous call stopped, but it stated its resume point only inside `(db_version > ? OR (db_version = ? AND seq >= ?))`, whose two arms carry distinct parameters. SQLite does not derive a range from that, so the scan over `cloudsync_changes` ran with an upper bound only and re-read the window from the beginning on every call, discarding rows until it reached the resume point — making a full drain quadratic in the number of chunks, and long enough on a large tenant to hit a server-side deadline and never complete. An explicit `db_version >= ?` is now stated alongside the disjunction; it selects exactly the same rows and lets the scan seek. Locally, draining a 188-chunk window went from 6110 ms to 210 ms, and the per-chunk cost no longer depends on how large the window is. PostgreSQL was never affected. - **A row rewritten in one statement no longer keeps the old blocks of a shorter value.** Writing a whole row rewrites its block column from the first position, so any block past the end of the new value stayed stored. On SQLite `INSERT OR REPLACE` skips the old row's delete trigger unless `recursive_triggers` is on, so those leftovers kept their metadata and were delivered as content: replacing `AAA\nBBB\nCCC` with `ZZZ` left `ZZZ\nBBB\nCCC` on the peers and on a later local read. On PostgreSQL they carried no metadata, so peers were unaffected, but they stayed in the blocks table for the life of the row. Blocks the new value does not cover are now retired with it. - **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. diff --git a/src/cloudsync.c b/src/cloudsync.c index 0407b036..c6da6f30 100644 --- a/src/cloudsync.c +++ b/src/cloudsync.c @@ -176,6 +176,8 @@ struct cloudsync_context { int64_t db_version; // version the DB would have if the transaction committed now int64_t pending_db_version; + // set when a payload group ends inside a transaction: the next change takes a new version + bool pending_closed; // used to set an order inside each transaction int seq; @@ -493,6 +495,8 @@ int64_t cloudsync_dbversion_next (cloudsync_context *data, int64_t merging_versi int64_t result = data->db_version + 1; if (result < data->pending_db_version) result = data->pending_db_version; + if (data->pending_closed && result <= data->pending_db_version) result = data->pending_db_version + 1; + data->pending_closed = false; if (merging_version != CLOUDSYNC_VALUE_NOTSET && result < merging_version) result = merging_version; data->pending_db_version = result; @@ -2648,6 +2652,7 @@ int cloudsync_commit_hook (void *ctx) { data->db_version = data->pending_db_version; data->pending_db_version = CLOUDSYNC_VALUE_NOTSET; + data->pending_closed = false; data->seq = 0; return DBRES_OK; @@ -2657,19 +2662,27 @@ void cloudsync_rollback_hook (void *ctx) { cloudsync_context *data = (cloudsync_context *)ctx; data->pending_db_version = CLOUDSYNC_VALUE_NOTSET; + data->pending_closed = false; data->seq = 0; } -// Runs the commit or rollback hook where the database does not: at the end of every -// PostgreSQL transaction, and between source db_versions of a payload applied inside one. -// It acts only when a db_version was taken, so the cached one survives read-only -// transactions. +// For a backend without commit and rollback hooks (PostgreSQL), called at the end of +// every transaction. It acts only on a transaction that took a db_version, so the cached +// db_version survives the many read-only transactions in between. void cloudsync_transaction_end (cloudsync_context *data, bool committed) { if (!data || data->pending_db_version == CLOUDSYNC_VALUE_NOTSET) return; if (committed) cloudsync_commit_hook(data); else cloudsync_rollback_hook(data); } +// Ends the pending db_version without committing it: the next change takes a new one with +// seq restarting at 0. The committed db_version is left alone, so a rollback restores it. +static void cloudsync_pending_version_close (cloudsync_context *data) { + if (data->pending_db_version == CLOUDSYNC_VALUE_NOTSET) return; + data->pending_closed = true; + data->seq = 0; +} + int cloudsync_begin_alter (cloudsync_context *data, const char *table_name) { // init cloudsync_settings if (cloudsync_context_init(data) == NULL) { @@ -4497,6 +4510,8 @@ int cloudsync_payload_apply (cloudsync_context *data, const char *payload, int b break; } int n = 0; + // each value takes its own db_version, as each group does in the row path + cloudsync_pending_version_close(data); rc = cloudsync_payload_apply_fragment_row(data, &row, checkpoint_db_version != CLOUDSYNC_CHECKPOINT_LAST_APPLIED, &n); // stop at the first error, as the row path does if (rc != DBRES_OK) break; @@ -4504,6 +4519,7 @@ int cloudsync_payload_apply (cloudsync_context *data, const char *payload, int b buffer += seek; buf_len -= seek; } + cloudsync_pending_version_close(data); if (clone) cloudsync_memory_free(clone); cloudsync_apply_stats_add(data, applied_rows); if (pnrows) *pnrows = applied_rows; @@ -4588,7 +4604,7 @@ int cloudsync_payload_apply (cloudsync_context *data, const char *payload, int b // Inside a caller's transaction (always on PostgreSQL) no commit hook runs between // groups: close the group's db_version here, or every source db_version would get // the same local one and their seqs would collide. A no-op after the RELEASE above. - if (db_version_changed) cloudsync_transaction_end(data, true); + if (db_version_changed) cloudsync_pending_version_close(data); if (!in_savepoint && db_version_changed && !database_in_transaction(data)) { rc = database_begin_savepoint(data, "cloudsync_payload_apply"); @@ -4663,6 +4679,8 @@ int cloudsync_payload_apply (cloudsync_context *data, const char *payload, int b in_savepoint = false; } } + // Local writes that follow in the same transaction take a db_version of their own. + cloudsync_pending_version_close(data); rc = fail_rc; if (rc != DBRES_OK) { diff --git a/test/postgresql/66_db_version_per_transaction.sql b/test/postgresql/66_db_version_per_transaction.sql index f6bffe5c..004a769d 100644 --- a/test/postgresql/66_db_version_per_transaction.sql +++ b/test/postgresql/66_db_version_per_transaction.sql @@ -86,6 +86,71 @@ SELECT :'pinned_list' = 'p1@10/0 p2@11/0 p3@12/0 p4@12/1 p5@13/0' AS pinned_ SELECT (:fail::int + 1) AS fail \gset \endif +-- A payload applied inside a transaction takes its own db_versions, apart from the local +-- writes before and after it, and a rollback gives them all back. Same expectations as +-- test_db_version_apply_in_transaction in test/review_regressions.c. +\connect postgres +\ir helper_psql_conn_setup.sql +DROP DATABASE IF EXISTS cloudsync_test_66_src; +DROP DATABASE IF EXISTS cloudsync_test_66_dst; +DROP DATABASE IF EXISTS cloudsync_test_66_rb; +CREATE DATABASE cloudsync_test_66_src; +CREATE DATABASE cloudsync_test_66_dst; +CREATE DATABASE cloudsync_test_66_rb; + +\connect cloudsync_test_66_src +\ir helper_psql_conn_setup.sql +CREATE EXTENSION IF NOT EXISTS cloudsync; +CREATE TABLE t(id TEXT PRIMARY KEY NOT NULL, v TEXT); +SELECT cloudsync_init('t') AS _init \gset +INSERT INTO t VALUES ('r1','x'),('r2','x'); +INSERT INTO t VALUES ('r3','x'); +SELECT '\x' || 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_66_dst +\ir helper_psql_conn_setup.sql +CREATE EXTENSION IF NOT EXISTS cloudsync; +CREATE TABLE t(id TEXT PRIMARY KEY NOT NULL, v TEXT); +SELECT cloudsync_init('t') AS _init \gset +BEGIN; +INSERT INTO t VALUES ('l1','x'); +SELECT cloudsync_payload_apply(decode(substr(:'payload', 3), 'hex')) AS _apply \gset +INSERT INTO t VALUES ('l2','x'); +COMMIT; +SELECT string_agg(cloudsync_pk_decode(pk,1) || '@' || db_version || '/' || seq, ' ' ORDER BY db_version, seq) AS apply_list +FROM cloudsync_changes \gset +SELECT :'apply_list' = 'l1@1/0 r1@2/0 r2@2/1 r3@3/0 l2@4/0' AS apply_ok \gset +\if :apply_ok +\echo [PASS] (:testid) a payload applied inside a transaction keeps its db_versions apart from local writes +\else +\echo [FAIL] (:testid) payload applied inside a transaction: :apply_list +SELECT (:fail::int + 1) AS fail \gset +\endif + +\connect cloudsync_test_66_rb +\ir helper_psql_conn_setup.sql +CREATE EXTENSION IF NOT EXISTS cloudsync; +CREATE TABLE t(id TEXT PRIMARY KEY NOT NULL, v TEXT); +SELECT cloudsync_init('t') AS _init \gset +BEGIN; +SELECT cloudsync_payload_apply(decode(substr(:'payload', 3), 'hex')) AS _apply \gset +ROLLBACK; +SELECT cloudsync_db_version() AS rb_dbv \gset +INSERT INTO t VALUES ('l1','x'); +SELECT string_agg(cloudsync_pk_decode(pk,1) || '@' || db_version || '/' || seq, ' ' ORDER BY db_version, seq) AS rb_list +FROM cloudsync_changes \gset +SELECT :rb_dbv = 0 AND :'rb_list' = 'l1@1/0' AS rb_ok \gset +\if :rb_ok +\echo [PASS] (:testid) a rolled back payload apply takes no db_version +\else +\echo [FAIL] (:testid) after a rolled back apply: db_version :rb_dbv, :rb_list +SELECT (:fail::int + 1) AS fail \gset +\endif + \connect postgres \ir helper_psql_conn_setup.sql DROP DATABASE IF EXISTS cloudsync_test_66; +DROP DATABASE IF EXISTS cloudsync_test_66_src; +DROP DATABASE IF EXISTS cloudsync_test_66_dst; +DROP DATABASE IF EXISTS cloudsync_test_66_rb; diff --git a/test/review_regressions.c b/test/review_regressions.c index d80856c9..14e27fd1 100644 --- a/test/review_regressions.c +++ b/test/review_regressions.c @@ -395,6 +395,35 @@ static void test_db_version_per_transaction(void) { CHECK(scalar(db, "SELECT cloudsync_db_version()") == 8); CHECK(close_db(db) == SQLITE_OK); } +static void test_db_version_apply_in_transaction(void) { + // A payload applied inside the caller's transaction takes its own db_versions, apart + // from the local writes before and after it, and a rollback gives them all back. + const char *changes = "SELECT string_agg(cloudsync_pk_decode(pk,1) || '@' || db_version || '/' || seq, ' ' ORDER BY db_version, seq) FROM cloudsync_changes"; + const char *schema = "CREATE TABLE t(id TEXT PRIMARY KEY NOT NULL, v TEXT); SELECT cloudsync_init('t');"; + sqlite3 *source = open_db(); + CHECK(sql(source, schema) == SQLITE_OK); + CHECK(sql(source, "INSERT INTO t VALUES ('r1','x'),('r2','x')") == SQLITE_OK); + CHECK(sql(source, "INSERT INTO t VALUES ('r3','x')") == SQLITE_OK); + + sqlite3 *db = open_db(); + CHECK(sql(db, schema) == SQLITE_OK); + CHECK(sql(db, "BEGIN; INSERT INTO t VALUES ('l1','x');") == SQLITE_OK); + CHECK(apply_payload(source, db) == SQLITE_ROW); + CHECK(sql(db, "INSERT INTO t VALUES ('l2','x'); COMMIT;") == SQLITE_OK); + CHECK(text_is(db, changes, "l1@1/0 r1@2/0 r2@2/1 r3@3/0 l2@4/0")); + CHECK(close_db(db) == SQLITE_OK); + + db = open_db(); + CHECK(sql(db, schema) == SQLITE_OK); + CHECK(sql(db, "BEGIN;") == SQLITE_OK); + CHECK(apply_payload(source, db) == SQLITE_ROW); + CHECK(sql(db, "ROLLBACK;") == SQLITE_OK); + CHECK(scalar(db, "SELECT cloudsync_db_version()") == 0); + CHECK(sql(db, "INSERT INTO t VALUES ('l1','x')") == SQLITE_OK); + CHECK(text_is(db, changes, "l1@1/0")); + CHECK(close_db(db) == SQLITE_OK); + CHECK(close_db(source) == SQLITE_OK); +} static void test_block_write_errors(void) { for (int update = 0; update < 2; update++) { sqlite3 *db = open_db(); @@ -677,6 +706,7 @@ int main(void) { test_batched_update_missing_row(); test_fragment_retention(); test_db_version_per_transaction(); + test_db_version_apply_in_transaction(); test_block_write_errors(); test_block_materialize_errors(); test_block_migration_orphan();