From e143be3ac446d06fb6412e881c827ae9046f0353 Mon Sep 17 00:00:00 2001 From: Andrea Donetti Date: Thu, 19 Mar 2026 15:58:34 -0600 Subject: [PATCH 1/9] test: add "Payload Apply Lock Test" --- test/postgresql/39_concurrent_write_apply.sql | 179 ++++++++++++++ test/postgresql/full_test.sql | 1 + test/unit.c | 221 ++++++++++++++++++ 3 files changed, 401 insertions(+) create mode 100644 test/postgresql/39_concurrent_write_apply.sql diff --git a/test/postgresql/39_concurrent_write_apply.sql b/test/postgresql/39_concurrent_write_apply.sql new file mode 100644 index 0000000..3d34856 --- /dev/null +++ b/test/postgresql/39_concurrent_write_apply.sql @@ -0,0 +1,179 @@ +-- 'Test concurrent write lock during payload apply' +-- NOTE: The lock-contention portion requires dblink with table access. +-- On environments where dblink cannot lock the table (e.g. Supabase), +-- the lock test is skipped and only apply + consistency are verified. + +\set testid '39' +\ir helper_test_init.sql + +\connect postgres +\ir helper_psql_conn_setup.sql +DROP DATABASE IF EXISTS cloudsync_test_39_a; +DROP DATABASE IF EXISTS cloudsync_test_39_b; +CREATE DATABASE cloudsync_test_39_a; +CREATE DATABASE cloudsync_test_39_b; + +-- Setup db_a +\connect cloudsync_test_39_a +\ir helper_psql_conn_setup.sql +CREATE EXTENSION IF NOT EXISTS cloudsync; +CREATE TABLE concurrent_tbl (id TEXT PRIMARY KEY, val TEXT); +SELECT cloudsync_init('concurrent_tbl', 'CLS', true) AS _init_a \gset + +-- Setup db_b +\connect cloudsync_test_39_b +\ir helper_psql_conn_setup.sql +CREATE EXTENSION IF NOT EXISTS cloudsync; +CREATE TABLE concurrent_tbl (id TEXT PRIMARY KEY, val TEXT); +SELECT cloudsync_init('concurrent_tbl', 'CLS', true) AS _init_b \gset + +-- Insert row1 on db_a and sync to db_b +\connect cloudsync_test_39_a +INSERT INTO concurrent_tbl VALUES ('row1', 'val_a'); + +SELECT CASE WHEN payload IS NULL OR octet_length(payload) = 0 + THEN '' + ELSE '\x' || encode(payload, 'hex') + END AS payload_init, + (payload IS NOT NULL AND octet_length(payload) > 0) AS payload_init_ok +FROM ( + SELECT cloudsync_payload_encode(tbl, pk, col_name, col_value, col_version, + db_version, site_id, cl, seq) AS payload + FROM cloudsync_changes WHERE site_id = cloudsync_siteid() +) AS p \gset + +\connect cloudsync_test_39_b +\if :payload_init_ok +SELECT cloudsync_payload_apply(decode(substr(:'payload_init', 3), 'hex')) AS _apply_init \gset +\endif + +-- Update row1 on db_a +\connect cloudsync_test_39_a +UPDATE concurrent_tbl SET val = 'val_a_updated' WHERE id = 'row1'; + +SELECT CASE WHEN payload IS NULL OR octet_length(payload) = 0 + THEN '' + ELSE '\x' || encode(payload, 'hex') + END AS payload_upd, + (payload IS NOT NULL AND octet_length(payload) > 0) AS payload_upd_ok +FROM ( + SELECT cloudsync_payload_encode(tbl, pk, col_name, col_value, col_version, + db_version, site_id, cl, seq) AS payload + FROM cloudsync_changes WHERE site_id = cloudsync_siteid() +) AS p \gset + +-- Try to set up dblink and acquire a table lock +\connect cloudsync_test_39_b +CREATE EXTENSION IF NOT EXISTS dblink; + +SELECT dblink_connect('locker', 'dbname=cloudsync_test_39_b') AS _conn \gset +SELECT dblink_exec('locker', 'BEGIN') AS _begin \gset + +-- Try to acquire EXCLUSIVE lock — if this fails (e.g. permission denied on +-- Supabase), _lock won't be set and we skip the lock-contention test +\unset _lock +SELECT dblink_exec('locker', 'LOCK TABLE concurrent_tbl IN EXCLUSIVE MODE') AS _lock \gset + +\if :{?_lock} +-- ===== Lock acquired — run lock-contention test ===== + +BEGIN; +\set ON_ERROR_ROLLBACK on +SET LOCAL lock_timeout = '500ms'; + +\if :payload_upd_ok +SELECT cloudsync_payload_apply(decode(substr(:'payload_upd', 3), 'hex')) AS _blocked_apply \gset +\endif + +COMMIT; +\set ON_ERROR_ROLLBACK off + +-- row1 should still have the OLD value because the apply was blocked +SELECT val AS row1_val_check FROM concurrent_tbl WHERE id = 'row1' \gset +SELECT (:'row1_val_check' = 'val_a') AS blocked_ok \gset +\if :blocked_ok +\echo [PASS] (:testid) Apply correctly blocked by concurrent table lock +\else +\echo [FAIL] (:testid) Expected val_a (blocked), got :'row1_val_check' +SELECT (:fail::int + 1) AS fail \gset +\endif + +-- Release the table lock +SELECT dblink_exec('locker', 'COMMIT') AS _release \gset +SELECT dblink_disconnect('locker') AS _disconn \gset + +-- Retry apply — should succeed now +\if :payload_upd_ok +SELECT cloudsync_payload_apply(decode(substr(:'payload_upd', 3), 'hex')) AS _apply_retry \gset +\endif + +SELECT val AS row1_val FROM concurrent_tbl WHERE id = 'row1' \gset +SELECT (:'row1_val' = 'val_a_updated') AS retry_ok \gset +\if :retry_ok +\echo [PASS] (:testid) Apply succeeded after lock released +\else +\echo [FAIL] (:testid) Apply after unlock - expected val_a_updated, got :'row1_val' +SELECT (:fail::int + 1) AS fail \gset +\endif + +\else +-- ===== Lock failed — skip contention test, apply directly ===== +\echo [SKIP] (:testid) Lock-contention test skipped (dblink cannot lock table) + +-- Clean up the dblink connection (transaction is aborted) +SELECT dblink_exec('locker', 'ROLLBACK') AS _rollback \gset +SELECT dblink_disconnect('locker') AS _disconn \gset + +\if :payload_upd_ok +SELECT cloudsync_payload_apply(decode(substr(:'payload_upd', 3), 'hex')) AS _apply_direct \gset +\endif + +SELECT val AS row1_val FROM concurrent_tbl WHERE id = 'row1' \gset +SELECT (:'row1_val' = 'val_a_updated') AS direct_ok \gset +\if :direct_ok +\echo [PASS] (:testid) Apply succeeded (no lock contention) +\else +\echo [FAIL] (:testid) Apply failed - expected val_a_updated, got :'row1_val' +SELECT (:fail::int + 1) AS fail \gset +\endif + +\endif + +-- Full cross-sync for consistency +SELECT CASE WHEN payload IS NULL OR octet_length(payload) = 0 + THEN '' + ELSE '\x' || encode(payload, 'hex') + END AS payload_b_final, + (payload IS NOT NULL AND octet_length(payload) > 0) AS payload_b_final_ok +FROM ( + SELECT cloudsync_payload_encode(tbl, pk, col_name, col_value, col_version, + db_version, site_id, cl, seq) AS payload + FROM cloudsync_changes WHERE site_id = cloudsync_siteid() +) AS p \gset + +\connect cloudsync_test_39_a +\if :payload_b_final_ok +SELECT cloudsync_payload_apply(decode(substr(:'payload_b_final', 3), 'hex')) AS _apply_final \gset +\endif + +SELECT md5(COALESCE(string_agg(id || ':' || COALESCE(val, ''), ',' ORDER BY id), '')) AS hash_a +FROM concurrent_tbl \gset + +\connect cloudsync_test_39_b +SELECT md5(COALESCE(string_agg(id || ':' || COALESCE(val, ''), ',' ORDER BY id), '')) AS hash_b +FROM concurrent_tbl \gset + +SELECT (:'hash_a' = :'hash_b') AS consistency_ok \gset +\if :consistency_ok +\echo [PASS] (:testid) Cross-database consistency verified +\else +\echo [FAIL] (:testid) Consistency failed (hash_a=:'hash_a' hash_b=:'hash_b') +SELECT (:fail::int + 1) AS fail \gset +\endif + +-- Cleanup +\ir helper_test_cleanup.sql +\if :should_cleanup +DROP DATABASE IF EXISTS cloudsync_test_39_a; +DROP DATABASE IF EXISTS cloudsync_test_39_b; +\endif diff --git a/test/postgresql/full_test.sql b/test/postgresql/full_test.sql index d02440a..ba69198 100644 --- a/test/postgresql/full_test.sql +++ b/test/postgresql/full_test.sql @@ -46,6 +46,7 @@ \ir 36_block_lww_round3.sql \ir 37_block_lww_round4.sql \ir 38_block_lww_round5.sql +\ir 39_concurrent_write_apply.sql -- 'Test summary' \echo '\nTest summary:' diff --git a/test/unit.c b/test/unit.c index 0487f9d..8b78dc4 100644 --- a/test/unit.c +++ b/test/unit.c @@ -6436,6 +6436,226 @@ bool do_test_merge_json_columns (int nclients, bool print_result, bool cleanup_d } // Test concurrent merge attempts +bool do_test_payload_apply_concurrent_write (int nclients, bool print_result, bool cleanup_databases) { + sqlite3 *db[MAX_SIMULATED_CLIENTS] = {NULL}; + sqlite3 *db_target2 = NULL; + sqlite3_stmt *select_stmt = NULL; + sqlite3_stmt *apply_stmt = NULL; + bool result = false; + int rc = SQLITE_OK; + + memset(db, 0, sizeof(sqlite3 *) * MAX_SIMULATED_CLIENTS); + if (nclients < 2) nclients = 2; + if (nclients > 2) nclients = 2; // this test uses exactly 2 databases + + time_t timestamp = time(NULL); + int saved_counter = test_counter; + + // create two file-based databases: db[0]=src, db[1]=target + for (int i = 0; i < nclients; ++i) { + db[i] = do_create_database_file(i, timestamp, test_counter++); + if (db[i] == NULL) return false; + + rc = sqlite3_exec(db[i], "CREATE TABLE concurrent_tbl (id TEXT PRIMARY KEY, val TEXT);", NULL, NULL, NULL); + if (rc != SQLITE_OK) goto finalize; + + rc = sqlite3_exec(db[i], "SELECT cloudsync_init('concurrent_tbl', 'cls', 1);", NULL, NULL, NULL); + if (rc != SQLITE_OK) goto finalize; + } + + // insert data on src (db[0]) + rc = sqlite3_exec(db[0], "INSERT INTO concurrent_tbl VALUES ('row1', 'hello');", NULL, NULL, NULL); + if (rc != SQLITE_OK) goto finalize; + rc = sqlite3_exec(db[0], "INSERT INTO concurrent_tbl VALUES ('row2', 'world');", NULL, NULL, NULL); + if (rc != SQLITE_OK) goto finalize; + + // extract payload from db[0] + const char *encode_sql = "SELECT cloudsync_payload_encode(tbl, pk, col_name, col_value, col_version, db_version, site_id, cl, seq) FROM cloudsync_changes WHERE site_id=cloudsync_siteid();"; + rc = sqlite3_prepare_v2(db[0], encode_sql, -1, &select_stmt, NULL); + if (rc != SQLITE_OK) goto finalize; + + rc = sqlite3_step(select_stmt); + if (rc != SQLITE_ROW) goto finalize; + + const void *payload_data = sqlite3_column_blob(select_stmt, 0); + int payload_size = sqlite3_column_bytes(select_stmt, 0); + if (payload_data == NULL || payload_size == 0) goto finalize; + + // copy payload since we'll need it after finalizing select_stmt + void *payload_copy = malloc(payload_size); + if (payload_copy == NULL) goto finalize; + memcpy(payload_copy, payload_data, payload_size); + sqlite3_finalize(select_stmt); + select_stmt = NULL; + + // open second connection to same file as db[1] (target) + { + char buf[256]; + do_build_database_path(buf, 1, timestamp, saved_counter + 1); + rc = sqlite3_open(buf, &db_target2); + if (rc != SQLITE_OK) { + printf("Error opening db_target2: %s\n", sqlite3_errmsg(db_target2)); + free(payload_copy); + goto finalize; + } + sqlite3_exec(db_target2, "PRAGMA journal_mode=WAL;", NULL, NULL, NULL); + sqlite3_cloudsync_init(db_target2, NULL, NULL); + } + + // on db[1] (target): begin immediate to hold write lock + rc = sqlite3_exec(db[1], "BEGIN IMMEDIATE;", NULL, NULL, NULL); + if (rc != SQLITE_OK) { + printf("BEGIN IMMEDIATE failed: %s\n", sqlite3_errmsg(db[1])); + free(payload_copy); + goto finalize; + } + rc = sqlite3_exec(db[1], "INSERT INTO concurrent_tbl VALUES ('blocker', 'blocking');", NULL, NULL, NULL); + if (rc != SQLITE_OK) { + printf("Blocker INSERT failed: %s\n", sqlite3_errmsg(db[1])); + free(payload_copy); + goto finalize; + } + + // on db_target2: try to apply payload — should fail with BUSY + rc = sqlite3_prepare_v2(db_target2, "SELECT cloudsync_payload_decode(?);", -1, &apply_stmt, NULL); + if (rc != SQLITE_OK) { + printf("Prepare apply failed: %s\n", sqlite3_errmsg(db_target2)); + free(payload_copy); + goto finalize; + } + rc = sqlite3_bind_blob(apply_stmt, 1, payload_copy, payload_size, SQLITE_STATIC); + if (rc != SQLITE_OK) { + printf("Bind failed: %s\n", sqlite3_errmsg(db_target2)); + free(payload_copy); + goto finalize; + } + + // set a short busy timeout so it doesn't wait forever (0 = fail immediately) + sqlite3_busy_timeout(db_target2, 0); + + rc = sqlite3_step(apply_stmt); + if (rc == SQLITE_ROW || rc == SQLITE_DONE) { + printf("Expected BUSY error but apply succeeded (rc=%d)\n", rc); + free(payload_copy); + goto finalize; + } + + // verify we got a BUSY-related error + int errcode = sqlite3_errcode(db_target2); + if (errcode != SQLITE_BUSY && errcode != SQLITE_LOCKED) { + printf("Expected SQLITE_BUSY or SQLITE_LOCKED but got error %d: %s\n", errcode, sqlite3_errmsg(db_target2)); + free(payload_copy); + goto finalize; + } + + if (print_result) { + printf(" Step 1: Apply blocked as expected (errcode=%d: %s)\n", errcode, sqlite3_errmsg(db_target2)); + } + + // release the write lock on db[1] + rc = sqlite3_exec(db[1], "COMMIT;", NULL, NULL, NULL); + if (rc != SQLITE_OK) { + printf("COMMIT failed: %s\n", sqlite3_errmsg(db[1])); + free(payload_copy); + goto finalize; + } + + // retry: reset and step again — should succeed now + sqlite3_reset(apply_stmt); + sqlite3_busy_timeout(db_target2, 5000); // give it time now + + rc = sqlite3_step(apply_stmt); + if (rc != SQLITE_ROW) { + printf("Expected SQLITE_ROW on retry but got %d: %s\n", rc, sqlite3_errmsg(db_target2)); + free(payload_copy); + goto finalize; + } + + if (print_result) { + printf(" Step 2: Apply succeeded after lock released\n"); + } + + sqlite3_finalize(apply_stmt); + apply_stmt = NULL; + free(payload_copy); + payload_copy = NULL; + + // verify: db_target2 should have row1, row2 (from payload) + blocker (from db[1]) + { + sqlite3_stmt *count_stmt = NULL; + rc = sqlite3_prepare_v2(db_target2, "SELECT COUNT(*) FROM concurrent_tbl;", -1, &count_stmt, NULL); + if (rc != SQLITE_OK) goto finalize; + rc = sqlite3_step(count_stmt); + if (rc != SQLITE_ROW) { sqlite3_finalize(count_stmt); goto finalize; } + int count = sqlite3_column_int(count_stmt, 0); + sqlite3_finalize(count_stmt); + if (count != 3) { + printf("Expected 3 rows but got %d\n", count); + goto finalize; + } + if (print_result) { + printf(" Step 3: Target has %d rows (expected 3)\n", count); + } + } + + // full consistency: merge all databases using payload + // first close db_target2 to avoid lock conflicts during merge + close_db(db_target2); + db_target2 = NULL; + + // merge db[0] <-> db[1] in both directions + if (do_merge_using_payload(db[0], db[1], false, true) == false) { + printf("Merge src->target failed\n"); + goto finalize; + } + if (do_merge_using_payload(db[1], db[0], false, true) == false) { + printf("Merge target->src failed\n"); + goto finalize; + } + + // verify consistency + { + const char *sql = "SELECT * FROM concurrent_tbl ORDER BY id;"; + bool cmp = do_compare_queries(db[0], sql, db[1], sql, -1, -1, print_result); + if (!cmp) { + printf("Consistency check failed between src and target\n"); + goto finalize; + } + } + + if (print_result) { + printf(" Step 4: Full consistency verified\n"); + } + + result = true; + +finalize: + if (select_stmt) sqlite3_finalize(select_stmt); + if (apply_stmt) sqlite3_finalize(apply_stmt); + if (db_target2) close_db(db_target2); + for (int i = 0; i < nclients; ++i) { + if (rc != SQLITE_OK && db[i] && (sqlite3_errcode(db[i]) != SQLITE_OK)) + printf("do_test_payload_apply_concurrent_write error: %s\n", sqlite3_errmsg(db[i])); + if (db[i]) { + if (sqlite3_get_autocommit(db[i]) == 0) { + result = false; + printf("do_test_payload_apply_concurrent_write error: db %d is in transaction\n", i); + } + int counter = close_db(db[i]); + if (counter > 0) { + result = false; + printf("do_test_payload_apply_concurrent_write error: db %d has %d unterminated statements\n", i, counter); + } + } + if (cleanup_databases) { + char buf[256]; + do_build_database_path(buf, i, timestamp, saved_counter++); + file_delete_internal(buf); + } + } + return result; +} + bool do_test_merge_concurrent_attempts (int nclients, bool print_result, bool cleanup_databases) { sqlite3 *db[MAX_SIMULATED_CLIENTS] = {NULL}; bool result = false; @@ -10254,6 +10474,7 @@ int main (int argc, const char * argv[]) { result += test_report("Merge Index Consistency:", do_test_merge_index_consistency(2, print_result, cleanup_databases)); result += test_report("Merge JSON Columns:", do_test_merge_json_columns(2, print_result, cleanup_databases)); result += test_report("Merge Concurrent Attempts:", do_test_merge_concurrent_attempts(3, print_result, cleanup_databases)); + result += test_report("Payload Apply Lock Test:", do_test_payload_apply_concurrent_write(2, print_result, cleanup_databases)); result += test_report("Merge Composite PK 10 Clients:", do_test_merge_composite_pk_10_clients(10, print_result, cleanup_databases)); result += test_report("PriKey NULL Test:", do_test_prikey(2, print_result, cleanup_databases)); result += test_report("Test Double Init:", do_test_double_init(2, cleanup_databases)); From 44ea471ed535967779d4ea58311c3196470a879b Mon Sep 17 00:00:00 2001 From: Andrea Donetti Date: Thu, 19 Mar 2026 16:51:25 -0600 Subject: [PATCH 2/9] test: add edge-case tests for CRDT sync correctness and error handling Add 7 new SQLite unit tests and 7 new PostgreSQL test files covering: - DWS/AWS algorithm rejection (unsupported CRDT algos return clean errors) - Corrupted payload handling (empty, garbage, truncated, bit-flipped) - Payload apply idempotency (3x apply produces identical results) - Causal-length tie-breaking determinism (3-way concurrent update convergence) - Delete/resurrect with out-of-order payload delivery - Large composite primary key (5-column mixed-type PK roundtrip) - PostgreSQL-specific type roundtrips (JSONB, TIMESTAMPTZ, NUMERIC, BYTEA) - Schema hash mismatch detection (ALTER TABLE without cloudsync workflow) --- test/postgresql/40_unsupported_algorithms.sql | 80 +++ test/postgresql/41_corrupted_payload.sql | 126 ++++ test/postgresql/42_payload_idempotency.sql | 88 +++ .../43_delete_resurrect_ordering.sql | 147 +++++ test/postgresql/44_large_composite_pk.sql | 142 +++++ test/postgresql/45_pg_specific_types.sql | 176 ++++++ test/postgresql/46_schema_hash_mismatch.sql | 96 +++ test/postgresql/full_test.sql | 7 + test/unit.c | 555 ++++++++++++++++++ 9 files changed, 1417 insertions(+) create mode 100644 test/postgresql/40_unsupported_algorithms.sql create mode 100644 test/postgresql/41_corrupted_payload.sql create mode 100644 test/postgresql/42_payload_idempotency.sql create mode 100644 test/postgresql/43_delete_resurrect_ordering.sql create mode 100644 test/postgresql/44_large_composite_pk.sql create mode 100644 test/postgresql/45_pg_specific_types.sql create mode 100644 test/postgresql/46_schema_hash_mismatch.sql diff --git a/test/postgresql/40_unsupported_algorithms.sql b/test/postgresql/40_unsupported_algorithms.sql new file mode 100644 index 0000000..d312ff6 --- /dev/null +++ b/test/postgresql/40_unsupported_algorithms.sql @@ -0,0 +1,80 @@ +-- Test unsupported CRDT algorithms (DWS, AWS) +-- Verifies that cloudsync_init rejects DWS and AWS with clear errors +-- and that no metadata tables are created. + +\set testid '40' +\ir helper_test_init.sql + +\connect postgres +\ir helper_psql_conn_setup.sql +DROP DATABASE IF EXISTS cloudsync_test_40; +CREATE DATABASE cloudsync_test_40; + +\connect cloudsync_test_40 +\ir helper_psql_conn_setup.sql +CREATE EXTENSION IF NOT EXISTS cloudsync; + +CREATE TABLE test_dws (id TEXT PRIMARY KEY, val TEXT); +CREATE TABLE test_aws (id TEXT PRIMARY KEY, val TEXT); + +-- Test DWS rejection +DO $$ +BEGIN + PERFORM cloudsync_init('test_dws', 'dws', true); + RAISE EXCEPTION 'cloudsync_init with dws should have failed'; +EXCEPTION WHEN OTHERS THEN + IF SQLERRM NOT LIKE '%not yet supported%' THEN + RAISE EXCEPTION 'Unexpected error for dws: %', SQLERRM; + END IF; +END $$; + +-- Verify no companion table was created for DWS +SELECT COUNT(*) = 0 AS no_dws_meta +FROM information_schema.tables +WHERE table_name = 'test_dws_cloudsync' \gset +\if :no_dws_meta +\echo [PASS] (:testid) DWS rejected - no metadata table created +\else +\echo [FAIL] (:testid) DWS metadata table should not exist +SELECT (:fail::int + 1) AS fail \gset +\endif + +-- Test AWS rejection +DO $$ +BEGIN + PERFORM cloudsync_init('test_aws', 'aws', true); + RAISE EXCEPTION 'cloudsync_init with aws should have failed'; +EXCEPTION WHEN OTHERS THEN + IF SQLERRM NOT LIKE '%not yet supported%' THEN + RAISE EXCEPTION 'Unexpected error for aws: %', SQLERRM; + END IF; +END $$; + +-- Verify no companion table was created for AWS +SELECT COUNT(*) = 0 AS no_aws_meta +FROM information_schema.tables +WHERE table_name = 'test_aws_cloudsync' \gset +\if :no_aws_meta +\echo [PASS] (:testid) AWS rejected - no metadata table created +\else +\echo [FAIL] (:testid) AWS metadata table should not exist +SELECT (:fail::int + 1) AS fail \gset +\endif + +-- Verify CLS still works (sanity check) +SELECT cloudsync_init('test_dws', 'cls', true) AS _init_cls \gset +SELECT COUNT(*) = 1 AS cls_meta_ok +FROM information_schema.tables +WHERE table_name = 'test_dws_cloudsync' \gset +\if :cls_meta_ok +\echo [PASS] (:testid) CLS init works after DWS/AWS rejection +\else +\echo [FAIL] (:testid) CLS init should work +SELECT (:fail::int + 1) AS fail \gset +\endif + +-- Cleanup +\ir helper_test_cleanup.sql +\if :should_cleanup +DROP DATABASE IF EXISTS cloudsync_test_40; +\endif diff --git a/test/postgresql/41_corrupted_payload.sql b/test/postgresql/41_corrupted_payload.sql new file mode 100644 index 0000000..0e73c10 --- /dev/null +++ b/test/postgresql/41_corrupted_payload.sql @@ -0,0 +1,126 @@ +-- Test corrupted payload handling +-- Verifies that cloudsync_payload_apply rejects corrupted payloads +-- without crashing or corrupting state. + +\set testid '41' +\ir helper_test_init.sql + +\connect postgres +\ir helper_psql_conn_setup.sql +DROP DATABASE IF EXISTS cloudsync_test_41_src; +DROP DATABASE IF EXISTS cloudsync_test_41_dst; +CREATE DATABASE cloudsync_test_41_src; +CREATE DATABASE cloudsync_test_41_dst; + +-- Setup source database with data +\connect cloudsync_test_41_src +\ir helper_psql_conn_setup.sql +CREATE EXTENSION IF NOT EXISTS cloudsync; +CREATE TABLE test_tbl (id TEXT PRIMARY KEY, val TEXT); +SELECT cloudsync_init('test_tbl', 'CLS', true) AS _init_src \gset +INSERT INTO test_tbl VALUES ('id1', 'value1'); +INSERT INTO test_tbl VALUES ('id2', 'value2'); + +-- Get a valid payload +SELECT encode(cloudsync_payload_encode(tbl, pk, col_name, col_value, col_version, db_version, site_id, cl, seq), 'hex') AS valid_payload_hex +FROM cloudsync_changes +WHERE site_id = cloudsync_siteid() \gset + +-- Setup destination database +\connect cloudsync_test_41_dst +\ir helper_psql_conn_setup.sql +CREATE EXTENSION IF NOT EXISTS cloudsync; +CREATE TABLE test_tbl (id TEXT PRIMARY KEY, val TEXT); +SELECT cloudsync_init('test_tbl', 'CLS', true) AS _init_dst \gset + +-- Record initial state +SELECT COUNT(*) AS initial_count FROM test_tbl \gset + +-- Test 1: Empty blob (zero bytes) +DO $$ +BEGIN + PERFORM cloudsync_payload_apply(''::bytea); + -- If it returns without error with 0 rows, that's also acceptable +EXCEPTION WHEN OTHERS THEN + -- Expected: error on empty payload + NULL; +END $$; + +SELECT COUNT(*) AS count_after_empty FROM test_tbl \gset +SELECT (:count_after_empty::int = :initial_count::int) AS empty_blob_ok \gset +\if :empty_blob_ok +\echo [PASS] (:testid) Empty blob rejected - table unchanged +\else +\echo [FAIL] (:testid) Empty blob corrupted table state +SELECT (:fail::int + 1) AS fail \gset +\endif + +-- Test 2: Random garbage bytes +DO $$ +BEGIN + PERFORM cloudsync_payload_apply(decode('deadbeefcafebabe0102030405060708', 'hex')); +EXCEPTION WHEN OTHERS THEN + -- Expected: error on garbage payload + NULL; +END $$; + +SELECT COUNT(*) AS count_after_garbage FROM test_tbl \gset +SELECT (:count_after_garbage::int = :initial_count::int) AS garbage_ok \gset +\if :garbage_ok +\echo [PASS] (:testid) Garbage bytes rejected - table unchanged +\else +\echo [FAIL] (:testid) Garbage bytes corrupted table state +SELECT (:fail::int + 1) AS fail \gset +\endif + +-- Test 3: Truncated payload (first 10 bytes of valid payload) +-- Build truncated hex at top level using psql variable interpolation +SELECT substr(:'valid_payload_hex', 1, 20) AS truncated_hex \gset +SELECT cloudsync_payload_apply(decode(:'truncated_hex', 'hex')) AS _apply_truncated \gset +-- If the above errors, psql continues (ON_ERROR_STOP is off) + +SELECT COUNT(*) AS count_after_truncated FROM test_tbl \gset +SELECT (:count_after_truncated::int = :initial_count::int) AS truncated_ok \gset +\if :truncated_ok +\echo [PASS] (:testid) Truncated payload rejected - table unchanged +\else +\echo [FAIL] (:testid) Truncated payload corrupted table state +SELECT (:fail::int + 1) AS fail \gset +\endif + +-- Test 4: Valid payload with flipped byte in the middle +-- Compute corrupted payload at top level: flip one byte via XOR with FF +SELECT + substr(:'valid_payload_hex', 1, length(:'valid_payload_hex') / 2 - 1) + || lpad(to_hex(get_byte(decode(substr(:'valid_payload_hex', length(:'valid_payload_hex') / 2, 2), 'hex'), 0) # 255), 2, '0') + || substr(:'valid_payload_hex', length(:'valid_payload_hex') / 2 + 2) + AS corrupted_hex \gset +SELECT cloudsync_payload_apply(decode(:'corrupted_hex', 'hex')) AS _apply_corrupted \gset +-- If the above errors, psql continues (ON_ERROR_STOP is off) + +SELECT COUNT(*) AS count_after_flipped FROM test_tbl \gset +SELECT (:count_after_flipped::int = :initial_count::int) AS flipped_ok \gset +\if :flipped_ok +\echo [PASS] (:testid) Flipped-byte payload rejected - table unchanged +\else +\echo [FAIL] (:testid) Flipped-byte payload corrupted table state +SELECT (:fail::int + 1) AS fail \gset +\endif + +-- Test 5: Now apply the VALID payload to confirm it still works +SELECT cloudsync_payload_apply(decode(:'valid_payload_hex', 'hex')) AS valid_apply \gset +SELECT COUNT(*) AS count_after_valid FROM test_tbl \gset +SELECT (:count_after_valid::int = 2) AS valid_ok \gset +\if :valid_ok +\echo [PASS] (:testid) Valid payload applied successfully after corrupted attempts +\else +\echo [FAIL] (:testid) Valid payload failed after corrupted attempts - count: :count_after_valid +SELECT (:fail::int + 1) AS fail \gset +\endif + +-- Cleanup +\ir helper_test_cleanup.sql +\if :should_cleanup +DROP DATABASE IF EXISTS cloudsync_test_41_src; +DROP DATABASE IF EXISTS cloudsync_test_41_dst; +\endif diff --git a/test/postgresql/42_payload_idempotency.sql b/test/postgresql/42_payload_idempotency.sql new file mode 100644 index 0000000..43617dd --- /dev/null +++ b/test/postgresql/42_payload_idempotency.sql @@ -0,0 +1,88 @@ +-- Test payload apply idempotency +-- Applying the same payload multiple times must produce identical results. + +\set testid '42' +\ir helper_test_init.sql + +\connect postgres +\ir helper_psql_conn_setup.sql +DROP DATABASE IF EXISTS cloudsync_test_42_src; +DROP DATABASE IF EXISTS cloudsync_test_42_dst; +CREATE DATABASE cloudsync_test_42_src; +CREATE DATABASE cloudsync_test_42_dst; + +-- Setup source with data +\connect cloudsync_test_42_src +\ir helper_psql_conn_setup.sql +CREATE EXTENSION IF NOT EXISTS cloudsync; +CREATE TABLE test_tbl (id TEXT PRIMARY KEY, val TEXT, num INTEGER); +SELECT cloudsync_init('test_tbl', 'CLS', true) AS _init_src \gset +INSERT INTO test_tbl VALUES ('id1', 'hello', 10); +INSERT INTO test_tbl VALUES ('id2', 'world', 20); +UPDATE test_tbl SET val = 'hello_updated' WHERE id = 'id1'; + +-- Encode payload +SELECT encode(cloudsync_payload_encode(tbl, pk, col_name, col_value, col_version, db_version, site_id, cl, seq), 'hex') AS payload_hex +FROM cloudsync_changes +WHERE site_id = cloudsync_siteid() \gset + +-- Setup destination +\connect cloudsync_test_42_dst +\ir helper_psql_conn_setup.sql +CREATE EXTENSION IF NOT EXISTS cloudsync; +CREATE TABLE test_tbl (id TEXT PRIMARY KEY, val TEXT, num INTEGER); +SELECT cloudsync_init('test_tbl', 'CLS', true) AS _init_dst \gset + +-- Apply #1 +SELECT cloudsync_payload_apply(decode(:'payload_hex', 'hex')) AS apply_1 \gset +SELECT md5(COALESCE(string_agg(id || ':' || COALESCE(val, '') || ':' || COALESCE(num::text, ''), ',' ORDER BY id), '')) AS hash_1 +FROM test_tbl \gset +SELECT COUNT(*) AS count_1 FROM test_tbl \gset + +-- Apply #2 +SELECT cloudsync_payload_apply(decode(:'payload_hex', 'hex')) AS apply_2 \gset +SELECT md5(COALESCE(string_agg(id || ':' || COALESCE(val, '') || ':' || COALESCE(num::text, ''), ',' ORDER BY id), '')) AS hash_2 +FROM test_tbl \gset +SELECT COUNT(*) AS count_2 FROM test_tbl \gset + +-- Apply #3 +SELECT cloudsync_payload_apply(decode(:'payload_hex', 'hex')) AS apply_3 \gset +SELECT md5(COALESCE(string_agg(id || ':' || COALESCE(val, '') || ':' || COALESCE(num::text, ''), ',' ORDER BY id), '')) AS hash_3 +FROM test_tbl \gset +SELECT COUNT(*) AS count_3 FROM test_tbl \gset + +-- Verify row count stays constant +SELECT (:count_1::int = :count_2::int AND :count_2::int = :count_3::int) AS count_stable \gset +\if :count_stable +\echo [PASS] (:testid) Row count stable across 3 applies (:count_1 rows) +\else +\echo [FAIL] (:testid) Row count changed: :count_1 -> :count_2 -> :count_3 +SELECT (:fail::int + 1) AS fail \gset +\endif + +-- Verify data hash is identical after each apply +SELECT (:'hash_1' = :'hash_2' AND :'hash_2' = :'hash_3') AS hash_stable \gset +\if :hash_stable +\echo [PASS] (:testid) Data hash identical across 3 applies +\else +\echo [FAIL] (:testid) Data hash changed: :hash_1 -> :hash_2 -> :hash_3 +SELECT (:fail::int + 1) AS fail \gset +\endif + +-- Verify data values are correct +SELECT COUNT(*) = 1 AS data_ok +FROM test_tbl +WHERE id = 'id1' AND val = 'hello_updated' AND num = 10 \gset +\if :data_ok +\echo [PASS] (:testid) Data values correct after idempotent applies +\else +\echo [FAIL] (:testid) Data values incorrect +SELECT (:fail::int + 1) AS fail \gset +\endif + +-- Cleanup +\ir helper_test_cleanup.sql +\if :should_cleanup +DROP DATABASE IF EXISTS cloudsync_test_42_src; +DROP DATABASE IF EXISTS cloudsync_test_42_dst; +\endif diff --git a/test/postgresql/43_delete_resurrect_ordering.sql b/test/postgresql/43_delete_resurrect_ordering.sql new file mode 100644 index 0000000..1689ea5 --- /dev/null +++ b/test/postgresql/43_delete_resurrect_ordering.sql @@ -0,0 +1,147 @@ +-- Test delete/resurrect with out-of-order payload delivery +-- Verifies CRDT causal length parity handles resurrection correctly +-- even when payloads arrive in non-causal order. + +\set testid '43' +\ir helper_test_init.sql + +\connect postgres +\ir helper_psql_conn_setup.sql +DROP DATABASE IF EXISTS cloudsync_test_43_a; +DROP DATABASE IF EXISTS cloudsync_test_43_b; +DROP DATABASE IF EXISTS cloudsync_test_43_c; +CREATE DATABASE cloudsync_test_43_a; +CREATE DATABASE cloudsync_test_43_b; +CREATE DATABASE cloudsync_test_43_c; + +-- Setup all three databases +\connect cloudsync_test_43_a +\ir helper_psql_conn_setup.sql +CREATE EXTENSION IF NOT EXISTS cloudsync; +CREATE TABLE test_tbl (id TEXT PRIMARY KEY, val TEXT); +SELECT cloudsync_init('test_tbl', 'CLS', true) AS _init_a \gset + +\connect cloudsync_test_43_b +\ir helper_psql_conn_setup.sql +CREATE EXTENSION IF NOT EXISTS cloudsync; +CREATE TABLE test_tbl (id TEXT PRIMARY KEY, val TEXT); +SELECT cloudsync_init('test_tbl', 'CLS', true) AS _init_b \gset + +\connect cloudsync_test_43_c +\ir helper_psql_conn_setup.sql +CREATE EXTENSION IF NOT EXISTS cloudsync; +CREATE TABLE test_tbl (id TEXT PRIMARY KEY, val TEXT); +SELECT cloudsync_init('test_tbl', 'CLS', true) AS _init_c \gset + +-- Round 1: A inserts row, sync to all +\connect cloudsync_test_43_a +INSERT INTO test_tbl VALUES ('row1', 'original'); +SELECT CASE WHEN payload IS NULL OR octet_length(payload) = 0 + THEN '' + ELSE '\x' || encode(payload, 'hex') + END AS payload_a_r1, + (payload IS NOT NULL AND octet_length(payload) > 0) AS payload_a_r1_ok +FROM ( + SELECT cloudsync_payload_encode(tbl, pk, col_name, col_value, col_version, db_version, site_id, cl, seq) AS payload + FROM cloudsync_changes + WHERE site_id = cloudsync_siteid() +) AS p \gset + +\connect cloudsync_test_43_b +\if :payload_a_r1_ok +SELECT cloudsync_payload_apply(decode(substr(:'payload_a_r1', 3), 'hex')) AS _apply \gset +\endif + +\connect cloudsync_test_43_c +\if :payload_a_r1_ok +SELECT cloudsync_payload_apply(decode(substr(:'payload_a_r1', 3), 'hex')) AS _apply \gset +\endif + +-- Round 2: A deletes row (CL goes 1->2) +\connect cloudsync_test_43_a +DELETE FROM test_tbl WHERE id = 'row1'; +SELECT CASE WHEN payload IS NULL OR octet_length(payload) = 0 + THEN '' + ELSE '\x' || encode(payload, 'hex') + END AS payload_a_r2, + (payload IS NOT NULL AND octet_length(payload) > 0) AS payload_a_r2_ok +FROM ( + SELECT cloudsync_payload_encode(tbl, pk, col_name, col_value, col_version, db_version, site_id, cl, seq) AS payload + FROM cloudsync_changes + WHERE site_id = cloudsync_siteid() +) AS p \gset + +-- Sync delete to B +\connect cloudsync_test_43_b +\if :payload_a_r2_ok +SELECT cloudsync_payload_apply(decode(substr(:'payload_a_r2', 3), 'hex')) AS _apply \gset +\endif + +-- Round 3: B re-inserts (CL goes 2->3, resurrection) +\connect cloudsync_test_43_b +INSERT INTO test_tbl VALUES ('row1', 'resurrected_by_b'); +SELECT CASE WHEN payload IS NULL OR octet_length(payload) = 0 + THEN '' + ELSE '\x' || encode(payload, 'hex') + END AS payload_b_r3, + (payload IS NOT NULL AND octet_length(payload) > 0) AS payload_b_r3_ok +FROM ( + SELECT cloudsync_payload_encode(tbl, pk, col_name, col_value, col_version, db_version, site_id, cl, seq) AS payload + FROM cloudsync_changes + WHERE site_id = cloudsync_siteid() +) AS p \gset + +-- C receives payloads in REVERSE order: B's resurrection FIRST, then A's delete +\connect cloudsync_test_43_c +\if :payload_b_r3_ok +SELECT cloudsync_payload_apply(decode(substr(:'payload_b_r3', 3), 'hex')) AS _apply_b \gset +\endif +\if :payload_a_r2_ok +SELECT cloudsync_payload_apply(decode(substr(:'payload_a_r2', 3), 'hex')) AS _apply_a \gset +\endif + +-- A receives B's resurrection +\connect cloudsync_test_43_a +\if :payload_b_r3_ok +SELECT cloudsync_payload_apply(decode(substr(:'payload_b_r3', 3), 'hex')) AS _apply_b \gset +\endif + +-- Final convergence check: all three should have row1 with 'resurrected_by_b' +\connect cloudsync_test_43_a +SELECT md5(COALESCE(string_agg(id || ':' || COALESCE(val, ''), ',' ORDER BY id), '')) AS hash_a +FROM test_tbl \gset + +\connect cloudsync_test_43_b +SELECT md5(COALESCE(string_agg(id || ':' || COALESCE(val, ''), ',' ORDER BY id), '')) AS hash_b +FROM test_tbl \gset + +\connect cloudsync_test_43_c +SELECT md5(COALESCE(string_agg(id || ':' || COALESCE(val, ''), ',' ORDER BY id), '')) AS hash_c +FROM test_tbl \gset + +SELECT (:'hash_a' = :'hash_b' AND :'hash_b' = :'hash_c') AS all_converge \gset +\if :all_converge +\echo [PASS] (:testid) All 3 databases converge after out-of-order delete/resurrect +\else +\echo [FAIL] (:testid) Databases diverged - A: :hash_a, B: :hash_b, C: :hash_c +SELECT (:fail::int + 1) AS fail \gset +\endif + +-- Verify the row exists (resurrection won) +SELECT COUNT(*) = 1 AS row_exists +FROM test_tbl +WHERE id = 'row1' \gset +\if :row_exists +\echo [PASS] (:testid) Resurrected row exists on C (received out of order) +\else +\echo [FAIL] (:testid) Resurrected row missing on C +SELECT (:fail::int + 1) AS fail \gset +\endif + +-- Cleanup +\ir helper_test_cleanup.sql +\if :should_cleanup +DROP DATABASE IF EXISTS cloudsync_test_43_a; +DROP DATABASE IF EXISTS cloudsync_test_43_b; +DROP DATABASE IF EXISTS cloudsync_test_43_c; +\endif diff --git a/test/postgresql/44_large_composite_pk.sql b/test/postgresql/44_large_composite_pk.sql new file mode 100644 index 0000000..21da26d --- /dev/null +++ b/test/postgresql/44_large_composite_pk.sql @@ -0,0 +1,142 @@ +-- Test large composite primary key (5 columns) +-- Verifies pk_encode/pk_decode handles complex multi-column PKs correctly. + +\set testid '44' +\ir helper_test_init.sql + +\connect postgres +\ir helper_psql_conn_setup.sql +DROP DATABASE IF EXISTS cloudsync_test_44_a; +DROP DATABASE IF EXISTS cloudsync_test_44_b; +CREATE DATABASE cloudsync_test_44_a; +CREATE DATABASE cloudsync_test_44_b; + +-- Setup Database A +\connect cloudsync_test_44_a +\ir helper_psql_conn_setup.sql +CREATE EXTENSION IF NOT EXISTS cloudsync; + +CREATE TABLE composite_pk_tbl ( + pk_text1 TEXT NOT NULL, + pk_int1 INTEGER NOT NULL, + pk_text2 TEXT NOT NULL, + pk_int2 INTEGER NOT NULL, + pk_text3 TEXT NOT NULL, + data_col TEXT, + num_col INTEGER, + PRIMARY KEY (pk_text1, pk_int1, pk_text2, pk_int2, pk_text3) +); + +SELECT cloudsync_init('composite_pk_tbl', 'CLS', true) AS _init_a \gset + +INSERT INTO composite_pk_tbl VALUES ('alpha', 1, 'beta', 100, 'gamma', 'data_a1', 42); +INSERT INTO composite_pk_tbl VALUES ('alpha', 2, 'beta', 200, 'delta', 'data_a2', 84); +INSERT INTO composite_pk_tbl VALUES ('x', 999, 'y', -1, 'z', 'edge_case', 0); + +-- Setup Database B +\connect cloudsync_test_44_b +\ir helper_psql_conn_setup.sql +CREATE EXTENSION IF NOT EXISTS cloudsync; + +CREATE TABLE composite_pk_tbl ( + pk_text1 TEXT NOT NULL, + pk_int1 INTEGER NOT NULL, + pk_text2 TEXT NOT NULL, + pk_int2 INTEGER NOT NULL, + pk_text3 TEXT NOT NULL, + data_col TEXT, + num_col INTEGER, + PRIMARY KEY (pk_text1, pk_int1, pk_text2, pk_int2, pk_text3) +); + +SELECT cloudsync_init('composite_pk_tbl', 'CLS', true) AS _init_b \gset + +INSERT INTO composite_pk_tbl VALUES ('alpha', 1, 'beta', 100, 'gamma', 'data_b1', 99); +INSERT INTO composite_pk_tbl VALUES ('foo', 3, 'bar', 300, 'baz', 'data_b2', 77); + +-- Encode and exchange payloads +\connect cloudsync_test_44_a +\ir helper_psql_conn_setup.sql +SELECT cloudsync_init('composite_pk_tbl', 'CLS', true) AS _reinit \gset +SELECT encode(cloudsync_payload_encode(tbl, pk, col_name, col_value, col_version, db_version, site_id, cl, seq), 'hex') AS payload_a +FROM cloudsync_changes +WHERE site_id = cloudsync_siteid() \gset + +\connect cloudsync_test_44_b +\ir helper_psql_conn_setup.sql +SELECT cloudsync_init('composite_pk_tbl', 'CLS', true) AS _reinit \gset +SELECT encode(cloudsync_payload_encode(tbl, pk, col_name, col_value, col_version, db_version, site_id, cl, seq), 'hex') AS payload_b +FROM cloudsync_changes +WHERE site_id = cloudsync_siteid() \gset + +-- Apply A -> B +SELECT cloudsync_payload_apply(decode(:'payload_a', 'hex')) AS apply_a_to_b \gset + +-- Apply B -> A +\connect cloudsync_test_44_a +\ir helper_psql_conn_setup.sql +SELECT cloudsync_init('composite_pk_tbl', 'CLS', true) AS _reinit \gset +SELECT cloudsync_payload_apply(decode(:'payload_b', 'hex')) AS apply_b_to_a \gset + +-- Update a row on A +UPDATE composite_pk_tbl SET data_col = 'updated_on_a' WHERE pk_text1 = 'foo' AND pk_int1 = 3 AND pk_text2 = 'bar' AND pk_int2 = 300 AND pk_text3 = 'baz'; + +SELECT encode(cloudsync_payload_encode(tbl, pk, col_name, col_value, col_version, db_version, site_id, cl, seq), 'hex') AS payload_a2 +FROM cloudsync_changes +WHERE site_id = cloudsync_siteid() \gset + +\connect cloudsync_test_44_b +\ir helper_psql_conn_setup.sql +SELECT cloudsync_init('composite_pk_tbl', 'CLS', true) AS _reinit \gset +SELECT cloudsync_payload_apply(decode(:'payload_a2', 'hex')) AS apply_a2_to_b \gset + +-- Final hash comparison +SELECT md5(COALESCE(string_agg( + pk_text1 || ':' || pk_int1::text || ':' || pk_text2 || ':' || pk_int2::text || ':' || pk_text3 || ':' || + COALESCE(data_col, 'NULL') || ':' || COALESCE(num_col::text, 'NULL'), + '|' ORDER BY pk_text1, pk_int1, pk_text2, pk_int2, pk_text3 +), '')) AS hash_b FROM composite_pk_tbl \gset + +\connect cloudsync_test_44_a +\ir helper_psql_conn_setup.sql +SELECT md5(COALESCE(string_agg( + pk_text1 || ':' || pk_int1::text || ':' || pk_text2 || ':' || pk_int2::text || ':' || pk_text3 || ':' || + COALESCE(data_col, 'NULL') || ':' || COALESCE(num_col::text, 'NULL'), + '|' ORDER BY pk_text1, pk_int1, pk_text2, pk_int2, pk_text3 +), '')) AS hash_a FROM composite_pk_tbl \gset + +SELECT (:'hash_a' = :'hash_b') AS hashes_match \gset +\if :hashes_match +\echo [PASS] (:testid) Large composite PK (5 cols) roundtrip and update +\else +\echo [FAIL] (:testid) Hash mismatch - A: :hash_a, B: :hash_b +SELECT (:fail::int + 1) AS fail \gset +\endif + +-- Verify row count +SELECT COUNT(*) AS row_count FROM composite_pk_tbl \gset +SELECT (:row_count::int = 4) AS count_ok \gset +\if :count_ok +\echo [PASS] (:testid) Row count correct (4 rows with 5-col composite PK) +\else +\echo [FAIL] (:testid) Expected 4 rows, got :row_count +SELECT (:fail::int + 1) AS fail \gset +\endif + +-- Verify update propagated +SELECT COUNT(*) = 1 AS update_ok +FROM composite_pk_tbl +WHERE pk_text1 = 'foo' AND pk_int1 = 3 AND data_col = 'updated_on_a' \gset +\if :update_ok +\echo [PASS] (:testid) Update propagated correctly for composite PK row +\else +\echo [FAIL] (:testid) Update not propagated +SELECT (:fail::int + 1) AS fail \gset +\endif + +-- Cleanup +\ir helper_test_cleanup.sql +\if :should_cleanup +DROP DATABASE IF EXISTS cloudsync_test_44_a; +DROP DATABASE IF EXISTS cloudsync_test_44_b; +\endif diff --git a/test/postgresql/45_pg_specific_types.sql b/test/postgresql/45_pg_specific_types.sql new file mode 100644 index 0000000..25d96c9 --- /dev/null +++ b/test/postgresql/45_pg_specific_types.sql @@ -0,0 +1,176 @@ +-- Test PostgreSQL-specific type roundtrips +-- Covers JSONB, TIMESTAMPTZ, NUMERIC with precision, BYTEA + +\set testid '45' +\ir helper_test_init.sql + +\connect postgres +\ir helper_psql_conn_setup.sql +DROP DATABASE IF EXISTS cloudsync_test_45_a; +DROP DATABASE IF EXISTS cloudsync_test_45_b; +CREATE DATABASE cloudsync_test_45_a; +CREATE DATABASE cloudsync_test_45_b; + +-- Setup Database A +\connect cloudsync_test_45_a +\ir helper_psql_conn_setup.sql +CREATE EXTENSION IF NOT EXISTS cloudsync; + +CREATE TABLE typed_tbl ( + id TEXT PRIMARY KEY, + json_col JSONB, + ts_col TIMESTAMPTZ, + num_col NUMERIC(18, 6), + bin_col BYTEA +); + +SELECT cloudsync_init('typed_tbl', 'CLS', true) AS _init_a \gset + +INSERT INTO typed_tbl VALUES ( + 'row1', + '{"key": "value", "nested": {"arr": [1, 2, 3]}}', + '2025-01-15 10:30:00+00', + 123456.789012, + '\x48656c6c6f' +); + +INSERT INTO typed_tbl VALUES ( + 'row2', + '[1, "two", null, true, false]', + '2024-06-30 23:59:59.999999+05:30', + -999999.123456, + '\xdeadbeef' +); + +INSERT INTO typed_tbl VALUES ( + 'row3', + 'null', + '1970-01-01 00:00:00+00', + 0.000000, + '\x' +); + +-- Setup Database B +\connect cloudsync_test_45_b +\ir helper_psql_conn_setup.sql +CREATE EXTENSION IF NOT EXISTS cloudsync; + +CREATE TABLE typed_tbl ( + id TEXT PRIMARY KEY, + json_col JSONB, + ts_col TIMESTAMPTZ, + num_col NUMERIC(18, 6), + bin_col BYTEA +); + +SELECT cloudsync_init('typed_tbl', 'CLS', true) AS _init_b \gset + +-- Encode and apply A -> B +\connect cloudsync_test_45_a +\ir helper_psql_conn_setup.sql +SELECT cloudsync_init('typed_tbl', 'CLS', true) AS _reinit \gset +SELECT encode(cloudsync_payload_encode(tbl, pk, col_name, col_value, col_version, db_version, site_id, cl, seq), 'hex') AS payload_a +FROM cloudsync_changes +WHERE site_id = cloudsync_siteid() \gset + +\connect cloudsync_test_45_b +\ir helper_psql_conn_setup.sql +SELECT cloudsync_init('typed_tbl', 'CLS', true) AS _reinit \gset +SELECT cloudsync_payload_apply(decode(:'payload_a', 'hex')) AS apply_a_to_b \gset + +-- Verify row count +SELECT COUNT(*) AS count_b FROM typed_tbl \gset +SELECT (:count_b::int = 3) AS count_ok \gset +\if :count_ok +\echo [PASS] (:testid) All 3 rows synced to B +\else +\echo [FAIL] (:testid) Expected 3 rows, got :count_b +SELECT (:fail::int + 1) AS fail \gset +\endif + +-- Verify JSONB roundtrip +SELECT COUNT(*) = 1 AS jsonb_ok +FROM typed_tbl +WHERE id = 'row1' + AND json_col @> '{"key": "value"}' + AND json_col -> 'nested' -> 'arr' = '[1, 2, 3]'::jsonb \gset +\if :jsonb_ok +\echo [PASS] (:testid) JSONB roundtrip correct +\else +\echo [FAIL] (:testid) JSONB data mismatch +SELECT (:fail::int + 1) AS fail \gset +\endif + +-- Verify TIMESTAMPTZ roundtrip +SELECT COUNT(*) = 1 AS ts_ok +FROM typed_tbl +WHERE id = 'row1' + AND ts_col = '2025-01-15 10:30:00+00'::timestamptz \gset +\if :ts_ok +\echo [PASS] (:testid) TIMESTAMPTZ roundtrip correct +\else +\echo [FAIL] (:testid) TIMESTAMPTZ data mismatch +SELECT (:fail::int + 1) AS fail \gset +\endif + +-- Verify NUMERIC roundtrip +SELECT COUNT(*) = 1 AS num_ok +FROM typed_tbl +WHERE id = 'row1' + AND num_col = 123456.789012 \gset +\if :num_ok +\echo [PASS] (:testid) NUMERIC(18,6) roundtrip correct +\else +\echo [FAIL] (:testid) NUMERIC data mismatch +SELECT (:fail::int + 1) AS fail \gset +\endif + +-- Verify BYTEA roundtrip +SELECT COUNT(*) = 1 AS bytea_ok +FROM typed_tbl +WHERE id = 'row1' + AND bin_col = '\x48656c6c6f'::bytea \gset +\if :bytea_ok +\echo [PASS] (:testid) BYTEA roundtrip correct +\else +\echo [FAIL] (:testid) BYTEA data mismatch +SELECT (:fail::int + 1) AS fail \gset +\endif + +-- Verify full hash match +\connect cloudsync_test_45_a +\ir helper_psql_conn_setup.sql +SELECT md5(COALESCE(string_agg( + id || ':' || + COALESCE(json_col::text, 'NULL') || ':' || + COALESCE(ts_col::text, 'NULL') || ':' || + COALESCE(num_col::text, 'NULL') || ':' || + COALESCE(encode(bin_col, 'hex'), 'NULL'), + '|' ORDER BY id +), '')) AS hash_a FROM typed_tbl \gset + +\connect cloudsync_test_45_b +\ir helper_psql_conn_setup.sql +SELECT md5(COALESCE(string_agg( + id || ':' || + COALESCE(json_col::text, 'NULL') || ':' || + COALESCE(ts_col::text, 'NULL') || ':' || + COALESCE(num_col::text, 'NULL') || ':' || + COALESCE(encode(bin_col, 'hex'), 'NULL'), + '|' ORDER BY id +), '')) AS hash_b FROM typed_tbl \gset + +SELECT (:'hash_a' = :'hash_b') AS hash_match \gset +\if :hash_match +\echo [PASS] (:testid) Full data hash matches for PG-specific types +\else +\echo [FAIL] (:testid) Hash mismatch - A: :hash_a, B: :hash_b +SELECT (:fail::int + 1) AS fail \gset +\endif + +-- Cleanup +\ir helper_test_cleanup.sql +\if :should_cleanup +DROP DATABASE IF EXISTS cloudsync_test_45_a; +DROP DATABASE IF EXISTS cloudsync_test_45_b; +\endif diff --git a/test/postgresql/46_schema_hash_mismatch.sql b/test/postgresql/46_schema_hash_mismatch.sql new file mode 100644 index 0000000..20e6867 --- /dev/null +++ b/test/postgresql/46_schema_hash_mismatch.sql @@ -0,0 +1,96 @@ +-- Test schema hash mismatch during merge +-- Verifies detection when ALTER TABLE is done without cloudsync_begin/commit_alter. + +\set testid '46' +\ir helper_test_init.sql + +\connect postgres +\ir helper_psql_conn_setup.sql +DROP DATABASE IF EXISTS cloudsync_test_46_src; +DROP DATABASE IF EXISTS cloudsync_test_46_dst; +CREATE DATABASE cloudsync_test_46_src; +CREATE DATABASE cloudsync_test_46_dst; + +-- Setup source +\connect cloudsync_test_46_src +\ir helper_psql_conn_setup.sql +CREATE EXTENSION IF NOT EXISTS cloudsync; +CREATE TABLE test_tbl (id TEXT PRIMARY KEY, val TEXT); +SELECT cloudsync_init('test_tbl', 'CLS', true) AS _init_src \gset +INSERT INTO test_tbl VALUES ('id1', 'value1'); + +-- Setup destination with same schema +\connect cloudsync_test_46_dst +\ir helper_psql_conn_setup.sql +CREATE EXTENSION IF NOT EXISTS cloudsync; +CREATE TABLE test_tbl (id TEXT PRIMARY KEY, val TEXT); +SELECT cloudsync_init('test_tbl', 'CLS', true) AS _init_dst \gset + +-- Initial sync to get both in sync +\connect cloudsync_test_46_src +\ir helper_psql_conn_setup.sql +SELECT cloudsync_init('test_tbl', 'CLS', true) AS _reinit \gset +SELECT encode(cloudsync_payload_encode(tbl, pk, col_name, col_value, col_version, db_version, site_id, cl, seq), 'hex') AS payload_initial +FROM cloudsync_changes +WHERE site_id = cloudsync_siteid() \gset + +\connect cloudsync_test_46_dst +\ir helper_psql_conn_setup.sql +SELECT cloudsync_init('test_tbl', 'CLS', true) AS _reinit \gset +SELECT cloudsync_payload_apply(decode(:'payload_initial', 'hex')) AS _apply_initial \gset + +-- Now ALTER TABLE on destination WITHOUT using cloudsync_begin/commit_alter +ALTER TABLE test_tbl ADD COLUMN extra TEXT DEFAULT 'default'; + +-- Insert new data on source +\connect cloudsync_test_46_src +\ir helper_psql_conn_setup.sql +SELECT cloudsync_init('test_tbl', 'CLS', true) AS _reinit \gset +INSERT INTO test_tbl VALUES ('id2', 'value2'); + +SELECT encode(cloudsync_payload_encode(tbl, pk, col_name, col_value, col_version, db_version, site_id, cl, seq), 'hex') AS payload_post_alter +FROM cloudsync_changes +WHERE site_id = cloudsync_siteid() \gset + +-- Apply payload from pre-alter source to post-alter destination +-- This should detect schema mismatch +\connect cloudsync_test_46_dst +\ir helper_psql_conn_setup.sql + +-- Reinit to pick up new schema +SELECT cloudsync_init('test_tbl', 'CLS', true) AS _reinit_dst \gset + +-- The apply may error due to schema mismatch, or succeed silently. +-- Either outcome is acceptable — the key is no corruption. +\set apply_ok true +SELECT cloudsync_payload_apply(decode(:'payload_post_alter', 'hex')) AS _apply_mismatch \gset +-- If the above errors, psql continues (ON_ERROR_STOP is off) and apply_ok stays true. +-- The test just verifies integrity below. + +-- Verify database is in a consistent state (not corrupted) +SELECT COUNT(*) AS final_count FROM test_tbl \gset +SELECT (:final_count::int >= 1) AS state_ok \gset +\if :state_ok +\echo [PASS] (:testid) Database consistent after schema mismatch scenario (rows: :final_count) +\else +\echo [FAIL] (:testid) Database corrupted after schema mismatch +SELECT (:fail::int + 1) AS fail \gset +\endif + +-- Verify original data is intact +SELECT COUNT(*) = 1 AS original_ok +FROM test_tbl +WHERE id = 'id1' AND val = 'value1' \gset +\if :original_ok +\echo [PASS] (:testid) Original data intact after schema mismatch +\else +\echo [FAIL] (:testid) Original data corrupted +SELECT (:fail::int + 1) AS fail \gset +\endif + +-- Cleanup +\ir helper_test_cleanup.sql +\if :should_cleanup +DROP DATABASE IF EXISTS cloudsync_test_46_src; +DROP DATABASE IF EXISTS cloudsync_test_46_dst; +\endif diff --git a/test/postgresql/full_test.sql b/test/postgresql/full_test.sql index ba69198..2341687 100644 --- a/test/postgresql/full_test.sql +++ b/test/postgresql/full_test.sql @@ -47,6 +47,13 @@ \ir 37_block_lww_round4.sql \ir 38_block_lww_round5.sql \ir 39_concurrent_write_apply.sql +\ir 40_unsupported_algorithms.sql +\ir 41_corrupted_payload.sql +\ir 42_payload_idempotency.sql +\ir 43_delete_resurrect_ordering.sql +\ir 44_large_composite_pk.sql +\ir 45_pg_specific_types.sql +\ir 46_schema_hash_mismatch.sql -- 'Test summary' \echo '\nTest summary:' diff --git a/test/unit.c b/test/unit.c index 8b78dc4..b62fa60 100644 --- a/test/unit.c +++ b/test/unit.c @@ -10364,6 +10364,552 @@ bool do_test_block_lww_whitespace(int nclients, bool print_result, bool cleanup_ return false; } +// MARK: - New edge-case tests + +bool do_test_unsupported_algorithms (sqlite3 *db) { + // Test that DWS and AWS algorithms are rejected with an error + const char *sql; + int rc; + + // Create tables for the test + sql = "CREATE TABLE IF NOT EXISTS test_dws (id TEXT PRIMARY KEY, val TEXT);" + "CREATE TABLE IF NOT EXISTS test_aws (id TEXT PRIMARY KEY, val TEXT);"; + rc = sqlite3_exec(db, sql, NULL, NULL, NULL); + if (rc != SQLITE_OK) return false; + + // DWS should fail + sql = "SELECT cloudsync_init('test_dws', 'dws');"; + rc = sqlite3_exec(db, sql, NULL, NULL, NULL); + if (rc != SQLITE_ERROR) return false; + + // AWS should fail + sql = "SELECT cloudsync_init('test_aws', 'aws');"; + rc = sqlite3_exec(db, sql, NULL, NULL, NULL); + if (rc != SQLITE_ERROR) return false; + + // Verify no companion tables were created + sqlite3_stmt *stmt = NULL; + rc = sqlite3_prepare_v2(db, "SELECT COUNT(*) FROM sqlite_master WHERE name='test_dws_cloudsync' OR name='test_aws_cloudsync';", -1, &stmt, NULL); + if (rc != SQLITE_OK) return false; + if (sqlite3_step(stmt) != SQLITE_ROW) { sqlite3_finalize(stmt); return false; } + int count = sqlite3_column_int(stmt, 0); + sqlite3_finalize(stmt); + if (count != 0) return false; + + // CLS should still work on the same table + sql = "SELECT cloudsync_init('test_dws', 'cls');"; + rc = sqlite3_exec(db, sql, NULL, NULL, NULL); + if (rc != SQLITE_OK) return false; + + return true; +} + +bool do_test_corrupted_payload (int nclients, bool print_result, bool cleanup_databases) { + sqlite3 *db[2] = {NULL, NULL}; + bool result = false; + + time_t timestamp = time(NULL); + int saved_counter = test_counter; + + // Create source and destination databases + for (int i = 0; i < 2; i++) { + db[i] = do_create_database_file(i, timestamp, test_counter++); + if (!db[i]) return false; + + int rc = sqlite3_exec(db[i], "CREATE TABLE test_tbl (id TEXT PRIMARY KEY, val TEXT);", NULL, NULL, NULL); + if (rc != SQLITE_OK) goto finalize; + rc = sqlite3_exec(db[i], "SELECT cloudsync_init('test_tbl');", NULL, NULL, NULL); + if (rc != SQLITE_OK) goto finalize; + } + + // Insert data in source + sqlite3_exec(db[0], "INSERT INTO test_tbl VALUES ('id1', 'value1');", NULL, NULL, NULL); + sqlite3_exec(db[0], "INSERT INTO test_tbl VALUES ('id2', 'value2');", NULL, NULL, NULL); + + // Get valid payload as blob + sqlite3_stmt *enc_stmt = NULL; + int rc = sqlite3_prepare_v2(db[0], "SELECT cloudsync_payload_encode(tbl, pk, col_name, col_value, col_version, db_version, site_id, cl, seq) FROM cloudsync_changes WHERE site_id=cloudsync_siteid();", -1, &enc_stmt, NULL); + if (rc != SQLITE_OK) goto finalize; + + rc = sqlite3_step(enc_stmt); + if (rc != SQLITE_ROW) { sqlite3_finalize(enc_stmt); goto finalize; } + + int valid_len = sqlite3_column_bytes(enc_stmt, 0); + const void *valid_blob = sqlite3_column_blob(enc_stmt, 0); + if (!valid_blob || valid_len < 20) { sqlite3_finalize(enc_stmt); goto finalize; } + + // Copy valid payload + char *payload_copy = (char *)malloc(valid_len); + if (!payload_copy) { sqlite3_finalize(enc_stmt); goto finalize; } + memcpy(payload_copy, valid_blob, valid_len); + sqlite3_finalize(enc_stmt); + + // Test 1: Empty blob + { + sqlite3_stmt *dec_stmt = NULL; + rc = sqlite3_prepare_v2(db[1], "SELECT cloudsync_payload_decode(?);", -1, &dec_stmt, NULL); + if (rc == SQLITE_OK) { + sqlite3_bind_blob(dec_stmt, 1, "", 0, SQLITE_STATIC); + rc = sqlite3_step(dec_stmt); + // Should either error or return without inserting + sqlite3_finalize(dec_stmt); + } + } + + // Test 2: Random garbage + { + char garbage[16] = {0xDE, 0xAD, 0xBE, 0xEF, 0xCA, 0xFE, 0xBA, 0xBE, 0x01, 0x02, 0x03, 0x04, 0x05, 0x06, 0x07, 0x08}; + sqlite3_stmt *dec_stmt = NULL; + rc = sqlite3_prepare_v2(db[1], "SELECT cloudsync_payload_decode(?);", -1, &dec_stmt, NULL); + if (rc == SQLITE_OK) { + sqlite3_bind_blob(dec_stmt, 1, garbage, sizeof(garbage), SQLITE_STATIC); + rc = sqlite3_step(dec_stmt); + sqlite3_finalize(dec_stmt); + } + } + + // Test 3: Truncated payload (first 10 bytes) + { + sqlite3_stmt *dec_stmt = NULL; + rc = sqlite3_prepare_v2(db[1], "SELECT cloudsync_payload_decode(?);", -1, &dec_stmt, NULL); + if (rc == SQLITE_OK) { + sqlite3_bind_blob(dec_stmt, 1, payload_copy, 10, SQLITE_STATIC); + rc = sqlite3_step(dec_stmt); + sqlite3_finalize(dec_stmt); + } + } + + // Test 4: Valid payload with flipped byte in the middle + { + char *corrupted = (char *)malloc(valid_len); + memcpy(corrupted, payload_copy, valid_len); + corrupted[valid_len / 2] ^= 0xFF; + + sqlite3_stmt *dec_stmt = NULL; + rc = sqlite3_prepare_v2(db[1], "SELECT cloudsync_payload_decode(?);", -1, &dec_stmt, NULL); + if (rc == SQLITE_OK) { + sqlite3_bind_blob(dec_stmt, 1, corrupted, valid_len, SQLITE_STATIC); + rc = sqlite3_step(dec_stmt); + sqlite3_finalize(dec_stmt); + } + free(corrupted); + } + + // Verify destination table is still empty (no corrupted data inserted) + { + sqlite3_stmt *count_stmt = NULL; + rc = sqlite3_prepare_v2(db[1], "SELECT COUNT(*) FROM test_tbl;", -1, &count_stmt, NULL); + if (rc != SQLITE_OK) { free(payload_copy); goto finalize; } + if (sqlite3_step(count_stmt) != SQLITE_ROW) { sqlite3_finalize(count_stmt); free(payload_copy); goto finalize; } + int count = sqlite3_column_int(count_stmt, 0); + sqlite3_finalize(count_stmt); + if (count != 0) { printf("corrupted_payload: expected 0 rows but got %d\n", count); free(payload_copy); goto finalize; } + } + + // Test 5: Valid payload should still work + if (!do_merge_using_payload(db[0], db[1], false, true)) { + printf("corrupted_payload: valid payload failed after corrupted attempts\n"); + free(payload_copy); + goto finalize; + } + + { + sqlite3_stmt *count_stmt = NULL; + rc = sqlite3_prepare_v2(db[1], "SELECT COUNT(*) FROM test_tbl;", -1, &count_stmt, NULL); + if (rc != SQLITE_OK) { free(payload_copy); goto finalize; } + if (sqlite3_step(count_stmt) != SQLITE_ROW) { sqlite3_finalize(count_stmt); free(payload_copy); goto finalize; } + int count = sqlite3_column_int(count_stmt, 0); + sqlite3_finalize(count_stmt); + if (count != 2) { printf("corrupted_payload: expected 2 rows after valid apply but got %d\n", count); free(payload_copy); goto finalize; } + } + + free(payload_copy); + result = true; + +finalize: + for (int i = 0; i < 2; i++) { + if (db[i]) close_db(db[i]); + if (cleanup_databases) { + char buf[256]; + do_build_database_path(buf, i, timestamp, saved_counter++); + file_delete_internal(buf); + } + } + return result; +} + +bool do_test_payload_idempotency (int nclients, bool print_result, bool cleanup_databases) { + sqlite3 *db[2] = {NULL, NULL}; + bool result = false; + + time_t timestamp = time(NULL); + int saved_counter = test_counter; + + for (int i = 0; i < 2; i++) { + db[i] = do_create_database_file(i, timestamp, test_counter++); + if (!db[i]) return false; + + int rc = sqlite3_exec(db[i], "CREATE TABLE test_tbl (id TEXT PRIMARY KEY, val TEXT, num INTEGER);", NULL, NULL, NULL); + if (rc != SQLITE_OK) goto finalize; + rc = sqlite3_exec(db[i], "SELECT cloudsync_init('test_tbl');", NULL, NULL, NULL); + if (rc != SQLITE_OK) goto finalize; + } + + // Insert data on source + sqlite3_exec(db[0], "INSERT INTO test_tbl VALUES ('id1', 'hello', 10);", NULL, NULL, NULL); + sqlite3_exec(db[0], "INSERT INTO test_tbl VALUES ('id2', 'world', 20);", NULL, NULL, NULL); + sqlite3_exec(db[0], "UPDATE test_tbl SET val = 'hello_updated' WHERE id = 'id1';", NULL, NULL, NULL); + + // Apply payload 3 times and check after each + int prev_count = -1; + for (int apply = 0; apply < 3; apply++) { + if (!do_merge_using_payload(db[0], db[1], false, true)) { + printf("payload_idempotency: apply #%d failed\n", apply + 1); + goto finalize; + } + + // Check row count + sqlite3_stmt *stmt = NULL; + int rc = sqlite3_prepare_v2(db[1], "SELECT COUNT(*) FROM test_tbl;", -1, &stmt, NULL); + if (rc != SQLITE_OK) goto finalize; + sqlite3_step(stmt); + int count = sqlite3_column_int(stmt, 0); + sqlite3_finalize(stmt); + + if (count != 2) { + printf("payload_idempotency: expected 2 rows after apply #%d, got %d\n", apply + 1, count); + goto finalize; + } + + if (prev_count >= 0 && count != prev_count) { + printf("payload_idempotency: row count changed from %d to %d on apply #%d\n", prev_count, count, apply + 1); + goto finalize; + } + prev_count = count; + } + + // Verify data values are correct + { + sqlite3_stmt *stmt = NULL; + int rc = sqlite3_prepare_v2(db[1], "SELECT val FROM test_tbl WHERE id = 'id1';", -1, &stmt, NULL); + if (rc != SQLITE_OK) goto finalize; + if (sqlite3_step(stmt) != SQLITE_ROW) { sqlite3_finalize(stmt); goto finalize; } + const char *val = (const char *)sqlite3_column_text(stmt, 0); + if (!val || strcmp(val, "hello_updated") != 0) { + printf("payload_idempotency: expected 'hello_updated', got '%s'\n", val ? val : "NULL"); + sqlite3_finalize(stmt); + goto finalize; + } + sqlite3_finalize(stmt); + } + + // Compare source and target + result = do_compare_queries(db[0], "SELECT * FROM test_tbl ORDER BY id;", + db[1], "SELECT * FROM test_tbl ORDER BY id;", + -1, -1, print_result); + +finalize: + for (int i = 0; i < 2; i++) { + if (db[i]) close_db(db[i]); + if (cleanup_databases) { + char buf[256]; + do_build_database_path(buf, i, timestamp, saved_counter++); + file_delete_internal(buf); + } + } + return result; +} + +bool do_test_causal_length_tiebreak (int nclients, bool print_result, bool cleanup_databases) { + sqlite3 *db[3] = {NULL, NULL, NULL}; + bool result = false; + + time_t timestamp = time(NULL); + int saved_counter = test_counter; + + // Create 3 databases with the same table + for (int i = 0; i < 3; i++) { + db[i] = do_create_database_file(i, timestamp, test_counter++); + if (!db[i]) return false; + + int rc = sqlite3_exec(db[i], "CREATE TABLE test_tbl (id TEXT PRIMARY KEY, val TEXT);", NULL, NULL, NULL); + if (rc != SQLITE_OK) goto finalize; + rc = sqlite3_exec(db[i], "SELECT cloudsync_init('test_tbl');", NULL, NULL, NULL); + if (rc != SQLITE_OK) goto finalize; + } + + // Seed row on db[0] and sync to all + sqlite3_exec(db[0], "INSERT INTO test_tbl VALUES ('row1', 'seed');", NULL, NULL, NULL); + do_merge_using_payload(db[0], db[1], false, true); + do_merge_using_payload(db[0], db[2], false, true); + + // All 3 independently update the same row+column (producing equal CL) + sqlite3_exec(db[0], "UPDATE test_tbl SET val = 'value_from_db0' WHERE id = 'row1';", NULL, NULL, NULL); + sqlite3_exec(db[1], "UPDATE test_tbl SET val = 'value_from_db1' WHERE id = 'row1';", NULL, NULL, NULL); + sqlite3_exec(db[2], "UPDATE test_tbl SET val = 'value_from_db2' WHERE id = 'row1';", NULL, NULL, NULL); + + // Merge all pairs in both directions + sqlite3 *all_db[MAX_SIMULATED_CLIENTS] = {NULL}; + all_db[0] = db[0]; all_db[1] = db[1]; all_db[2] = db[2]; + if (!do_merge(all_db, 3, true)) { + printf("causal_length_tiebreak: merge failed\n"); + goto finalize; + } + + // All 3 must converge to the same value + const char *query = "SELECT val FROM test_tbl WHERE id = 'row1';"; + char *values[3] = {NULL, NULL, NULL}; + + for (int i = 0; i < 3; i++) { + sqlite3_stmt *stmt = NULL; + int rc = sqlite3_prepare_v2(db[i], query, -1, &stmt, NULL); + if (rc != SQLITE_OK) goto tiebreak_finalize; + if (sqlite3_step(stmt) != SQLITE_ROW) { sqlite3_finalize(stmt); goto tiebreak_finalize; } + const char *text = (const char *)sqlite3_column_text(stmt, 0); + values[i] = text ? strdup(text) : NULL; + sqlite3_finalize(stmt); + } + + // Check convergence + if (values[0] && values[1] && values[2] && + strcmp(values[0], values[1]) == 0 && strcmp(values[1], values[2]) == 0) { + result = true; + } else { + printf("causal_length_tiebreak: databases diverged: '%s', '%s', '%s'\n", + values[0] ? values[0] : "NULL", + values[1] ? values[1] : "NULL", + values[2] ? values[2] : "NULL"); + } + +tiebreak_finalize: + for (int i = 0; i < 3; i++) { + if (values[i]) free(values[i]); + } + +finalize: + for (int i = 0; i < 3; i++) { + if (db[i]) close_db(db[i]); + if (cleanup_databases) { + char buf[256]; + do_build_database_path(buf, i, timestamp, saved_counter++); + file_delete_internal(buf); + } + } + return result; +} + +bool do_test_delete_resurrect_ordering (int nclients, bool print_result, bool cleanup_databases) { + sqlite3 *db[3] = {NULL, NULL, NULL}; + bool result = false; + + time_t timestamp = time(NULL); + int saved_counter = test_counter; + + for (int i = 0; i < 3; i++) { + db[i] = do_create_database_file(i, timestamp, test_counter++); + if (!db[i]) return false; + + int rc = sqlite3_exec(db[i], "CREATE TABLE test_tbl (id TEXT PRIMARY KEY, val TEXT);", NULL, NULL, NULL); + if (rc != SQLITE_OK) goto finalize; + rc = sqlite3_exec(db[i], "SELECT cloudsync_init('test_tbl');", NULL, NULL, NULL); + if (rc != SQLITE_OK) goto finalize; + } + + // Site A: insert row, sync to B and C + sqlite3_exec(db[0], "INSERT INTO test_tbl VALUES ('row1', 'original');", NULL, NULL, NULL); + do_merge_using_payload(db[0], db[1], false, true); + do_merge_using_payload(db[0], db[2], false, true); + + // Site A: delete row (CL 1->2) + sqlite3_exec(db[0], "DELETE FROM test_tbl WHERE id = 'row1';", NULL, NULL, NULL); + + // Sync delete to B + do_merge_using_payload(db[0], db[1], true, true); + + // Site B: re-insert (CL 2->3, resurrection) + sqlite3_exec(db[1], "INSERT INTO test_tbl VALUES ('row1', 'resurrected_by_b');", NULL, NULL, NULL); + + // Site C receives payloads in REVERSE order: B's resurrection first, then A's delete + do_merge_using_payload(db[1], db[2], true, true); + do_merge_using_payload(db[0], db[2], true, true); + + // Site A receives B's resurrection + do_merge_using_payload(db[1], db[2], true, true); + do_merge_using_payload(db[1], db[0], true, true); + + // All 3 should converge: row exists + const char *query = "SELECT * FROM test_tbl ORDER BY id;"; + result = do_compare_queries(db[0], query, db[1], query, -1, -1, print_result); + if (result) result = do_compare_queries(db[0], query, db[2], query, -1, -1, print_result); + + // Verify the row exists (resurrection should win) + if (result) { + sqlite3_stmt *stmt = NULL; + int rc = sqlite3_prepare_v2(db[2], "SELECT COUNT(*) FROM test_tbl WHERE id = 'row1';", -1, &stmt, NULL); + if (rc == SQLITE_OK && sqlite3_step(stmt) == SQLITE_ROW) { + int count = sqlite3_column_int(stmt, 0); + if (count != 1) { + printf("delete_resurrect_ordering: expected row1 to exist on db[2], count=%d\n", count); + result = false; + } + } + if (stmt) sqlite3_finalize(stmt); + } + +finalize: + for (int i = 0; i < 3; i++) { + if (db[i]) close_db(db[i]); + if (cleanup_databases) { + char buf[256]; + do_build_database_path(buf, i, timestamp, saved_counter++); + file_delete_internal(buf); + } + } + return result; +} + +bool do_test_large_composite_pk (int nclients, bool print_result, bool cleanup_databases) { + sqlite3 *db[2] = {NULL, NULL}; + bool result = false; + + time_t timestamp = time(NULL); + int saved_counter = test_counter; + + for (int i = 0; i < 2; i++) { + db[i] = do_create_database_file(i, timestamp, test_counter++); + if (!db[i]) return false; + + int rc = sqlite3_exec(db[i], + "CREATE TABLE cpk_tbl (" + " pk_text1 TEXT NOT NULL," + " pk_int1 INTEGER NOT NULL," + " pk_text2 TEXT NOT NULL," + " pk_int2 INTEGER NOT NULL," + " pk_text3 TEXT NOT NULL," + " data_col TEXT," + " num_col INTEGER," + " PRIMARY KEY (pk_text1, pk_int1, pk_text2, pk_int2, pk_text3)" + ");", NULL, NULL, NULL); + if (rc != SQLITE_OK) goto finalize; + rc = sqlite3_exec(db[i], "SELECT cloudsync_init('cpk_tbl');", NULL, NULL, NULL); + if (rc != SQLITE_OK) goto finalize; + } + + // Insert data on both sides + sqlite3_exec(db[0], "INSERT INTO cpk_tbl VALUES ('alpha', 1, 'beta', 100, 'gamma', 'data_a1', 42);", NULL, NULL, NULL); + sqlite3_exec(db[0], "INSERT INTO cpk_tbl VALUES ('alpha', 2, 'beta', 200, 'delta', 'data_a2', 84);", NULL, NULL, NULL); + sqlite3_exec(db[0], "INSERT INTO cpk_tbl VALUES ('x', 999, 'y', -1, 'z', 'edge_case', 0);", NULL, NULL, NULL); + + sqlite3_exec(db[1], "INSERT INTO cpk_tbl VALUES ('alpha', 1, 'beta', 100, 'gamma', 'data_b1', 99);", NULL, NULL, NULL); + sqlite3_exec(db[1], "INSERT INTO cpk_tbl VALUES ('foo', 3, 'bar', 300, 'baz', 'data_b2', 77);", NULL, NULL, NULL); + + // Merge both directions + if (!do_merge_using_payload(db[0], db[1], false, true)) goto finalize; + if (!do_merge_using_payload(db[1], db[0], false, true)) goto finalize; + + // Update on db[0] and sync + sqlite3_exec(db[0], "UPDATE cpk_tbl SET data_col = 'updated_on_a' WHERE pk_text1 = 'foo' AND pk_int1 = 3 AND pk_text2 = 'bar' AND pk_int2 = 300 AND pk_text3 = 'baz';", NULL, NULL, NULL); + if (!do_merge_using_payload(db[0], db[1], true, true)) goto finalize; + + // Compare + result = do_compare_queries(db[0], "SELECT * FROM cpk_tbl ORDER BY pk_text1, pk_int1, pk_text2, pk_int2, pk_text3;", + db[1], "SELECT * FROM cpk_tbl ORDER BY pk_text1, pk_int1, pk_text2, pk_int2, pk_text3;", + -1, -1, print_result); + + // Verify row count + if (result) { + sqlite3_stmt *stmt = NULL; + int rc = sqlite3_prepare_v2(db[1], "SELECT COUNT(*) FROM cpk_tbl;", -1, &stmt, NULL); + if (rc == SQLITE_OK && sqlite3_step(stmt) == SQLITE_ROW) { + int count = sqlite3_column_int(stmt, 0); + if (count != 4) { + printf("large_composite_pk: expected 4 rows, got %d\n", count); + result = false; + } + } + if (stmt) sqlite3_finalize(stmt); + } + +finalize: + for (int i = 0; i < 2; i++) { + if (db[i]) close_db(db[i]); + if (cleanup_databases) { + char buf[256]; + do_build_database_path(buf, i, timestamp, saved_counter++); + file_delete_internal(buf); + } + } + return result; +} + +bool do_test_schema_hash_mismatch (int nclients, bool print_result, bool cleanup_databases) { + sqlite3 *db[2] = {NULL, NULL}; + bool result = false; + + time_t timestamp = time(NULL); + int saved_counter = test_counter; + + for (int i = 0; i < 2; i++) { + db[i] = do_create_database_file(i, timestamp, test_counter++); + if (!db[i]) return false; + + int rc = sqlite3_exec(db[i], "CREATE TABLE test_tbl (id TEXT PRIMARY KEY, val TEXT);", NULL, NULL, NULL); + if (rc != SQLITE_OK) goto finalize; + rc = sqlite3_exec(db[i], "SELECT cloudsync_init('test_tbl');", NULL, NULL, NULL); + if (rc != SQLITE_OK) goto finalize; + } + + // Initial sync + sqlite3_exec(db[0], "INSERT INTO test_tbl VALUES ('id1', 'value1');", NULL, NULL, NULL); + if (!do_merge_using_payload(db[0], db[1], false, true)) goto finalize; + + // ALTER TABLE on destination WITHOUT cloudsync_begin/commit_alter + int rc = sqlite3_exec(db[1], "ALTER TABLE test_tbl ADD COLUMN extra TEXT;", NULL, NULL, NULL); + if (rc != SQLITE_OK) goto finalize; + + // Re-init to pick up changed schema + rc = sqlite3_exec(db[1], "SELECT cloudsync_init('test_tbl');", NULL, NULL, NULL); + if (rc != SQLITE_OK) goto finalize; + + // Insert new data on source + sqlite3_exec(db[0], "INSERT INTO test_tbl VALUES ('id2', 'value2');", NULL, NULL, NULL); + + // Apply payload from pre-alter source to post-alter destination + // This should fail due to schema hash mismatch + bool merge_result = do_merge_using_payload(db[0], db[1], true, false); + if (merge_result) { + // If merge succeeded despite schema mismatch, it means the extension + // accepted the fewer-columns payload — verify data isn't corrupted + } + + // Verify original data is intact regardless + { + sqlite3_stmt *stmt = NULL; + rc = sqlite3_prepare_v2(db[1], "SELECT COUNT(*) FROM test_tbl WHERE id = 'id1' AND val = 'value1';", -1, &stmt, NULL); + if (rc != SQLITE_OK) goto finalize; + if (sqlite3_step(stmt) != SQLITE_ROW) { sqlite3_finalize(stmt); goto finalize; } + int count = sqlite3_column_int(stmt, 0); + sqlite3_finalize(stmt); + if (count != 1) { + printf("schema_hash_mismatch: original data corrupted\n"); + goto finalize; + } + } + + result = true; + +finalize: + for (int i = 0; i < 2; i++) { + if (db[i]) close_db(db[i]); + if (cleanup_databases) { + char buf[256]; + do_build_database_path(buf, i, timestamp, saved_counter++); + file_delete_internal(buf); + } + } + return result; +} + int test_report(const char *description, bool result){ printf("%-30s %s\n", description, (result) ? "OK" : "FAILED"); return result ? 0 : 1; @@ -10393,6 +10939,7 @@ int main (int argc, const char * argv[]) { result += test_report("DBUtils Test:", do_test_dbutils()); result += test_report("Minor Test:", do_test_others(db)); result += test_report("Test Error Cases:", do_test_error_cases(db)); + result += test_report("Unsupported Algos Test:", do_test_unsupported_algorithms(db)); result += test_report("Null PK Insert Test:", do_test_null_prikey_insert(db)); result += test_report("Test Single PK:", do_test_single_pk(print_result)); @@ -10528,6 +11075,14 @@ int main (int argc, const char * argv[]) { result += test_report("Test Block LWW LongLine:", do_test_block_lww_long_line(2, print_result, cleanup_databases)); result += test_report("Test Block LWW Whitespace:", do_test_block_lww_whitespace(2, print_result, cleanup_databases)); + // edge-case tests + result += test_report("Corrupted Payload Test:", do_test_corrupted_payload(2, print_result, cleanup_databases)); + result += test_report("Payload Idempotency Test:", do_test_payload_idempotency(2, print_result, cleanup_databases)); + result += test_report("CL Tiebreak Test:", do_test_causal_length_tiebreak(3, print_result, cleanup_databases)); + result += test_report("Delete/Resurrect Order:", do_test_delete_resurrect_ordering(3, print_result, cleanup_databases)); + result += test_report("Large Composite PK Test:", do_test_large_composite_pk(2, print_result, cleanup_databases)); + result += test_report("Schema Hash Mismatch:", do_test_schema_hash_mismatch(2, print_result, cleanup_databases)); + finalize: if (rc != SQLITE_OK) printf("%s (%d)\n", (db) ? sqlite3_errmsg(db) : "N/A", rc); close_db(db); From 9882a84d711f5c324bd688d8ba55714234bc3b9d Mon Sep 17 00:00:00 2001 From: Andrea Donetti Date: Mon, 23 Mar 2026 13:24:43 -0600 Subject: [PATCH 3/9] test: improved stress test command --- .../commands/stress-test-sync-sqlitecloud.md | 50 +++++++++++++------ 1 file changed, 35 insertions(+), 15 deletions(-) diff --git a/.claude/commands/stress-test-sync-sqlitecloud.md b/.claude/commands/stress-test-sync-sqlitecloud.md index 22102cb..f7f1c75 100644 --- a/.claude/commands/stress-test-sync-sqlitecloud.md +++ b/.claude/commands/stress-test-sync-sqlitecloud.md @@ -19,7 +19,11 @@ Ask the user for the following configuration using a single question set: - Medium: 10K rows, 10 iterations, 4 concurrent databases - Large: 100K rows, 50 iterations, 4 concurrent databases (Jim's original scenario) - Custom: let the user specify rows, iterations, and number of concurrent databases -4. **RLS mode** — with RLS (requires user tokens) or without RLS +4. **Operations per iteration** — how many UPDATE and DELETE operations to perform each iteration: + - `NUM_UPDATES`: number of UPDATE operations per iteration (default: 1). Each UPDATE runs `UPDATE SET value = value + 1;` affecting all rows. + - `NUM_DELETES`: number of DELETE operations per iteration (default: 1). Each DELETE runs `DELETE FROM
WHERE rowid IN (SELECT rowid FROM
ORDER BY RANDOM() LIMIT 10);` removing 10 random rows. Set to 0 to skip deletes entirely. + - Propose defaults of 1 update and 1 delete. The user can set 0 deletes for update-only tests. +5. **RLS mode** — with RLS (requires user tokens) or without RLS 5. **Table schema** — offer simple default or custom: ```sql CREATE TABLE test_sync (id TEXT PRIMARY KEY, user_id TEXT NOT NULL DEFAULT '', name TEXT, value INTEGER); @@ -34,6 +38,8 @@ Save these as variables: - `ROWS` (number of rows per iteration) - `ITERATIONS` (number of delete/insert/update cycles) - `NUM_DBS` (number of concurrent databases) +- `NUM_UPDATES` (number of UPDATE operations per iteration, default 1) +- `NUM_DELETES` (number of DELETE operations per iteration, default 1; 0 to skip) ### Step 2: Setup SQLiteCloud Database and Table @@ -106,15 +112,15 @@ Create a bash script at `/tmp/stress_test_concurrent.sh` that: 2. **Defines a worker function** that runs in a subshell for each database: - Each worker logs all output to `/tmp/sync_concurrent_.log` - Each iteration does: - a. **UPDATE all/some rows** (e.g., `UPDATE
SET value = value + 1;`) - b. **DELETE a few rows** (e.g., `DELETE FROM
WHERE rowid IN (SELECT rowid FROM
ORDER BY RANDOM() LIMIT 10);`) + a. **UPDATE** — run `UPDATE
SET value = value + 1;` repeated `NUM_UPDATES` times (skip if 0) + b. **DELETE** — run `DELETE FROM
WHERE rowid IN (SELECT rowid FROM
ORDER BY RANDOM() LIMIT 10);` repeated `NUM_DELETES` times (skip if 0) c. **Sync using the 3-step send/check/check pattern:** 1. `SELECT cloudsync_network_send_changes();` — send local changes to the server 2. `SELECT cloudsync_network_check_changes();` — ask the server to prepare a payload of remote changes 3. Sleep 1 second (outside sqlite3, between two separate sqlite3 invocations) 4. `SELECT cloudsync_network_check_changes();` — download the prepared payload, if any - Each sqlite3 session must: `.load` the extension, call `cloudsync_network_init()`/`cloudsync_network_init_custom()`, `cloudsync_network_set_apikey()`/`cloudsync_network_set_token()` (depending on RLS mode), do the work, call `cloudsync_terminate()` - - **Timing**: Log the wall-clock execution time (in milliseconds) for each `cloudsync_network_send_changes()`, `cloudsync_network_check_changes()` call. Use bash `date +%s%3N` before and after each sqlite3 invocation that calls a network function, and compute the delta. Log lines like: `[DB][iter ] send_changes: 123ms`, `[DB][iter ] check_changes_1: 45ms`, `[DB][iter ] check_changes_2: 67ms` + - **Timing**: Log the wall-clock execution time (in milliseconds) for each `cloudsync_network_send_changes()`, `cloudsync_network_check_changes()` call. Define a `now_ms()` helper function at the top of the script and use it before and after each sqlite3 invocation that calls a network function, computing the delta. On **macOS**, `date` does not support `%3N` (nanoseconds) — use `python3 -c 'import time; print(int(time.time()*1000))'` instead. On **Linux**, `date +%s%3N` works fine. The script should detect the platform and define `now_ms()` accordingly. Log lines like: `[DB][iter ] send_changes: 123ms`, `[DB][iter ] check_changes_1: 45ms`, `[DB][iter ] check_changes_2: 67ms` - Include labeled output lines like `[DB][iter ] updated count=, deleted count=` for grep-ability 3. **Launches all workers in parallel** using `&` and collects PIDs @@ -151,21 +157,29 @@ After the test completes, provide a detailed breakdown: After all workers have terminated, perform a **final sync on every local database** to ensure all databases converge to the same state. Then verify data integrity. -1. **Final sync loop** (max 10 retries): Repeat the following until all local databases have the same row count, or the retry limit is reached: +**IMPORTANT — RLS mode changes what "convergence" means:** When RLS is enabled, each user can only see their own rows. Databases belonging to different users will have different row counts and different data — this is correct behavior. All convergence and integrity checks must therefore be scoped **per user group** (i.e., only compare databases that share the same userId/token). + +1. **Final sync loop** (max 10 retries): Repeat the following until convergence is achieved within each user group, or the retry limit is reached: a. For each local database (sequentially): - Load the extension, call `cloudsync_network_init`/`cloudsync_network_init_custom`, authenticate with `cloudsync_network_set_apikey`/`cloudsync_network_set_token` - Run `SELECT cloudsync_network_sync(100, 10);` to sync remaining changes - Call `cloudsync_terminate()` b. After syncing all databases, query `SELECT COUNT(*) FROM
` on each database - c. If all row counts are identical, convergence is achieved — break out of the loop - d. Otherwise, log the round number and the distinct row counts, then repeat from (a) - e. If the retry limit is reached without convergence, report it as a failure + c. **If RLS is disabled:** Check that all databases have the same row count. If so, convergence is achieved — break. + d. **If RLS is enabled:** Group databases by userId. Within each user group, check that all databases have the same row count. Convergence is achieved when every user group is internally consistent — break. Different user groups are expected to have different row counts. + e. Otherwise, log the round number and the distinct row counts (per group if RLS), then repeat from (a) + f. If the retry limit is reached without convergence, report it as a failure -2. **Row count verification**: Report the final row counts. All databases should have the same number of rows. Also check SQLiteCloud (as admin) for total row count. +2. **Row count verification**: + - **If RLS is disabled:** Report the final row counts. All databases should have the same number of rows. + - **If RLS is enabled:** Report row counts grouped by user. All databases within the same user group should have identical row counts. Different user groups may differ. Also verify that each database only contains rows matching its userId. + - In both cases, also check SQLiteCloud (as admin) for total row count. -3. **Row content verification**: Pick one random row ID from the first database (`SELECT id FROM
ORDER BY RANDOM() LIMIT 1;`). Then query that same row (`SELECT id, user_id, name, value FROM
WHERE id = '';`) on **every** local database. Compare the results — all databases must return identical column values for that row. Report the row ID, the expected values, and any mismatches. +3. **Row content verification**: + - **If RLS is disabled:** Pick one random row ID from the first database. Query that row on every local database. All must return identical values. + - **If RLS is enabled:** For each user group, pick one random row ID from the first database in that group. Query that row on all databases in the same user group. All databases in the group must return identical values. Do NOT expect databases from other user groups to have this row — they should return empty (RLS blocks cross-user access). -4. If RLS is enabled, verify no cross-user data leakage. +4. **RLS cross-user leak check** (RLS mode only): For a sample of databases (e.g., one per user group), verify that `SELECT COUNT(*) FROM
WHERE user_id != ''` returns 0. Report any cross-user data leakage as a test failure. ## Output Format @@ -194,15 +208,21 @@ If errors are found, include: The test **PASSES** if: 1. All workers complete all iterations 2. Zero `error`, `locked`, `SQLITE_BUSY`, or HTTP 500 responses in any log -3. After the final sync, all local databases have the same row count -4. A randomly selected row has identical content across all local databases +3. After the final sync, databases converge: + - **Without RLS:** all local databases have the same row count + - **With RLS:** all databases within each user group have the same row count (different user groups may differ) +4. Row content is consistent: + - **Without RLS:** a randomly selected row has identical content across all local databases + - **With RLS:** a randomly selected row has identical content across all databases in the same user group; databases from other user groups correctly return empty for that row +5. **With RLS:** no cross-user data leakage (each database contains only rows matching its userId) The test **FAILS** if: 1. Any worker crashes or fails to complete 2. Any `database is locked` or `SQLITE_BUSY` errors appear 3. Server returns 500 errors under concurrent load -4. Row counts differ across local databases after the final sync loop exhausts all retries -5. Row content differs across local databases (data corruption) +4. Row counts differ within the comparison scope (all DBs without RLS, same-user DBs with RLS) after the final sync loop exhausts all retries +5. Row content differs within the comparison scope (data corruption) +6. **With RLS:** any database contains rows belonging to a different userId (cross-user data leakage) ## Important Notes From e2e9f8eefd5a0db7df496bef87999a3b2f6de18e Mon Sep 17 00:00:00 2001 From: Andrea Donetti Date: Mon, 23 Mar 2026 13:25:57 -0600 Subject: [PATCH 4/9] feat(ci): add PostgreSQL extension builds for Linux, macOS, and Windows Add a postgres-build matrix job that compiles the PostgreSQL extension for 5 platform/arch combinations (linux-x86_64, linux-arm64, macos-arm64, macos-x86_64, windows-x86_64) and includes them as release assets. Add postgres-package Makefile target and Windows platform support (PG_EXTENSION_LIB, -lpostgres linking). --- .github/workflows/main.yml | 69 +++++++++++++++++++++++++++++++++++++- docker/Makefile.postgresql | 41 +++++++++++++++++----- src/cloudsync.h | 2 +- 3 files changed, 101 insertions(+), 11 deletions(-) diff --git a/.github/workflows/main.yml b/.github/workflows/main.yml index 21c9466..52b4801 100644 --- a/.github/workflows/main.yml +++ b/.github/workflows/main.yml @@ -265,10 +265,77 @@ jobs: docker cp test/postgresql cloudsync-postgres:/tmp/cloudsync/test/postgresql docker exec cloudsync-postgres psql -U postgres -d postgres -f /tmp/cloudsync/test/postgresql/full_test.sql + postgres-build: + runs-on: ${{ matrix.os }} + name: postgresql-${{ matrix.name }}-${{ matrix.arch }} build + timeout-minutes: 15 + strategy: + fail-fast: false + matrix: + include: + - os: ubuntu-22.04 + arch: x86_64 + name: linux + - os: ubuntu-22.04-arm + arch: arm64 + name: linux + - os: macos-15 + arch: arm64 + name: macos + - os: macos-13 + arch: x86_64 + name: macos + - os: windows-2022 + arch: x86_64 + name: windows + + steps: + + - uses: actions/checkout@v4.2.2 + with: + submodules: true + + - name: linux install postgresql dev headers + if: matrix.name == 'linux' + run: | + sudo sh -c 'echo "deb http://apt.postgresql.org/pub/repos/apt $(lsb_release -cs)-pgdg main" > /etc/apt/sources.list.d/pgdg.list' + curl -fsSL https://www.postgresql.org/media/keys/ACCC4CF8.asc | sudo gpg --dearmor -o /etc/apt/trusted.gpg.d/postgresql.gpg + sudo apt-get update + sudo apt-get install -y postgresql-server-dev-17 + + - name: macos install postgresql + if: matrix.name == 'macos' + run: brew install postgresql@17 + + - uses: msys2/setup-msys2@v2.27.0 + if: matrix.name == 'windows' + with: + msystem: ucrt64 + install: mingw-w64-ucrt-x86_64-gcc make mingw-w64-ucrt-x86_64-postgresql + + - name: build and package postgresql extension (linux) + if: matrix.name == 'linux' + run: make postgres-package + + - name: build and package postgresql extension (macos) + if: matrix.name == 'macos' + run: make postgres-package PG_CONFIG=$(brew --prefix postgresql@17)/bin/pg_config + + - name: build and package postgresql extension (windows) + if: matrix.name == 'windows' + shell: msys2 {0} + run: make postgres-package + + - uses: actions/upload-artifact@v4.6.2 + with: + name: cloudsync-postgresql-${{ matrix.name }}-${{ matrix.arch }} + path: dist/postgresql/ + if-no-files-found: error + release: runs-on: ubuntu-22.04 name: release - needs: [build, postgres-test] + needs: [build, postgres-test, postgres-build] if: github.ref == 'refs/heads/main' env: diff --git a/docker/Makefile.postgresql b/docker/Makefile.postgresql index 70b3da9..2fca61d 100644 --- a/docker/Makefile.postgresql +++ b/docker/Makefile.postgresql @@ -14,6 +14,17 @@ PG_INCLUDEDIR := $(shell $(PG_CONFIG) --includedir-server 2>/dev/null) EXTENSION = cloudsync EXTVERSION = 1.0 +# Platform-specific PostgreSQL settings +ifeq ($(OS),Windows_NT) + PG_EXTENSION_LIB = $(EXTENSION).dll + PG_CFLAGS = -Wall -Wextra -Wno-unused-parameter -std=c11 -O2 + PG_LDFLAGS = -shared -L$(shell $(PG_CONFIG) --libdir) -lpostgres +else + PG_EXTENSION_LIB = $(EXTENSION).so + PG_CFLAGS = -fPIC -Wall -Wextra -Wno-unused-parameter -std=c11 -O2 + PG_LDFLAGS = -shared +endif + # Source files - core platform-agnostic code PG_CORE_SRC = \ src/cloudsync.c \ @@ -38,15 +49,16 @@ PG_OBJS = $(PG_ALL_SRC:.c=.o) # Compiler flags # Define POSIX macros as compiler flags to ensure they're defined before any includes PG_CPPFLAGS = -I$(PG_INCLUDEDIR) -Isrc -Isrc/postgresql -Imodules/fractional-indexing -DCLOUDSYNC_POSTGRESQL_BUILD -D_POSIX_C_SOURCE=200809L -D_GNU_SOURCE -PG_CFLAGS = -fPIC -Wall -Wextra -Wno-unused-parameter -std=c11 -O2 PG_DEBUG ?= 0 ifeq ($(PG_DEBUG),1) +ifeq ($(OS),Windows_NT) +PG_CFLAGS = -Wall -Wextra -Wno-unused-parameter -std=c11 -g -O0 -fno-omit-frame-pointer +else PG_CFLAGS = -fPIC -Wall -Wextra -Wno-unused-parameter -std=c11 -g -O0 -fno-omit-frame-pointer endif -PG_LDFLAGS = -shared +endif # Output files -PG_EXTENSION_SO = $(EXTENSION).so PG_EXTENSION_SQL = src/postgresql/$(EXTENSION)--$(EXTVERSION).sql PG_EXTENSION_CONTROL = docker/postgresql/$(EXTENSION).control @@ -54,7 +66,7 @@ PG_EXTENSION_CONTROL = docker/postgresql/$(EXTENSION).control # PostgreSQL Build Targets # ============================================================================ -.PHONY: postgres-check postgres-build postgres-install postgres-clean postgres-test \ +.PHONY: postgres-check postgres-build postgres-install postgres-package postgres-clean postgres-test \ postgres-docker-build postgres-docker-build-asan postgres-docker-run postgres-docker-run-asan postgres-docker-stop postgres-docker-rebuild \ postgres-docker-debug-build postgres-docker-debug-run postgres-docker-debug-rebuild \ postgres-docker-shell postgres-dev-rebuild postgres-help unittest-pg \ @@ -78,16 +90,16 @@ postgres-build: postgres-check echo " CC $$src"; \ $(CC) $(PG_CPPFLAGS) $(PG_CFLAGS) -c $$src -o $${src%.c}.o || exit 1; \ done - @echo "Linking $(PG_EXTENSION_SO)..." - $(CC) $(PG_LDFLAGS) -o $(PG_EXTENSION_SO) $(PG_OBJS) - @echo "Build complete: $(PG_EXTENSION_SO)" + @echo "Linking $(PG_EXTENSION_LIB)..." + $(CC) $(PG_LDFLAGS) -o $(PG_EXTENSION_LIB) $(PG_OBJS) + @echo "Build complete: $(PG_EXTENSION_LIB)" # Install extension to PostgreSQL postgres-install: postgres-build @echo "Installing CloudSync extension to PostgreSQL..." @echo "Installing shared library to $(PG_PKGLIBDIR)/" install -d $(PG_PKGLIBDIR) - install -m 755 $(PG_EXTENSION_SO) $(PG_PKGLIBDIR)/ + install -m 755 $(PG_EXTENSION_LIB) $(PG_PKGLIBDIR)/ @echo "Installing SQL script to $(PG_SHAREDIR)/extension/" install -d $(PG_SHAREDIR)/extension install -m 644 $(PG_EXTENSION_SQL) $(PG_SHAREDIR)/extension/ @@ -98,10 +110,21 @@ postgres-install: postgres-build @echo "To use the extension, run in psql:" @echo " CREATE EXTENSION $(EXTENSION);" +# Package extension files for distribution +PG_DIST_DIR = dist/postgresql + +postgres-package: postgres-build + @echo "Packaging PostgreSQL extension..." + @mkdir -p $(PG_DIST_DIR) + cp $(PG_EXTENSION_LIB) $(PG_DIST_DIR)/ + cp $(PG_EXTENSION_SQL) $(PG_DIST_DIR)/ + cp $(PG_EXTENSION_CONTROL) $(PG_DIST_DIR)/ + @echo "Package ready in $(PG_DIST_DIR)/" + # Clean PostgreSQL build artifacts postgres-clean: @echo "Cleaning PostgreSQL build artifacts..." - rm -f $(PG_OBJS) $(PG_EXTENSION_SO) + rm -f $(PG_OBJS) $(PG_EXTENSION_LIB) @echo "Clean complete" # Test extension (requires running PostgreSQL) diff --git a/src/cloudsync.h b/src/cloudsync.h index 7e233c8..3b62a27 100644 --- a/src/cloudsync.h +++ b/src/cloudsync.h @@ -18,7 +18,7 @@ extern "C" { #endif -#define CLOUDSYNC_VERSION "0.9.202" +#define CLOUDSYNC_VERSION "0.9.203" #define CLOUDSYNC_MAX_TABLENAME_LEN 512 #define CLOUDSYNC_VALUE_NOTSET -1 From 25fe161a8be4215953da37ec383f23c988722d89 Mon Sep 17 00:00:00 2001 From: Andrea Donetti Date: Mon, 23 Mar 2026 13:52:27 -0600 Subject: [PATCH 5/9] fix(ci): resolve PostgreSQL extension build failures on macOS and Windows macOS: install gettext for libintl.h and pass include path via PG_EXTRA_CFLAGS. Windows: skip POSIX defines (_POSIX_C_SOURCE, _GNU_SOURCE) so PG headers use Winsock instead of netinet/in.h. --- .github/workflows/main.yml | 4 ++-- docker/Makefile.postgresql | 7 ++++++- 2 files changed, 8 insertions(+), 3 deletions(-) diff --git a/.github/workflows/main.yml b/.github/workflows/main.yml index 52b4801..a3ea7a8 100644 --- a/.github/workflows/main.yml +++ b/.github/workflows/main.yml @@ -305,7 +305,7 @@ jobs: - name: macos install postgresql if: matrix.name == 'macos' - run: brew install postgresql@17 + run: brew install postgresql@17 gettext - uses: msys2/setup-msys2@v2.27.0 if: matrix.name == 'windows' @@ -319,7 +319,7 @@ jobs: - name: build and package postgresql extension (macos) if: matrix.name == 'macos' - run: make postgres-package PG_CONFIG=$(brew --prefix postgresql@17)/bin/pg_config + run: make postgres-package PG_CONFIG=$(brew --prefix postgresql@17)/bin/pg_config PG_EXTRA_CFLAGS="-I$(brew --prefix gettext)/include" - name: build and package postgresql extension (windows) if: matrix.name == 'windows' diff --git a/docker/Makefile.postgresql b/docker/Makefile.postgresql index 2fca61d..514f1e6 100644 --- a/docker/Makefile.postgresql +++ b/docker/Makefile.postgresql @@ -48,7 +48,12 @@ PG_OBJS = $(PG_ALL_SRC:.c=.o) # Compiler flags # Define POSIX macros as compiler flags to ensure they're defined before any includes -PG_CPPFLAGS = -I$(PG_INCLUDEDIR) -Isrc -Isrc/postgresql -Imodules/fractional-indexing -DCLOUDSYNC_POSTGRESQL_BUILD -D_POSIX_C_SOURCE=200809L -D_GNU_SOURCE +# On Windows, skip POSIX defines — PG headers use Winsock instead of netinet/in.h +PG_EXTRA_CFLAGS ?= +PG_CPPFLAGS = -I$(PG_INCLUDEDIR) -Isrc -Isrc/postgresql -Imodules/fractional-indexing -DCLOUDSYNC_POSTGRESQL_BUILD $(PG_EXTRA_CFLAGS) +ifneq ($(OS),Windows_NT) +PG_CPPFLAGS += -D_POSIX_C_SOURCE=200809L -D_GNU_SOURCE +endif PG_DEBUG ?= 0 ifeq ($(PG_DEBUG),1) ifeq ($(OS),Windows_NT) From 8105742acfaf9e47947180297686b263866d83f6 Mon Sep 17 00:00:00 2001 From: Andrea Donetti Date: Mon, 23 Mar 2026 14:15:22 -0600 Subject: [PATCH 6/9] fix(ci): resolve PostgreSQL extension build on macOS, drop Windows MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit macOS: guard Security.h include behind !CLOUDSYNC_POSTGRESQL_BUILD to avoid type conflicts (Size, uint64) with PostgreSQL headers, use getentropy instead. Windows: dropped from postgres-build matrix — MSYS2 PostgreSQL headers expect Unix includes incompatible with MinGW. --- .github/workflows/main.yml | 15 --------------- src/utils.c | 7 +++++-- 2 files changed, 5 insertions(+), 17 deletions(-) diff --git a/.github/workflows/main.yml b/.github/workflows/main.yml index a3ea7a8..32dd318 100644 --- a/.github/workflows/main.yml +++ b/.github/workflows/main.yml @@ -285,10 +285,6 @@ jobs: - os: macos-13 arch: x86_64 name: macos - - os: windows-2022 - arch: x86_64 - name: windows - steps: - uses: actions/checkout@v4.2.2 @@ -307,12 +303,6 @@ jobs: if: matrix.name == 'macos' run: brew install postgresql@17 gettext - - uses: msys2/setup-msys2@v2.27.0 - if: matrix.name == 'windows' - with: - msystem: ucrt64 - install: mingw-w64-ucrt-x86_64-gcc make mingw-w64-ucrt-x86_64-postgresql - - name: build and package postgresql extension (linux) if: matrix.name == 'linux' run: make postgres-package @@ -321,11 +311,6 @@ jobs: if: matrix.name == 'macos' run: make postgres-package PG_CONFIG=$(brew --prefix postgresql@17)/bin/pg_config PG_EXTRA_CFLAGS="-I$(brew --prefix gettext)/include" - - name: build and package postgresql extension (windows) - if: matrix.name == 'windows' - shell: msys2 {0} - run: make postgres-package - - uses: actions/upload-artifact@v4.6.2 with: name: cloudsync-postgresql-${{ matrix.name }}-${{ matrix.arch }} diff --git a/src/utils.c b/src/utils.c index 9fbe12a..fff6cdd 100644 --- a/src/utils.c +++ b/src/utils.c @@ -18,7 +18,7 @@ #define file_close _close #else #include -#if defined(__APPLE__) +#if defined(__APPLE__) && !defined(CLOUDSYNC_POSTGRESQL_BUILD) #include #elif !defined(__ANDROID__) #include @@ -57,9 +57,12 @@ int cloudsync_uuid_v7 (uint8_t value[UUID_LEN]) { // fill the buffer with high-quality random data #ifdef _WIN32 if (BCryptGenRandom(NULL, (BYTE*)value, UUID_LEN, BCRYPT_USE_SYSTEM_PREFERRED_RNG) != STATUS_SUCCESS) return -1; - #elif defined(__APPLE__) + #elif defined(__APPLE__) && !defined(CLOUDSYNC_POSTGRESQL_BUILD) // Use SecRandomCopyBytes for macOS/iOS if (SecRandomCopyBytes(kSecRandomDefault, UUID_LEN, value) != errSecSuccess) return -1; + #elif defined(__APPLE__) && defined(CLOUDSYNC_POSTGRESQL_BUILD) + // PostgreSQL build: use getentropy to avoid Security.framework type conflicts + if (getentropy(value, UUID_LEN) != 0) return -1; #elif defined(__ANDROID__) //arc4random_buf doesn't have a return value to check for success arc4random_buf(value, UUID_LEN); From 5037e0eb53f32d8fb0468e86372e32828a64c950 Mon Sep 17 00:00:00 2001 From: Andrea Donetti Date: Mon, 23 Mar 2026 14:43:30 -0600 Subject: [PATCH 7/9] fix(ci): use _DARWIN_C_SOURCE on macOS for PostgreSQL extension build Replace _GNU_SOURCE with _DARWIN_C_SOURCE on macOS to expose preadv/pwritev declarations required by PostgreSQL headers. _GNU_SOURCE is now Linux-only. --- docker/Makefile.postgresql | 8 +++++++- 1 file changed, 7 insertions(+), 1 deletion(-) diff --git a/docker/Makefile.postgresql b/docker/Makefile.postgresql index 514f1e6..fec88da 100644 --- a/docker/Makefile.postgresql +++ b/docker/Makefile.postgresql @@ -52,7 +52,13 @@ PG_OBJS = $(PG_ALL_SRC:.c=.o) PG_EXTRA_CFLAGS ?= PG_CPPFLAGS = -I$(PG_INCLUDEDIR) -Isrc -Isrc/postgresql -Imodules/fractional-indexing -DCLOUDSYNC_POSTGRESQL_BUILD $(PG_EXTRA_CFLAGS) ifneq ($(OS),Windows_NT) -PG_CPPFLAGS += -D_POSIX_C_SOURCE=200809L -D_GNU_SOURCE +PG_CPPFLAGS += -D_POSIX_C_SOURCE=200809L +UNAME_S := $(shell uname -s) +ifeq ($(UNAME_S),Darwin) +PG_CPPFLAGS += -D_DARWIN_C_SOURCE +else +PG_CPPFLAGS += -D_GNU_SOURCE +endif endif PG_DEBUG ?= 0 ifeq ($(PG_DEBUG),1) From 7fb4e61bb49cc3788e97806186b9c1c11fc0776a Mon Sep 17 00:00:00 2001 From: Andrea Donetti Date: Mon, 23 Mar 2026 15:00:51 -0600 Subject: [PATCH 8/9] fix(ci): add -undefined dynamic_lookup for macOS PostgreSQL extension linking macOS linker requires explicit handling of undefined symbols. PostgreSQL extensions resolve symbols at load time, so -undefined dynamic_lookup is needed. Also consolidate UNAME_S detection to a single location. --- docker/Makefile.postgresql | 13 +++++++++++-- 1 file changed, 11 insertions(+), 2 deletions(-) diff --git a/docker/Makefile.postgresql b/docker/Makefile.postgresql index fec88da..fea571a 100644 --- a/docker/Makefile.postgresql +++ b/docker/Makefile.postgresql @@ -14,6 +14,11 @@ PG_INCLUDEDIR := $(shell $(PG_CONFIG) --includedir-server 2>/dev/null) EXTENSION = cloudsync EXTVERSION = 1.0 +# Detect OS for platform-specific settings +ifneq ($(OS),Windows_NT) +UNAME_S := $(shell uname -s) +endif + # Platform-specific PostgreSQL settings ifeq ($(OS),Windows_NT) PG_EXTENSION_LIB = $(EXTENSION).dll @@ -22,7 +27,12 @@ ifeq ($(OS),Windows_NT) else PG_EXTENSION_LIB = $(EXTENSION).so PG_CFLAGS = -fPIC -Wall -Wextra -Wno-unused-parameter -std=c11 -O2 - PG_LDFLAGS = -shared + ifeq ($(UNAME_S),Darwin) + # macOS: allow undefined symbols resolved at load time by PostgreSQL + PG_LDFLAGS = -shared -undefined dynamic_lookup + else + PG_LDFLAGS = -shared + endif endif # Source files - core platform-agnostic code @@ -53,7 +63,6 @@ PG_EXTRA_CFLAGS ?= PG_CPPFLAGS = -I$(PG_INCLUDEDIR) -Isrc -Isrc/postgresql -Imodules/fractional-indexing -DCLOUDSYNC_POSTGRESQL_BUILD $(PG_EXTRA_CFLAGS) ifneq ($(OS),Windows_NT) PG_CPPFLAGS += -D_POSIX_C_SOURCE=200809L -UNAME_S := $(shell uname -s) ifeq ($(UNAME_S),Darwin) PG_CPPFLAGS += -D_DARWIN_C_SOURCE else From c25bd248e7baff3bb3b1fdc817db918a256fc8b3 Mon Sep 17 00:00:00 2001 From: Andrea Donetti Date: Mon, 23 Mar 2026 15:12:55 -0600 Subject: [PATCH 9/9] fix(ci): use macos-15 for PostgreSQL x86_64 build with cross-compilation macos-13 runners are deprecated. Cross-compile x86_64 on macos-15 ARM runner via -arch x86_64 passed through PG_EXTRA_CFLAGS. Also pass PG_EXTRA_CFLAGS to linker for arch flag propagation. --- .github/workflows/main.yml | 4 ++-- docker/Makefile.postgresql | 2 +- 2 files changed, 3 insertions(+), 3 deletions(-) diff --git a/.github/workflows/main.yml b/.github/workflows/main.yml index 32dd318..6a79be6 100644 --- a/.github/workflows/main.yml +++ b/.github/workflows/main.yml @@ -282,7 +282,7 @@ jobs: - os: macos-15 arch: arm64 name: macos - - os: macos-13 + - os: macos-15 arch: x86_64 name: macos steps: @@ -309,7 +309,7 @@ jobs: - name: build and package postgresql extension (macos) if: matrix.name == 'macos' - run: make postgres-package PG_CONFIG=$(brew --prefix postgresql@17)/bin/pg_config PG_EXTRA_CFLAGS="-I$(brew --prefix gettext)/include" + run: make postgres-package PG_CONFIG=$(brew --prefix postgresql@17)/bin/pg_config PG_EXTRA_CFLAGS="-I$(brew --prefix gettext)/include ${{ matrix.arch == 'x86_64' && '-arch x86_64' || '' }}" - uses: actions/upload-artifact@v4.6.2 with: diff --git a/docker/Makefile.postgresql b/docker/Makefile.postgresql index fea571a..8e1a514 100644 --- a/docker/Makefile.postgresql +++ b/docker/Makefile.postgresql @@ -111,7 +111,7 @@ postgres-build: postgres-check $(CC) $(PG_CPPFLAGS) $(PG_CFLAGS) -c $$src -o $${src%.c}.o || exit 1; \ done @echo "Linking $(PG_EXTENSION_LIB)..." - $(CC) $(PG_LDFLAGS) -o $(PG_EXTENSION_LIB) $(PG_OBJS) + $(CC) $(PG_LDFLAGS) $(PG_EXTRA_CFLAGS) -o $(PG_EXTENSION_LIB) $(PG_OBJS) @echo "Build complete: $(PG_EXTENSION_LIB)" # Install extension to PostgreSQL