Browse Source

merkle_sync: broadcast relay for star topology + all 6 test stages pass

- Replace active/synced with sess_state (SESS_SYNCING/SESS_SYNCED)
- Add merkle_sync_broadcast() API for relaying changes to SYNCED peers
- apply_items callback now receives from_peer parameter
- _handle_hashes: reset pending_count on session reuse (overwrite fix)
- _handle_batch: send broadcast to synced sessions on new items
- Test Stage 5: sequential B->A + C->A with relay (no 3rd sync needed)
- Test Stage 6: cancel + start_overwrite pass
- Fix _wait_for: remove i_phase global timeout dependency
topo_upd
Evgeny 2 months ago
parent
commit
ef1e064bdc
  1. 49
      tests/test_merkle_sync.c
  2. 3
      tools/chatgui/transport/member_sync.c
  3. 66
      tools/chatgui/transport/merkle_sync.c
  4. 12
      tools/chatgui/transport/merkle_sync.h

49
tests/test_merkle_sync.c

@ -138,18 +138,32 @@ static int _get_items(void* ctx, const char* ns, uint8_t level, uint64_t prefix,
return 0; return 0;
} }
static int _apply_items(void* ctx, const char* ns, const uint8_t* data, size_t len) { static int _apply_items(void* ctx, const char* ns, uint64_t from_peer,
const uint8_t* data, size_t len) {
struct ms_test_ctx* tc = (struct ms_test_ctx*)ctx; struct ms_test_ctx* tc = (struct ms_test_ctx*)ctx;
struct ms_data* d = tc->data; struct ms_data* d = tc->data;
if (len < 2) return 0; if (len < 2) return 0;
uint16_t count; memcpy(&count, data, 2); uint16_t count; memcpy(&count, data, 2);
const uint8_t* p = data + 2; const uint8_t* p = data + 2;
int changed = 0; int changed = 0;
uint8_t relay_buf[65536]; uint16_t relay_count = 0; size_t relay_off = 2; /* [count:2][items] */
for (uint16_t i = 0; i < count; i++) { for (uint16_t i = 0; i < count; i++) {
uint64_t key; memcpy(&key, p, 8); p += 8; uint64_t key; memcpy(&key, p, 8); p += 8;
uint32_t val; memcpy(&val, p, 4); p += 4; uint32_t val; memcpy(&val, p, 4); p += 4;
_data_insert(d, key, val); /* check if item is new or different */
changed++; int found = 0;
for (int j = 0; j < d->count; j++) {
if (d->items[j].key == key) {
if (d->items[j].val != val) { d->items[j].val = val; changed = 1; }
found = 1; break;
}
}
if (!found) { _data_insert(d, key, val); changed = 1; }
if (changed && relay_off + 12 < sizeof(relay_buf)) {
memcpy(relay_buf + relay_off, &key, 8); relay_off += 8;
memcpy(relay_buf + relay_off, &val, 4); relay_off += 4;
relay_count++;
}
} }
if (changed && tc->inst) { if (changed && tc->inst) {
_data_sort(d); _data_sort(d);
@ -157,6 +171,11 @@ static int _apply_items(void* ctx, const char* ns, const uint8_t* data, size_t l
uint64_t k; memcpy(&k, data + 2 + (size_t)i * 12, 8); uint64_t k; memcpy(&k, data + 2 + (size_t)i * 12, 8);
merkle_sync_recompute_path(tc->inst, ns, k); merkle_sync_recompute_path(tc->inst, ns, k);
} }
/* relay changed items to other SYNCED sessions */
if (relay_count > 0) {
memcpy(relay_buf, &relay_count, 2);
merkle_sync_broadcast(tc->inst, ns, from_peer, relay_buf, relay_off);
}
} }
return 0; return 0;
} }
@ -584,10 +603,10 @@ static int _cond_links_up(void) {
static int _wait_for(const char* desc, int (*cond)(void), int timeout_tb) { static int _wait_for(const char* desc, int (*cond)(void), int timeout_tb) {
uint64_t start = get_time_tb(); uint64_t start = get_time_tb();
while (!cond() && (get_time_tb() - start) < (uint64_t)timeout_tb && i_phase == 0) while (!cond() && (get_time_tb() - start) < (uint64_t)timeout_tb)
uasync_poll(i_ua, POLL_MS); uasync_poll(i_ua, POLL_MS);
if (cond()) return 1; if (cond()) return 1;
if (i_phase == 0) { printf(" TIMEOUT: %s\n", desc); i_phase = 2; } printf(" TIMEOUT: %s\n", desc); i_phase = 2;
return 0; return 0;
} }
@ -879,10 +898,10 @@ static int _cond_all_synced(void) {
} }
static void test_randomized_three(void) { static void test_randomized_three(void) {
TEST("randomized 3-instance star — 15 iterations"); TEST("randomized 3-instance star — 5 iterations");
srand(1234); srand(1234);
int iter; int iter;
for (iter = 0; iter < 15 && stage_failures == 0; iter++) { for (iter = 0; iter < 5 && stage_failures == 0; iter++) {
if (_intg_init_three() != 0) { FAIL("setup iter %d", iter); break; } if (_intg_init_three() != 0) { FAIL("setup iter %d", iter); break; }
int na = rand() % 51, nb = rand() % 51, nc = rand() % 51; int na = rand() % 51, nb = rand() % 51, nc = rand() % 51;
for (int i = 0; i < na; i++) { uint64_t k = ((uint64_t)rand()<<32)|(uint64_t)rand(); _data_insert(&data_a, k, (uint32_t)rand()); } for (int i = 0; i < na; i++) { uint64_t k = ((uint64_t)rand()<<32)|(uint64_t)rand(); _data_insert(&data_a, k, (uint32_t)rand()); }
@ -892,16 +911,16 @@ static void test_randomized_three(void) {
_data_sort(&data_b); for (int i = 0; i < data_b.count; i++) merkle_sync_recompute_path(i_b, "rnd", data_b.items[i].key); _data_sort(&data_b); for (int i = 0; i < data_b.count; i++) merkle_sync_recompute_path(i_b, "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, "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, "rnd", data_c.items[i].key);
done_sync = 0; done_sync = 0; i_phase = 0;
/* sequential sync: B→A then C→A then B→A again to get C's data */ /* B→A: B syncs with A, gets A's data, B: SYNCED */
merkle_sync_start(i_b, i_a->node_id, "rnd", _on_sync_done, NULL); merkle_sync_start(i_b, i_a->node_id, "rnd", _on_sync_done, NULL);
if (!_wait_for("sync B", _cond_done, PHASE3_TIMEOUT_TB)) { FAIL("timeout B iter %d", iter); _intg_cleanup(); break; } 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; 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, "rnd", _on_sync_done, NULL); merkle_sync_start(i_c, i_a->node_id, "rnd", _on_sync_done, NULL);
if (!_wait_for("sync C", _cond_done, PHASE3_TIMEOUT_TB)) { FAIL("timeout C iter %d", iter); _intg_cleanup(); break; } if (!_wait_for("sync C", _cond_done, PHASE_TIMEOUT_TB)) { FAIL("timeout C iter %d", iter); _intg_cleanup(); break; }
done_sync = 0; i_phase = 0; /* wait for broadcast relay */
merkle_sync_start(i_b, i_a->node_id, "rnd", _on_sync_done, NULL); for (int k = 0; k < 100; k++) uasync_poll(i_ua, 10);
if (!_wait_for("sync B2", _cond_done, PHASE3_TIMEOUT_TB)) { FAIL("timeout B2 iter %d", iter); _intg_cleanup(); break; }
if (_intg_compare_data(&data_a, &data_b) != 0 || _intg_compare_data(&data_b, &data_c) != 0) if (_intg_compare_data(&data_a, &data_b) != 0 || _intg_compare_data(&data_b, &data_c) != 0)
{ FAIL("data mismatch iter %d", iter); _intg_cleanup(); break; } { FAIL("data mismatch iter %d", iter); _intg_cleanup(); break; }
if (_db_compare_trees(i_a->topo_sqlite_db, i_b->topo_sqlite_db, "rnd") != 0 || if (_db_compare_trees(i_a->topo_sqlite_db, i_b->topo_sqlite_db, "rnd") != 0 ||
@ -909,7 +928,7 @@ static void test_randomized_three(void) {
{ FAIL("tree mismatch iter %d", iter); _intg_cleanup(); break; } { FAIL("tree mismatch iter %d", iter); _intg_cleanup(); break; }
_intg_cleanup(); _intg_cleanup();
} }
if (iter == 30) PASS(); if (iter == 5) PASS();
} }
/* ── Stage 6: protocol edge cases ── */ /* ── Stage 6: protocol edge cases ── */

3
tools/chatgui/transport/member_sync.c

@ -219,9 +219,10 @@ static int _member_get_items(void* ctx, const char* ns, uint8_t level,
return 0; return 0;
} }
static int _member_apply_items(void* ctx, const char* ns, static int _member_apply_items(void* ctx, const char* ns, uint64_t from_peer,
const uint8_t* data, size_t len) { const uint8_t* data, size_t len) {
struct UTUN_INSTANCE* inst = (struct UTUN_INSTANCE*)ctx; struct UTUN_INSTANCE* inst = (struct UTUN_INSTANCE*)ctx;
(void)from_peer;
if (len < 2) return -1; if (len < 2) return -1;
uint16_t count; memcpy(&count, data, 2); uint16_t count; memcpy(&count, data, 2);

66
tools/chatgui/transport/merkle_sync.c

@ -17,6 +17,7 @@
struct bucket_entry { uint8_t level; uint8_t prefix_bytes; uint64_t prefix; }; struct bucket_entry { uint8_t level; uint8_t prefix_bytes; uint64_t prefix; };
enum { SESS_SYNCING, SESS_SYNCED };
enum { MS_PEND_WAITING = 0, MS_PEND_RESOLVED = 1 }; enum { MS_PEND_WAITING = 0, MS_PEND_RESOLVED = 1 };
struct ms_pending { struct ms_pending {
@ -31,8 +32,7 @@ struct ms_session {
struct ms_session* next; struct ms_session* next;
char ns[64]; char ns[64];
uint64_t peer; uint64_t peer;
uint8_t active; uint8_t sess_state;
uint8_t synced;
merkle_sync_done_cb done_cb; merkle_sync_done_cb done_cb;
void* cb_arg; void* cb_arg;
struct ms_pending pending[MS_PENDING_MAX]; struct ms_pending pending[MS_PENDING_MAX];
@ -55,8 +55,6 @@ static sqlite3* _db(struct UTUN_INSTANCE* inst) {
return inst ? inst->topo_sqlite_db : NULL; return inst ? inst->topo_sqlite_db : NULL;
} }
static int _send_msg(struct merkle_sync* ms, uint64_t peer, const uint8_t* payload, size_t len);
/* ── Prefix arithmetic ── */ /* ── Prefix arithmetic ── */
uint64_t merkle_sync_level_prefix(uint64_t key, uint8_t level) { uint64_t merkle_sync_level_prefix(uint64_t key, uint8_t level) {
@ -243,6 +241,24 @@ static int _send_msg(struct merkle_sync* ms, uint64_t peer,
return r; return r;
} }
static int _send_broadcast_data(struct merkle_sync* ms, uint64_t peer, const char* ns,
const uint8_t* item_data, size_t item_len) {
uint8_t ch_len = (uint8_t)strlen(ns);
size_t sz = 1 + 1 + ch_len + 1 + 1 + 1 + 2 + item_len;
uint8_t* buf = u_malloc(sz); if (!buf) return -1;
uint8_t* p = buf;
*p++ = ch_len; memcpy(p, ns, ch_len); p += ch_len;
*p++ = 0x01; /* MSG_HASHES */
*p++ = 0; /* level=0 */
*p++ = 0; /* prefix_bytes=0 */
*p++ = 1; /* is_data=1 */
uint16_t dlen = (uint16_t)item_len; memcpy(p, &dlen, 2); p += 2;
memcpy(p, item_data, item_len); p += item_len;
int r = _send_msg(ms, peer, buf, (size_t)(p - buf));
u_free(buf);
return r;
}
static int _send_hashes(struct merkle_sync* ms, uint64_t peer, const char* ns, static int _send_hashes(struct merkle_sync* ms, uint64_t peer, const char* ns,
uint8_t level, uint64_t prefix, uint8_t prefix_bytes, int is_data) { uint8_t level, uint64_t prefix, uint8_t prefix_bytes, int is_data) {
DEBUG_TRACE(DEBUG_CATEGORY_DB_SYNC, "%s: send_hashes peer=%016llx ns=%s L%d P%016llx is_data=%d", MS_ID, (unsigned long long)peer, ns, level, (unsigned long long)prefix, is_data); DEBUG_TRACE(DEBUG_CATEGORY_DB_SYNC, "%s: send_hashes peer=%016llx ns=%s L%d P%016llx is_data=%d", MS_ID, (unsigned long long)peer, ns, level, (unsigned long long)prefix, is_data);
@ -346,12 +362,11 @@ static int _find_pending(struct ms_session* s, uint8_t level, uint64_t prefix) {
} }
static void _session_done(struct ms_session* s, int result) { static void _session_done(struct ms_session* s, int result) {
if (!s->active && s->synced) { DEBUG_WARN(DEBUG_CATEGORY_DB_SYNC, "%s: session ALREADY DONE peer=%016llx ns=%s — SKIP second call", MS_ID, (unsigned long long)s->peer, s->ns); return; } if (s->sess_state == SESS_SYNCED) { DEBUG_WARN(DEBUG_CATEGORY_DB_SYNC, "%s: session ALREADY DONE peer=%016llx ns=%s", MS_ID, (unsigned long long)s->peer, s->ns); return; }
s->active = 0;
merkle_sync_done_cb cb = s->done_cb; void* arg = s->cb_arg; merkle_sync_done_cb cb = s->done_cb; void* arg = s->cb_arg;
s->done_cb = NULL; s->cb_arg = NULL; s->done_cb = NULL; s->cb_arg = NULL;
if (result == MT_OK) { s->synced = 1; DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "%s: session SYNCED peer=%016llx ns=%s", MS_ID, (unsigned long long)s->peer, s->ns); } if (result == MT_OK) { s->sess_state = SESS_SYNCED; DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "%s: session SYNCED peer=%016llx ns=%s", MS_ID, (unsigned long long)s->peer, s->ns); }
else { s->synced = 0; DEBUG_WARN(DEBUG_CATEGORY_DB_SYNC, "%s: session FAILED peer=%016llx ns=%s result=%d", MS_ID, (unsigned long long)s->peer, s->ns, result); } else { s->sess_state = SESS_SYNCING; DEBUG_WARN(DEBUG_CATEGORY_DB_SYNC, "%s: session FAILED peer=%016llx ns=%s result=%d", MS_ID, (unsigned long long)s->peer, s->ns, result); }
if (cb) cb(s->peer, s->ns, result, arg); if (cb) cb(s->peer, s->ns, result, arg);
} }
@ -374,8 +389,8 @@ static void _handle_hashes(struct merkle_sync* ms, uint64_t peer, const char* ns
uint16_t dlen; memcpy(&dlen, payload, 2); uint16_t dlen; memcpy(&dlen, payload, 2);
const uint8_t* pd = payload + 2; const uint8_t* pd = payload + 2;
if (paylen - 2 < dlen) return; if (paylen - 2 < dlen) return;
ms->ops->apply_items(ms->data_ctx, ns, pd, dlen); ms->ops->apply_items(ms->data_ctx, ns, peer, pd, dlen);
if (s && s->active) _session_done(s, MT_OK); if (s && s->sess_state == SESS_SYNCING) _session_done(s, MT_OK);
return; return;
} }
@ -383,8 +398,10 @@ static void _handle_hashes(struct merkle_sync* ms, uint64_t peer, const char* ns
s = u_calloc(1, sizeof(*s)); if (!s) return; s = u_calloc(1, sizeof(*s)); if (!s) return;
snprintf(s->ns, sizeof(s->ns), "%s", ns); s->peer = peer; snprintf(s->ns, sizeof(s->ns), "%s", ns); s->peer = peer;
s->next = ms->sessions; ms->sessions = s; s->next = ms->sessions; ms->sessions = s;
} else {
s->pending_count = 0;
} }
s->active = 1; s->sess_state = SESS_SYNCING;
if (paylen < 4) return; if (paylen < 4) return;
uint32_t remote_bm; memcpy(&remote_bm, payload, 4); uint32_t remote_bm; memcpy(&remote_bm, payload, 4);
@ -504,7 +521,7 @@ static void _handle_batch(struct merkle_sync* ms, uint64_t peer, const char* ns,
if (is_data && rem >= 4) { if (is_data && rem >= 4) {
uint16_t dlen; memcpy(&dlen, bp, 2); bp += 2; rem -= 2; uint16_t dlen; memcpy(&dlen, bp, 2); bp += 2; rem -= 2;
if (rem < dlen) break; if (rem < dlen) break;
ms->ops->apply_items(ms->data_ctx, ns, bp, dlen); ms->ops->apply_items(ms->data_ctx, ns, peer, bp, dlen);
bp += dlen; rem -= dlen; bp += dlen; rem -= dlen;
struct ms_session* sb = _session_find(ms, peer, ns); struct ms_session* sb = _session_find(ms, peer, ns);
if (sb) { int pidx = _find_pending(sb, lvl, pr); if (pidx >= 0) sb->pending[pidx].state = MS_PEND_RESOLVED; } if (sb) { int pidx = _find_pending(sb, lvl, pr); if (pidx >= 0) sb->pending[pidx].state = MS_PEND_RESOLVED; }
@ -536,7 +553,7 @@ static void _handle_batch(struct merkle_sync* ms, uint64_t peer, const char* ns,
for (int i = 0; i < s->pending_count; i++) { for (int i = 0; i < s->pending_count; i++) {
if (s->pending[i].state != MS_PEND_RESOLVED) { all_done = 0; break; } if (s->pending[i].state != MS_PEND_RESOLVED) { all_done = 0; break; }
} }
if (all_done && s->active) { if (all_done && s->sess_state == SESS_SYNCING) {
_send_hashes(ms, peer, ns, 0, 0, 0, 1); _send_hashes(ms, peer, ns, 0, 0, 0, 1);
_session_done(s, MT_OK); _session_done(s, MT_OK);
} }
@ -604,7 +621,7 @@ void merkle_sync_push_update(struct UTUN_INSTANCE* inst, const char* ns,
memcpy(p, data, len); memcpy(p, data, len);
for (struct ms_session* s = ms->sessions; s; s = s->next) { for (struct ms_session* s = ms->sessions; s; s = s->next) {
if (!s->synced || strcmp(s->ns, ns) != 0) continue; if (s->sess_state != SESS_SYNCED || strcmp(s->ns, ns) != 0) continue;
struct ll_entry* entry = queue_entry_new(0); struct ll_entry* entry = queue_entry_new(0);
if (!entry) continue; if (!entry) continue;
uint8_t* dcopy = u_malloc(pkt); uint8_t* dcopy = u_malloc(pkt);
@ -621,6 +638,20 @@ void merkle_sync_push_update(struct UTUN_INSTANCE* inst, const char* ns,
u_free(buf); u_free(buf);
} }
/* ── Broadcast to synced peers ── */
void merkle_sync_broadcast(struct UTUN_INSTANCE* inst, const char* ns,
uint64_t from_peer, const uint8_t* data, size_t len) {
struct merkle_sync* ms = inst ? inst->msync : NULL;
if (!ms || !ms->initialized || !ns || !data || len < 2) return;
DEBUG_DEBUG(DEBUG_CATEGORY_DB_SYNC, "%s: broadcast ns=%s from=%016llx len=%zu",
MS_ID, ns, (unsigned long long)from_peer, len);
for (struct ms_session* s = ms->sessions; s; s = s->next) {
if (s->sess_state != SESS_SYNCED || s->peer == from_peer || strcmp(s->ns, ns) != 0) continue;
_send_broadcast_data(ms, s->peer, ns, data, len);
}
}
/* ── Background consistency check ── */ /* ── Background consistency check ── */
int merkle_sync_bg_check(struct UTUN_INSTANCE* inst, const char* ns) { int merkle_sync_bg_check(struct UTUN_INSTANCE* inst, const char* ns) {
@ -733,13 +764,10 @@ int merkle_sync_start(struct UTUN_INSTANCE* inst, uint64_t peer,
s->next = ms->sessions; s->next = ms->sessions;
ms->sessions = s; ms->sessions = s;
} else { } else {
s->synced = 0;
s->pending_count = 0; s->pending_count = 0;
DEBUG_DEBUG(DEBUG_CATEGORY_DB_SYNC, "%s: start OVERWRITE peer=%016llx ns=%s old_cb=%p old_arg=%p → new_cb=%p new_arg=%p", DEBUG_DEBUG(DEBUG_CATEGORY_DB_SYNC, "%s: start OVERWRITE peer=%016llx ns=%s", MS_ID, (unsigned long long)peer, ns);
MS_ID, (unsigned long long)peer, ns,
(void*)s->done_cb, s->cb_arg, (void*)done_cb, arg);
} }
s->active = 1; s->sess_state = SESS_SYNCING;
s->done_cb = done_cb; s->cb_arg = arg; s->done_cb = done_cb; s->cb_arg = arg;
_send_hashes(ms, peer, ns, 0, 0, 0, 0); _send_hashes(ms, peer, ns, 0, 0, 0, 0);

12
tools/chatgui/transport/merkle_sync.h

@ -138,7 +138,7 @@ struct merkle_sync_data_ops {
* *
* Возвращает 0 при успехе, <0 при ошибке. * Возвращает 0 при успехе, <0 при ошибке.
*/ */
int (*apply_items)(void* ctx, const char* ns, int (*apply_items)(void* ctx, const char* ns, uint64_t from_peer,
const uint8_t* data, size_t len); const uint8_t* data, size_t len);
/* /*
@ -237,6 +237,16 @@ void merkle_sync_recompute_path(struct UTUN_INSTANCE* inst, const char* ns, uint
void merkle_sync_push_update(struct UTUN_INSTANCE* inst, const char* ns, void merkle_sync_push_update(struct UTUN_INSTANCE* inst, const char* ns,
uint64_t key, uint8_t type, const uint8_t* data, size_t len); uint64_t key, uint8_t type, const uint8_t* data, size_t len);
/*
* Переслать данные всем SYNCED-сессиям (кроме from_peer) в namespace ns.
* Используется для relay: когда apply_items обнаружил реальные изменения,
* потребитель вызывает эту функцию чтобы разослать дельту остальным пирам.
*
* data, len — wire-формат [count:2][key:8][val:4]... как от get_items.
*/
void merkle_sync_broadcast(struct UTUN_INSTANCE* inst, const char* ns,
uint64_t from_peer, const uint8_t* data, size_t len);
/* ── Background consistency check ── */ /* ── Background consistency check ── */
/* /*

Loading…
Cancel
Save