diff --git a/API.md b/API.md index 6322be78..a132a956 100644 --- a/API.md +++ b/API.md @@ -662,6 +662,8 @@ On PostgreSQL, apply chunks as individual statements from the transport/client l - Monolithic payloads generated by [`cloudsync_payload_encode()`](#cloudsync_payload_encodetbl-pk-col_name-col_value-col_version-db_version-site_id-cl-seq). - Chunk-fragment payloads generated by [`cloudsync_payload_chunks()`](#cloudsync_payload_chunkssince_db_version-filter_site_id-until_db_version-exclude_filter_site_id). +On PostgreSQL, concurrent merges of the same row are serialized before reading its clocks under `READ COMMITTED` (including `READ UNCOMMITTED`, which PostgreSQL treats identically). This also covers direct inserts into `cloudsync_changes`. Locks last until the caller's transaction ends and use a fixed pool of 256 keys per database, so unrelated rows can also wait for each other. `SERIALIZABLE` relies on PostgreSQL's conflict detection; retry the whole transaction on serialization failure (`40001`) or deadlock (`40P01`). Merging under `REPEATABLE READ` is refused with `0A000`, because waiting cannot refresh that transaction's snapshot. Transactions merging multiple rows can deadlock when acquiring locks in different orders. + When a v3 fragment payload is received, CloudSync stores the fragment in an internal table and returns after applying zero or more completed values. Once the final fragment for a value is received, the completed value is validated and applied. Fragments can arrive in any order, and duplicate fragment delivery is idempotent. Applying a fragment never moves the receive checkpoint. On PostgreSQL, pieces of one value applied by concurrent transactions wait for each other under `READ COMMITTED`, and fail with a retryable serialization error under `SERIALIZABLE` when they conflict; a fragment is refused under `REPEATABLE READ`, where a transaction could miss a piece committed while it waited. **Parameters:** diff --git a/docker/README.md b/docker/README.md index 95777d65..02f262fd 100644 --- a/docker/README.md +++ b/docker/README.md @@ -247,6 +247,26 @@ EXECUTE FUNCTION bump_app_schema_version(); ## Development Workflow +### Reproducible PostgreSQL tests + +From the repository root, run: + +```bash +./scripts/test-postgres-docker.sh +# Issue #70 regression only (serial/concurrent merge, rollback, isolation, lock bound) +./scripts/test-postgres-docker.sh 67_concurrent_merge.sql +# Select another PostgreSQL image tag +POSTGRES_TAG=15-bookworm ./scripts/test-postgres-docker.sh +POSTGRES_TAG=18-bookworm ./scripts/test-postgres-docker.sh +``` + +The script builds the extension from the current source, creates an isolated container, +runs psql with `ON_ERROR_STOP`, and removes the container and its volumes on exit. +It does not publish ports or use an existing database. The image remains cached. +The concurrency test uses dblink and checks that the second session is waiting on +a lock before allowing the first to commit. On the code before the issue #70 fix, +it fails with `expected higher/3, got lower/2`. + ### 1. Make Changes Edit source files in `src/postgresql/` or `src/` (shared code). diff --git a/scripts/test-postgres-docker.sh b/scripts/test-postgres-docker.sh new file mode 100755 index 00000000..dab8a92a --- /dev/null +++ b/scripts/test-postgres-docker.sh @@ -0,0 +1,31 @@ +#!/usr/bin/env bash +# Disposable PostgreSQL build and tests; no host ports or persistent volumes. +set -euo pipefail +cd "$(dirname "$0")/.." +postgres_tag=${POSTGRES_TAG:-17} +test_file=${1:-full_test.sql} +if [[ ! -f "test/postgresql/$test_file" || "$test_file" == */* ]]; then + echo "Expected a SQL file name in test/postgresql" >&2 + exit 2 +fi +image="sqlite-sync-test:${postgres_tag}" +container="sqlite-sync-test-$$" +cleanup() { docker rm -f -v "$container" >/dev/null 2>&1 || true; } +trap cleanup EXIT + +docker build --build-arg "POSTGRES_TAG=$postgres_tag" -t "$image" -f docker/postgresql/Dockerfile . +docker run -d --name "$container" -e POSTGRES_PASSWORD=postgres \ + -v "$PWD/test:/tests:ro" "$image" >/dev/null +for ((i=0; i<60; i++)); do + if docker exec "$container" pg_isready -U postgres -d postgres >/dev/null 2>&1; then + # The image entrypoint briefly runs a temporary server during initialization. + if docker exec "$container" psql -h 127.0.0.1 -U postgres -d postgres -c 'SELECT 1' >/dev/null 2>&1; then + docker exec "$container" psql -U postgres -d postgres -v ON_ERROR_STOP=1 \ + -f "/tests/postgresql/$test_file" + exit 0 + fi + fi + sleep 1 +done +docker logs "$container" >&2 +exit 1 diff --git a/src/cloudsync.c b/src/cloudsync.c index c6da6f30..3c7ddcd1 100644 --- a/src/cloudsync.c +++ b/src/cloudsync.c @@ -2197,6 +2197,11 @@ int table_col_index (cloudsync_table_context *table, const char *col_name) { } int merge_insert (cloudsync_context *data, cloudsync_table_context *table, const char *insert_pk, int insert_pk_len, int64_t insert_cl, const char *insert_name, dbvalue_t *insert_value, int64_t insert_col_version, int64_t insert_db_version, const char *insert_site_id, int insert_site_id_len, int64_t insert_seq, int64_t *rowid) { + // Hold through the transaction, including deferred column writes. Lock before + // reading any row clocks, even when the row does not exist yet. + int lock_rc = database_merge_lock(data, table->meta_ref, insert_pk, insert_pk_len); + if (lock_rc != DBRES_OK) return lock_rc; + // Handle DWS and AWS algorithms here // Delete-Wins Set (DWS): table_algo_crdt_dws // Add-Wins Set (AWS): table_algo_crdt_aws diff --git a/src/database.h b/src/database.h index 50cb621d..9e96c600 100644 --- a/src/database.h +++ b/src/database.h @@ -98,6 +98,7 @@ int database_commit_savepoint (cloudsync_context *data, const char *savepoint_na int database_rollback_savepoint (cloudsync_context *data, const char *savepoint_name); bool database_in_transaction (cloudsync_context *data); int database_fragment_lock (cloudsync_context *data, const char *value_id); +int database_merge_lock (cloudsync_context *data, const char *table_ref, const void *pk, int pklen); int database_errcode (cloudsync_context *data); const char *database_errmsg (cloudsync_context *data); void database_log_warning (cloudsync_context *data, const char *message); diff --git a/src/postgresql/database_postgresql.c b/src/postgresql/database_postgresql.c index da4e55d6..708c3ba5 100644 --- a/src/postgresql/database_postgresql.c +++ b/src/postgresql/database_postgresql.c @@ -23,6 +23,7 @@ // PostgreSQL SPI and other headers #include "access/xact.h" #include "catalog/pg_type.h" +#include "common/hashfn.h" #include "executor/spi.h" #include "funcapi.h" #include "utils/array.h" @@ -1212,6 +1213,31 @@ int database_fragment_lock (cloudsync_context *data, const char *value_id) { return (rc == DBRES_ROW) ? DBRES_OK : cloudsync_set_error(data, "cloudsync_payload_apply: unable to lock a fragmented value", rc); } +// Serialize merge decisions for every column of a row, including tombstones and +// blocks. READ COMMITTED takes a fresh snapshot on the subsequent clock reads. +// SERIALIZABLE detects stale decisions itself; REPEATABLE READ cannot refresh its +// snapshot after waiting and is refused, as for fragmented values above. +int database_merge_lock (cloudsync_context *data, const char *table_ref, const void *pk, int pklen) { + if (IsolationIsSerializable()) return DBRES_OK; + if (IsolationUsesXactSnapshot()) { + int rc = cloudsync_set_error(data, "cloudsync merge cannot run under REPEATABLE READ, use READ COMMITTED or SERIALIZABLE", DBRES_MISUSE); + cloudsync_set_sqlstate(data, ERRCODE_FEATURE_NOT_SUPPORTED); + return rc; + } + + // A fixed pool bounds lock-table use even for very large imports. Collisions + // only serialize unrelated rows. Include the qualified metadata table so all + // callers resolving the same table agree, regardless of their search_path. + uint32 bucket = (hash_bytes((const unsigned char *)table_ref, strlen(table_ref)) ^ + hash_bytes((const unsigned char *)pk, pklen)) & 255; + dbvm_t *vm = NULL; + int rc = databasevm_prepare(data, "SELECT pg_advisory_xact_lock(1129530963, $1::integer);", &vm, 0); + if (rc == DBRES_OK) rc = databasevm_bind_int(vm, 1, bucket); + if (rc == DBRES_OK) rc = databasevm_step(vm); + if (vm) databasevm_finalize(vm); + return (rc == DBRES_ROW) ? DBRES_OK : rc; +} + bool database_table_exists (cloudsync_context *data, const char *name, const char *schema) { return database_system_exists(data, name, "table", false, schema); } diff --git a/src/sqlite/database_sqlite.c b/src/sqlite/database_sqlite.c index cc1e0e78..359bbfbd 100644 --- a/src/sqlite/database_sqlite.c +++ b/src/sqlite/database_sqlite.c @@ -603,6 +603,11 @@ int database_fragment_lock (cloudsync_context *data, const char *value_id) { return DBRES_OK; } +int database_merge_lock (cloudsync_context *data, const char *table_ref, const void *pk, int pklen) { + // SQLite already serializes writers, including the clock reads and writes. + return DBRES_OK; +} + bool database_table_exists (cloudsync_context *data, const char *name, const char *schema) { UNUSED_PARAMETER(schema); return database_system_exists(data, name, "table"); diff --git a/test/postgresql/67_concurrent_merge.sql b/test/postgresql/67_concurrent_merge.sql new file mode 100644 index 00000000..e734a108 --- /dev/null +++ b/test/postgresql/67_concurrent_merge.sql @@ -0,0 +1,221 @@ +-- Issue #70: two merge decisions must not use the same stale row clocks. +-- Run with psql -v ON_ERROR_STOP=1 -f test/postgresql/67_concurrent_merge.sql. +-- dblink observes the waiter before releasing x; no timing-dependent overlap. +\set ON_ERROR_STOP on +\set testid '67-concurrent-merge' +\ir helper_test_init.sql +\connect postgres +\ir helper_psql_conn_setup.sql +DROP DATABASE IF EXISTS cloudsync_test_67; +CREATE DATABASE cloudsync_test_67; +\connect cloudsync_test_67 +\ir helper_psql_conn_setup.sql +CREATE EXTENSION cloudsync; +CREATE EXTENSION dblink; +CREATE TABLE t(id TEXT PRIMARY KEY, v TEXT); +SELECT cloudsync_init('t') AS _init \gset +CREATE TABLE transport(name TEXT PRIMARY KEY, payload BYTEA); +INSERT INTO t VALUES ('row', 'higher'); +INSERT INTO transport SELECT 'higher', cloudsync_payload_encode(tbl, pk, col_name, + col_value, 3, 3, decode(repeat('01',16),'hex'), 1, 0) FROM cloudsync_changes; +UPDATE t SET v='lower'; +INSERT INTO transport SELECT 'lower', cloudsync_payload_encode(tbl, pk, col_name, + col_value, 2, 2, decode(repeat('02',16),'hex'), 1, 0) FROM cloudsync_changes; +TRUNCATE t, t_cloudsync; +CREATE FUNCTION apply_value(n TEXT) RETURNS INT LANGUAGE plpgsql AS $$ +DECLARE data BYTEA; BEGIN + SELECT payload INTO STRICT data FROM transport WHERE name=n; + RETURN cloudsync_payload_apply(data); +END $$; +CREATE FUNCTION assert_winner(label TEXT) RETURNS VOID LANGUAGE plpgsql AS $$ +DECLARE val TEXT; ver BIGINT; BEGIN + SELECT v INTO val FROM t WHERE id='row'; + SELECT col_version INTO ver FROM t_cloudsync WHERE pk=cloudsync_pk_encode('row'::text) AND col_name='v'; + IF val IS DISTINCT FROM 'higher' OR ver IS DISTINCT FROM 3::bigint THEN + RAISE EXCEPTION '%: expected higher/3, got %/%', label, val, ver; + END IF; +END $$; +SELECT apply_value('lower'); +SELECT apply_value('higher'); +SELECT assert_winner('serial lower then higher'); +TRUNCATE t, t_cloudsync; +SELECT apply_value('higher'); +SELECT apply_value('lower'); +SELECT assert_winner('serial higher then lower'); +SELECT dblink_connect('x', format('dbname=%s user=%s application_name=cloudsync_67_x',current_database(),current_user)); +SELECT dblink_connect('y', format('dbname=%s user=%s application_name=cloudsync_67_y',current_database(),current_user)); +-- Initialize both worker contexts outside the contending transactions. +SELECT * FROM dblink('x', 'SELECT apply_value(''higher'')') AS r(n INT); +SELECT * FROM dblink('y', 'SELECT apply_value(''higher'')') AS r(n INT); +SELECT dblink_exec('x', 'SET statement_timeout=''15s'''); +SELECT dblink_exec('y', 'SET statement_timeout=''15s'''); +CREATE FUNCTION wait_for_y() RETURNS VOID LANGUAGE plpgsql AS $$ +BEGIN + FOR i IN 1..500 LOOP + PERFORM pg_stat_clear_snapshot(); + IF EXISTS (SELECT FROM pg_stat_activity WHERE application_name='cloudsync_67_y' AND wait_event_type='Lock') THEN RETURN; END IF; + IF dblink_is_busy('y')=0 THEN RAISE EXCEPTION 'y finished without waiting'; END IF; + PERFORM pg_sleep(0.01); + END LOOP; + RAISE EXCEPTION 'timeout waiting for y'; +END $$; +-- Fresh and existing rows; both arrival orders; rollback must release the lock too. +-- seeded=False, first=higher, rollback=False +TRUNCATE t, t_cloudsync; + +SELECT dblink_exec('x','BEGIN'); +SELECT * FROM dblink('x','SELECT apply_value(''higher'')') AS r(n INT); +SELECT dblink_send_query('y','SELECT apply_value(''lower'')'); +SELECT wait_for_y(); +SELECT dblink_exec('x','COMMIT'); +SELECT * FROM dblink_get_result('y') AS r(n INT); +SELECT * FROM dblink_get_result('y') AS r(n INT); + +SELECT assert_winner('seeded=False first=higher rollback=False'); + +-- seeded=False, first=higher, rollback=True +TRUNCATE t, t_cloudsync; + +SELECT dblink_exec('x','BEGIN'); +SELECT * FROM dblink('x','SELECT apply_value(''higher'')') AS r(n INT); +SELECT dblink_send_query('y','SELECT apply_value(''lower'')'); +SELECT wait_for_y(); +SELECT dblink_exec('x','ROLLBACK'); +SELECT * FROM dblink_get_result('y') AS r(n INT); +SELECT * FROM dblink_get_result('y') AS r(n INT); +SELECT apply_value('higher'); +SELECT assert_winner('seeded=False first=higher rollback=True'); + +-- seeded=False, first=lower, rollback=False +TRUNCATE t, t_cloudsync; + +SELECT dblink_exec('x','BEGIN'); +SELECT * FROM dblink('x','SELECT apply_value(''lower'')') AS r(n INT); +SELECT dblink_send_query('y','SELECT apply_value(''higher'')'); +SELECT wait_for_y(); +SELECT dblink_exec('x','COMMIT'); +SELECT * FROM dblink_get_result('y') AS r(n INT); +SELECT * FROM dblink_get_result('y') AS r(n INT); + +SELECT assert_winner('seeded=False first=lower rollback=False'); + +-- seeded=False, first=lower, rollback=True +TRUNCATE t, t_cloudsync; + +SELECT dblink_exec('x','BEGIN'); +SELECT * FROM dblink('x','SELECT apply_value(''lower'')') AS r(n INT); +SELECT dblink_send_query('y','SELECT apply_value(''higher'')'); +SELECT wait_for_y(); +SELECT dblink_exec('x','ROLLBACK'); +SELECT * FROM dblink_get_result('y') AS r(n INT); +SELECT * FROM dblink_get_result('y') AS r(n INT); +SELECT apply_value('lower'); +SELECT assert_winner('seeded=False first=lower rollback=True'); + +-- seeded=True, first=higher, rollback=False +TRUNCATE t, t_cloudsync; +SELECT apply_value('lower'); +SELECT dblink_exec('x','BEGIN'); +SELECT * FROM dblink('x','SELECT apply_value(''higher'')') AS r(n INT); +SELECT dblink_send_query('y','SELECT apply_value(''lower'')'); +SELECT wait_for_y(); +SELECT dblink_exec('x','COMMIT'); +SELECT * FROM dblink_get_result('y') AS r(n INT); +SELECT * FROM dblink_get_result('y') AS r(n INT); + +SELECT assert_winner('seeded=True first=higher rollback=False'); + +-- seeded=True, first=higher, rollback=True +TRUNCATE t, t_cloudsync; +SELECT apply_value('lower'); +SELECT dblink_exec('x','BEGIN'); +SELECT * FROM dblink('x','SELECT apply_value(''higher'')') AS r(n INT); +SELECT dblink_send_query('y','SELECT apply_value(''lower'')'); +SELECT wait_for_y(); +SELECT dblink_exec('x','ROLLBACK'); +SELECT * FROM dblink_get_result('y') AS r(n INT); +SELECT * FROM dblink_get_result('y') AS r(n INT); +SELECT apply_value('higher'); +SELECT assert_winner('seeded=True first=higher rollback=True'); + +-- seeded=True, first=lower, rollback=False +TRUNCATE t, t_cloudsync; +SELECT apply_value('lower'); +SELECT dblink_exec('x','BEGIN'); +SELECT * FROM dblink('x','SELECT apply_value(''lower'')') AS r(n INT); +SELECT dblink_send_query('y','SELECT apply_value(''higher'')'); +SELECT wait_for_y(); +SELECT dblink_exec('x','COMMIT'); +SELECT * FROM dblink_get_result('y') AS r(n INT); +SELECT * FROM dblink_get_result('y') AS r(n INT); + +SELECT assert_winner('seeded=True first=lower rollback=False'); + +-- seeded=True, first=lower, rollback=True +TRUNCATE t, t_cloudsync; +SELECT apply_value('lower'); +SELECT dblink_exec('x','BEGIN'); +SELECT * FROM dblink('x','SELECT apply_value(''lower'')') AS r(n INT); +SELECT dblink_send_query('y','SELECT apply_value(''higher'')'); +SELECT wait_for_y(); +SELECT dblink_exec('x','ROLLBACK'); +SELECT * FROM dblink_get_result('y') AS r(n INT); +SELECT * FROM dblink_get_result('y') AS r(n INT); +SELECT apply_value('lower'); +SELECT assert_winner('seeded=True first=lower rollback=True'); +-- REPEATABLE READ cannot refresh a snapshot after a lock wait. Preserve SQLSTATE. +CREATE FUNCTION apply_checked(n TEXT) RETURNS TEXT LANGUAGE plpgsql AS $$ +BEGIN + PERFORM apply_value(n); + RETURN 'ok'; +EXCEPTION WHEN OTHERS THEN RETURN SQLSTATE; +END $$; +BEGIN ISOLATION LEVEL REPEATABLE READ; +DO $$ BEGIN + IF apply_checked('lower') <> '0A000' THEN RAISE EXCEPTION 'expected REPEATABLE READ refusal'; END IF; +END $$; +ROLLBACK; + +-- SERIALIZABLE must abort the stale writer. Retrying then converges. +TRUNCATE t, t_cloudsync; +SELECT dblink_exec('x','BEGIN ISOLATION LEVEL SERIALIZABLE'); +SELECT * FROM dblink('x','SELECT apply_value(''higher'')') AS r(n INT); +SELECT dblink_exec('y','BEGIN ISOLATION LEVEL SERIALIZABLE'); +SELECT dblink_send_query('y','SELECT apply_checked(''lower'')'); +SELECT wait_for_y(); +SELECT dblink_exec('x','COMMIT'); +SELECT state = '40001' AS serialization_ok FROM dblink_get_result('y') AS r(state TEXT) \gset +SELECT * FROM dblink_get_result('y') AS r(state TEXT); +SELECT dblink_exec('y','ROLLBACK'); +\if :serialization_ok +\else +DO $$ BEGIN RAISE EXCEPTION 'expected SQLSTATE 40001'; END $$; +\endif +SELECT apply_value('lower'); +SELECT assert_winner('SERIALIZABLE retry'); + +-- The direct cloudsync_changes API uses the same protection. Applying many rows +-- in one transaction must not allocate one advisory lock for each distinct PK. +CREATE TEMP TABLE encoded_value AS SELECT col_value FROM cloudsync_changes WHERE tbl='t' AND col_name='v'; +BEGIN; +INSERT INTO cloudsync_changes(tbl,pk,col_name,col_value,col_version,db_version,site_id,cl,seq) +SELECT 't',cloudsync_pk_encode('bulk-' || i), 'v', e.col_value, 3, 3, + decode(repeat('01',16),'hex'), 1, i FROM generate_series(1,10000) AS g(i) CROSS JOIN encoded_value e; +DO $$ DECLARE n INT; BEGIN + SELECT count(*) INTO n FROM pg_locks WHERE pid=pg_backend_pid() AND locktype='advisory' + AND classid=1129530963 AND objsubid=2; + IF n=0 OR n>256 THEN RAISE EXCEPTION 'unbounded or missing row locks: %',n; END IF; + IF (SELECT count(*) FROM t WHERE id LIKE 'bulk-%')<>10000 THEN RAISE EXCEPTION 'bulk rows missing'; END IF; +END $$; +ROLLBACK; +DO $$ BEGIN + IF EXISTS (SELECT FROM pg_locks WHERE pid=pg_backend_pid() AND locktype='advisory' AND classid=1129530963) + THEN RAISE EXCEPTION 'row locks survived rollback'; END IF; +END $$; +\echo [PASS] (:testid) isolation errors, retry and bounded locks for 10000 rows + +SELECT dblink_disconnect('x'); +SELECT dblink_disconnect('y'); +\echo [PASS] (:testid) serial and concurrent merges converge on higher/3 +\connect postgres +DROP DATABASE cloudsync_test_67; diff --git a/test/postgresql/full_test.sql b/test/postgresql/full_test.sql index fe21bd0f..a91ad568 100644 --- a/test/postgresql/full_test.sql +++ b/test/postgresql/full_test.sql @@ -73,6 +73,7 @@ \ir 63_deep_savepoints.sql \ir 64_block_rewrite_leftovers.sql \ir 66_db_version_per_transaction.sql +\ir 67_concurrent_merge.sql -- 'Test summary' \echo '\nTest summary:'