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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 2 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -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, 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.
Expand Down
33 changes: 33 additions & 0 deletions src/cloudsync.c
Original file line number Diff line number Diff line change
Expand Up @@ -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;

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

Expand Down Expand Up @@ -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;
Expand All @@ -2657,6 +2662,24 @@ 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;
}

// 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;
}

Expand Down Expand Up @@ -4487,13 +4510,16 @@ 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;
applied_rows += n;
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;
Expand Down Expand Up @@ -4575,6 +4601,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_pending_version_close(data);

if (!in_savepoint && db_version_changed && !database_in_transaction(data)) {
rc = database_begin_savepoint(data, "cloudsync_payload_apply");
if (rc != DBRES_OK) {
Expand Down Expand Up @@ -4648,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) {
Expand Down
1 change: 1 addition & 0 deletions src/cloudsync.h
Original file line number Diff line number Diff line change
Expand Up @@ -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);
Expand Down
26 changes: 26 additions & 0 deletions src/postgresql/cloudsync_postgresql.c
Original file line number Diff line number Diff line change
Expand Up @@ -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) {
Expand All @@ -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) {
Expand Down
6 changes: 4 additions & 2 deletions src/postgresql/sql_postgresql.c
Original file line number Diff line number Diff line change
Expand Up @@ -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) "
Expand Down
156 changes: 156 additions & 0 deletions test/postgresql/66_db_version_per_transaction.sql
Original file line number Diff line number Diff line change
@@ -0,0 +1,156 @@
-- 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

-- 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;
1 change: 1 addition & 0 deletions test/postgresql/full_test.sql
Original file line number Diff line number Diff line change
Expand Up @@ -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:'
Expand Down
Loading
Loading