// test_db_sync.c — sync protocol test: peer_empty, hash_MATCH, divergence merge #include #include #include #include #include #include #include "../lib/platform_compat.h" #include "test_utils.h" #ifdef _WIN32 #include #include #else #include #endif #include "etcp.h" #include "etcp_connections.h" #include "../src/config_parser.h" #include "../src/config_updater.h" #include "../src/utun_instance.h" #include "routing.h" #include "../src/tun_if.h" #include "secure_channel.h" #include "secure_channel.h" #include "../src/db_sync.h" #include "../lib/u_async.h" #include "../lib/debug_config.h" #include "../lib/mem.h" #define TEST_TIMEOUT_TB 120000 // 12s #define PHASE_TIMEOUT_TB 100000 // 10s per phase #define POLL_INTERVAL_MS 1 static struct UTUN_INSTANCE* inst_a = NULL; static struct UTUN_INSTANCE* inst_b = NULL; static struct UTUN_INSTANCE* inst_c = NULL; static struct DB_SYNC_INSTANCE* si_a = NULL; static struct DB_SYNC_INSTANCE* si_b = NULL; static struct DB_SYNC_INSTANCE* si_c = NULL; static struct UASYNC* ua = NULL; static int test_phase = 0; static void* timeout_id = NULL; static char temp_dir[] = "/tmp/utun_dbsync_XXXXXX"; static char config_a[256], config_b[256], config_c[256]; static int port_a_srv, port_b_srv, port_c_srv; static int write_file(const char* path, const char* fmt, ...) { va_list ap; FILE* f = fopen(path, "w"); if (!f) return -1; va_start(ap, fmt); vfprintf(f, fmt, ap); va_end(ap); fclose(f); return 0; } static char* get_pubkey(const char* path) { struct utun_config* cfg = parse_config(path); if (!cfg) return NULL; char* pub = strdup(cfg->global.my_public_key_hex); free_config(cfg); return pub; } static int create_temp_configs(void) { if (test_mkdtemp(temp_dir) != 0) { fprintf(stderr, "mkdtemp failed\n"); return -1; } int base = 42000 + (getpid() % 15000); port_a_srv = base; port_b_srv = base + 1; port_c_srv = base + 2; snprintf(config_a, sizeof(config_a), "%s/a.conf", temp_dir); snprintf(config_b, sizeof(config_b), "%s/b.conf", temp_dir); snprintf(config_c, sizeof(config_c), "%s/c.conf", temp_dir); char db_path[320]; snprintf(db_path, sizeof(db_path), "%s/db_a", temp_dir); utun_mkdir(db_path, 0755); snprintf(db_path, sizeof(db_path), "%s/db_b", temp_dir); utun_mkdir(db_path, 0755); snprintf(db_path, sizeof(db_path), "%s/db_c", temp_dir); utun_mkdir(db_path, 0755); write_file(config_a, "[global]\nmy_node_id=0xAAAAAAAAAAAAAAAA\ntun_ip=10.200.0.1/24\n" "tun_ifname=tun200\nkeepalive_adaptive=0\ndb_path=%s/db_a\ndb_sync_enabled=1\n\n" "[server: srv_a]\naddr=127.0.0.1:%d\ntype=public\n\n" "[allowed_keys]\nallow_all=1\n", temp_dir, port_a_srv); config_ensure_keys_and_node_id(config_a); char* pub_a = get_pubkey(config_a); if (!pub_a) { fprintf(stderr, "pub_a fail\n"); return -1; } write_file(config_b, "[global]\nmy_node_id=0xBBBBBBBBBBBBBBBB\ntun_ip=10.200.0.2/24\n" "tun_ifname=tun201\nkeepalive_adaptive=0\ndb_path=%s/db_b\ndb_sync_enabled=1\n\n" "[server: srv_b]\naddr=127.0.0.1:%d\ntype=public\n\n" "[client: to_a]\nkeepalive=1\npeer_public_key=%s\nlink=srv_b:127.0.0.1:%d\n\n" "[allowed_keys]\nallow_all=1\n", temp_dir, port_b_srv, pub_a, port_a_srv); write_file(config_c, "[global]\nmy_node_id=0xCCCCCCCCCCCCCCCC\ntun_ip=10.200.0.3/24\n" "tun_ifname=tun202\nkeepalive_adaptive=0\ndb_path=%s/db_c\ndb_sync_enabled=1\n\n" "[server: srv_c]\naddr=127.0.0.1:%d\ntype=public\n\n" "[client: to_a]\nkeepalive=1\npeer_public_key=%s\nlink=srv_c:127.0.0.1:%d\n\n" "[allowed_keys]\nallow_all=1\n", temp_dir, port_c_srv, pub_a, port_a_srv); free(pub_a); config_ensure_keys_and_node_id(config_b); config_ensure_keys_and_node_id(config_c); return 0; } static void cleanup_temp_configs(void) { unlink(config_a); unlink(config_b); unlink(config_c); char pa[320]; snprintf(pa, sizeof(pa), "%s/db_a/chats.db", temp_dir); unlink(pa); snprintf(pa, sizeof(pa), "%s/db_a/chats.db-wal", temp_dir); unlink(pa); snprintf(pa, sizeof(pa), "%s/db_a/chats.db-shm", temp_dir); unlink(pa); snprintf(pa, sizeof(pa), "%s/db_b/chats.db", temp_dir); unlink(pa); snprintf(pa, sizeof(pa), "%s/db_b/chats.db-wal", temp_dir); unlink(pa); snprintf(pa, sizeof(pa), "%s/db_b/chats.db-shm", temp_dir); unlink(pa); snprintf(pa, sizeof(pa), "%s/db_c/chats.db", temp_dir); unlink(pa); snprintf(pa, sizeof(pa), "%s/db_c/chats.db-wal", temp_dir); unlink(pa); snprintf(pa, sizeof(pa), "%s/db_c/chats.db-shm", temp_dir); unlink(pa); snprintf(pa, sizeof(pa), "%s/db_a", temp_dir); test_rmdir(pa); snprintf(pa, sizeof(pa), "%s/db_b", temp_dir); test_rmdir(pa); snprintf(pa, sizeof(pa), "%s/db_c", temp_dir); test_rmdir(pa); test_rmdir(temp_dir); } static void test_timeout(void* arg) { (void)arg; test_phase = 2; } static int wait_for(const char* desc, int (*cond)(void), int timeout_tb) { uint64_t start = get_time_tb(); while (!cond() && (get_time_tb() - start) < (uint64_t)timeout_tb && test_phase == 0) uasync_poll(ua, POLL_INTERVAL_MS); if (cond()) return 1; if (test_phase == 0) { fprintf(stderr, "TIMEOUT: %s\n", desc); test_phase = 2; } return 0; } static int cond_links_init(void) { if (!inst_a || !inst_b) return 0; struct ll_entry* e = inst_a->connections->head; int links = 0; while (e) { struct conn_queue_entry* ce = (struct conn_queue_entry*)e->data; struct ETCP_LINK* l = ce->conn->links; while (l) { if (l->initialized) links++; l = l->next; } e = e->next; } return links >= (inst_c ? 2 : 1); } static uint32_t ca_target, cb_target, cc_target; static int _cond_ca(void) { return si_a && db_sync_count(si_a) == ca_target; } static int _cond_cb(void) { return si_b && db_sync_count(si_b) == cb_target; } static int _cond_cc(void) { return si_c && db_sync_count(si_c) == cc_target; } static int _cond_both(void) { return _cond_ca() && _cond_cb(); } static int _cond_all(void) { return _cond_ca() && _cond_cb() && _cond_cc(); } static int insert_many(struct DB_SYNC_INSTANCE* si, struct UTUN_INSTANCE* inst, int start, int count) { char buf[128]; for (int i = start; i < start + count && test_phase == 0; i++) { snprintf(buf, sizeof(buf), "{\"idx\":%d,\"val\":\"data_%d\"}", i, i); uint64_t ts = db_sync_next_timestamp(si); uint8_t sig_msg[256]; size_t off = 0; memcpy(sig_msg + off, &ts, 8); off += 8; size_t jl = strlen(buf); memcpy(sig_msg + off, buf, jl); off += jl; uint8_t sig[64]; if (sc_ed25519_sign(inst->my_ed25519_privkey, sig_msg, off, sig) != SC_OK) { fprintf(stderr,"sign fail\n"); return -1; } if (db_sync_insert_signed(si, buf, jl, sig, 64, ts) < 0) { fprintf(stderr,"insert fail\n"); return -1; } } return 0; } static void remove_si(struct DB_SYNC_INSTANCE** psi) { if (*psi) { db_sync_instance_remove(*psi); *psi = NULL; } } int main(void) { printf("=== test_db_sync ===\n"); debug_config_init(); debug_set_level(DEBUG_LEVEL_ERROR); debug_set_category_level(DEBUG_CATEGORY_DB_SYNC, DEBUG_LEVEL_DEBUG); if (create_temp_configs() != 0) { cleanup_temp_configs(); return 1; } utun_instance_set_tun_init_enabled(0); ua = uasync_create(); if (!ua) { cleanup_temp_configs(); return 1; } inst_a = utun_instance_create(ua, config_a); inst_b = utun_instance_create(ua, config_b); if (!inst_a || !inst_b || utun_instance_init(inst_a) != 0 || utun_instance_init(inst_b) != 0) { fprintf(stderr, "init fail\n"); cleanup_temp_configs(); return 1; } timeout_id = uasync_set_timeout(ua, TEST_TIMEOUT_TB, NULL, test_timeout, "global_timeout"); if (!wait_for("links up", cond_links_init, PHASE_TIMEOUT_TB)) { test_phase = 2; goto done; } // =================================================================== // Phase 1: peer_empty. add_quiet (no auto-sync) → insert A=5, // reinitiate → INIT_SYNC(5) → B INIT_RESP(ch8=0,sc=0) → A sends all 5. // =================================================================== printf("Phase 1: peer_empty\n"); si_a = db_sync_instance_add(inst_a, "test", 1, 0); si_b = db_sync_instance_add(inst_b, "test", 1, 0); if (!si_a || !si_b) { test_phase = 2; goto done; } if (insert_many(si_a, inst_a, 0, 5) != 0) { test_phase = 2; goto done; } db_sync_reinitiate(si_b, inst_a->node_id); cb_target = 5; if (!wait_for("B=5", _cond_cb, PHASE_TIMEOUT_TB)) { test_phase = 2; goto done; } if (db_sync_count(si_b) != 5 || db_sync_chain_verify(si_a) || db_sync_chain_verify(si_b)) { test_phase = 2; goto done; } printf(" PASS\n"); // =================================================================== // Phase 2: hash_MATCH tail-send. Insert 3 more on A (PUSH disabled — // si_a from phase 1 has sync_state=2? No, reinitiate set it to 1, then // SYNC_DONE set it to 2. Disable PUSH, insert, reinitiate. // =================================================================== printf("Phase 2: hash_MATCH tail-send\n"); db_sync_peer_set_state(si_a, inst_b->node_id, 0); if (insert_many(si_a, inst_a, 5, 3) != 0) { test_phase = 2; goto done; } db_sync_reinitiate(si_b, inst_a->node_id); cb_target = 8; if (!wait_for("B=8", _cond_cb, PHASE_TIMEOUT_TB)) { test_phase = 2; goto done; } if (db_sync_count(si_b) != 8 || db_sync_chain_verify(si_a) || db_sync_chain_verify(si_b)) { test_phase = 2; goto done; } printf(" PASS\n"); // =================================================================== // Phase 3: divergence merge. remove si, add_quiet → insert A=3 B=2 → // reinitiate both → divergence → REFINE merge to 5. // =================================================================== printf("Phase 3: divergence merge\n"); remove_si(&si_a); remove_si(&si_b); si_a = db_sync_instance_add(inst_a, "test2", 2, 0); si_b = db_sync_instance_add(inst_b, "test2", 2, 0); if (!si_a || !si_b) { test_phase = 2; goto done; } if (insert_many(si_a, inst_a, 0, 3) != 0 || insert_many(si_b, inst_b, 10, 2) != 0) { test_phase = 2; goto done; } db_sync_reinitiate(si_b, inst_a->node_id); ca_target = 5; cb_target = 5; if (!wait_for("both=5", _cond_both, PHASE_TIMEOUT_TB)) { test_phase = 2; goto done; } if (db_sync_chain_verify(si_a) || db_sync_chain_verify(si_b)) { test_phase = 2; goto done; } printf(" PASS\n"); // =================================================================== // Phase 4: divergence merge. remove si, add_quiet → insert A=2 B=2 → // reinitiate both → merge to 4. // =================================================================== printf("Phase 4: divergence merge 2\n"); remove_si(&si_a); remove_si(&si_b); si_a = db_sync_instance_add(inst_a, "test_div", 30, 0); si_b = db_sync_instance_add(inst_b, "test_div", 30, 0); if (!si_a || !si_b) { test_phase = 2; goto done; } if (insert_many(si_a, inst_a, 0, 2) != 0 || insert_many(si_b, inst_b, 10, 2) != 0) { test_phase = 2; goto done; } db_sync_reinitiate(si_b, inst_a->node_id); ca_target = 4; cb_target = 4; if (!wait_for("both=4", _cond_both, PHASE_TIMEOUT_TB)) { test_phase = 2; goto done; } if (db_sync_chain_verify(si_a) || db_sync_chain_verify(si_b)) { test_phase = 2; goto done; } printf(" PASS\n"); // =================================================================== // Phase 5: PUSH not-at-tail. Use db_sync_instance_add (normal, with auto-sync) // so PUSH fires. Insert early timestamp record → PUSH → cascade_from. // =================================================================== printf("Phase 5: PUSH not-at-tail\n"); remove_si(&si_a); remove_si(&si_b); si_a = db_sync_instance_add(inst_a, "push_t", 40, 1); si_b = db_sync_instance_add(inst_b, "push_t", 40, 1); if (!si_a || !si_b) { test_phase = 2; goto done; } if (insert_many(si_a, inst_a, 0, 2) != 0) { test_phase = 2; goto done; } cb_target = 2; if (!wait_for("seed=2", _cond_cb, PHASE_TIMEOUT_TB)) { test_phase = 2; goto done; } { char ebuf[64]; snprintf(ebuf, sizeof(ebuf), "{\"idx\":\"early\",\"val\":\"not_at_tail\"}"); uint64_t early_ts = 500; uint8_t msg[256]; size_t moff = 0; memcpy(msg + moff, &early_ts, 8); moff += 8; size_t jl = strlen(ebuf); memcpy(msg + moff, ebuf, jl); moff += jl; uint8_t sig[64]; if (sc_ed25519_sign(inst_a->my_ed25519_privkey, msg, moff, sig) != SC_OK || db_sync_insert_signed(si_a, ebuf, jl, sig, 64, early_ts) != 0) { test_phase = 2; goto done; } } ca_target = 3; cb_target = 3; if (!wait_for("both=3", _cond_both, PHASE_TIMEOUT_TB)) { test_phase = 2; goto done; } if (db_sync_chain_verify(si_a) || db_sync_chain_verify(si_b)) { test_phase = 2; goto done; } printf(" PASS\n"); // =================================================================== // Phase 6: triple star — cascade notification. // =================================================================== printf("Phase 6: triple star\n"); remove_si(&si_a); remove_si(&si_b); inst_c = utun_instance_create(ua, config_c); if (!inst_c || utun_instance_init(inst_c) != 0) { fprintf(stderr, "inst_c fail\n"); test_phase = 2; goto done; } if (!wait_for("C links up", cond_links_init, PHASE_TIMEOUT_TB)) { test_phase = 2; goto done; } si_a = db_sync_instance_add(inst_a, "triple", 70, 0); si_b = db_sync_instance_add(inst_b, "triple", 70, 0); si_c = db_sync_instance_add(inst_c, "triple", 70, 0); if (!si_a || !si_b || !si_c) { test_phase = 2; goto done; } if (insert_many(si_a, inst_a, 0, 2) != 0 || insert_many(si_b, inst_b, 10, 1) != 0 || insert_many(si_c, inst_c, 20, 1) != 0) { test_phase = 2; goto done; } db_sync_reinitiate(si_b, inst_a->node_id); db_sync_reinitiate(si_c, inst_a->node_id); ca_target = 4; cb_target = 4; cc_target = 4; if (!wait_for("all=4", _cond_all, PHASE_TIMEOUT_TB)) { test_phase = 2; goto done; } if (db_sync_chain_verify(si_a) || db_sync_chain_verify(si_b) || db_sync_chain_verify(si_c)) { test_phase = 2; goto done; } printf(" merge=4 PASS\n"); // 6b: PUSH at tail on A → delivered to B and C via PUSH (sync_state >=1 after merge). if (insert_many(si_a, inst_a, 30, 1) != 0) { test_phase = 2; goto done; } ca_target = 5; cb_target = 5; cc_target = 5; if (!wait_for("all=5 (tail PUSH)", _cond_all, PHASE_TIMEOUT_TB)) { test_phase = 2; goto done; } if (db_sync_chain_verify(si_a) || db_sync_chain_verify(si_b) || db_sync_chain_verify(si_c)) { test_phase = 2; goto done; } printf(" tail PUSH=5 PASS\n"); // 6c: PUSH not-at-tail on A → delivered to B and C → cascade_from, synced_pos reset. { char ebuf[64]; snprintf(ebuf, sizeof(ebuf), "{\"idx\":\"mid\",\"val\":\"triple_insert\"}"); uint64_t early_ts = 500; uint8_t msg[256]; size_t moff = 0; memcpy(msg + moff, &early_ts, 8); moff += 8; size_t jl = strlen(ebuf); memcpy(msg + moff, ebuf, jl); moff += jl; uint8_t sig[64]; if (sc_ed25519_sign(inst_a->my_ed25519_privkey, msg, moff, sig) != SC_OK || db_sync_insert_signed(si_a, ebuf, jl, sig, 64, early_ts) != 0) { test_phase = 2; goto done; } } ca_target = 6; cb_target = 6; cc_target = 6; if (!wait_for("all=6 (mid PUSH cascade)", _cond_all, PHASE_TIMEOUT_TB)) { test_phase = 2; goto done; } if (db_sync_chain_verify(si_a) || db_sync_chain_verify(si_b) || db_sync_chain_verify(si_c)) { test_phase = 2; goto done; } printf(" mid PUSH cascade=6 PASS\n"); // =================================================================== // Phase 7: randomized triple — cold start divergence, live PUSH, // disconnect/reconnect with additional random inserts. // =================================================================== printf("Phase 7: randomized triple\n"); remove_si(&si_a); remove_si(&si_b); remove_si(&si_c); srand((unsigned)time(NULL)); printf("seed=%u\n", (unsigned)time(NULL)); si_a = db_sync_instance_add(inst_a, "rnd_triple", 80, 0); si_b = db_sync_instance_add(inst_b, "rnd_triple", 80, 0); si_c = db_sync_instance_add(inst_c, "rnd_triple", 80, 0); if (!si_a || !si_b || !si_c) { test_phase = 2; goto done; } // -- Round 1: cold start divergence merge with 0..20 random records per peer -- int r1_a = rand() % 2, r1_b = rand() % 2, r1_c = rand() % 2; int total = r1_a + r1_b + r1_c; printf(" R1: A=%d B=%d C=%d (total=%d)\n", r1_a, r1_b, r1_c, total); if ((r1_a > 0 && insert_many(si_a, inst_a, 0, r1_a) != 0) || (r1_b > 0 && insert_many(si_b, inst_b, 100, r1_b) != 0) || (r1_c > 0 && insert_many(si_c, inst_c, 200, r1_c) != 0)) { test_phase = 2; goto done; } db_sync_reinitiate(si_b, inst_a->node_id); db_sync_reinitiate(si_c, inst_a->node_id); ca_target = cb_target = cc_target = (uint32_t)total; if (!wait_for("R1 converge", _cond_all, PHASE_TIMEOUT_TB)) { test_phase = 2; goto done; } { uint64_t ha, hb, hc; if (db_sync_last_chain_hash8(si_a, &ha) || db_sync_last_chain_hash8(si_b, &hb) || db_sync_last_chain_hash8(si_c, &hc) || ha != hb || hb != hc) { test_phase = 2; goto done; } if (db_sync_chain_verify(si_a) || db_sync_chain_verify(si_b) || db_sync_chain_verify(si_c)) { test_phase = 2; goto done; } } printf(" R1 PASS (%u records)\n", ca_target); // -- Round 2: live PUSH — 20 records on random peer, verify convergence -- { int r2_peer = rand() % 3; struct DB_SYNC_INSTANCE* pick = (r2_peer == 0) ? si_a : (r2_peer == 1) ? si_b : si_c; struct UTUN_INSTANCE* pick_inst = (r2_peer == 0) ? inst_a : (r2_peer == 1) ? inst_b : inst_c; const char* pick_name = (r2_peer == 0) ? "A" : (r2_peer == 1) ? "B" : "C"; printf(" R2: +5 on peer %s\n", pick_name); if (insert_many(pick, pick_inst, 300, 5) != 0) { test_phase = 2; goto done; } ca_target = cb_target = cc_target = (uint32_t)(total + 5); if (!wait_for("R2 converge", _cond_all, PHASE_TIMEOUT_TB)) { test_phase = 2; goto done; } uint64_t ha, hb, hc; if (db_sync_last_chain_hash8(si_a, &ha) || db_sync_last_chain_hash8(si_b, &hb) || db_sync_last_chain_hash8(si_c, &hc) || ha != hb || hb != hc) { test_phase = 2; goto done; } printf(" R2 PASS (%u records)\n", ca_target); total += 5; } // -- Round 3: simulate disconnect (reset sync state), add random records, reinitiate, sync -- printf(" R3: disconnect/reset sync state, insert random, reinitiate\n"); db_sync_peer_set_state(si_a, inst_b->node_id, 0); db_sync_peer_set_state(si_a, inst_c->node_id, 0); db_sync_peer_set_state(si_b, inst_a->node_id, 0); db_sync_peer_set_state(si_c, inst_a->node_id, 0); int r3_a = rand() % 2, r3_b = rand() % 2, r3_c = rand() % 2; total += r3_a + r3_b + r3_c; printf(" R3: A=%d B=%d C=%d (total=%d)\n", r3_a, r3_b, r3_c, total); if ((r3_a > 0 && insert_many(si_a, inst_a, 400, r3_a) != 0) || (r3_b > 0 && insert_many(si_b, inst_b, 500, r3_b) != 0) || (r3_c > 0 && insert_many(si_c, inst_c, 600, r3_c) != 0)) { test_phase = 2; goto done; } db_sync_reinitiate(si_b, inst_a->node_id); db_sync_reinitiate(si_c, inst_a->node_id); ca_target = cb_target = cc_target = (uint32_t)total; if (!wait_for("R3 converge", _cond_all, PHASE_TIMEOUT_TB)) { test_phase = 2; goto done; } { uint64_t ha, hb, hc; if (db_sync_last_chain_hash8(si_a, &ha) || db_sync_last_chain_hash8(si_b, &hb) || db_sync_last_chain_hash8(si_c, &hc) || ha != hb || hb != hc) { test_phase = 2; goto done; } } printf(" R3 PASS (%u records)\n", ca_target); done: remove_si(&si_a); remove_si(&si_b); remove_si(&si_c); if (timeout_id && ua) { uasync_cancel_timeout(ua, timeout_id); timeout_id = NULL; } if (inst_c) { inst_c->running = 0; utun_instance_destroy(inst_c); inst_c = NULL; } if (inst_b) { inst_b->running = 0; utun_instance_destroy(inst_b); inst_b = NULL; } if (inst_a) { inst_a->running = 0; utun_instance_destroy(inst_a); inst_a = NULL; } if (ua) { uasync_destroy(ua, 0); ua = NULL; } cleanup_temp_configs(); if (test_phase == 0) test_phase = 1; printf("=== %s ===\n", test_phase == 1 ? "PASS" : "FAIL"); return test_phase == 1 ? 0 : 1; }