From ef1e064bdc5aa48a59cfd58b90e692bab53a45b5 Mon Sep 17 00:00:00 2001 From: Evgeny Date: Sun, 26 Jul 2026 18:20:10 +0300 Subject: [PATCH] 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 --- tests/test_merkle_sync.c | 49 ++++++++++++++------ tools/chatgui/transport/member_sync.c | 5 +- tools/chatgui/transport/merkle_sync.c | 66 +++++++++++++++++++-------- tools/chatgui/transport/merkle_sync.h | 12 ++++- 4 files changed, 95 insertions(+), 37 deletions(-) diff --git a/tests/test_merkle_sync.c b/tests/test_merkle_sync.c index f353ceb9..31454174 100644 --- a/tests/test_merkle_sync.c +++ b/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; } -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_data* d = tc->data; if (len < 2) return 0; uint16_t count; memcpy(&count, data, 2); const uint8_t* p = data + 2; 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++) { uint64_t key; memcpy(&key, p, 8); p += 8; uint32_t val; memcpy(&val, p, 4); p += 4; - _data_insert(d, key, val); - changed++; + /* check if item is new or different */ + 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) { _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); 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; } @@ -584,10 +603,10 @@ static int _cond_links_up(void) { 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 && i_phase == 0) + while (!cond() && (get_time_tb() - start) < (uint64_t)timeout_tb) uasync_poll(i_ua, POLL_MS); 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; } @@ -879,10 +898,10 @@ static int _cond_all_synced(void) { } static void test_randomized_three(void) { - TEST("randomized 3-instance star — 15 iterations"); + TEST("randomized 3-instance star — 5 iterations"); srand(1234); 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; } 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()); } @@ -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_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; - /* sequential sync: B→A then C→A then B→A again to get C's data */ + 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, "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; + /* 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); - if (!_wait_for("sync C", _cond_done, PHASE3_TIMEOUT_TB)) { FAIL("timeout C iter %d", iter); _intg_cleanup(); break; } - done_sync = 0; i_phase = 0; - merkle_sync_start(i_b, i_a->node_id, "rnd", _on_sync_done, NULL); - if (!_wait_for("sync B2", _cond_done, PHASE3_TIMEOUT_TB)) { FAIL("timeout B2 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; } + /* wait for broadcast relay */ + for (int k = 0; k < 100; k++) uasync_poll(i_ua, 10); 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; } 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; } _intg_cleanup(); } - if (iter == 30) PASS(); + if (iter == 5) PASS(); } /* ── Stage 6: protocol edge cases ── */ diff --git a/tools/chatgui/transport/member_sync.c b/tools/chatgui/transport/member_sync.c index e1341492..3bc03ee8 100644 --- a/tools/chatgui/transport/member_sync.c +++ b/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; } -static int _member_apply_items(void* ctx, const char* ns, - const uint8_t* data, size_t len) { +static int _member_apply_items(void* ctx, const char* ns, uint64_t from_peer, + const uint8_t* data, size_t len) { struct UTUN_INSTANCE* inst = (struct UTUN_INSTANCE*)ctx; + (void)from_peer; if (len < 2) return -1; uint16_t count; memcpy(&count, data, 2); diff --git a/tools/chatgui/transport/merkle_sync.c b/tools/chatgui/transport/merkle_sync.c index 9d868307..116e1de3 100644 --- a/tools/chatgui/transport/merkle_sync.c +++ b/tools/chatgui/transport/merkle_sync.c @@ -17,6 +17,7 @@ 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 }; struct ms_pending { @@ -31,8 +32,7 @@ struct ms_session { struct ms_session* next; char ns[64]; uint64_t peer; - uint8_t active; - uint8_t synced; + uint8_t sess_state; merkle_sync_done_cb done_cb; void* cb_arg; 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; } -static int _send_msg(struct merkle_sync* ms, uint64_t peer, const uint8_t* payload, size_t len); - /* ── Prefix arithmetic ── */ 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; } +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, 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); @@ -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) { - 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; } - s->active = 0; + 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; } merkle_sync_done_cb cb = s->done_cb; void* arg = s->cb_arg; 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); } - 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); } + 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->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); } @@ -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); const uint8_t* pd = payload + 2; if (paylen - 2 < dlen) return; - ms->ops->apply_items(ms->data_ctx, ns, pd, dlen); - if (s && s->active) _session_done(s, MT_OK); + ms->ops->apply_items(ms->data_ctx, ns, peer, pd, dlen); + if (s && s->sess_state == SESS_SYNCING) _session_done(s, MT_OK); 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; snprintf(s->ns, sizeof(s->ns), "%s", ns); s->peer = peer; s->next = ms->sessions; ms->sessions = s; + } else { + s->pending_count = 0; } - s->active = 1; + s->sess_state = SESS_SYNCING; if (paylen < 4) return; 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) { uint16_t dlen; memcpy(&dlen, bp, 2); bp += 2; rem -= 2; 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; 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; } @@ -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++) { 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); _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); 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); if (!entry) continue; 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); } +/* ── 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 ── */ 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; ms->sessions = s; } else { - s->synced = 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", - MS_ID, (unsigned long long)peer, ns, - (void*)s->done_cb, s->cb_arg, (void*)done_cb, arg); + DEBUG_DEBUG(DEBUG_CATEGORY_DB_SYNC, "%s: start OVERWRITE peer=%016llx ns=%s", MS_ID, (unsigned long long)peer, ns); } - s->active = 1; + s->sess_state = SESS_SYNCING; s->done_cb = done_cb; s->cb_arg = arg; _send_hashes(ms, peer, ns, 0, 0, 0, 0); diff --git a/tools/chatgui/transport/merkle_sync.h b/tools/chatgui/transport/merkle_sync.h index 19c4a59d..6ab16e18 100644 --- a/tools/chatgui/transport/merkle_sync.h +++ b/tools/chatgui/transport/merkle_sync.h @@ -138,7 +138,7 @@ struct merkle_sync_data_ops { * * Возвращает 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); /* @@ -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, 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 ── */ /*