diff --git a/Makefile b/Makefile index 4c119a14..78a61ff0 100644 --- a/Makefile +++ b/Makefile @@ -368,6 +368,17 @@ chunk-bench: $(TARGET) $(DIST_DIR)/chunk_bench$(EXE) if [ -n "$(CHUNK_BENCH_VERBOSE)" ]; then export CHUNK_BENCH_VERBOSE="$(CHUNK_BENCH_VERBOSE)"; fi; \ ./$(DIST_DIR)/chunk_bench$(EXE) +# Concurrent cloudsync_payload_apply stress test across transaction modes. Not part +# of `make test`: it needs threads and reports machine-dependent timings. Rebuilt on +# every run so APPLY_STRESS_SQLITE_DIR can point at another sqlite3.c (one with +# BEGIN CONCURRENT, for example). +APPLY_STRESS_SQLITE_DIR ?= $(SQLITE_DIR) + +apply-stress: $(TARGET) + $(CC) -O2 -Wall -Wno-unused-parameter -I$(APPLY_STRESS_SQLITE_DIR) -DSQLITE_THREADSAFE=1 -DHAVE_USLEEP=1 \ + $(TEST_DIR)/stress/apply_stress.c $(APPLY_STRESS_SQLITE_DIR)/sqlite3.c -o $(DIST_DIR)/apply_stress$(EXE) -lpthread -lm + ./$(DIST_DIR)/apply_stress$(EXE) + OPENSSL_TARBALL = $(OPENSSL_DIR)/$(OPENSSL_VERSION).tar.gz $(OPENSSL_TARBALL): @@ -597,4 +608,4 @@ help: # Include PostgreSQL extension targets include docker/Makefile.postgresql -.PHONY: all clean test unittest e2e extension help version xcframework aar +.PHONY: all clean test unittest e2e apply-stress extension help version xcframework aar diff --git a/test/stress/apply_stress.c b/test/stress/apply_stress.c new file mode 100644 index 00000000..9f38939c --- /dev/null +++ b/test/stress/apply_stress.c @@ -0,0 +1,570 @@ +// +// apply_stress.c +// cloudsync +// +// Concurrent cloudsync_payload_apply stress test. Several devices produce payloads +// that each span several db_versions; one thread per device then applies its +// payloads to a shared WAL database file on its own connection, all at once, the +// way a sync server applies uploads from many devices. Each mode wraps the apply +// differently: +// +// autocommit — SELECT cloudsync_payload_apply(?) on its own: the apply commits +// once per source db_version through its own savepoints. +// deferred — BEGIN; apply; COMMIT. What the SQLite Cloud server does by default: +// it runs every write command inside an implicit BEGIN TRANSACTION. +// immediate — BEGIN IMMEDIATE; apply; COMMIT. The write lock is taken before any +// read, so contention is a lock wait that busy_timeout absorbs. +// concurrent — BEGIN CONCURRENT; apply; COMMIT. The server's mode when concurrent +// transactions are enabled. Needs a SQLite built from the begin-concurrent +// branch; skipped otherwise. +// +// Retries follow the server: an apply request makes up to 1 + APPLY_STRESS_RETRIES +// attempts, and a COMMIT that returns BUSY is retried 5 more times (10..50 ms) before +// the transaction is rolled back. A request that runs out of attempts fails and the +// device sends it again after APPLY_STRESS_REQUEST_BACKOFF_MS. +// +// Every mode must converge to the same content as a reference database that applied +// every payload serially, and in immediate mode the apply statement itself must never +// return BUSY. Failed requests, retries, BUSY codes, wall time and process CPU time +// are reported for comparison. +// +// Not built or run by CI: it needs threads and its timings are machine-dependent. +// Run it by hand with `make apply-stress`; APPLY_STRESS_SQLITE_DIR= builds it +// against another sqlite3.c (for example the server's, for BEGIN CONCURRENT). +// +// Env: APPLY_STRESS_DEVICES (default 8), APPLY_STRESS_PAYLOADS (payloads per device, +// default 20), APPLY_STRESS_VERSIONS (db_versions per payload, default 5), +// APPLY_STRESS_ROWS (rows per db_version, default 4), +// APPLY_STRESS_SHARED (1 = devices write the same keys, default 0), +// APPLY_STRESS_MODE (comma-separated modes, default all four), +// APPLY_STRESS_BUSY_TIMEOUT_MS (default 5000), +// APPLY_STRESS_RETRIES (retries per request, default 6), +// APPLY_STRESS_BACKOFF_MS (max jittered sleep between retries, default 5), +// APPLY_STRESS_REQUEST_BACKOFF_MS (sleep before re-sending a failed request, +// default 100), APPLY_STRESS_MAX_REQUESTS (per payload, default 1000), +// APPLY_STRESS_EXT (default ./dist/cloudsync). +// + +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include "sqlite3.h" + +#define DB_PATH "dist/apply-stress.sqlite" +#define EXT_PATH "./dist/cloudsync" +#define MAX_CODES 16 +#define COMMIT_RETRIES 5 + +typedef enum { + MODE_AUTOCOMMIT, + MODE_DEFERRED, + MODE_IMMEDIATE, + MODE_CONCURRENT, + MODE_COUNT +} apply_mode; + +static const char *mode_names[MODE_COUNT] = {"autocommit", "deferred", "immediate", "concurrent"}; +static const char *mode_begin[MODE_COUNT] = {NULL, "BEGIN", "BEGIN IMMEDIATE", "BEGIN CONCURRENT"}; + +typedef enum { + ATTEMPT_OK, + ATTEMPT_BUSY, + ATTEMPT_FATAL +} attempt_result; + +typedef struct { + char *blob; + int size; +} payload_t; + +typedef struct { + int code; + int count; + char message[160]; +} code_count; + +typedef struct { + // input + int device; + apply_mode mode; + payload_t *payloads; + int npayloads; + // output + int applied; + int attempts; + int failed_requests; + int busy_begin; + int busy_apply; + int busy_commit; + int fatal; + char fatal_message[256]; + code_count codes[MAX_CODES]; + int ncodes; +} worker_t; + +static struct { + int devices, payloads, versions, rows; + bool shared; + int busy_timeout_ms, retries, backoff_ms, request_backoff_ms, max_requests; + const char *ext_path; +} cfg; + +static pthread_mutex_t gate_lock = PTHREAD_MUTEX_INITIALIZER; +static pthread_cond_t gate_cond = PTHREAD_COND_INITIALIZER; +static bool gate_open = false; + +static int env_int(const char *name, int dflt) { + const char *v = getenv(name); + if (!v || !*v) return dflt; + char *end = NULL; + long p = strtol(v, &end, 10); + if (!end || *end != '\0' || p < 0) return dflt; + return (int)p; +} + +static const char *env_str(const char *name, const char *dflt) { + const char *v = getenv(name); + return (v && *v) ? v : dflt; +} + +static double monotonic_ms(void) { + struct timespec ts; + clock_gettime(CLOCK_MONOTONIC, &ts); + return ((double)ts.tv_sec * 1000.0) + ((double)ts.tv_nsec / 1000000.0); +} + +static double cpu_ms(void) { + struct rusage ru; + getrusage(RUSAGE_SELF, &ru); + return (ru.ru_utime.tv_sec + ru.ru_stime.tv_sec) * 1000.0 + (ru.ru_utime.tv_usec + ru.ru_stime.tv_usec) / 1000.0; +} + +static int db_exec(sqlite3 *db, const char *sql) { + char *err = NULL; + int rc = sqlite3_exec(db, sql, NULL, NULL, &err); + if (rc != SQLITE_OK) { + fprintf(stderr, "exec failed: %s: %s\n", sql, err ? err : sqlite3_errmsg(db)); + sqlite3_free(err); + } + return rc; +} + +static sqlite3 *open_db(const char *path) { + sqlite3 *db = NULL; + if (sqlite3_open(path, &db) != SQLITE_OK) { + fprintf(stderr, "open failed: %s: %s\n", path, db ? sqlite3_errmsg(db) : "out of memory"); + sqlite3_close(db); + return NULL; + } + sqlite3_extended_result_codes(db, 1); + sqlite3_enable_load_extension(db, 1); + char *err = NULL; + if (sqlite3_load_extension(db, cfg.ext_path, NULL, &err) != SQLITE_OK) { + fprintf(stderr, "load_extension failed: %s: %s\n", cfg.ext_path, err ? err : sqlite3_errmsg(db)); + sqlite3_free(err); + sqlite3_close(db); + return NULL; + } + return db; +} + +static const char *schema_sql = + "CREATE TABLE t (id TEXT PRIMARY KEY NOT NULL, dev INTEGER, n INTEGER, body TEXT);" + "SELECT cloudsync_init('t');"; + +// Builds one device's payloads: each covers cfg.versions local transactions. +static int build_device_payloads(int device, payload_t *out) { + sqlite3 *db = open_db(":memory:"); + if (!db) return 1; + int rc = db_exec(db, schema_sql); + sqlite3_stmt *ins = NULL, *enc = NULL; + if (rc == SQLITE_OK) rc = sqlite3_prepare_v2(db, + "INSERT INTO t (id, dev, n, body) VALUES (?1, ?2, ?3, ?4) " + "ON CONFLICT(id) DO UPDATE SET dev=excluded.dev, n=excluded.n, body=excluded.body", -1, &ins, NULL); + if (rc == SQLITE_OK) rc = sqlite3_prepare_v2(db, + "SELECT cloudsync_payload_encode(tbl, pk, col_name, col_value, col_version, db_version, site_id, cl, seq) " + "FROM cloudsync_changes WHERE db_version > ?1", -1, &enc, NULL); + int64_t last = 0; + int n = 0; + for (int p = 0; rc == SQLITE_OK && p < cfg.payloads; p++) { + for (int v = 0; rc == SQLITE_OK && v < cfg.versions; v++) { + rc = db_exec(db, "BEGIN"); + for (int r = 0; rc == SQLITE_OK && r < cfg.rows; r++, n++) { + char id[64], body[64]; + if (cfg.shared) snprintf(id, sizeof(id), "row-%d", (v * cfg.rows) + r); + else snprintf(id, sizeof(id), "d%d-p%d-v%d-r%d", device, p, v, r); + snprintf(body, sizeof(body), "device %d write %d", device, n); + sqlite3_bind_text(ins, 1, id, -1, SQLITE_TRANSIENT); + sqlite3_bind_int(ins, 2, device); + sqlite3_bind_int(ins, 3, n); + sqlite3_bind_text(ins, 4, body, -1, SQLITE_TRANSIENT); + rc = sqlite3_step(ins); + rc = (rc == SQLITE_DONE) ? SQLITE_OK : rc; + sqlite3_reset(ins); + } + if (rc == SQLITE_OK) rc = db_exec(db, "COMMIT"); + } + if (rc != SQLITE_OK) break; + sqlite3_bind_int64(enc, 1, last); + rc = sqlite3_step(enc); + if (rc == SQLITE_ROW && sqlite3_column_type(enc, 0) == SQLITE_BLOB) { + out[p].size = sqlite3_column_bytes(enc, 0); + out[p].blob = malloc((size_t)out[p].size); + if (out[p].blob) memcpy(out[p].blob, sqlite3_column_blob(enc, 0), (size_t)out[p].size); + rc = out[p].blob ? SQLITE_OK : SQLITE_NOMEM; + } else { + fprintf(stderr, "device %d payload %d: encode failed: %s\n", device, p, sqlite3_errmsg(db)); + rc = SQLITE_ERROR; + } + sqlite3_reset(enc); + sqlite3_stmt *mx = NULL; + if (rc == SQLITE_OK && sqlite3_prepare_v2(db, "SELECT max(db_version) FROM cloudsync_changes", -1, &mx, NULL) == SQLITE_OK) { + if (sqlite3_step(mx) == SQLITE_ROW) last = sqlite3_column_int64(mx, 0); + } + sqlite3_finalize(mx); + } + if (rc != SQLITE_OK) fprintf(stderr, "device %d: %s\n", device, sqlite3_errmsg(db)); + sqlite3_finalize(ins); + sqlite3_finalize(enc); + sqlite3_close(db); + return rc == SQLITE_OK ? 0 : 1; +} + +// Ordered content of t, used to compare a stressed database with the reference. +static char *table_content(sqlite3 *db) { + sqlite3_stmt *stmt = NULL; + char *out = NULL; + if (sqlite3_prepare_v2(db, + "SELECT count(*) || ':' || ifnull(group_concat(id || '|' || dev || '|' || n || '|' || body, ','), '') " + "FROM (SELECT * FROM t ORDER BY id)", -1, &stmt, NULL) == SQLITE_OK && + sqlite3_step(stmt) == SQLITE_ROW) { + const char *s = (const char *)sqlite3_column_text(stmt, 0); + out = s ? strdup(s) : NULL; + } + sqlite3_finalize(stmt); + return out; +} + +static bool is_busy(int rc, const char *msg) { + if ((rc & 0xff) == SQLITE_BUSY || (rc & 0xff) == SQLITE_LOCKED) return true; + return msg && strstr(msg, "database is locked") != NULL; +} + +static void record_code(worker_t *w, int rc, const char *msg) { + for (int i = 0; i < w->ncodes; i++) { + if (w->codes[i].code == rc) { w->codes[i].count++; return; } + } + if (w->ncodes == MAX_CODES) return; + code_count *c = &w->codes[w->ncodes++]; + c->code = rc; + c->count = 1; + snprintf(c->message, sizeof(c->message), "%s", msg ? msg : ""); +} + +static void sleep_jittered(unsigned int *seed, int max_ms) { + if (max_ms <= 0) return; + usleep((useconds_t)(rand_r(seed) % (unsigned)(max_ms * 1000)) + 100); +} + +static void rollback_if_open(sqlite3 *db) { + if (!sqlite3_get_autocommit(db)) sqlite3_exec(db, "ROLLBACK", NULL, NULL, NULL); +} + +// One attempt to apply a payload: an error that is not BUSY is fatal. +static attempt_result apply_once(worker_t *w, sqlite3 *db, sqlite3_stmt *apply, const payload_t *payload) { + const char *begin = mode_begin[w->mode]; + int rc; + if (begin) { + rc = sqlite3_exec(db, begin, NULL, NULL, NULL); + if (rc != SQLITE_OK) { + record_code(w, sqlite3_extended_errcode(db), sqlite3_errmsg(db)); + if (!is_busy(rc, sqlite3_errmsg(db))) goto fatal; + rollback_if_open(db); + w->busy_begin++; + return ATTEMPT_BUSY; + } + } + + sqlite3_bind_blob(apply, 1, payload->blob, payload->size, SQLITE_STATIC); + rc = sqlite3_step(apply); + if (rc != SQLITE_ROW) { + int code = sqlite3_extended_errcode(db); + char msg[256]; + snprintf(msg, sizeof(msg), "%s", sqlite3_errmsg(db)); + sqlite3_reset(apply); + record_code(w, code, msg); + rollback_if_open(db); + if (!is_busy(code, msg)) { + snprintf(w->fatal_message, sizeof(w->fatal_message), "apply: %s (%d)", msg, code); + return ATTEMPT_FATAL; + } + w->busy_apply++; + return ATTEMPT_BUSY; + } + sqlite3_reset(apply); + if (!begin) return ATTEMPT_OK; + + rc = sqlite3_exec(db, "COMMIT", NULL, NULL, NULL); + for (int i = 0; is_busy(rc, NULL) && i < COMMIT_RETRIES; i++) { + usleep((useconds_t)(10000 * (i + 1))); + rc = sqlite3_exec(db, "COMMIT", NULL, NULL, NULL); + } + if (rc == SQLITE_OK) return ATTEMPT_OK; + record_code(w, sqlite3_extended_errcode(db), sqlite3_errmsg(db)); + if (!is_busy(rc, sqlite3_errmsg(db))) goto fatal; + rollback_if_open(db); + w->busy_commit++; + return ATTEMPT_BUSY; + +fatal: + snprintf(w->fatal_message, sizeof(w->fatal_message), "%s", sqlite3_errmsg(db)); + rollback_if_open(db); + return ATTEMPT_FATAL; +} + +static void *worker_main(void *arg) { + worker_t *w = (worker_t *)arg; + unsigned int seed = (unsigned int)(w->device * 7919 + 17); + sqlite3 *db = open_db(DB_PATH); + sqlite3_stmt *apply = NULL; + if (!db || sqlite3_prepare_v2(db, "SELECT cloudsync_payload_apply(?1)", -1, &apply, NULL) != SQLITE_OK) { + w->fatal = 1; + snprintf(w->fatal_message, sizeof(w->fatal_message), "setup: %s", db ? sqlite3_errmsg(db) : "open failed"); + goto done; + } + sqlite3_busy_timeout(db, cfg.busy_timeout_ms); + + pthread_mutex_lock(&gate_lock); + while (!gate_open) pthread_cond_wait(&gate_cond, &gate_lock); + pthread_mutex_unlock(&gate_lock); + + for (int p = 0; p < w->npayloads && !w->fatal; p++) { + bool applied = false; + for (int request = 0; !applied && !w->fatal; request++) { + if (request == cfg.max_requests) { + snprintf(w->fatal_message, sizeof(w->fatal_message), "payload %d: gave up after %d failed requests", p, request); + w->fatal = 1; + break; + } + if (request > 0) sleep_jittered(&seed, cfg.request_backoff_ms); + for (int attempt = 0; attempt <= cfg.retries; attempt++) { + if (attempt > 0) sleep_jittered(&seed, cfg.backoff_ms); + w->attempts++; + attempt_result res = apply_once(w, db, apply, &w->payloads[p]); + if (res == ATTEMPT_OK) { applied = true; break; } + if (res == ATTEMPT_FATAL) { w->fatal = 1; break; } + } + if (!applied && !w->fatal) w->failed_requests++; + } + if (applied) w->applied++; + } + +done: + sqlite3_finalize(apply); + sqlite3_close(db); + return NULL; +} + +static void remove_db_files(void) { + remove(DB_PATH); + remove(DB_PATH "-wal"); + remove(DB_PATH "-shm"); +} + +static bool begin_concurrent_supported(void) { + sqlite3 *db = NULL; + sqlite3_stmt *stmt = NULL; + bool ok = sqlite3_open(":memory:", &db) == SQLITE_OK && + sqlite3_prepare_v2(db, "BEGIN CONCURRENT", -1, &stmt, NULL) == SQLITE_OK; + sqlite3_finalize(stmt); + sqlite3_close(db); + return ok; +} + +// Runs one mode against a fresh database. Returns the number of failed checks. +static int run_mode(apply_mode mode, payload_t **payloads, const char *reference) { + remove_db_files(); + sqlite3 *setup = open_db(DB_PATH); + if (!setup) return 1; + if (db_exec(setup, "PRAGMA journal_mode=WAL; PRAGMA synchronous=NORMAL;") != SQLITE_OK || + db_exec(setup, schema_sql) != SQLITE_OK) { + sqlite3_close(setup); + return 1; + } + sqlite3_close(setup); + + worker_t *workers = calloc((size_t)cfg.devices, sizeof(worker_t)); + pthread_t *threads = calloc((size_t)cfg.devices, sizeof(pthread_t)); + if (!workers || !threads) { free(workers); free(threads); return 1; } + gate_open = false; + for (int d = 0; d < cfg.devices; d++) { + workers[d].device = d; + workers[d].mode = mode; + workers[d].payloads = payloads[d]; + workers[d].npayloads = cfg.payloads; + pthread_create(&threads[d], NULL, worker_main, &workers[d]); + } + + // Let every thread open its connection before the clock starts. + usleep(200 * 1000); + double cpu0 = cpu_ms(), t0 = monotonic_ms(); + pthread_mutex_lock(&gate_lock); + gate_open = true; + pthread_cond_broadcast(&gate_cond); + pthread_mutex_unlock(&gate_lock); + for (int d = 0; d < cfg.devices; d++) pthread_join(threads[d], NULL); + double wall = monotonic_ms() - t0, cpu = cpu_ms() - cpu0; + + int failures = 0, applied = 0, attempts = 0, failed_requests = 0, busy_begin = 0, busy_apply = 0, busy_commit = 0; + code_count totals[MAX_CODES]; + int ntotals = 0; + printf("\n[%s]\n", mode_names[mode]); + for (int d = 0; d < cfg.devices; d++) { + worker_t *w = &workers[d]; + applied += w->applied; + attempts += w->attempts; + failed_requests += w->failed_requests; + busy_begin += w->busy_begin; + busy_apply += w->busy_apply; + busy_commit += w->busy_commit; + if (w->fatal) { + printf(" FAIL device %d: %s\n", d, w->fatal_message); + failures++; + } + for (int i = 0; i < w->ncodes; i++) { + int j = 0; + while (j < ntotals && totals[j].code != w->codes[i].code) j++; + if (j == ntotals) { + if (ntotals == MAX_CODES) continue; + totals[ntotals++] = w->codes[i]; + } else { + totals[j].count += w->codes[i].count; + } + } + } + + int expected = cfg.devices * cfg.payloads; + printf(" payloads applied %d / %d\n", applied, expected); + printf(" apply attempts %d (%d retries)\n", attempts, attempts - applied); + printf(" failed requests %d\n", failed_requests); + printf(" BUSY at BEGIN %d\n", busy_begin); + printf(" BUSY in apply %d\n", busy_apply); + printf(" BUSY at COMMIT %d\n", busy_commit); + for (int i = 0; i < ntotals; i++) + printf(" code %-5d x%-6d %s\n", totals[i].code, totals[i].count, totals[i].message); + printf(" wall time %.0f ms\n", wall); + printf(" process CPU time %.0f ms (%.2f cores busy)\n", cpu, wall > 0 ? cpu / wall : 0); + printf(" CPU per payload %.2f ms\n", applied > 0 ? cpu / applied : 0); + + if (applied != expected) failures++; + if (mode == MODE_IMMEDIATE && busy_apply > 0) { + printf(" FAIL the apply statement returned BUSY while holding the write lock\n"); + failures++; + } + + sqlite3 *check = open_db(DB_PATH); + char *content = check ? table_content(check) : NULL; + if (!content || strcmp(content, reference) != 0) { + printf(" FAIL content differs from the serial reference\n"); + failures++; + } else { + printf(" content matches the serial reference\n"); + } + free(content); + sqlite3_close(check); + free(workers); + free(threads); + return failures; +} + +static bool mode_selected(const char *list, const char *name) { + size_t len = strlen(name); + for (const char *p = list; (p = strstr(p, name)) != NULL; p += len) { + bool starts = (p == list || p[-1] == ','); + bool ends = (p[len] == '\0' || p[len] == ','); + if (starts && ends) return true; + } + return false; +} + +int main(void) { + cfg.devices = env_int("APPLY_STRESS_DEVICES", 8); + cfg.payloads = env_int("APPLY_STRESS_PAYLOADS", 20); + cfg.versions = env_int("APPLY_STRESS_VERSIONS", 5); + cfg.rows = env_int("APPLY_STRESS_ROWS", 4); + cfg.shared = env_int("APPLY_STRESS_SHARED", 0) != 0; + cfg.busy_timeout_ms = env_int("APPLY_STRESS_BUSY_TIMEOUT_MS", 5000); + cfg.retries = env_int("APPLY_STRESS_RETRIES", 6); + cfg.backoff_ms = env_int("APPLY_STRESS_BACKOFF_MS", 5); + cfg.request_backoff_ms = env_int("APPLY_STRESS_REQUEST_BACKOFF_MS", 100); + cfg.max_requests = env_int("APPLY_STRESS_MAX_REQUESTS", 1000); + cfg.ext_path = env_str("APPLY_STRESS_EXT", EXT_PATH); + const char *modes = env_str("APPLY_STRESS_MODE", "autocommit,deferred,immediate,concurrent"); + if (cfg.devices < 1 || cfg.payloads < 1 || cfg.versions < 1 || cfg.rows < 1 || cfg.max_requests < 1) { + fprintf(stderr, "devices, payloads, versions, rows and max requests must be positive\n"); + return 1; + } + if (sqlite3_threadsafe() == 0) { + fprintf(stderr, "SQLite was built without thread safety\n"); + return 1; + } + + printf("SQLite %s, devices %d, payloads/device %d, db_versions/payload %d, rows/db_version %d, %s keys\n", + sqlite3_libversion(), cfg.devices, cfg.payloads, cfg.versions, cfg.rows, cfg.shared ? "shared" : "distinct"); + printf("busy_timeout %d ms, %d retries per request (backoff up to %d ms), request re-send after up to %d ms\n", + cfg.busy_timeout_ms, cfg.retries, cfg.backoff_ms, cfg.request_backoff_ms); + + payload_t **payloads = calloc((size_t)cfg.devices, sizeof(payload_t *)); + if (!payloads) return 1; + for (int d = 0; d < cfg.devices; d++) { + payloads[d] = calloc((size_t)cfg.payloads, sizeof(payload_t)); + if (!payloads[d] || build_device_payloads(d, payloads[d]) != 0) return 1; + } + + // Apply order does not matter to a converged CRDT, so a serial apply is the oracle. + sqlite3 *ref = open_db(":memory:"); + if (!ref || db_exec(ref, schema_sql) != SQLITE_OK) return 1; + sqlite3_stmt *apply = NULL; + if (sqlite3_prepare_v2(ref, "SELECT cloudsync_payload_apply(?1)", -1, &apply, NULL) != SQLITE_OK) return 1; + for (int d = 0; d < cfg.devices; d++) { + for (int p = 0; p < cfg.payloads; p++) { + sqlite3_bind_blob(apply, 1, payloads[d][p].blob, payloads[d][p].size, SQLITE_STATIC); + if (sqlite3_step(apply) != SQLITE_ROW) { + fprintf(stderr, "reference apply failed: %s\n", sqlite3_errmsg(ref)); + return 1; + } + sqlite3_reset(apply); + } + } + sqlite3_finalize(apply); + char *reference = table_content(ref); + sqlite3_close(ref); + if (!reference) return 1; + + int failures = 0; + for (int m = 0; m < MODE_COUNT; m++) { + if (!mode_selected(modes, mode_names[m])) continue; + if (m == MODE_CONCURRENT && !begin_concurrent_supported()) { + printf("\n[%s]\n skipped: this SQLite build has no BEGIN CONCURRENT (set APPLY_STRESS_SQLITE_DIR)\n", mode_names[m]); + continue; + } + failures += run_mode((apply_mode)m, payloads, reference); + } + + free(reference); + for (int d = 0; d < cfg.devices; d++) { + for (int p = 0; p < cfg.payloads; p++) free(payloads[d][p].blob); + free(payloads[d]); + } + free(payloads); + remove_db_files(); + + printf("\n%s\n", failures ? "FAILED" : "PASSED"); + return failures ? 1 : 0; +}