From be10714972e21406c490c1f392e10bd5b200dd1d Mon Sep 17 00:00:00 2001 From: Evgeny Date: Sat, 25 Jul 2026 23:00:42 +0300 Subject: [PATCH] merkle_sync: fix prefix reconstruction, remove timeouts/retries MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - _get_level_hashes: remove (0x1F << parent_shift) from range_end so bitmap only includes children of the actual parent prefix (not all parents). Add level 0 virtual root that queries all level-1 entries. - merkle_sync_start: send HASHES at level 0 (virtual root) instead of level 1/prefix 0. Root-level REQUEST prefix reconstruction now uses the correct level-1 prefix from DB. - _handle_batch: use merkle_sync_prefix_bytes(next_lvl) for synthetic child HASHES instead of current level's pb_i (could be wrong size). - Remove timeout/retry mechanism: delete _session_timeout_cb, _session_start_timer, timer/retries fields. ETCP provides reliable transport — no need for application-layer retries. --- tools/chatgui/transport/merkle_sync.c | 65 +++++++-------------------- tools/chatgui/transport/merkle_sync.h | 4 +- 2 files changed, 17 insertions(+), 52 deletions(-) diff --git a/tools/chatgui/transport/merkle_sync.c b/tools/chatgui/transport/merkle_sync.c index b8e15390..2490f896 100644 --- a/tools/chatgui/transport/merkle_sync.c +++ b/tools/chatgui/transport/merkle_sync.c @@ -12,7 +12,6 @@ #include #define MS_ID "merkle_sync" -#define MS_SYNC_TIMEOUT_MS 10000 #define MS_BG_INTERVAL_MS 100 static struct merkle_sync* g_merkle = NULL; @@ -25,8 +24,6 @@ struct ms_session { uint64_t peer; uint8_t active; uint8_t synced; - void* timer; - uint8_t retries; merkle_sync_done_cb done_cb; void* cb_arg; }; @@ -168,11 +165,19 @@ static int _get_level_hashes(struct merkle_sync* ms, const char* ns, if (level >= MT_MAX_LEVEL) return 0; sqlite3* db = _db(ms->inst); if (!db) return -1; - int parent_shift = 63 - (int)level * 5; if (parent_shift < 0) parent_shift = 0; - int next_shift = 63 - ((int)level + 1) * 5; if (next_shift < 0) next_shift = 0; - uint64_t range_end = prefix | (0x1FULL << parent_shift) | (0x1FULL << (next_shift > 0 ? next_shift : 0)); + int next_shift; + + if (level == 0) { + next_shift = 63 - 5; + } else { + next_shift = 63 - ((int)level + 1) * 5; if (next_shift < 0) next_shift = 0; + } + + uint64_t range_end = (level == 0) ? UINT64_MAX + : prefix | (0x1FULL << (next_shift > 0 ? next_shift : 0)); sqlite3_stmt* stmt = NULL; + int query_level = (int)(level + 1); if (sqlite3_prepare_v2(db, "SELECT prefix64, hash FROM merkle_tree_hash" " WHERE namespace=? AND level=? AND prefix64>=? AND prefix64<=?" @@ -181,7 +186,7 @@ static int _get_level_hashes(struct merkle_sync* ms, const char* ns, return -1; } sqlite3_bind_text(stmt, 1, ns, -1, SQLITE_STATIC); - sqlite3_bind_int(stmt, 2, (int)(level + 1)); + sqlite3_bind_int(stmt, 2, query_level); sqlite3_bind_int64(stmt, 3, (sqlite3_int64)prefix); sqlite3_bind_int64(stmt, 4, (sqlite3_int64)range_end); while (sqlite3_step(stmt) == SQLITE_ROW) { @@ -315,34 +320,6 @@ static int _send_batch(struct merkle_sync* ms, uint64_t peer, const char* ns, /* ── Sessions ── */ -static void _session_start_timer(struct ms_session* s); - -static void _session_timeout_cb(void* arg) { - struct ms_session* s = (struct ms_session*)arg; - if (!s || !s->active || !g_merkle || !g_merkle->initialized) return; - DEBUG_TRACE(DEBUG_CATEGORY_DB_SYNC, "%s: timeout peer=%016llx ns=%s retry=%d/3", MS_ID, (unsigned long long)s->peer, s->ns, s->retries); - s->timer = NULL; - s->retries++; - if (s->retries > 3) { - DEBUG_TRACE(DEBUG_CATEGORY_DB_SYNC, "%s: → max retries, done(ERR_TIMEOUT)", MS_ID); - DEBUG_WARN(DEBUG_CATEGORY_DB_SYNC, "%s: sync timeout peer=%016llx ns=%s", MS_ID, - (unsigned long long)s->peer, s->ns); - s->active = 0; - if (s->done_cb) s->done_cb(s->peer, s->ns, MT_ERR_TIMEOUT, s->cb_arg); - return; - } - DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "%s: retry %d peer=%016llx ns=%s", MS_ID, - s->retries, (unsigned long long)s->peer, s->ns); - _send_hashes(g_merkle, s->peer, s->ns, 1, 0, 1, 0); - _session_start_timer(s); -} - -static void _session_start_timer(struct ms_session* s) { - if (!g_merkle || !g_merkle->inst) return; - s->timer = uasync_set_timeout(g_merkle->inst->ua, (uint32_t)(MS_SYNC_TIMEOUT_MS * 10), - s, _session_timeout_cb, "ms_sync"); -} - static struct ms_session* _session_find(struct merkle_sync* ms, uint64_t peer, const char* ns) { for (struct ms_session* s = ms->sessions; s; s = s->next) if (s->peer == peer && strcmp(s->ns, ns) == 0) return s; @@ -351,7 +328,6 @@ static struct ms_session* _session_find(struct merkle_sync* ms, uint64_t peer, c static void _session_done(struct ms_session* s, int result) { s->active = 0; - if (s->timer) { uasync_cancel_timeout(g_merkle->inst->ua, s->timer); s->timer = 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 (s->done_cb) s->done_cb(s->peer, s->ns, result, s->cb_arg); @@ -370,7 +346,6 @@ static void _handle_hashes(struct merkle_sync* ms, uint64_t peer, const char* ns DEBUG_TRACE(DEBUG_CATEGORY_DB_SYNC, "%s: handle_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); struct ms_session* s = _session_find(ms, peer, ns); - if (s && s->timer) { uasync_cancel_timeout(g_merkle->inst->ua, s->timer); s->timer = NULL; } if (is_data) { if (paylen < 2) return; @@ -447,7 +422,6 @@ static void _handle_hashes(struct merkle_sync* ms, uint64_t peer, const char* ns _send_msg(ms, peer, rbuf, (size_t)(wr - rbuf)); u_free(rbuf); } - if (s) _session_start_timer(s); } else { if (s) _session_done(s, MT_OK); } @@ -472,11 +446,6 @@ static void _handle_request(struct merkle_sync* ms, uint64_t peer, const char* n } if (bc > 0) { _send_batch(ms, peer, ns, buckets, bc); - struct ms_session* s = _session_find(ms, peer, ns); - if (s && s->active) { - if (s->timer) { uasync_cancel_timeout(ms->inst->ua, s->timer); s->timer = NULL; } - _session_start_timer(s); - } } } @@ -506,7 +475,7 @@ static void _handle_batch(struct merkle_sync* ms, uint64_t peer, const char* ns, uint8_t sub_pl[4096]; size_t sub_len = 0; uint8_t next_lvl = (uint8_t)(lvl < MT_MAX_LEVEL ? lvl + 1 : lvl); sub_pl[sub_len++] = next_lvl; - sub_pl[sub_len++] = pb_i; + sub_pl[sub_len++] = merkle_sync_prefix_bytes(next_lvl); _prefix_write(sub_pl + sub_len, pr, pb_i); sub_len += pb_i; sub_pl[sub_len++] = 0; /* is_data=0 */ uint32_t bm; memcpy(&bm, bp, 4); bp += 4; rem -= 4; @@ -690,7 +659,6 @@ void merkle_sync_destroy(struct UTUN_INSTANCE* inst) { struct ms_session* s = ms->sessions; while (s) { struct ms_session* next = s->next; - if (s->timer) uasync_cancel_timeout(inst->ua, s->timer); u_free(s); s = next; } u_free(ms); @@ -715,14 +683,12 @@ int merkle_sync_start(struct UTUN_INSTANCE* inst, uint64_t peer, ms->sessions = s; } else { /* replace callback if session already exists */ - if (s->timer) { uasync_cancel_timeout(inst->ua, s->timer); s->timer = NULL; } s->synced = 0; } - s->active = 1; s->retries = 0; + s->active = 1; s->done_cb = done_cb; s->cb_arg = arg; - _send_hashes(ms, peer, ns, 1, 0, 1, 0); - _session_start_timer(s); + _send_hashes(ms, peer, ns, 0, 0, 0, 0); return 0; } @@ -734,7 +700,6 @@ void merkle_sync_cancel(struct UTUN_INSTANCE* inst, uint64_t peer, const char* n while (*p) { struct ms_session* s = *p; if (s->peer == peer && strcmp(s->ns, ns) == 0) { - if (s->timer) { uasync_cancel_timeout(inst->ua, s->timer); s->timer = NULL; } *p = s->next; u_free(s); return; } diff --git a/tools/chatgui/transport/merkle_sync.h b/tools/chatgui/transport/merkle_sync.h index 8aa88fe9..19c4a59d 100644 --- a/tools/chatgui/transport/merkle_sync.h +++ b/tools/chatgui/transport/merkle_sync.h @@ -85,8 +85,8 @@ struct UTUN_INSTANCE; * * ── Сессии ── * - * Таймаут 10 секунд, 3 ретрая. При таймауте сессия перезапускает - * MSG_HASHES с level=1. После 3 ретраев — done_cb с MT_ERR_TIMEOUT. + * Протокол работает поверх надёжного транспорта (ETCP). + * Таймаутов и ретраев нет — при ошибке сессия завершается. */ /* ── Data model callbacks ── */