|
|
|
|
@ -112,19 +112,20 @@ static int _bucket_hash(void* ctx, const char* ns, uint8_t level, uint64_t prefi
|
|
|
|
|
return count; |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
static int _get_page(void* ctx, const char* ns, uint64_t prefix, int after_valid, uint64_t after, |
|
|
|
|
uint8_t* buf, size_t* len, uint64_t* next, int* more) { |
|
|
|
|
static int _get_page(void* ctx, const char* ns, uint64_t prefix, int after_valid, const merkle_key_t* after, |
|
|
|
|
uint8_t* buf, size_t* len, merkle_key_t* next, int* more) { |
|
|
|
|
(void)ns; |
|
|
|
|
struct ms_data* d = ((struct ms_test_ctx*)ctx)->data; |
|
|
|
|
uint16_t count = 0; size_t off = 2; |
|
|
|
|
*next = 0; *more = 0; |
|
|
|
|
memset(next, 0, sizeof(*next)); *more = 0; |
|
|
|
|
uint64_t after_id; memcpy(&after_id, after->bytes, 8); after_id = be64toh(after_id); |
|
|
|
|
for (int i = 0; i < d->count; i++) { |
|
|
|
|
if (merkle_sync_level_prefix(d->items[i].key, MT_MAX_LEVEL) != prefix |
|
|
|
|
|| (after_valid && d->items[i].key <= after)) continue; |
|
|
|
|
|| (after_valid && d->items[i].key <= after_id)) continue; |
|
|
|
|
if (off + 12 > *len) { if (!count) return -1; *more = 1; break; } |
|
|
|
|
memcpy(buf + off, &d->items[i].key, 8); off += 8; |
|
|
|
|
memcpy(buf + off, &d->items[i].val, 4); off += 4; |
|
|
|
|
*next = d->items[i].key; count++; |
|
|
|
|
uint64_t next_id = htobe64(d->items[i].key); memcpy(next->bytes, &next_id, 8); count++; |
|
|
|
|
} |
|
|
|
|
memcpy(buf, &count, 2); *len = off; return 0; |
|
|
|
|
} |
|
|
|
|
@ -153,7 +154,7 @@ static int _apply_items(void* ctx, const char* ns, uint64_t from_peer, const uin
|
|
|
|
|
_data_sort(d); |
|
|
|
|
/* Страница всегда относится к одному листу: пересчитываем его один раз. */ |
|
|
|
|
uint64_t key; memcpy(&key, data + 2, 8); |
|
|
|
|
if (merkle_sync_recompute_path(tc->inst, ns, key) < 0) return -1; |
|
|
|
|
if (merkle_sync_recompute_path(tc->inst, MT_DATA_MEMBERS, ns, key) < 0) return -1; |
|
|
|
|
} |
|
|
|
|
return 0; |
|
|
|
|
} |
|
|
|
|
@ -164,7 +165,7 @@ static const struct merkle_sync_data_ops g_test_ops = {
|
|
|
|
|
|
|
|
|
|
static const uint8_t* _test_hash(struct UTUN_INSTANCE* inst, const char* ns, uint8_t level, uint64_t prefix) { |
|
|
|
|
static uint8_t hash[MT_HASH_SIZE]; |
|
|
|
|
if (merkle_sync_read_hash(inst, ns, level, prefix, hash) < 0) abort(); |
|
|
|
|
if (merkle_sync_read_hash(inst, MT_DATA_MEMBERS, ns, level, prefix, hash) < 0) abort(); |
|
|
|
|
return hash; |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
@ -261,9 +262,10 @@ static void test_level_prefix_bucket_partition(void) {
|
|
|
|
|
static void _ensure_table(sqlite3* db) { if (merkle_tree_init(db) < 0) abort(); } |
|
|
|
|
|
|
|
|
|
static int _db_count_rows(sqlite3* db, const char* ns) { |
|
|
|
|
char tree_ns[25]; snprintf(tree_ns, sizeof(tree_ns), "0:%s", ns); |
|
|
|
|
sqlite3_stmt* st = NULL; |
|
|
|
|
sqlite3_prepare_v2(db, "SELECT COUNT(*) FROM merkle_tree_hash WHERE namespace=?", -1, &st, NULL); |
|
|
|
|
sqlite3_bind_text(st, 1, ns, -1, SQLITE_STATIC); |
|
|
|
|
sqlite3_bind_text(st, 1, tree_ns, -1, SQLITE_STATIC); |
|
|
|
|
int n = 0; |
|
|
|
|
if (sqlite3_step(st) == SQLITE_ROW) n = sqlite3_column_int(st, 0); |
|
|
|
|
sqlite3_finalize(st); |
|
|
|
|
@ -271,13 +273,14 @@ static int _db_count_rows(sqlite3* db, const char* ns) {
|
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
static int _db_compare_trees(sqlite3* da, sqlite3* db, const char* ns) { |
|
|
|
|
char tree_ns[25]; snprintf(tree_ns, sizeof(tree_ns), "0:%s", ns); |
|
|
|
|
sqlite3_stmt* sa = NULL, *sb = NULL; |
|
|
|
|
sqlite3_prepare_v2(da, "SELECT level,prefix64,hash FROM merkle_tree_hash WHERE namespace=? ORDER BY level,prefix64", |
|
|
|
|
-1, &sa, NULL); |
|
|
|
|
sqlite3_prepare_v2(db, "SELECT level,prefix64,hash FROM merkle_tree_hash WHERE namespace=? ORDER BY level,prefix64", |
|
|
|
|
-1, &sb, NULL); |
|
|
|
|
sqlite3_bind_text(sa, 1, ns, -1, SQLITE_STATIC); |
|
|
|
|
sqlite3_bind_text(sb, 1, ns, -1, SQLITE_STATIC); |
|
|
|
|
sqlite3_bind_text(sa, 1, tree_ns, -1, SQLITE_STATIC); |
|
|
|
|
sqlite3_bind_text(sb, 1, tree_ns, -1, SQLITE_STATIC); |
|
|
|
|
int mismatch = 0; |
|
|
|
|
while (1) { |
|
|
|
|
int ra = sqlite3_step(sa), rb = sqlite3_step(sb); |
|
|
|
|
@ -306,7 +309,7 @@ static void test_tree_empty_ns(void) {
|
|
|
|
|
inst.ua = uasync_create(); |
|
|
|
|
if (merkle_sync_init(&inst, 0x72, &g_test_ops, &ctx) != 0) abort(); |
|
|
|
|
|
|
|
|
|
merkle_sync_recompute_path(&inst, MS_NS_TEST, 0x1234ULL); |
|
|
|
|
merkle_sync_recompute_path(&inst, MT_DATA_MEMBERS, MT_DATA_MEMBERS, MT_DATA_MEMBERS, MT_DATA_MEMBERS, MT_DATA_MEMBERS, MT_DATA_MEMBERS, MT_DATA_MEMBERS, MT_DATA_MEMBERS, MT_DATA_MEMBERS, MS_NS_TEST, 0x1234ULL); |
|
|
|
|
if (_db_count_rows(db, MS_NS_TEST) == 0) PASS(); else FAIL("got %d rows", _db_count_rows(db, MS_NS_TEST)); |
|
|
|
|
|
|
|
|
|
merkle_sync_destroy(&inst); uasync_destroy(inst.ua, 0); sqlite3_close(db); |
|
|
|
|
@ -470,7 +473,7 @@ static void test_tree_consistency(void) {
|
|
|
|
|
sqlite3_prepare_v2(db, |
|
|
|
|
"SELECT level, prefix64 FROM merkle_tree_hash WHERE namespace=? ORDER BY level, prefix64", |
|
|
|
|
-1, &st, NULL); |
|
|
|
|
sqlite3_bind_text(st, 1, "ns", -1, SQLITE_STATIC); |
|
|
|
|
sqlite3_bind_text(st, 1, "0:ns", -1, SQLITE_STATIC); |
|
|
|
|
int ok = 1; |
|
|
|
|
while (sqlite3_step(st) == SQLITE_ROW && ok) { |
|
|
|
|
uint8_t lv = (uint8_t)sqlite3_column_int(st, 0); |
|
|
|
|
@ -681,7 +684,7 @@ static void _intg_ins_many(struct ms_data* d, struct UTUN_INSTANCE* inst, int ba
|
|
|
|
|
} |
|
|
|
|
_data_sort(d); |
|
|
|
|
for (int i = 0; i < d->count; i++) |
|
|
|
|
merkle_sync_recompute_path(inst, MS_NS_TEST, d->items[i].key); |
|
|
|
|
merkle_sync_recompute_path(inst, MT_DATA_MEMBERS, MS_NS_TEST, d->items[i].key); |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
static int _intg_compare_data(struct ms_data* a, struct ms_data* b) { |
|
|
|
|
@ -708,7 +711,7 @@ static void test_peer_empty(void) {
|
|
|
|
|
if (_intg_init_two() != 0) { FAIL("setup failed"); return; } |
|
|
|
|
_intg_ins_many(&data_b, i_b, 0, 5); |
|
|
|
|
done_sync = 0; |
|
|
|
|
merkle_sync_start(i_a, i_b->node_id, MS_NS_TEST, _on_sync_done, NULL); |
|
|
|
|
merkle_sync_start(i_a, MT_DATA_MEMBERS, MT_DATA_MEMBERS, MT_DATA_MEMBERS, MT_DATA_MEMBERS, MT_DATA_MEMBERS, MT_DATA_MEMBERS, MT_DATA_MEMBERS, MT_DATA_MEMBERS, i_b->node_id, MS_NS_TEST, _on_sync_done, NULL); |
|
|
|
|
if (!_wait_for("sync done", _cond_done, PHASE_TIMEOUT_TB)) |
|
|
|
|
{ FAIL("sync timeout"); _intg_cleanup(); return; } |
|
|
|
|
printf(" A data=%d B data=%d, A tree=%d B tree=%d\n", |
|
|
|
|
@ -778,8 +781,8 @@ static void test_randomized_two(void) {
|
|
|
|
|
_data_insert(&data_b, k, (uint32_t)rand()); |
|
|
|
|
} |
|
|
|
|
_data_sort(&data_a); _data_sort(&data_b); |
|
|
|
|
for (int i = 0; i < data_a.count; i++) merkle_sync_recompute_path(i_a, MS_NS_RND, data_a.items[i].key); |
|
|
|
|
for (int i = 0; i < data_b.count; i++) merkle_sync_recompute_path(i_b, MS_NS_RND, data_b.items[i].key); |
|
|
|
|
for (int i = 0; i < data_a.count; i++) merkle_sync_recompute_path(i_a, MT_DATA_MEMBERS, MT_DATA_MEMBERS, MS_NS_RND, data_a.items[i].key); |
|
|
|
|
for (int i = 0; i < data_b.count; i++) merkle_sync_recompute_path(i_b, MT_DATA_MEMBERS, MT_DATA_MEMBERS, MT_DATA_MEMBERS, MT_DATA_MEMBERS, MS_NS_RND, data_b.items[i].key); |
|
|
|
|
|
|
|
|
|
done_sync = 0; |
|
|
|
|
merkle_sync_start(i_a, i_b->node_id, MS_NS_RND, _on_sync_done, NULL); |
|
|
|
|
@ -877,15 +880,15 @@ static void test_randomized_three(void) {
|
|
|
|
|
for (int i = 0; i < nc; i++) { uint64_t k = ((uint64_t)rand()<<32)|(uint64_t)rand(); _data_insert(&data_c, k, (uint32_t)rand()); } |
|
|
|
|
_data_sort(&data_a); for (int i = 0; i < data_a.count; i++) merkle_sync_recompute_path(i_a, MS_NS_RND, data_a.items[i].key); |
|
|
|
|
_data_sort(&data_b); for (int i = 0; i < data_b.count; i++) merkle_sync_recompute_path(i_b, MS_NS_RND, data_b.items[i].key); |
|
|
|
|
_data_sort(&data_c); for (int i = 0; i < data_c.count; i++) merkle_sync_recompute_path(i_c, MS_NS_RND, data_c.items[i].key); |
|
|
|
|
_data_sort(&data_c); for (int i = 0; i < data_c.count; i++) merkle_sync_recompute_path(i_c, MT_DATA_MEMBERS, MS_NS_RND, data_c.items[i].key); |
|
|
|
|
|
|
|
|
|
done_sync = 0; i_phase = 0; |
|
|
|
|
/* B→A: B syncs with A, gets A's data, B: SYNCED */ |
|
|
|
|
merkle_sync_start(i_b, i_a->node_id, MS_NS_RND, _on_sync_done, NULL); |
|
|
|
|
merkle_sync_start(i_b, MT_DATA_MEMBERS, i_a->node_id, MS_NS_RND, _on_sync_done, NULL); |
|
|
|
|
if (!_wait_for("sync B", _cond_done, PHASE_TIMEOUT_TB)) { FAIL("timeout B iter %d", iter); _intg_cleanup(); break; } |
|
|
|
|
done_sync = 0; i_phase = 0; |
|
|
|
|
/* C→A: C syncs with A, gets A∪B, A broadcasts C's new items to B (relay) */ |
|
|
|
|
merkle_sync_start(i_c, i_a->node_id, MS_NS_RND, _on_sync_done, NULL); |
|
|
|
|
merkle_sync_start(i_c, MT_DATA_MEMBERS, i_a->node_id, MS_NS_RND, _on_sync_done, NULL); |
|
|
|
|
if (!_wait_for("sync C", _cond_done, PHASE_TIMEOUT_TB)) { FAIL("timeout C iter %d", iter); _intg_cleanup(); break; } |
|
|
|
|
/* Автоматический следующий раунд A↔B завершается по условию, не по числу poll(). */ |
|
|
|
|
if (!_wait_for("automatic convergence of all peers", _cond_all_synced, PHASE3_TIMEOUT_TB)) { |
|
|
|
|
@ -909,7 +912,7 @@ static void test_cancel(void) {
|
|
|
|
|
_intg_ins_many(&data_b, i_b, 0, 5); |
|
|
|
|
done_sync = 0; |
|
|
|
|
merkle_sync_start(i_a, i_b->node_id, MS_NS_TEST, _on_sync_done, NULL); |
|
|
|
|
merkle_sync_cancel(i_a, i_b->node_id, MS_NS_TEST); |
|
|
|
|
merkle_sync_cancel(i_a, MT_DATA_MEMBERS, i_b->node_id, MS_NS_TEST); |
|
|
|
|
for (int i = 0; i < 100 && i_phase == 0; i++) uasync_poll(i_ua, 10); |
|
|
|
|
if (!done_sync) PASS(); else FAIL("done_cb was called"); |
|
|
|
|
_intg_cleanup(); |
|
|
|
|
@ -1014,8 +1017,8 @@ static void _str_spam_cb(void* arg) {
|
|
|
|
|
/* spammer: bump own member, sync to hub */ |
|
|
|
|
uint64_t key = ((uint64_t)idx) << 60; |
|
|
|
|
_data_insert(&gs->data[idx], key, (uint32_t)(gs->data[idx].count + 1)); |
|
|
|
|
merkle_sync_recompute_path(gs->inst[idx], MS_NS_STRESS, key); |
|
|
|
|
merkle_sync_start(gs->inst[idx], gs->inst[STRESS_HUB]->node_id, MS_NS_STRESS, NULL, NULL); |
|
|
|
|
merkle_sync_recompute_path(gs->inst[idx], MT_DATA_MEMBERS, MS_NS_STRESS, key); |
|
|
|
|
merkle_sync_start(gs->inst[idx], MT_DATA_MEMBERS, MT_DATA_MEMBERS, gs->inst[STRESS_HUB]->node_id, MS_NS_STRESS, NULL, NULL); |
|
|
|
|
|
|
|
|
|
_str_spam_schedule(gs, idx); |
|
|
|
|
} |
|
|
|
|
@ -1142,12 +1145,12 @@ static int _str_init(void) {
|
|
|
|
|
for (int i = 1; i < STRESS_N; i++) { |
|
|
|
|
uint64_t key = ((uint64_t)i) << 60; |
|
|
|
|
_data_insert(&s->data[i], key, (uint32_t)i); |
|
|
|
|
merkle_sync_recompute_path(s->inst[i], MS_NS_STRESS, key); |
|
|
|
|
merkle_sync_recompute_path(s->inst[i], MT_DATA_MEMBERS, MS_NS_STRESS, key); |
|
|
|
|
} |
|
|
|
|
/* initial sync: each spoke → hub */ |
|
|
|
|
s->sync_pending = STRESS_N - 1; |
|
|
|
|
for (int i = 1; i < STRESS_N; i++) |
|
|
|
|
merkle_sync_start(s->inst[i], s->inst[STRESS_HUB]->node_id, MS_NS_STRESS, _str_sync_done, s); |
|
|
|
|
merkle_sync_start(s->inst[i], MT_DATA_MEMBERS, MT_DATA_MEMBERS, s->inst[STRESS_HUB]->node_id, MS_NS_STRESS, _str_sync_done, s); |
|
|
|
|
_str_wait_for("initial sync", _str_cond_sync_done, STRESS_SYNC_TB); |
|
|
|
|
|
|
|
|
|
/* observers: sync with hub too */ |
|
|
|
|
|