diff --git a/src/db_sync.c b/src/db_sync.c index 27b1bd2b..e87127ad 100644 --- a/src/db_sync.c +++ b/src/db_sync.c @@ -738,19 +738,6 @@ static void db_handle_init_resp(struct DB_SYNC_INSTANCE* si, uint64_t my_dh = 0; db_datahash_at(si, tp, &my_dh); - if (sc == 0 && my_dh == peer_dh) { - struct SI_PEER* sp = si_peer_find(si, src); - uint32_t mc = db_count(si); - if (sp) { - sp->synced_pos = tp; - sp->sync_state = 2; - } - DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, - "sync complete with %016llx tp=%u my_mc=%u", - (unsigned long long)src, tp, mc); - return; - } - if (peer_dh == 0 && sc == 0) { uint32_t mc = db_count(si); uint32_t batches = (mc + DB_SEND_DATA_MAX - 1) / DB_SEND_DATA_MAX; @@ -821,12 +808,62 @@ static void db_handle_init_resp(struct DB_SYNC_INSTANCE* si, if (my_dh == peer_dh) { uint32_t mc = db_count(si); struct SI_PEER* sp = si_peer_find(si, src); + uint32_t tail = tp + 1; + if (mc > tail) { + DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, + "dh match at tp=%u my_mc=%u — sending tail [%u,%u) to %016llx", + tp, mc, tail, mc, (unsigned long long)src); + uint32_t sent = tail; + while (sent < mc) { + uint32_t b = mc - sent; + if (b > DB_SEND_DATA_MAX) b = DB_SEND_DATA_MAX; + uint8_t sdbuf[8192]; + uint32_t off = 0; + sdbuf[off++] = DB_MSG_SEND_DATA; + memcpy(sdbuf + off, &sent, 4); off += 4; + uint16_t rc = 0; + uint16_t* rcp = (uint16_t*)(sdbuf + off); off += 2; + sqlite3_stmt* stmt; + si_prep(si, &stmt, + "SELECT id,timestamp,datahash,data,author_signature" + " FROM \"%s\" ORDER BY timestamp,datahash" + " LIMIT ? OFFSET ?"); + sqlite3_bind_int64(stmt, 1, (sqlite3_int64)b); + sqlite3_bind_int64(stmt, 2, (sqlite3_int64)sent); + while (sqlite3_step(stmt) == SQLITE_ROW) { + uint64_t rid = (uint64_t)sqlite3_column_int64(stmt, 0); + uint64_t rts = (uint64_t)sqlite3_column_int64(stmt, 1); + uint64_t rdh = (uint64_t)sqlite3_column_int64(stmt, 2); + const uint8_t* rd = (const uint8_t*)sqlite3_column_blob(stmt, 3); + uint32_t rdl = (uint32_t)sqlite3_column_bytes(stmt, 3); + if (!rd) rdl = 0; + const uint8_t* rsig = (const uint8_t*)sqlite3_column_blob(stmt, 4); + int rsl = sqlite3_column_bytes(stmt, 4); + if (!rsig) rsl = 0; + int rec_sz = 28 + rdl + 1 + rsl; + if (off + rec_sz > (int)sizeof(sdbuf)) break; + memcpy(sdbuf + off, &rid, 8); off += 8; + memcpy(sdbuf + off, &rts, 8); off += 8; + memcpy(sdbuf + off, &rdh, 8); off += 8; + memcpy(sdbuf + off, &rdl, 4); off += 4; + if (rdl > 0) { memcpy(sdbuf + off, rd, rdl); off += rdl; } + sdbuf[off++] = (uint8_t)rsl; + if (rsl > 0) { memcpy(sdbuf + off, rsig, rsl); off += rsl; } + rc++; + } + sqlite3_finalize(stmt); + *rcp = rc; + sent += rc; + db_sync_send(si, src, sdbuf, off); + if (rc == 0) break; + } + } if (sp) { sp->synced_pos = tp; sp->sync_state = 2; } DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, - "dh match with %016llx at tp=%u my_mc=%u", + "sync complete with %016llx tp=%u my_mc=%u", (unsigned long long)src, tp, mc); return; } diff --git a/tests/test_db_sync.c b/tests/test_db_sync.c index e60d3845..8ccdc81f 100644 --- a/tests/test_db_sync.c +++ b/tests/test_db_sync.c @@ -1,4 +1,4 @@ -// test_db_sync.c — всестороннее тестирование модуля db_sync +// test_db_sync.c — тестирование db_sync: dh_match tail-send и peer empty #include #include #include @@ -27,19 +27,16 @@ #include "../lib/debug_config.h" #include "../lib/mem.h" -#define TEST_TIMEOUT_TB 300000 // 30s total -#define PHASE_TIMEOUT_TB 150000 // 15s per phase +#define TEST_TIMEOUT_TB 600000 // 60s +#define PHASE_TIMEOUT_TB 250000 // 25s #define POLL_INTERVAL_MS 5 -#define NODE_ID_A 0xAAAAAAAAAAAAAAAAULL -#define NODE_ID_B 0xBBBBBBBBBBBBBBBBULL - static struct UTUN_INSTANCE* inst_a = NULL; static struct UTUN_INSTANCE* inst_b = NULL; static struct DB_SYNC_INSTANCE* si_a = NULL; static struct DB_SYNC_INSTANCE* si_b = NULL; static struct UASYNC* ua = NULL; -static int test_phase = 0; // 0=running, 1=success, 2=failure +static int test_phase = 0; static void* timeout_id = NULL; static char temp_dir[] = "/tmp/utun_dbsync_XXXXXX"; @@ -47,8 +44,7 @@ static char config_a[256], config_b[256]; static int port_a_srv, port_b_srv; static int write_file(const char* path, const char* fmt, ...) { - va_list ap; - FILE* f = fopen(path, "w"); + 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; @@ -66,52 +62,24 @@ static int create_temp_configs(void) { port_a_srv = base; port_b_srv = base + 1; snprintf(config_a, sizeof(config_a), "%s/a.conf", temp_dir); snprintf(config_b, sizeof(config_b), "%s/b.conf", temp_dir); - - // Create SQLite db directories (parent dir for sync file) 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); if (write_file(config_a, - "[global]\n" - "my_node_id=0xAAAAAAAAAAAAAAAA\n" - "tun_ip=10.200.0.1/24\n" - "tun_ifname=tun200\n" - "keepalive_adaptive=0\n" - "db_path=%s/db_a\n" - "db_sync_enabled=1\n" - "\n" - "[server: srv_a]\n" - "addr=127.0.0.1:%d\n" - "type=public\n" - "\n" - "[allowed_keys]\n" - "allow_all=1\n", - temp_dir, port_a_srv) != 0) return -1; + "[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) != 0) return -1; if (config_ensure_keys_and_node_id(config_a) != 0) return -1; - char* pub_a = get_pubkey(config_a); - if (!pub_a) return -1; + char* pub_a = get_pubkey(config_a); if (!pub_a) return -1; if (write_file(config_b, - "[global]\n" - "my_node_id=0xBBBBBBBBBBBBBBBB\n" - "tun_ip=10.200.0.2/24\n" - "tun_ifname=tun201\n" - "keepalive_adaptive=0\n" - "db_path=%s/db_b\n" - "db_sync_enabled=1\n" - "\n" - "[server: srv_b]\n" - "addr=127.0.0.1:%d\n" - "type=public\n" - "\n" - "[client: to_a]\n" - "keepalive=1\n" - "peer_public_key=%s\n" - "link=srv_b:127.0.0.1:%d\n" - "\n" - "[allowed_keys]\n" - "allow_all=1\n", + "[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) != 0) { free(pub_a); return -1; } free(pub_a); if (config_ensure_keys_and_node_id(config_b) != 0) return -1; @@ -143,103 +111,112 @@ static int wait_for(const char* desc, int (*cond)(void), int timeout_tb) { return 0; } -static void sleep_tb(int tb) { - uint64_t end = get_time_tb() + (uint64_t)tb; - while (get_time_tb() < end && test_phase == 0) uasync_poll(ua, POLL_INTERVAL_MS); -} - -// ---- Condition functions ---- static int cond_links_init(void) { if (!inst_a || !inst_b) return 0; - struct ll_entry* entry = inst_a->connections->head; - while (entry) { struct conn_queue_entry* ce = (struct conn_queue_entry*)entry->data; - struct ETCP_CONN* ca = ce->conn; - struct ETCP_LINK* l = ca->links; while (l) { if (l->initialized) return 1; l = l->next; } - entry = entry->next; } + struct ll_entry* e = inst_a->connections->head; + 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) return 1; l = l->next; } + e = e->next; } return 0; } - -static int cond_count_a(uint32_t expected) { - if (!si_a) return 0; - return db_sync_count(si_a) == expected; -} -static int cond_count_b(uint32_t expected) { - if (!si_b) return 0; - return db_sync_count(si_b) == expected; -} -static uint32_t count_a_target, count_b_target; -static int _cond_count_a(void) { return cond_count_a(count_a_target); } -static int _cond_count_b(void) { return cond_count_b(count_b_target); } +static uint32_t ca_target, cb_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 insert_many(struct DB_SYNC_INSTANCE* si, 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\",\"pad\":\"%s\"}", i, i, + snprintf(buf, sizeof(buf), "{\"idx\":%d,\"val\":\"data_%d\",\"pad\":\"%.50s\"}", i, i, "xxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxx"); - int ret = db_sync_insert(si, buf); - if (ret < 0) { fprintf(stderr, "insert_many failed at %d ret=%d\n", i, ret); return -1; } + if (db_sync_insert(si, buf) < 0) { fprintf(stderr,"insert_many fail %d\n", i); return -1; } } return 0; } -static int insert_many_batch(struct DB_SYNC_INSTANCE* si, int start, int count) { - char buf[256]; - for (int i = start; i < start + count && test_phase == 0; i++) { - snprintf(buf, sizeof(buf), "{\"n\":%d,\"text\":\"record_number_%d_abcdefghijklmnopqrstuvwxyz\"}", i, i); - db_sync_insert(si, buf); - } - return 0; +static void remove_si(struct DB_SYNC_INSTANCE** psi) { + if (*psi) { db_sync_instance_remove(*psi); *psi = NULL; } } // ---- Main test ---- int main(void) { printf("=== test_db_sync ===\n"); - debug_config_init(); - debug_set_level(DEBUG_LEVEL_ERROR); // quiet - + debug_config_init(); debug_set_level(DEBUG_LEVEL_ERROR); if (create_temp_configs() != 0) { fprintf(stderr, "config creation failed\n"); return 1; } utun_instance_set_tun_init_enabled(0); ua = uasync_create(); - if (!ua) { fprintf(stderr, "uasync_create failed\n"); cleanup_temp_configs(); return 1; } + 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) { fprintf(stderr, "instance create failed\n"); cleanup_temp_configs(); return 1; } - - if (utun_instance_init(inst_a) != 0 || utun_instance_init(inst_b) != 0) { - fprintf(stderr, "instance init failed\n"); cleanup_temp_configs(); return 1; + if (!inst_a || !inst_b || utun_instance_init(inst_a) != 0 || utun_instance_init(inst_b) != 0) { + fprintf(stderr, "instance create/init failed\n"); cleanup_temp_configs(); return 1; } + timeout_id = uasync_set_timeout(ua, TEST_TIMEOUT_TB, NULL, test_timeout, "global_timeout"); + + // =================================================================== + // Phase 1: late B — A=50, B=0. B создан ПОСЛЕ conn_up (A не инициирует). + // Только B инициирует: divergence → REFINE → SEND_DATA(1) → SYNC_DONE + // mismatch → A ретраит → dh_match(tp=0,sc=0) → должен отправить хвост [1..49]. + // =================================================================== + printf("Phase 1: late B — dh_match at tp=0 (sc=0) + tail-send 49 records\n"); + if (!wait_for("links up", cond_links_init, PHASE_TIMEOUT_TB)) { test_phase=2; goto done; } si_a = db_sync_instance_add(inst_a, "test", 1); - si_b = db_sync_instance_add(inst_b, "test", 1); - if (!si_a || !si_b) { fprintf(stderr, "db_sync_instance_add failed\n"); return 1; } + if (!si_a || insert_many(si_a, 0, 50) != 0) { test_phase=2; goto done; } + if (db_sync_count(si_a) != 50) { fprintf(stderr, "FAIL: A!=50\n"); test_phase=2; goto done; } + printf(" A has 50 records, creating B now (B empty, A conn_up already fired)\n"); - timeout_id = uasync_set_timeout(ua, TEST_TIMEOUT_TB, NULL, test_timeout, "test_timeout"); + si_b = db_sync_instance_add(inst_b, "test", 1); + if (!si_b || db_sync_count(si_b) != 0) { test_phase=2; goto done; } - // ===== Phase 1: базовый CRUD ===== - printf("Phase 1: basic CRUD...\n"); - if (db_sync_count(si_a) != 0) { fprintf(stderr, "FAIL: initial count not 0 (got %u)\n", db_sync_count(si_a)); test_phase=2; } + cb_target = 50; + if (!wait_for("B count=50 (dh_match tp=0 tail)", _cond_cb, PHASE_TIMEOUT_TB)) { test_phase=2; goto done; } + printf(" Phase 1 PASS: B=%u A=%u\n", db_sync_count(si_b), db_sync_count(si_a)); - if (db_sync_insert(si_a, "{\"key\":\"val1\"}") != 0) { fprintf(stderr,"FAIL: insert 1\n"); test_phase=2; } - if (db_sync_insert(si_a, "{\"key\":\"val2\"}") != 0) { fprintf(stderr,"FAIL: insert 2\n"); test_phase=2; } - if (db_sync_insert(si_a, "{\"key\":\"val3\"}") != 0) { fprintf(stderr,"FAIL: insert 3\n"); test_phase=2; } - if (db_sync_count(si_a) != 3) { fprintf(stderr,"FAIL: count not 3 (got %u)\n", db_sync_count(si_a)); test_phase=2; } - // Dedup by (timestamp,datahash): same content at different time = new record - if (db_sync_insert(si_a, "{\"key\":\"val4\"}") != 0) { fprintf(stderr,"FAIL: insert 4\n"); test_phase=2; } - if (db_sync_count(si_a) != 4) { fprintf(stderr,"FAIL: count not 4 (got %u)\n", db_sync_count(si_a)); test_phase=2; } - if (test_phase == 0) printf("Phase 1: PASS (count=4)\n"); + // =================================================================== + // Phase 2: recreate A after adding 30 more — A=80, B=50. + // A инициирует: dh_match(tp=49, sc>0, sparse hashes) — должен отправить хвост [50..79]. + // =================================================================== + printf("Phase 2: recreate A — dh_match at tp=49 (sc>0 sparse) + tail-send 30 records\n"); - // ===== Phase 2: initial sync A↔B ===== - printf("Phase 2: initial sync...\n"); - if (test_phase == 0) { count_b_target = 4; if (!wait_for("B count=4", _cond_count_b, PHASE_TIMEOUT_TB)) test_phase=2; } - if (test_phase == 0) printf("Phase 2: PASS (B synced %u records)\n", db_sync_count(si_b)); + if (insert_many(si_a, 50, 30) != 0) { test_phase=2; goto done; } + if (db_sync_count(si_a) != 80) { fprintf(stderr, "FAIL: A!=80\n"); test_phase=2; goto done; } - // ===== Cleanup ===== + remove_si(&si_a); + si_a = db_sync_instance_add(inst_a, "test", 1); + if (!si_a || db_sync_count(si_a) != 80) { test_phase=2; goto done; } + printf(" A recreated (80 records), B has 50 — A will initiate with tail\n"); + + cb_target = 80; + if (!wait_for("B count=80 (dh_match tp=49 sparse)", _cond_cb, PHASE_TIMEOUT_TB)) { test_phase=2; goto done; } + printf(" Phase 2 PASS: B=%u A=%u\n", db_sync_count(si_b), db_sync_count(si_a)); + + // =================================================================== + // Phase 3: fresh instances, A=30, B=0, both created simultaneously. + // A инициирует: peer empty (peer_dh=0, sc=0) — отправляет все 30. + // B тоже инициирует но расходится через divergence+reinit+dh_match(tp=0) — + // но A уже всё отправил, так что B и так получит через A's sync. + // =================================================================== + printf("Phase 3: fresh instances A=30 B=0, both initiate — peer empty\n"); + + remove_si(&si_a); + remove_si(&si_b); + si_a = db_sync_instance_add(inst_a, "test2", 2); + si_b = db_sync_instance_add(inst_b, "test2", 2); + if (!si_a || !si_b) { test_phase=2; goto done; } + + if (insert_many(si_a, 0, 30) != 0 || db_sync_count(si_a) != 30) { test_phase=2; goto done; } + + cb_target = 30; + if (!wait_for("B count=30 (peer empty)", _cond_cb, PHASE_TIMEOUT_TB)) { test_phase=2; goto done; } + printf(" Phase 3 PASS: B=%u A=%u\n", db_sync_count(si_b), db_sync_count(si_a)); + +done: + remove_si(&si_a); + remove_si(&si_b); if (timeout_id && ua) { uasync_cancel_timeout(ua, timeout_id); timeout_id = NULL; } - if (si_a) { db_sync_instance_remove(si_a); si_a = NULL; } - if (si_b) { db_sync_instance_remove(si_b); si_b = NULL; } if (inst_a) { inst_a->running = 0; utun_instance_destroy(inst_a); inst_a = NULL; } if (inst_b) { inst_b->running = 0; utun_instance_destroy(inst_b); inst_b = NULL; } if (ua) { uasync_destroy(ua, 0); ua = NULL; }