You can not select more than 25 topics
Topics must start with a letter or number, can include dashes ('-') and can be up to 35 characters long.
235 lines
10 KiB
235 lines
10 KiB
// test_db_sync.c — тестирование db_sync: dh_match tail-send и peer empty |
|
#include <stdio.h> |
|
#include <stdlib.h> |
|
#include <string.h> |
|
#include <stdarg.h> |
|
#include <time.h> |
|
#include <sys/stat.h> |
|
#include "../lib/platform_compat.h" |
|
#include "test_utils.h" |
|
#ifdef _WIN32 |
|
#include <windows.h> |
|
#include <direct.h> |
|
#else |
|
#include <unistd.h> |
|
#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 600000 // 60s |
|
#define PHASE_TIMEOUT_TB 250000 // 25s |
|
#define POLL_INTERVAL_MS 5 |
|
|
|
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; |
|
static void* timeout_id = NULL; |
|
|
|
static char temp_dir[] = "/tmp/utun_dbsync_XXXXXX"; |
|
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"); |
|
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; |
|
snprintf(config_a, sizeof(config_a), "%s/a.conf", temp_dir); |
|
snprintf(config_b, sizeof(config_b), "%s/b.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); |
|
|
|
if (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) != 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; |
|
|
|
if (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) != 0) { free(pub_a); return -1; } |
|
free(pub_a); |
|
if (config_ensure_keys_and_node_id(config_b) != 0) return -1; |
|
return 0; |
|
} |
|
|
|
static void cleanup_temp_configs(void) { |
|
unlink(config_a); unlink(config_b); |
|
char pa[320]; |
|
snprintf(pa, sizeof(pa), "%s/db_a/sync", temp_dir); unlink(pa); |
|
snprintf(pa, sizeof(pa), "%s/db_a/sync-wal", temp_dir); unlink(pa); |
|
snprintf(pa, sizeof(pa), "%s/db_a/sync-shm", temp_dir); unlink(pa); |
|
snprintf(pa, sizeof(pa), "%s/db_b/sync", temp_dir); unlink(pa); |
|
snprintf(pa, sizeof(pa), "%s/db_b/sync-wal", temp_dir); unlink(pa); |
|
snprintf(pa, sizeof(pa), "%s/db_b/sync-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); |
|
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; |
|
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 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, 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\",\"pad\":\"%.50s\"}", i, i, |
|
"xxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxx"); |
|
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 %d\n", i); return -1; } |
|
if (db_sync_insert_signed(si, buf, jl, sig, 64, ts) < 0) { fprintf(stderr,"insert_many fail %d\n", i); return -1; } |
|
} |
|
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); |
|
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) { 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, "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); |
|
if (!si_a || insert_many(si_a, inst_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"); |
|
|
|
si_b = db_sync_instance_add(inst_b, "test", 1); |
|
if (!si_b || db_sync_count(si_b) != 0) { test_phase=2; goto done; } |
|
|
|
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)); |
|
|
|
// =================================================================== |
|
// 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"); |
|
|
|
if (insert_many(si_a, inst_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; } |
|
|
|
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, inst_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 (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; } |
|
cleanup_temp_configs(); |
|
|
|
if (test_phase == 0) test_phase = 1; |
|
printf("=== %s ===\n", test_phase == 1 ? "PASS" : "FAIL"); |
|
return test_phase == 1 ? 0 : 1; |
|
}
|
|
|