Skip to content
Open
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 API.md
Original file line number Diff line number Diff line change
Expand Up @@ -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:**
Expand Down
20 changes: 20 additions & 0 deletions docker/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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).
Expand Down
31 changes: 31 additions & 0 deletions scripts/test-postgres-docker.sh
Original file line number Diff line number Diff line change
@@ -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
5 changes: 5 additions & 0 deletions src/cloudsync.c
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
1 change: 1 addition & 0 deletions src/database.h
Original file line number Diff line number Diff line change
Expand Up @@ -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);
Expand Down
26 changes: 26 additions & 0 deletions src/postgresql/database_postgresql.c
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand Down Expand Up @@ -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);
}
Expand Down
5 changes: 5 additions & 0 deletions src/sqlite/database_sqlite.c
Original file line number Diff line number Diff line change
Expand Up @@ -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");
Expand Down
221 changes: 221 additions & 0 deletions test/postgresql/67_concurrent_merge.sql
Original file line number Diff line number Diff line change
@@ -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;
1 change: 1 addition & 0 deletions test/postgresql/full_test.sql
Original file line number Diff line number Diff line change
Expand Up @@ -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:'
Expand Down
Loading