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

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
6 changes: 6 additions & 0 deletions .github/workflows/main.yml
Original file line number Diff line number Diff line change
Expand Up @@ -26,6 +26,12 @@ jobs:
container: ${{ matrix.container && matrix.container || '' }}
name: ${{ matrix.name }}${{ matrix.arch && format('-{0}', matrix.arch) || '' }} build${{ matrix.arch != 'arm64-v8a' && matrix.arch != 'armeabi-v7a' && matrix.name != 'ios-sim' && matrix.name != 'ios' && matrix.name != 'mac-catalyst' && matrix.name != 'apple-xcframework' && matrix.name != 'android-aar' && ( matrix.name != 'macos' || matrix.arch != 'x86_64' ) && ' + test' || ''}}
timeout-minutes: 20
# Only this leg receives the shared chunked tenant. Lock it across branches,
# and queue competing runs instead of canceling a pending PR's validation.
concurrency:
group: ${{ (matrix.name == 'linux' && matrix.arch == 'x86_64') && 'cloudsync-chunked-test-tenant' || format('build-{0}-{1}-{2}', github.run_id, matrix.name, matrix.arch) }}
cancel-in-progress: false
queue: max
strategy:
fail-fast: false
matrix:
Expand Down
5 changes: 4 additions & 1 deletion Makefile
Original file line number Diff line number Diff line change
Expand Up @@ -270,6 +270,8 @@ $(TEST_TARGET): $(TEST_OBJ)
# Object files
$(BUILD_RELEASE)/%.o: %.c
$(CC) $(CFLAGS) -O3 -fPIC -c $< -o $@
$(BUILD_TEST)/integration_bootstrap.o: $(TEST_DIR)/integration.c

$(BUILD_TEST)/sqlite3.o: $(SQLITE_DIR)/sqlite3.c
$(CC) $(CFLAGS) -DSQLITE_DQS=0 -DSQLITE_CORE -c $< -o $@
$(BUILD_TEST)/%.o: %.c
Expand Down Expand Up @@ -297,9 +299,10 @@ ifneq ($(COVERAGE),false)
endif

# Run only unit tests
unittest: $(TARGET) $(DIST_DIR)/unit$(EXE) $(DIST_DIR)/review_regressions$(EXE)
unittest: $(TARGET) $(DIST_DIR)/unit$(EXE) $(DIST_DIR)/review_regressions$(EXE) $(DIST_DIR)/integration_bootstrap$(EXE)
@./$(DIST_DIR)/unit$(EXE)
@./$(DIST_DIR)/review_regressions$(EXE)
@./$(DIST_DIR)/integration_bootstrap$(EXE)

# Run the SQLite unit and regression suites on a real big-endian host (s390x) under QEMU
# emulation. The payload and primary-key encodings are byte-order sensitive; this is the
Expand Down
11 changes: 11 additions & 0 deletions docs/internal/block-group-atomicity.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,11 @@
# Atomic application of mixed ordinary and block columns

Before materializing a block column, payload apply flushes the row's ordinary columns so the base row exists. Previously this flush also released the PK group's savepoint. A later block-table write error therefore left the ordinary columns and their metadata committed while the block column was missing or stale.

The pre-block flush now writes pending ordinary columns without releasing the group's savepoint or advancing the applied-row count. All blocks and ordinary columns remain inside that boundary until the PK, table or source database version changes. A rejected block rolls back that entire group, including metadata, blocks and resurrection changes. Previously completed groups retain their existing semantics; a failed apply does not advance the receive checkpoint. Fixing the underlying error and replaying the payload restores the complete row.

## Validation

`test_block_group_atomicity` in `test/review_regressions.c` runs 120 failure-and-retry cases across inserts, updates and resurrected rows, 3–10 blocks, rejection of first/middle/last blocks, and both autocommit and caller-owned transactions. It compares both directions of SQL EXCEPT for the base table, metadata and block table, checks the unchanged checkpoint and caller transaction, and verifies successful retry and duplicate delivery.

The core suite and audit regressions pass under AddressSanitizer and UndefinedBehaviorSanitizer, with SQLite itself instrumented and zero outstanding SQLite memory. Restoring the original pre-block flush fails 360 assertions in the new tests. The independent branch also passes all 521 PostgreSQL 18.6 checks for shared-code compatibility. These tests use local databases; they do not validate a deployed cloud server.
17 changes: 14 additions & 3 deletions docs/internal/cloud-e2e.md
Original file line number Diff line number Diff line change
Expand Up @@ -18,9 +18,20 @@ gh run watch RUN_ID --repo sqliteai/sqlite-sync --exit-status

A push already starts that workflow, so do not dispatch a second run for the same
commit unnecessarily. The workflow cancels older runs of the same branch. Avoid
running another branch or a local process against the shared chunked tenant at the
same time: the negative-cache test requires an idle, exclusive tenant, and the
workflow's concurrency group is per branch, not per tenant.
running a local process or an older workflow against the shared chunked tenant at
the same time: the negative-cache test requires an idle, exclusive tenant. The
Linux x86_64 job now takes a shared job-level concurrency lock across branches.
`queue: max` retains competing validations instead of replacing a pending run;
other matrix jobs remain parallel. See [GitHub's concurrency documentation](https://docs.github.com/en/actions/how-tos/write-workflows/choose-when-workflows-run/control-workflow-concurrency).

Fresh receivers use bounded polling: HTTP 202 while a download is prepared is not
proof of failure or successful synchronization. Bootstrap checks require received
rows and the expected fixture data, and reject SQL/protocol errors immediately.
The negative-cache idle assertion remains strict: any unexpected rows still fail.
`test/integration_bootstrap.c`, included in `make unittest`, exercises delayed and
partial delivery, exhaustion of the retry budget, absent data, absent received
rows, protocol failures, malformed JSON and SQL errors without a cloud connection.
It also asserts one sync call per attempt and zero outstanding SQLite memory.

Inspect the **linux-x86_64 build + test** job. Only that matrix leg receives
`INTEGRATION_TEST_CHUNKED_DATABASE_ID`. A green job alone is insufficient: optional
Expand Down
6 changes: 4 additions & 2 deletions src/cloudsync.c
Original file line number Diff line number Diff line change
Expand Up @@ -4604,9 +4604,11 @@ int cloudsync_payload_apply (cloudsync_context *data, const char *payload, int b
last_tbl_len = decoded_context.tbl_len;

if (batch.count > 0 && cloudsync_payload_row_is_block(&decoded_context)) {
fail_rc = cloudsync_payload_group_flush(data, &batch);
// Materialization needs the ordinary columns to exist, but they and
// every block still belong to the same atomic PK group. Keep its
// savepoint open until the actual PK/table/db_version boundary.
fail_rc = merge_flush_pending(data);
if (fail_rc != DBRES_OK) break;
applied = (int)i;
}

cloudsync_payload_group_open(data, &batch);
Expand Down
38 changes: 35 additions & 3 deletions test/integration.c
Original file line number Diff line number Diff line change
Expand Up @@ -232,6 +232,37 @@ int db_select_receive (sqlite3 *db, const char *sql, int *chunks, int *complete,
return sqlite3_finalize(stmt);
}

// A fresh site's first check may return HTTP 202 while its download is prepared.
// Require actual received rows AND the expected data, allowing bounded empty polls.
// Materialize the scalar result so JSON projections cannot invoke sync repeatedly.
int db_sync_await(sqlite3 *db, const char *expected_sql, int max_attempts, int delay_ms) {
bool received = false;
for (int attempt = 0; attempt < max_attempts; attempt++) {
int rows = 0, valid = 0;
char error[512];
int rc = db_select_receive(db,
"WITH result AS MATERIALIZED (SELECT cloudsync_network_sync(250,10) AS j) "
"SELECT j ->> '$.receive.rows', "
"coalesce((j ->> '$.send.status') <> 'error' AND "
"json_type(j,'$.receive.rows') = 'integer', 0), "
"coalesce(j ->> '$.receive.error', j ->> '$.send.lastFailure', j ->> '$.receive.lastFailure') "
"FROM result;", &rows, &valid, error, sizeof(error));
if (rc != SQLITE_OK) return rc;
if (!valid || rows < 0 || error[0]) {
printf("Error: bootstrap sync failed: %s\n", error[0] ? error : "invalid sync status");
return SQLITE_ERROR;
}
received = received || rows > 0;
int ready = 0;
rc = db_select_int(db, expected_sql, &ready);
if (rc != SQLITE_OK) return rc;
if (received && ready) return SQLITE_OK;
if (attempt + 1 < max_attempts) sqlite3_sleep(delay_ms);
}
printf("Error: bootstrap sync did not deliver the expected data after %d attempts\n", max_attempts);
return SQLITE_ERROR;
}

int db_expect_min (sqlite3 *db, const char *sql, int expect_min) {
int value = 0;
int rc = db_select_int(db, sql, &value);
Expand Down Expand Up @@ -587,7 +618,7 @@ int test_init (const char *db_path, int init) {
snprintf(sql, sizeof(sql), "INSERT INTO users (id, name) VALUES ('%s', '%s');", value, value);
rc = db_exec(db, sql); RCHECK
rc = db_expect_int(db, "SELECT COUNT(*) as count FROM users;", 1); RCHECK
rc = db_expect_gt0(db, "SELECT cloudsync_network_sync(250,10) ->> '$.receive.rows';"); RCHECK
rc = db_sync_await(db, "SELECT count(*) > 0 FROM activities;", 120, 250); RCHECK
rc = db_expect_gt0(db, "SELECT COUNT(*) as count FROM users;"); RCHECK
rc = db_expect_gt0(db, "SELECT COUNT(*) as count FROM activities;"); RCHECK
rc = db_expect_int(db, "SELECT COUNT(*) as count FROM workouts;", 0); RCHECK
Expand Down Expand Up @@ -689,7 +720,8 @@ int test_enable_disable(const char *db_path) {
rc = db_exec(db2, set_apikey2); RCHECK
}

rc = db_expect_gt0(db2, "SELECT cloudsync_network_sync(250,10) ->> '$.receive.rows';"); RCHECK
snprintf(sql, sizeof(sql), "SELECT COUNT(*) = 1 FROM users WHERE name='%s-should-sync';", value);
rc = db_sync_await(db2, sql, 120, 250); RCHECK

snprintf(sql, sizeof(sql), "SELECT COUNT(*) FROM users WHERE name='%s';", value);
rc = db_expect_int(db2, sql, 0); RCHECK
Expand Down Expand Up @@ -799,7 +831,7 @@ int test_token_auth (void) {
char sql[256];
snprintf(sql, sizeof(sql), "INSERT INTO users (id, name) VALUES ('%s', '%s');", value, value);
rc = db_exec(db, sql); RCHECK
rc = db_expect_gt0(db, "SELECT cloudsync_network_sync(250,10) ->> '$.receive.rows';"); RCHECK
rc = db_sync_await(db, "SELECT count(*) > 0 FROM activities;", 120, 250); RCHECK
rc = db_exec(db, "SELECT cloudsync_terminate();");

ABORT_TEST
Expand Down
74 changes: 74 additions & 0 deletions test/integration_bootstrap.c
Original file line number Diff line number Diff line change
@@ -0,0 +1,74 @@
// Exercise the integration wait policy without contacting the cloud.
#define main integration_main
#include "integration.c"
#undef main

static int failures;
#define CHECK(x) do { if (!(x)) { fprintf(stderr, "%s:%d: %s\n", __FILE__, __LINE__, #x); failures++; } } while (0)

typedef struct {
const char *json;
int ready;
} sync_reply;
typedef struct {
const sync_reply *replies;
int count;
int calls;
} sync_script;

static void mock_sync(sqlite3_context *ctx, int argc, sqlite3_value **argv) {
sync_script *script = sqlite3_user_data(ctx);
CHECK(argc == 2 && sqlite3_value_int(argv[0]) == 250 && sqlite3_value_int(argv[1]) == 10);
if (script->calls >= script->count) {
sqlite3_result_error(ctx, "unexpected extra sync", -1);
return;
}
sync_reply reply = script->replies[script->calls++];
if (reply.ready) CHECK(sqlite3_exec(sqlite3_context_db_handle(ctx), "UPDATE fixture SET ready=1", NULL, NULL, NULL) == SQLITE_OK);
if (reply.json) sqlite3_result_text(ctx, reply.json, -1, SQLITE_STATIC);
else sqlite3_result_error(ctx, "injected SQL error", -1);
}

static void run_case(const sync_reply *replies, int count, int expected_rc, int expected_calls) {
sqlite3 *db = NULL;
CHECK(sqlite3_open(":memory:", &db) == SQLITE_OK);
CHECK(sqlite3_exec(db, "CREATE TABLE fixture(ready); INSERT INTO fixture VALUES(0)", NULL, NULL, NULL) == SQLITE_OK);
sync_script script = {replies, count, 0};
CHECK(sqlite3_create_function(db, "cloudsync_network_sync", 2, SQLITE_UTF8, &script, mock_sync, NULL, NULL) == SQLITE_OK);
CHECK(db_sync_await(db, "SELECT ready FROM fixture", count, 0) == expected_rc);
CHECK(script.calls == expected_calls);
CHECK(sqlite3_close(db) == SQLITE_OK);
}

int main(void) {
const char *empty = "{\"send\":{\"status\":\"ok\"},\"receive\":{\"rows\":0}}";
const char *rows = "{\"send\":{\"status\":\"ok\"},\"receive\":{\"rows\":3}}";
sync_reply delayed[] = {{empty,0},{empty,0},{rows,1}};
sync_reply partial[] = {{rows,0},{empty,0},{empty,1}};
sync_reply never[] = {{empty,0},{empty,0},{empty,0}};
sync_reply no_data[] = {{rows,0},{rows,0}};
sync_reply no_rows[] = {{empty,1},{empty,1}};
sync_reply immediate[] = {{rows,1}};
for (int i = 0; i < 100; i++) {
run_case(delayed, 3, SQLITE_OK, 3);
run_case(partial, 3, SQLITE_OK, 3);
run_case(immediate, 1, SQLITE_OK, 1);
}
run_case(never, 3, SQLITE_ERROR, 3);
run_case(no_data, 2, SQLITE_ERROR, 2);
run_case(no_rows, 2, SQLITE_ERROR, 2);
const char *errors[] = {
"{\"send\":{\"status\":\"error\"},\"receive\":{\"rows\":3}}",
"{\"send\":{\"status\":\"ok\",\"lastFailure\":{\"message\":\"denied\"}},\"receive\":{\"rows\":3}}",
"{\"send\":{\"status\":\"ok\"},\"receive\":{\"rows\":3,\"error\":\"denied\"}}",
"{\"send\":{\"status\":\"ok\"},\"receive\":{\"rows\":3,\"lastFailure\":{\"message\":\"denied\"}}}",
"{\"receive\":{\"rows\":3}}", "{}", "null", "{invalid", NULL
};
for (unsigned i = 0; i < sizeof(errors)/sizeof(errors[0]); i++) {
sync_reply error[] = {{errors[i],1},{rows,1}};
run_case(error, 2, SQLITE_ERROR, 1);
}
CHECK(sqlite3_memory_used() == 0);
printf("Integration bootstrap policy: %d failures\n", failures);
return failures ? 1 : 0;
}
63 changes: 63 additions & 0 deletions test/review_regressions.c
Original file line number Diff line number Diff line change
Expand Up @@ -433,6 +433,68 @@ static bool fail_alloc(void) {
}
static void *fault_malloc(int size) { return fail_alloc() ? NULL : memory.xMalloc(size); }
static void *fault_realloc(void *ptr, int size) { return fail_alloc() ? NULL : memory.xRealloc(ptr, size); }
static void test_block_group_atomicity(void) {
const char *schema = "CREATE TABLE t(id TEXT PRIMARY KEY NOT NULL, owner TEXT, body TEXT);"
"SELECT cloudsync_init('t'); SELECT cloudsync_set_column('t','body','algo','block');"
"CREATE TABLE caller_work(value TEXT);";
// Insert, update and resurrection; every block position; caller and internal
// transactions. Repeat with varying payload sizes to exercise batch reuse.
for (int trial = 0; trial < 120; trial++) {
int mode = trial % 3, blocks = 3 + (trial / 3) % 8;
int denied = (trial * 7) % blocks;
bool caller = (trial / 24) % 2;
sqlite3 *source = open_db(), *target = open_db();
CHECK(sql(source, schema) == SQLITE_OK && sql(target, schema) == SQLITE_OK);
if (mode) {
CHECK(sql(source, "INSERT INTO t VALUES('a','old','old body');") == SQLITE_OK);
CHECK(apply_payload(source, target) == SQLITE_ROW);
if (mode == 2) {
CHECK(sql(source, "DELETE FROM t;") == SQLITE_OK);
CHECK(apply_payload(source, target) == SQLITE_ROW);
}
}
char body[1024] = {0}, stmt[2048];
for (int j = 0; j < blocks; j++) {
char part[64];
snprintf(part, sizeof(part), "%sblock-%d-trial-%d", j ? "\n" : "", j, trial);
strcat(body, part);
}
snprintf(stmt, sizeof(stmt), "INSERT INTO t VALUES('a','new','%s') ON CONFLICT(id) DO UPDATE SET owner=excluded.owner,body=excluded.body;", body);
CHECK(sql(source, stmt) == SQLITE_OK);
CHECK(sql(target, "CREATE TEMP TABLE before_meta AS SELECT * FROM t_cloudsync;"
"CREATE TEMP TABLE before_blocks AS SELECT * FROM t_cloudsync_blocks;"
"CREATE TEMP TABLE before_data AS SELECT * FROM t;") == SQLITE_OK);
int64_t checkpoint = scalar(target, "SELECT coalesce((SELECT value FROM cloudsync_settings WHERE key='check_dbversion'),0)");
snprintf(stmt, sizeof(stmt), "CREATE TRIGGER deny_block BEFORE INSERT ON t_cloudsync_blocks "
"WHEN NEW.col_value='block-%d-trial-%d' BEGIN SELECT RAISE(ABORT,'block rejected'); END", denied, trial);
CHECK(sql(target, stmt) == SQLITE_OK);
if (caller) CHECK(sql(target, "BEGIN; INSERT INTO caller_work VALUES('kept'); SAVEPOINT caller_sp;") == SQLITE_OK);
int rejected_rc = apply_payload(source, target);
CHECK(rejected_rc != SQLITE_ROW);
CHECK(strstr(sqlite3_errmsg(target), "block rejected") != NULL);
CHECK(sqlite3_get_autocommit(target) == !caller);
CHECK(scalar(target, "SELECT count(*) FROM (SELECT * FROM t EXCEPT SELECT * FROM before_data)") == 0);
CHECK(scalar(target, "SELECT count(*) FROM (SELECT * FROM before_data EXCEPT SELECT * FROM t)") == 0);
CHECK(scalar(target, "SELECT count(*) FROM (SELECT * FROM t_cloudsync EXCEPT SELECT * FROM before_meta)") == 0);
CHECK(scalar(target, "SELECT count(*) FROM (SELECT * FROM before_meta EXCEPT SELECT * FROM t_cloudsync)") == 0);
CHECK(scalar(target, "SELECT count(*) FROM (SELECT * FROM t_cloudsync_blocks EXCEPT SELECT * FROM before_blocks)") == 0);
CHECK(scalar(target, "SELECT count(*) FROM (SELECT * FROM before_blocks EXCEPT SELECT * FROM t_cloudsync_blocks)") == 0);
CHECK(scalar(target, "SELECT coalesce((SELECT value FROM cloudsync_settings WHERE key='check_dbversion'),0)") == checkpoint);
if (caller) {
CHECK(scalar(target, "SELECT count(*) FROM caller_work") == 1);
CHECK(sql(target, "RELEASE caller_sp;") == SQLITE_OK);
}
CHECK(sql(target, "DROP TRIGGER deny_block;") == SQLITE_OK);
CHECK(apply_payload(source, target) == SQLITE_ROW);
snprintf(stmt, sizeof(stmt), "SELECT count(*) FROM t WHERE owner='new' AND body='%s'", body);
CHECK(scalar(target, stmt) == 1);
CHECK(apply_payload(source, target) == SQLITE_ROW);
CHECK(scalar(target, stmt) == 1);
if (caller) CHECK(sql(target, "COMMIT") == SQLITE_OK);
CHECK(close_db(source) == SQLITE_OK && close_db(target) == SQLITE_OK);
}
}

static void test_block_oom(void) {
block_init_allocator();
for (int kind = 0; kind < 5; kind++) {
Expand Down Expand Up @@ -521,6 +583,7 @@ int main(void) {
test_block_materialize_errors();
test_block_migration_orphan();
test_block_not_null_payload();
test_block_group_atomicity();
test_refill_error();
test_block_oom();
cloudsync_memory_finalize();
Expand Down
Loading