diff --git a/tools/chatgui/CMakeLists.txt b/tools/chatgui/CMakeLists.txt index 48819429..cc3168b7 100644 --- a/tools/chatgui/CMakeLists.txt +++ b/tools/chatgui/CMakeLists.txt @@ -75,6 +75,7 @@ add_executable(chatgui transport/chat_core.c transport/chat_sync.c transport/member_sync.c + transport/merkle_sync.c db/db_manager.cpp ../../lib/sqlite3.c resources/chatgui.qrc @@ -82,7 +83,7 @@ add_executable(chatgui target_include_directories(chatgui PRIVATE ${CMAKE_SOURCE_DIR}/../../lib ${CMAKE_SOURCE_DIR}/../../src ${CMAKE_SOURCE_DIR}/db) target_compile_definitions(chatgui PRIVATE SQLITE_THREADSAFE=1 USE_SQLITE) -set_source_files_properties(../../lib/sqlite3.c transport/chat_core.c transport/chat_sync.c transport/member_sync.c PROPERTIES LANGUAGE C) +set_source_files_properties(../../lib/sqlite3.c transport/chat_core.c transport/chat_sync.c transport/member_sync.c transport/merkle_sync.c PROPERTIES LANGUAGE C) if(WIN32) target_link_libraries(chatgui PRIVATE ${QT_LIBS} rlottie::rlottie ZLIB::ZLIB OpenSSL::Crypto ZXing::ZXing utun pthread) else() diff --git a/tools/chatgui/libutun/CMakeLists.txt b/tools/chatgui/libutun/CMakeLists.txt index b524ba8b..128409b2 100644 --- a/tools/chatgui/libutun/CMakeLists.txt +++ b/tools/chatgui/libutun/CMakeLists.txt @@ -65,6 +65,7 @@ set(UTUN_COMMON_SOURCES ${SRC_DIR}/db_sync.c ${SRC_DIR}/conn_mgr.c ${TRANSPORT_DIR}/chat_sync.c + ${TRANSPORT_DIR}/merkle_sync.c ${TRANSPORT_DIR}/member_sync.c ${SRC_DIR}/routing.c ${SRC_DIR}/tun_if.c diff --git a/tools/chatgui/transport/chat_sync.c b/tools/chatgui/transport/chat_sync.c index e16491dd..c1a59e22 100644 --- a/tools/chatgui/transport/chat_sync.c +++ b/tools/chatgui/transport/chat_sync.c @@ -374,6 +374,13 @@ static void cs_join_timeout_cb(void* arg) { cs->pending_invite_ch_id = 0; } +static void _on_member_sync_done(uint64_t peer, const char* ns, int result, void* arg) { + struct channel_cache* ch = (struct channel_cache*)arg; + if (result == MT_OK && ch) ch->synced = CS_SYNC_DONE; + DEBUG_INFO(DEBUG_CATEGORY_CONNECTIVITY, "%s: member_sync %s ns=%s peer=%016llx", + CS_ID, result == MT_OK ? "OK" : "FAIL", ns, (unsigned long long)peer); +} + static void cs_on_conn_up(struct ETCP_CONN* conn, void* arg) { (void)arg; if (!conn || !g_cs) return; @@ -405,7 +412,7 @@ static void cs_on_conn_up(struct ETCP_CONN* conn, void* arg) { memcpy(msg + 1, &ch->msg_count, 4); memcpy(msg + 5, ch->last_chain_hash, 32); cs_send(g_cs, ch->channel_id, peer, msg, 37); - member_sync_start(g_cs->inst, peer, ch->channel_id); + member_sync_start(g_cs->inst, peer, ch->channel_id, _on_member_sync_done, ch); } member_sync_set_online(g_cs->inst, peer, 1); } @@ -985,7 +992,7 @@ static void cs_handle_welcome(struct chat_sync* cs, uint64_t peer, sqlite3_bind_int(stmt, 4, (int)port); sqlite3_step(stmt); sqlite3_finalize(stmt); } - member_sync_cancel_peer(g_cs->inst, peer); + member_sync_cancel(g_cs->inst, peer, ch_id); member_sync_set_online(g_cs->inst, peer, 0); } } diff --git a/tools/chatgui/transport/member_sync.c b/tools/chatgui/transport/member_sync.c index 2bfb6e7e..b1b225fe 100644 --- a/tools/chatgui/transport/member_sync.c +++ b/tools/chatgui/transport/member_sync.c @@ -2,79 +2,37 @@ #include "topo_node_sqlite.h" #include "../../../src/utun_instance.h" -#include "../../../src/etcp_router.h" -#include "../../../src/etcp_api.h" -#include "../../../src/etcp.h" #include "../../../src/topo_group.h" #include "../../../lib/debug_config.h" #include "../../../lib/mem.h" -#include "../../../lib/u_async.h" #include #include -#include #include #define MS_ID "member_sync" -#define MS_SYNC_TIMEOUT_MS 10000 -#define MS_BG_INTERVAL_MS 100 -static struct member_sync* g_ms = NULL; +struct addr_item { uint8_t family; uint8_t addr[16]; uint16_t port; }; -struct bucket_entry { - uint8_t level; - uint8_t prefix_bytes; - uint64_t prefix; -}; - -struct addr_item { - uint8_t family; - uint8_t addr[16]; - uint16_t port; -}; - -struct ms_session { - struct ms_session* next; - char ch_id[64]; - uint64_t peer; - uint8_t active; - void* timer; - uint32_t started_tb; - uint8_t retries; -}; - -struct member_sync { - struct UTUN_INSTANCE* inst; - struct ms_session* sessions; - uint8_t initialized; - void* bg_timer; -}; +/* ── DB access ── */ static sqlite3* _db(struct UTUN_INSTANCE* inst) { return inst && inst->topo_groups ? inst->topo_groups->topo_sqlite_db : NULL; } -static uint64_t _level_prefix(uint64_t member_id, uint8_t level) { - int shift = 63 - (int)level * 5; - if (shift < 0) shift = 0; - return (member_id >> shift) << shift; -} - -static uint8_t _prefix_bytes(uint8_t level) { - int bits = level * 5; - return (uint8_t)((bits + 7) / 8); -} - -static void _prefix_write(uint8_t* out, uint64_t prefix, uint8_t pb) { - for (int i = (int)pb - 1; i >= 0; i--) - out[pb - 1 - i] = (uint8_t)(prefix >> (i * 8)); +static void _peers_table(const char* ch_id, char* buf, size_t sz) { + size_t i = 0; + while (*ch_id && i < sz - 1) { + char c = *ch_id++; + if ((c >= 'a' && c <= 'z') || (c >= 'A' && c <= 'Z') || (c >= '0' && c <= '9') || c == '_') + buf[i++] = c; else buf[i++] = '_'; + } + buf[i] = '\0'; + char tbl[128]; snprintf(tbl, sizeof(tbl), "peers_%s", buf); + snprintf(buf, sz, "%s", tbl); } -static uint64_t _prefix_read(const uint8_t* in, uint8_t pb) { - uint64_t v = 0; - for (uint8_t i = 0; i < pb && i < 8; i++) v = (v << 8) | in[i]; - return v; -} +/* ── Member hash (identical to old _compute_member_hash) ── */ static int _addr_cmp(const void* a, const void* b) { const struct addr_item* ia = (const struct addr_item*)a; @@ -88,7 +46,7 @@ static int _addr_cmp(const void* a, const void* b) { static void _compute_member_hash(uint64_t node_id, const uint8_t* x25519, const uint8_t* ed25519, const uint8_t* join_sig, const uint8_t* addrs_data, int addr_count, int online, - uint8_t hash_out[MS_HASH_SIZE]) { + uint8_t hash_out[MT_HASH_SIZE]) { EVP_MD_CTX* ctx = EVP_MD_CTX_new(); EVP_DigestInit_ex(ctx, EVP_sha256(), NULL); EVP_DigestUpdate(ctx, &node_id, 8); @@ -102,8 +60,7 @@ static void _compute_member_hash(uint64_t node_id, const uint8_t* x25519, struct addr_item items[256]; int n = 0; const uint8_t* p = addrs_data; for (int i = 0; i < addr_count && n < 256; i++) { - uint8_t fam = *p++; - items[n].family = fam; + uint8_t fam = *p++; items[n].family = fam; int ip_len = fam == 4 ? 4 : 16; memcpy(items[n].addr, p, (size_t)ip_len); p += ip_len; items[n].port = ((uint16_t)p[0] << 8) | p[1]; p += 2; @@ -124,28 +81,18 @@ static void _compute_member_hash(uint64_t node_id, const uint8_t* x25519, EVP_MD_CTX_free(ctx); } -static void _peers_table(const char* ch_id, char* buf, size_t sz) { - size_t i = 0; - while (*ch_id && i < sz - 1) { - char c = *ch_id++; - if ((c >= 'a' && c <= 'z') || (c >= 'A' && c <= 'Z') || (c >= '0' && c <= '9') || c == '_') - buf[i++] = c; - else - buf[i++] = '_'; - } - buf[i] = '\0'; - char tbl[128]; snprintf(tbl, sizeof(tbl), "peers_%s", buf); - snprintf(buf, sz, "%s", tbl); -} +/* ── merkle_sync_data_ops implementation ── */ -static int _recompute_bucket(struct UTUN_INSTANCE* inst, const char* ch_id, - uint8_t level, uint64_t prefix64) { +static int _member_update_bucket_hash(void* ctx, const char* ns, uint8_t level, + uint64_t prefix64, EVP_MD_CTX* sha_ctx) { + struct UTUN_INSTANCE* inst = (struct UTUN_INSTANCE*)ctx; sqlite3* db = _db(inst); if (!db) return -1; - char peers_tbl[128]; _peers_table(ch_id, peers_tbl, sizeof(peers_tbl)); - char sql[512]; + char peers_tbl[128]; _peers_table(ns, peers_tbl, sizeof(peers_tbl)); + int mask_shift = 63 - (int)level * 5; uint64_t mask = (mask_shift >= 0) ? (~0ULL << mask_shift) : UINT64_MAX; - snprintf(sql, sizeof(sql), + + char sql[512]; snprintf(sql, sizeof(sql), "SELECT p.node_id, n.x25519_pubkey, n.ed25519_pubkey, p.join_sig, n.online" " FROM \"%s\" p JOIN nodes n ON p.node_id=n.node_id" " WHERE (p.node_id & %lld) == %lld ORDER BY p.node_id ASC", @@ -154,10 +101,7 @@ static int _recompute_bucket(struct UTUN_INSTANCE* inst, const char* ch_id, sqlite3_stmt* stmt = NULL; if (sqlite3_prepare_v2(db, sql, -1, &stmt, NULL) != SQLITE_OK) return -1; - EVP_MD_CTX* ctx = EVP_MD_CTX_new(); - EVP_DigestInit_ex(ctx, EVP_sha256(), NULL); int count = 0; - while (sqlite3_step(stmt) == SQLITE_ROW) { uint64_t nid = (uint64_t)sqlite3_column_int64(stmt, 0); const uint8_t* x25 = (const uint8_t*)sqlite3_column_blob(stmt, 1); @@ -165,214 +109,26 @@ static int _recompute_bucket(struct UTUN_INSTANCE* inst, const char* ch_id, const uint8_t* sig = (const uint8_t*)sqlite3_column_blob(stmt, 3); int online = sqlite3_column_int(stmt, 4); if (!x25 || !ed || !sig) continue; - - uint8_t mh[MS_HASH_SIZE]; + uint8_t mh[MT_HASH_SIZE]; _compute_member_hash(nid, x25, ed, sig, NULL, 0, online, mh); - EVP_DigestUpdate(ctx, mh, MS_HASH_SIZE); + EVP_DigestUpdate(sha_ctx, mh, MT_HASH_SIZE); count++; } sqlite3_finalize(stmt); - - if (count == 0) { - EVP_MD_CTX_free(ctx); - snprintf(sql, sizeof(sql), - "DELETE FROM member_tree_hash WHERE channel_id=? AND level=? AND prefix64=?"); - sqlite3_stmt* ds = NULL; - if (sqlite3_prepare_v2(db, sql, -1, &ds, NULL) == SQLITE_OK) { - sqlite3_bind_text(ds, 1, ch_id, -1, SQLITE_STATIC); - sqlite3_bind_int(ds, 2, level); - sqlite3_bind_int64(ds, 3, (sqlite3_int64)prefix64); - sqlite3_step(ds); sqlite3_finalize(ds); - } - return 0; - } - - uint8_t hash[MS_HASH_SIZE]; - EVP_DigestFinal_ex(ctx, hash, NULL); - EVP_MD_CTX_free(ctx); - - snprintf(sql, sizeof(sql), - "INSERT OR REPLACE INTO member_tree_hash(channel_id, level, prefix64, hash, member_count)" - " VALUES(?,?,?,?,?)"); - sqlite3_stmt* is = NULL; - if (sqlite3_prepare_v2(db, sql, -1, &is, NULL) == SQLITE_OK) { - sqlite3_bind_text(is, 1, ch_id, -1, SQLITE_STATIC); - sqlite3_bind_int(is, 2, level); - sqlite3_bind_int64(is, 3, (sqlite3_int64)prefix64); - sqlite3_bind_blob(is, 4, hash, MS_HASH_SIZE, SQLITE_STATIC); - sqlite3_bind_int(is, 5, count); - sqlite3_step(is); sqlite3_finalize(is); - } - return 0; -} - -static void _recompute_path(struct UTUN_INSTANCE* inst, const char* ch_id, uint64_t member_id) { - for (uint8_t level = 1; level <= MS_MAX_LEVEL; level++) - _recompute_bucket(inst, ch_id, level, _level_prefix(member_id, level)); -} - -static void _ensure_table(void) { - sqlite3* db = g_ms ? _db(g_ms->inst) : NULL; if (!db) return; - sqlite3_exec(db, - "CREATE TABLE IF NOT EXISTS member_tree_hash (" - " channel_id TEXT NOT NULL," - " level INTEGER NOT NULL CHECK(level BETWEEN 1 AND 5)," - " prefix64 INTEGER NOT NULL," - " hash BLOB NOT NULL," - " member_count INTEGER NOT NULL," - " PRIMARY KEY (channel_id, level, prefix64))", - NULL, NULL, NULL); -} - -const uint8_t* member_sync_get_hash(struct UTUN_INSTANCE* inst, const char* ch_id, - uint8_t level, uint64_t prefix64) { - static uint8_t zero[MS_HASH_SIZE]; /* returns static — not thread-safe but single-threaded uasync */ - sqlite3* db = _db(inst); if (!db) { memset(zero, 0, MS_HASH_SIZE); return zero; } - sqlite3_stmt* stmt = NULL; - if (sqlite3_prepare_v2(db, - "SELECT hash FROM member_tree_hash WHERE channel_id=? AND level=? AND prefix64=?", - -1, &stmt, NULL) != SQLITE_OK) { memset(zero, 0, MS_HASH_SIZE); return zero; } - sqlite3_bind_text(stmt, 1, ch_id, -1, SQLITE_STATIC); - sqlite3_bind_int(stmt, 2, level); - sqlite3_bind_int64(stmt, 3, (sqlite3_int64)prefix64); - memset(zero, 0, MS_HASH_SIZE); - if (sqlite3_step(stmt) == SQLITE_ROW) { - const void* h = sqlite3_column_blob(stmt, 0); - if (h) memcpy(zero, h, MS_HASH_SIZE); - } - sqlite3_finalize(stmt); - return zero; -} - -int member_sync_put(struct UTUN_INSTANCE* inst, const char* ch_id, - uint64_t member_id, const uint8_t* x25519, - const uint8_t* ed25519, const uint8_t* join_sig, - const uint8_t* addrs_data, int addr_count) { - if (!inst || !ch_id) return -1; - sqlite3* db = _db(inst); if (!db) return -1; - _ensure_table(); - - sqlite3_stmt* ns = NULL; - sqlite3_prepare_v2(db, - "INSERT OR REPLACE INTO nodes(node_id, name, x25519_pubkey, ed25519_pubkey, online)" - " VALUES(?,COALESCE((SELECT name FROM nodes WHERE node_id=?),''),?,?," - " COALESCE((SELECT online FROM nodes WHERE node_id=?),0))", - -1, &ns, NULL); - if (ns) { - sqlite3_bind_int64(ns, 1, (sqlite3_int64)member_id); - sqlite3_bind_int64(ns, 2, (sqlite3_int64)member_id); - sqlite3_bind_blob(ns, 3, x25519, 32, SQLITE_STATIC); - sqlite3_bind_blob(ns, 4, ed25519, 32, SQLITE_STATIC); - sqlite3_bind_int64(ns, 5, (sqlite3_int64)member_id); - sqlite3_step(ns); sqlite3_finalize(ns); - } - - topo_node_sqlite_member_put(db, ch_id, member_id, join_sig, NULL); - - if (addrs_data && addr_count > 0) { - sqlite3_exec(db, "BEGIN", NULL, NULL, NULL); - sqlite3_stmt* ds = NULL; - sqlite3_prepare_v2(db, "DELETE FROM node_addresses WHERE node_id=? AND is_nat=0", -1, &ds, NULL); - if (ds) { sqlite3_bind_int64(ds, 1, (sqlite3_int64)member_id); sqlite3_step(ds); sqlite3_finalize(ds); } - - sqlite3_stmt* as = NULL; - sqlite3_prepare_v2(db, - "INSERT INTO node_addresses(node_id,family,protocol,address,port,is_nat)" - " VALUES(?,?,1,?,?,0)", -1, &as, NULL); - if (as) { - const uint8_t* p = addrs_data; - for (int i = 0; i < addr_count; i++) { - uint8_t fam = *p++; int ip_sz = fam == 4 ? 4 : 16; - sqlite3_bind_int64(as, 1, (sqlite3_int64)member_id); - sqlite3_bind_int(as, 2, fam); - sqlite3_bind_blob(as, 3, p, ip_sz, SQLITE_STATIC); - p += ip_sz; - uint16_t port = ((uint16_t)p[0] << 8) | p[1]; p += 2; - sqlite3_bind_int(as, 4, port); - sqlite3_step(as); sqlite3_reset(as); - } - sqlite3_finalize(as); - } - sqlite3_exec(db, "COMMIT", NULL, NULL, NULL); - } - - _recompute_path(inst, ch_id, member_id); - return 0; -} - -int member_sync_del(struct UTUN_INSTANCE* inst, const char* ch_id, uint64_t member_id) { - if (!inst || !ch_id) return -1; - sqlite3* db = _db(inst); if (!db) return -1; - topo_node_sqlite_member_del(db, ch_id, member_id); - _recompute_path(inst, ch_id, member_id); - return 0; -} - -void member_sync_set_online(struct UTUN_INSTANCE* inst, uint64_t node_id, int online) { - if (!inst) return; - sqlite3* db = _db(inst); if (!db) return; - topo_node_sqlite_node_set_online(db, node_id, online); - if (!g_ms || !g_ms->initialized) return; - sqlite3_stmt* stmt = NULL; - if (sqlite3_prepare_v2(db, "SELECT channel_id FROM channels", -1, &stmt, NULL) != SQLITE_OK) return; - while (sqlite3_step(stmt) == SQLITE_ROW) { - const char* ch_id = (const char*)sqlite3_column_text(stmt, 0); - if (!ch_id) continue; - char peers_tbl[128]; _peers_table(ch_id, peers_tbl, sizeof(peers_tbl)); - char buf[256]; snprintf(buf, sizeof(buf), - "SELECT 1 FROM \"%s\" WHERE node_id=?", peers_tbl); - sqlite3_stmt* cs = NULL; - if (sqlite3_prepare_v2(db, buf, -1, &cs, NULL) == SQLITE_OK) { - sqlite3_bind_int64(cs, 1, (sqlite3_int64)node_id); - if (sqlite3_step(cs) == SQLITE_ROW) _recompute_path(inst, ch_id, node_id); - sqlite3_finalize(cs); - } - } - sqlite3_finalize(stmt); + return count; } -int member_sync_count(struct UTUN_INSTANCE* inst, const char* ch_id) { - sqlite3* db = _db(inst); - if (!db || !ch_id) return 0; - char peers_tbl[128]; _peers_table(ch_id, peers_tbl, sizeof(peers_tbl)); - char sql[256]; snprintf(sql, sizeof(sql), "SELECT COUNT(*) FROM \"%s\"", peers_tbl); - sqlite3_stmt* stmt = NULL; - if (sqlite3_prepare_v2(db, sql, -1, &stmt, NULL) != SQLITE_OK) return 0; - int c = 0; - if (sqlite3_step(stmt) == SQLITE_ROW) c = sqlite3_column_int(stmt, 0); - sqlite3_finalize(stmt); - return c; -} - -int member_sync_get_level_hashes(struct UTUN_INSTANCE* inst, const char* ch_id, - uint8_t level, uint64_t prefix, uint8_t prefix_bytes, - uint32_t* bitmap, uint8_t hashes[MS_BUCKETS][MS_HASH_SIZE]) { - (void)prefix_bytes; - memset(hashes, 0, sizeof(uint8_t) * MS_BUCKETS * MS_HASH_SIZE); - *bitmap = 0; - if (level >= MS_MAX_LEVEL) return 0; - int next_shift = 63 - ((int)level + 1) * 5; - - for (int i = 0; i < MS_BUCKETS; i++) { - uint64_t child_prefix = prefix | ((uint64_t)i << (next_shift > 0 ? next_shift : 0)); - const uint8_t* h = member_sync_get_hash(inst, ch_id, (uint8_t)(level + 1), child_prefix); - int empty = 1; - for (int j = 0; j < MS_HASH_SIZE; j++) { if (h[j] != 0) { empty = 0; break; } } - if (!empty) { *bitmap |= (1u << i); memcpy(hashes[i], h, MS_HASH_SIZE); } - } - return 0; -} - -int member_sync_get_bucket_members(struct UTUN_INSTANCE* inst, const char* ch_id, - uint8_t level, uint64_t prefix, uint8_t prefix_bytes, - uint8_t* buf, size_t* len) { - (void)prefix_bytes; +static int _member_get_items(void* ctx, const char* ns, uint8_t level, + uint64_t prefix, uint8_t prefix_bytes, + uint8_t* buf, size_t* len) { + struct UTUN_INSTANCE* inst = (struct UTUN_INSTANCE*)ctx; sqlite3* db = _db(inst); if (!db || !buf || !len) return -1; - char peers_tbl[128]; _peers_table(ch_id, peers_tbl, sizeof(peers_tbl)); - char sql[512]; + char peers_tbl[128]; _peers_table(ns, peers_tbl, sizeof(peers_tbl)); + int mask_shift = 63 - (int)level * 5; uint64_t mask = (mask_shift >= 0) ? (~0ULL << mask_shift) : UINT64_MAX; - snprintf(sql, sizeof(sql), + + char sql[512]; snprintf(sql, sizeof(sql), "SELECT p.node_id, n.x25519_pubkey, n.ed25519_pubkey, p.join_sig, n.online" " FROM \"%s\" p JOIN nodes n ON p.node_id=n.node_id" " WHERE (p.node_id & %lld) == %lld ORDER BY p.node_id ASC", @@ -430,514 +186,191 @@ int member_sync_get_bucket_members(struct UTUN_INSTANCE* inst, const char* ch_id return 0; } -static int _send_msg(struct UTUN_INSTANCE* inst, uint64_t peer, - const uint8_t* payload, size_t len) { - if (!inst || len < 1) return -1; - uint8_t* buf = u_malloc(1 + len); - if (!buf) return -1; - buf[0] = ETCP_RT_ID_MEMBER_SYNC; - memcpy(buf + 1, payload, len); - struct ll_entry* entry = queue_entry_new(0); - if (!entry) { u_free(buf); return -1; } - entry->dgram = buf; entry->len = 1 + len; - int r = etcp_route_send(inst, peer, entry, 0); - if (r != 0) { u_free(buf); queue_entry_free(entry); } - return r; -} - -static int _popcount_u32(uint32_t v) { - return __builtin_popcount(v); -} -#ifndef __has_builtin -#define __has_builtin(x) 0 -#endif -#if !defined(__GNUC__) && !defined(__clang__) -static int _popcount_u32(uint32_t v) { - v = v - ((v >> 1) & 0x55555555); - v = (v & 0x33333333) + ((v >> 2) & 0x33333333); - return ((v + (v >> 4) & 0x0F0F0F0F) * 0x01010101) >> 24; -} -#endif - -static int _send_hashes(struct UTUN_INSTANCE* inst, uint64_t peer, const char* ch_id, - uint8_t level, uint64_t prefix, uint8_t prefix_bytes, int is_data) { - uint8_t ch_len = (uint8_t)strlen(ch_id); - size_t max_sz = 1 + 1 + ch_len + 1 + 1 + prefix_bytes + 1 + 4 + MS_BUCKETS * MS_HASH_SIZE + 2 + 65536; - uint8_t* buf = u_malloc(max_sz); - if (!buf) return -1; - uint8_t* p = buf; - *p++ = MS_MSG_HASHES; - *p++ = ch_len; memcpy(p, ch_id, ch_len); p += ch_len; - *p++ = level; - *p++ = prefix_bytes; _prefix_write(p, prefix, prefix_bytes); p += prefix_bytes; - *p++ = (uint8_t)(is_data ? 1 : 0); - - if (is_data) { - size_t mlen = 65536; - uint8_t* mbuf = u_malloc(mlen); - if (mbuf) { - if (member_sync_get_bucket_members(inst, ch_id, level, prefix, prefix_bytes, mbuf, &mlen) == 0) - { memcpy(p, mbuf, mlen); p += mlen; } - else { uint16_t zero = 0; memcpy(p, &zero, 2); p += 2; } - u_free(mbuf); +static int _member_apply_items(void* ctx, const char* ns, + const uint8_t* data, size_t len) { + struct UTUN_INSTANCE* inst = (struct UTUN_INSTANCE*)ctx; + if (len < 2) return -1; + + uint16_t count; memcpy(&count, data, 2); + const uint8_t* mp = data + 2; size_t mrem = len - 2; + + for (uint16_t i = 0; i < count && mrem >= 137; i++) { + uint64_t nid; memcpy(&nid, mp, 8); mp += 8; mrem -= 8; + const uint8_t* x25 = mp; mp += 32; mrem -= 32; + const uint8_t* ed = mp; mp += 32; mrem -= 32; + const uint8_t* sig = mp; mp += 64; mrem -= 64; + uint8_t online = *mp++; mrem--; + uint8_t ac = *mp++; mrem--; + const uint8_t* addrs = mp; + int consumed = 0; + for (int a = 0; a < (int)ac && mrem >= (size_t)(1 + consumed); a++) { + uint8_t fam = mp[consumed]; consumed++; + int sz = fam == 4 ? 4 : 16; + consumed += sz + 2; } - } else { - uint32_t bitmap; uint8_t hashes[MS_BUCKETS][MS_HASH_SIZE]; - member_sync_get_level_hashes(inst, ch_id, level, prefix, prefix_bytes, &bitmap, hashes); - memcpy(p, &bitmap, 4); p += 4; - for (int i = 0; i < MS_BUCKETS; i++) - if (bitmap & (1u << i)) { memcpy(p, hashes[i], MS_HASH_SIZE); p += MS_HASH_SIZE; } + member_sync_put(inst, ns, nid, x25, ed, sig, addrs, (int)ac); + if (online) member_sync_set_online(inst, nid, 1); + mp += consumed; mrem -= (size_t)consumed; } - - int r = _send_msg(inst, peer, buf, (size_t)(p - buf)); - u_free(buf); - return r; + return 0; } -static int _send_batch(struct UTUN_INSTANCE* inst, uint64_t peer, const char* ch_id, - struct bucket_entry* buckets, int count) { - uint8_t ch_len = (uint8_t)strlen(ch_id); - size_t max_sz = 1 + 1 + ch_len + 1 + (size_t)count * (1 + 1 + 8 + 1 + 4 + MS_BUCKETS * MS_HASH_SIZE + 65536); - uint8_t* buf = u_malloc(max_sz); - if (!buf) return -1; - uint8_t* p = buf; - *p++ = MS_MSG_BATCH; - *p++ = ch_len; memcpy(p, ch_id, ch_len); p += ch_len; - *p++ = (uint8_t)count; - - for (int i = 0; i < count; i++) { - *p++ = buckets[i].level; - *p++ = buckets[i].prefix_bytes; - _prefix_write(p, buckets[i].prefix, buckets[i].prefix_bytes); p += buckets[i].prefix_bytes; - - int bcount = 0; - uint64_t mask = 0; - int mask_shift = 63 - (int)buckets[i].level * 5; - if (mask_shift >= 0) mask = ~0ULL << mask_shift; else mask = UINT64_MAX; - char peers_tbl[128]; _peers_table(ch_id, peers_tbl, sizeof(peers_tbl)); - sqlite3* db = _db(inst); - if (db) { - char sql[256]; snprintf(sql, sizeof(sql), - "SELECT COUNT(*) FROM \"%s\" p WHERE (p.node_id & %lld)==%lld", - peers_tbl, (long long)mask, (long long)buckets[i].prefix); - sqlite3_stmt* cs = NULL; - if (sqlite3_prepare_v2(db, sql, -1, &cs, NULL) == SQLITE_OK) { - if (sqlite3_step(cs) == SQLITE_ROW) bcount = sqlite3_column_int(cs, 0); - sqlite3_finalize(cs); - } - } +static const struct merkle_sync_data_ops g_member_ops = { + .update_bucket_hash = _member_update_bucket_hash, + .get_items = _member_get_items, + .apply_items = _member_apply_items, +}; - int is_terminal = (buckets[i].level >= MS_MAX_LEVEL || bcount < 8); - *p++ = (uint8_t)(is_terminal ? 1 : 0); - - if (is_terminal) { - size_t mlen = 65536; - uint8_t* mbuf = u_malloc(mlen); - if (mbuf) { - if (member_sync_get_bucket_members(inst, ch_id, buckets[i].level, - buckets[i].prefix, buckets[i].prefix_bytes, mbuf, &mlen) == 0) - { memcpy(p, mbuf, mlen); p += mlen; } - else { uint16_t z = 0; memcpy(p, &z, 2); p += 2; } - u_free(mbuf); - } - } else { - uint32_t bm; uint8_t hs[MS_BUCKETS][MS_HASH_SIZE]; - member_sync_get_level_hashes(inst, ch_id, buckets[i].level, - buckets[i].prefix, buckets[i].prefix_bytes, &bm, hs); - memcpy(p, &bm, 4); p += 4; - for (int j = 0; j < MS_BUCKETS; j++) - if (bm & (1u << j)) { memcpy(p, hs[j], MS_HASH_SIZE); p += MS_HASH_SIZE; } +/* ── node_updated callback ── */ + +static void _on_node_updated(struct UTUN_INSTANCE* inst, uint64_t node_id, + const uint8_t* x25519, const uint8_t* ed25519) { + (void)x25519; (void)ed25519; + if (!inst) return; + sqlite3* db = _db(inst); if (!db) return; + sqlite3_stmt* cs = NULL; + if (sqlite3_prepare_v2(db, "SELECT channel_id FROM channels", -1, &cs, NULL) != SQLITE_OK) return; + while (sqlite3_step(cs) == SQLITE_ROW) { + const char* ch = (const char*)sqlite3_column_text(cs, 0); + if (!ch) continue; + char peers_tbl[128]; _peers_table(ch, peers_tbl, sizeof(peers_tbl)); + char buf[256]; snprintf(buf, sizeof(buf), "SELECT 1 FROM \"%s\" WHERE node_id=?", peers_tbl); + sqlite3_stmt* ps = NULL; + if (sqlite3_prepare_v2(db, buf, -1, &ps, NULL) == SQLITE_OK) { + sqlite3_bind_int64(ps, 1, (sqlite3_int64)node_id); + if (sqlite3_step(ps) == SQLITE_ROW) + merkle_sync_recompute_path(inst, ch, node_id); + sqlite3_finalize(ps); } } - - int r = _send_msg(inst, peer, buf, (size_t)(p - buf)); - u_free(buf); - return r; + sqlite3_finalize(cs); } -/* ── 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_ms || !g_ms->initialized) return; - s->timer = NULL; - s->retries++; - if (s->retries > 3) { - DEBUG_WARN(DEBUG_CATEGORY_CONNECTIVITY, "%s: sync timeout peer=%016llx ch=%s", MS_ID, - (unsigned long long)s->peer, s->ch_id); - s->active = 0; - return; - } - DEBUG_INFO(DEBUG_CATEGORY_CONNECTIVITY, "%s: retry %d peer=%016llx ch=%s", MS_ID, - s->retries, (unsigned long long)s->peer, s->ch_id); - _send_hashes(g_ms->inst, s->peer, s->ch_id, 1, 0, 1, 0); - _session_start_timer(s); -} +/* ── Public API ── */ -static void _session_start_timer(struct ms_session* s) { - if (!g_ms || !g_ms->inst) return; - s->timer = uasync_set_timeout(g_ms->inst->ua, (uint32_t)(MS_SYNC_TIMEOUT_MS * 10), - s, _session_timeout_cb, "ms_sync"); +int member_sync_init(struct UTUN_INSTANCE* inst) { + if (!inst) return -1; + int rc = merkle_sync_init(inst, 0x31, &g_member_ops, inst); + if (rc != 0) return rc; + topo_groups_set_node_updated_cb(inst->topo_groups, _on_node_updated); + DEBUG_INFO(DEBUG_CATEGORY_CONNECTIVITY, "%s: initialized", MS_ID); + return 0; } -static struct ms_session* _session_find(uint64_t peer, const char* ch_id) { - for (struct ms_session* s = g_ms ? g_ms->sessions : NULL; s; s = s->next) - if (s->peer == peer && strcmp(s->ch_id, ch_id) == 0) return s; - return NULL; +void member_sync_destroy(struct UTUN_INSTANCE* inst) { + if (!inst) return; + topo_groups_set_node_updated_cb(inst->topo_groups, NULL); + merkle_sync_destroy(inst); } -void member_sync_start(struct UTUN_INSTANCE* inst, uint64_t peer, const char* ch_id) { - if (!inst || !ch_id || !g_ms || !g_ms->initialized) return; - _ensure_table(); - - struct ms_session* s = _session_find(peer, ch_id); - if (!s) { - s = u_calloc(1, sizeof(*s)); - if (!s) return; - snprintf(s->ch_id, sizeof(s->ch_id), "%s", ch_id); - s->peer = peer; - s->next = g_ms->sessions; - g_ms->sessions = s; - } - s->active = 1; s->retries = 0; - - _send_hashes(inst, peer, ch_id, 1, 0, 1, 0); - _session_start_timer(s); +int member_sync_start(struct UTUN_INSTANCE* inst, uint64_t peer, + const char* ch_id, merkle_sync_done_cb done_cb, void* arg) { + return merkle_sync_start(inst, peer, ch_id, done_cb, arg); } -void member_sync_cancel_peer(struct UTUN_INSTANCE* inst, uint64_t peer) { - (void)inst; - struct ms_session** p = g_ms ? &g_ms->sessions : NULL; - while (p && *p) { - struct ms_session* s = *p; - if (s->peer == peer) { - if (s->timer) { uasync_cancel_timeout(g_ms->inst->ua, s->timer); s->timer = NULL; } - *p = s->next; u_free(s); - } else { - p = &(*p)->next; - } - } +void member_sync_cancel(struct UTUN_INSTANCE* inst, uint64_t peer, const char* ch_id) { + merkle_sync_cancel(inst, peer, ch_id); } -/* ── recv handlers ── */ - -static void _handle_hashes(struct UTUN_INSTANCE* inst, uint64_t peer, const char* ch_id, - const uint8_t* pl, size_t plen) { - if (plen < 3) return; - uint8_t level = pl[0]; uint8_t pb = pl[1]; - uint64_t prefix = _prefix_read(pl + 2, pb); - uint8_t is_data = pl[2 + pb]; - const uint8_t* payload = pl + 3 + pb; - size_t paylen = plen - 3 - pb; - - struct ms_session* s = _session_find(peer, ch_id); - if (s && s->timer) { uasync_cancel_timeout(g_ms->inst->ua, s->timer); s->timer = NULL; } - - if (is_data) { - if (paylen < 2) return; - uint16_t count; memcpy(&count, payload, 2); - const uint8_t* mp = payload + 2; size_t mrem = paylen - 2; - for (uint16_t i = 0; i < count && mrem >= 137; i++) { - uint64_t nid; memcpy(&nid, mp, 8); mp += 8; mrem -= 8; - const uint8_t* x25 = mp; mp += 32; mrem -= 32; - const uint8_t* ed = mp; mp += 32; mrem -= 32; - const uint8_t* sig = mp; mp += 64; mrem -= 64; - uint8_t online = *mp++; mrem--; - uint8_t ac = *mp++; mrem--; - const uint8_t* addrs = mp; - int consumed = 0; - for (int a = 0; a < (int)ac && mrem >= (size_t)(1 + consumed); a++) { - uint8_t fam = mp[consumed]; consumed++; - int sz = fam == 4 ? 4 : 16; - consumed += sz + 2; - } - member_sync_put(inst, ch_id, nid, x25, ed, sig, addrs, (int)ac); - mp += consumed; mrem -= (size_t)consumed; - } - } else { - if (!s) { s = u_calloc(1, sizeof(*s)); if (!s) return; - snprintf(s->ch_id, sizeof(s->ch_id), "%s", ch_id); s->peer = peer; - s->next = g_ms->sessions; g_ms->sessions = s; } - s->active = 1; - - if (paylen < 4) return; - uint32_t remote_bm; memcpy(&remote_bm, payload, 4); - const uint8_t* rp = payload + 4; - - uint32_t local_bm; uint8_t lh[MS_BUCKETS][MS_HASH_SIZE]; - member_sync_get_level_hashes(inst, ch_id, level, prefix, pb, &local_bm, lh); - - uint32_t differs = remote_bm ^ local_bm; - for (int i = 0; i < MS_BUCKETS; i++) { - if (!(local_bm & (1u << i)) && !(remote_bm & (1u << i))) continue; - if ((local_bm & (1u << i)) && (remote_bm & (1u << i))) { - const uint8_t* rh_ptr = rp; - int rh_idx = 0; - for (int j = 0; j < i; j++) if (remote_bm & (1u << j)) rh_idx++; - if (memcmp(rh_ptr + (size_t)rh_idx * MS_HASH_SIZE, lh[i], MS_HASH_SIZE) == 0) - differs &= ~(1u << i); - } - } - - if (differs == 0) { s->active = 0; return; } - - int next_shift = 63 - ((int)level + 1) * 5; - struct bucket_entry requests[MS_BUCKETS]; int rcount = 0; - for (int i = 0; i < MS_BUCKETS && rcount < MS_MAX_BATCH; i++) { - if (!(differs & (1u << i))) continue; - uint8_t nl = (uint8_t)(level + 1); - if (nl > MS_MAX_LEVEL) nl = MS_MAX_LEVEL; - requests[rcount].level = nl; - requests[rcount].prefix_bytes = _prefix_bytes(nl); - requests[rcount].prefix = prefix | ((uint64_t)i << (next_shift > 0 ? next_shift : 0)); - rcount++; - } +int member_sync_put(struct UTUN_INSTANCE* inst, const char* ch_id, + uint64_t member_id, const uint8_t* x25519, + const uint8_t* ed25519, const uint8_t* join_sig, + const uint8_t* addrs_data, int addr_count) { + if (!inst || !ch_id) return -1; + sqlite3* db = _db(inst); if (!db) return -1; - if (rcount > 0) { - uint8_t ch_len = (uint8_t)strlen(ch_id); - size_t rs = 1 + 1 + ch_len + 1; - for (int i = 0; i < rcount; i++) rs += 1 + 1 + requests[i].prefix_bytes + 1; - uint8_t* rbuf = u_malloc(rs); - if (rbuf) { - uint8_t* wr = rbuf; - *wr++ = MS_MSG_REQUEST; *wr++ = ch_len; - memcpy(wr, ch_id, ch_len); wr += ch_len; - *wr++ = (uint8_t)rcount; - for (int i = 0; i < rcount; i++) { - *wr++ = requests[i].level; - *wr++ = requests[i].prefix_bytes; - _prefix_write(wr, requests[i].prefix, requests[i].prefix_bytes); - wr += requests[i].prefix_bytes; - int cnt = 0; sqlite3* db = _db(inst); - if (db) { - char peers_tbl[128]; _peers_table(ch_id, peers_tbl, sizeof(peers_tbl)); - char sql[256]; snprintf(sql, sizeof(sql), - "SELECT COUNT(*) FROM \"%s\" WHERE (node_id & " - "(CASE WHEN %d>=0 THEN ~0<<%d ELSE ~0 END)) == %lld", - peers_tbl, 63-(int)requests[i].level*5, 63-(int)requests[i].level*5, - (long long)requests[i].prefix); - sqlite3_stmt* cst = NULL; - if (sqlite3_prepare_v2(db, sql, -1, &cst, NULL) == SQLITE_OK) { - if (sqlite3_step(cst) == SQLITE_ROW) cnt = sqlite3_column_int(cst, 0); - sqlite3_finalize(cst); - } - } - *wr++ = (uint8_t)(cnt < 8 || requests[i].level >= MS_MAX_LEVEL ? 1 : 0); - } - _send_msg(inst, peer, rbuf, (size_t)(wr - rbuf)); - u_free(rbuf); - } - } + sqlite3_stmt* ns = NULL; + sqlite3_prepare_v2(db, + "INSERT OR REPLACE INTO nodes(node_id, name, x25519_pubkey, ed25519_pubkey, online)" + " VALUES(?,COALESCE((SELECT name FROM nodes WHERE node_id=?),''),?,?," + " COALESCE((SELECT online FROM nodes WHERE node_id=?),0))", + -1, &ns, NULL); + if (ns) { + sqlite3_bind_int64(ns, 1, (sqlite3_int64)member_id); + sqlite3_bind_int64(ns, 2, (sqlite3_int64)member_id); + sqlite3_bind_blob(ns, 3, x25519, 32, SQLITE_STATIC); + sqlite3_bind_blob(ns, 4, ed25519, 32, SQLITE_STATIC); + sqlite3_bind_int64(ns, 5, (sqlite3_int64)member_id); + sqlite3_step(ns); sqlite3_finalize(ns); } -} -static void _handle_request(struct UTUN_INSTANCE* inst, uint64_t peer, const char* ch_id, - const uint8_t* pl, size_t plen) { - if (plen < 1) return; - uint8_t count = pl[0]; const uint8_t* bp = pl + 1; size_t off = 0; - - struct bucket_entry buckets[MS_MAX_BATCH]; int bc = 0; - for (uint8_t i = 0; i < count && bc < MS_MAX_BATCH; i++) { - if (off + 2 > plen - 1) break; - uint8_t lvl = bp[off++]; - uint8_t pb = bp[off++]; - if (off + pb > plen - 1) break; - uint64_t pr = _prefix_read(bp + off, pb); off += pb; - buckets[bc].level = lvl; - buckets[bc].prefix_bytes = pb; - buckets[bc].prefix = pr; - bc++; - off++; /* skip full_data */ - } + topo_node_sqlite_member_put(db, ch_id, member_id, join_sig, NULL); - if (bc > 0) _send_batch(inst, peer, ch_id, buckets, bc); -} + if (addrs_data && addr_count > 0) { + sqlite3_exec(db, "BEGIN", NULL, NULL, NULL); + sqlite3_stmt* ds = NULL; + sqlite3_prepare_v2(db, "DELETE FROM node_addresses WHERE node_id=? AND is_nat=0", -1, &ds, NULL); + if (ds) { sqlite3_bind_int64(ds, 1, (sqlite3_int64)member_id); sqlite3_step(ds); sqlite3_finalize(ds); } -static void _handle_batch(struct UTUN_INSTANCE* inst, uint64_t peer, const char* ch_id, - const uint8_t* pl, size_t plen) { - if (plen < 1) return; - uint8_t count = pl[0]; const uint8_t* bp = pl + 1; size_t rem = plen - 1; - - for (uint8_t i = 0; i < count && rem >= 3; i++) { - uint8_t lvl = bp[0]; uint8_t pb = bp[1]; rem -= 2; bp += 2; - if (rem < pb + 1) break; - uint64_t pr = _prefix_read(bp, pb); bp += pb; rem -= pb; - uint8_t is_data = *bp++; rem--; - - if (is_data && rem >= 2) { - uint16_t mc; memcpy(&mc, bp, 2); bp += 2; rem -= 2; - for (uint16_t j = 0; j < mc && rem >= 137; j++) { - uint64_t nid; memcpy(&nid, bp, 8); bp += 8; rem -= 8; - const uint8_t* x25 = bp; bp += 32; rem -= 32; - const uint8_t* ed = bp; bp += 32; rem -= 32; - const uint8_t* sig = bp; bp += 64; rem -= 64; - uint8_t online = *bp++; rem--; - uint8_t ac = *bp++; rem--; - const uint8_t* addrs = bp; - int consumed = 0; - for (int a = 0; a < (int)ac && rem >= (size_t)(1 + consumed); a++) { - uint8_t fam = bp[consumed]; consumed++; - int sz = fam == 4 ? 4 : 16; - consumed += sz + 2; - } - member_sync_put(inst, ch_id, nid, x25, ed, sig, addrs, (int)ac); - bp += consumed; rem -= (size_t)consumed; - } - } else if (!is_data && rem >= 4) { - uint8_t sub_pl[4096]; size_t sub_len = 0; - uint8_t next_lvl = (uint8_t)(lvl < MS_MAX_LEVEL ? lvl + 1 : lvl); - sub_pl[sub_len++] = next_lvl; - sub_pl[sub_len++] = pb; - _prefix_write(sub_pl + sub_len, pr, pb); sub_len += pb; - sub_pl[sub_len++] = 0; - uint32_t bm; memcpy(&bm, bp, 4); bp += 4; rem -= 4; - memcpy(sub_pl + sub_len, &bm, 4); sub_len += 4; - int nh = _popcount_u32(bm); - if (rem >= (size_t)nh * MS_HASH_SIZE) { - memcpy(sub_pl + sub_len, bp, (size_t)nh * MS_HASH_SIZE); - sub_len += (size_t)nh * MS_HASH_SIZE; - bp += (size_t)nh * MS_HASH_SIZE; rem -= (size_t)nh * MS_HASH_SIZE; + sqlite3_stmt* as = NULL; + sqlite3_prepare_v2(db, + "INSERT INTO node_addresses(node_id,family,protocol,address,port,is_nat)" + " VALUES(?,?,1,?,?,0)", -1, &as, NULL); + if (as) { + const uint8_t* p = addrs_data; + for (int i = 0; i < addr_count; i++) { + uint8_t fam = *p++; int ip_sz = fam == 4 ? 4 : 16; + sqlite3_bind_int64(as, 1, (sqlite3_int64)member_id); + sqlite3_bind_int(as, 2, fam); + sqlite3_bind_blob(as, 3, p, ip_sz, SQLITE_STATIC); + p += ip_sz; + uint16_t port = ((uint16_t)p[0] << 8) | p[1]; p += 2; + sqlite3_bind_int(as, 4, port); + sqlite3_step(as); sqlite3_reset(as); } - _handle_hashes(inst, peer, ch_id, sub_pl, sub_len); + sqlite3_finalize(as); } - } -} - -static void _recv_cb(struct ETCP_CONN* conn, struct ll_entry* entry) { - if (!entry || entry->len < 4) { - if (entry) { if (entry->dgram) u_free(entry->dgram); queue_entry_free(entry); } - return; - } - if (!g_ms || !g_ms->initialized) { u_free(entry->dgram); queue_entry_free(entry); return; } - uint64_t peer = conn ? conn->peer_node_id : 0; - const uint8_t* d = entry->dgram; - size_t dlen = entry->len; - - uint8_t ch_len = d[1]; - if (dlen < (size_t)(2 + ch_len + 1)) { u_free(entry->dgram); queue_entry_free(entry); return; } - char ch_id[64]; memcpy(ch_id, d + 2, ch_len); ch_id[ch_len] = '\0'; - uint8_t type = d[2 + ch_len]; - const uint8_t* pl = d + 3 + ch_len; - size_t plen = dlen - 3 - ch_len; - - switch (type) { - case MS_MSG_HASHES: _handle_hashes(g_ms->inst, peer, ch_id, pl, plen); break; - case MS_MSG_REQUEST: _handle_request(g_ms->inst, peer, ch_id, pl, plen); break; - case MS_MSG_BATCH: _handle_batch(g_ms->inst, peer, ch_id, pl, plen); break; + sqlite3_exec(db, "COMMIT", NULL, NULL, NULL); } - u_free(entry->dgram); queue_entry_free(entry); + merkle_sync_recompute_path(inst, ch_id, member_id); + return 0; } -/* ── bg_check ── */ - -static void _bg_timer_cb(void* arg) { - struct member_sync* ms = (struct member_sync*)arg; - if (!ms || !ms->initialized || !ms->inst) return; - sqlite3* db = _db(ms->inst); if (!db) return; - - sqlite3_stmt* cs = NULL; - sqlite3_prepare_v2(db, "SELECT channel_id FROM channels", -1, &cs, NULL); - if (!cs) return; - while (sqlite3_step(cs) == SQLITE_ROW) { - const char* ch_id = (const char*)sqlite3_column_text(cs, 0); - if (!ch_id) continue; - int rc = member_sync_bg_check(ms->inst, ch_id); - if (rc < 0) continue; - break; - } - sqlite3_finalize(cs); - - ms->bg_timer = uasync_set_timeout(ms->inst->ua, (uint32_t)(MS_BG_INTERVAL_MS * 10), - ms, _bg_timer_cb, "ms_bg"); +int member_sync_del(struct UTUN_INSTANCE* inst, const char* ch_id, uint64_t member_id) { + if (!inst || !ch_id) return -1; + sqlite3* db = _db(inst); if (!db) return -1; + topo_node_sqlite_member_del(db, ch_id, member_id); + merkle_sync_recompute_path(inst, ch_id, member_id); + return 0; } -int member_sync_bg_check(struct UTUN_INSTANCE* inst, const char* ch_id) { - sqlite3* db = _db(inst); if (!db || !ch_id) return -1; - +int member_sync_count(struct UTUN_INSTANCE* inst, const char* ch_id) { + sqlite3* db = _db(inst); + if (!db || !ch_id) return 0; + char peers_tbl[128]; _peers_table(ch_id, peers_tbl, sizeof(peers_tbl)); + char sql[256]; snprintf(sql, sizeof(sql), "SELECT COUNT(*) FROM \"%s\"", peers_tbl); sqlite3_stmt* stmt = NULL; - if (sqlite3_prepare_v2(db, - "SELECT prefix64 FROM member_tree_hash WHERE channel_id=? AND level=? LIMIT 1", - -1, &stmt, NULL) != SQLITE_OK) return -1; - sqlite3_bind_text(stmt, 1, ch_id, -1, SQLITE_STATIC); - sqlite3_bind_int(stmt, 2, MS_MAX_LEVEL); - - if (sqlite3_step(stmt) != SQLITE_ROW) { sqlite3_finalize(stmt); return -1; } - uint64_t pref = (uint64_t)sqlite3_column_int64(stmt, 0); + if (sqlite3_prepare_v2(db, sql, -1, &stmt, NULL) != SQLITE_OK) return 0; + int c = 0; + if (sqlite3_step(stmt) == SQLITE_ROW) c = sqlite3_column_int(stmt, 0); sqlite3_finalize(stmt); - - uint8_t old_hash[MS_HASH_SIZE]; - memcpy(old_hash, member_sync_get_hash(inst, ch_id, MS_MAX_LEVEL, pref), MS_HASH_SIZE); - - _recompute_bucket(inst, ch_id, MS_MAX_LEVEL, pref); - - const uint8_t* new_hash = member_sync_get_hash(inst, ch_id, MS_MAX_LEVEL, pref); - if (memcmp(old_hash, new_hash, MS_HASH_SIZE) != 0) { - for (uint8_t lv = (uint8_t)(MS_MAX_LEVEL - 1); lv >= 1; lv--) - _recompute_bucket(inst, ch_id, lv, _level_prefix(pref, lv)); - return 1; - } - return 0; + return c; } -/* called when topo_group stores updated node info (pubkeys/addresses) to DB */ -static void _on_node_updated(struct UTUN_INSTANCE* inst, uint64_t node_id, - const uint8_t* x25519, const uint8_t* ed25519) { - (void)x25519; (void)ed25519; - if (!inst || !g_ms || !g_ms->initialized) return; +void member_sync_set_online(struct UTUN_INSTANCE* inst, uint64_t node_id, int online) { + if (!inst) return; sqlite3* db = _db(inst); if (!db) return; - sqlite3_stmt* cs = NULL; - if (sqlite3_prepare_v2(db, "SELECT channel_id FROM channels", -1, &cs, NULL) != SQLITE_OK) return; - while (sqlite3_step(cs) == SQLITE_ROW) { - const char* ch = (const char*)sqlite3_column_text(cs, 0); - if (!ch) continue; - char peers_tbl[128]; _peers_table(ch, peers_tbl, sizeof(peers_tbl)); + topo_node_sqlite_node_set_online(db, node_id, online); + sqlite3_stmt* stmt = NULL; + if (sqlite3_prepare_v2(db, "SELECT channel_id FROM channels", -1, &stmt, NULL) != SQLITE_OK) return; + while (sqlite3_step(stmt) == SQLITE_ROW) { + const char* ch_id = (const char*)sqlite3_column_text(stmt, 0); + if (!ch_id) continue; + char peers_tbl[128]; _peers_table(ch_id, peers_tbl, sizeof(peers_tbl)); char buf[256]; snprintf(buf, sizeof(buf), "SELECT 1 FROM \"%s\" WHERE node_id=?", peers_tbl); - sqlite3_stmt* ps = NULL; - if (sqlite3_prepare_v2(db, buf, -1, &ps, NULL) == SQLITE_OK) { - sqlite3_bind_int64(ps, 1, (sqlite3_int64)node_id); - if (sqlite3_step(ps) == SQLITE_ROW) _recompute_path(inst, ch, node_id); - sqlite3_finalize(ps); + sqlite3_stmt* cs = NULL; + if (sqlite3_prepare_v2(db, buf, -1, &cs, NULL) == SQLITE_OK) { + sqlite3_bind_int64(cs, 1, (sqlite3_int64)node_id); + if (sqlite3_step(cs) == SQLITE_ROW) + merkle_sync_recompute_path(inst, ch_id, node_id); + sqlite3_finalize(cs); } } - sqlite3_finalize(cs); -} - -int member_sync_init(struct UTUN_INSTANCE* inst) { - if (!inst) return -1; - struct member_sync* ms = u_calloc(1, sizeof(*ms)); - if (!ms) return -1; - ms->inst = inst; - ms->initialized = 1; - g_ms = ms; - - _ensure_table(); - - etcp_router_bind(inst, ETCP_RT_ID_MEMBER_SYNC, _recv_cb); - - topo_groups_set_node_updated_cb(inst->topo_groups, _on_node_updated); - - ms->bg_timer = uasync_set_timeout(inst->ua, (uint32_t)(MS_BG_INTERVAL_MS * 10), - ms, _bg_timer_cb, "ms_bg"); - - DEBUG_INFO(DEBUG_CATEGORY_CONNECTIVITY, "%s: initialized", MS_ID); - return 0; + sqlite3_finalize(stmt); } -void member_sync_destroy(struct UTUN_INSTANCE* inst) { - if (!g_ms || !inst) return; - g_ms->initialized = 0; - - etcp_router_bind(inst, ETCP_RT_ID_MEMBER_SYNC, NULL); - - if (g_ms->bg_timer) { uasync_cancel_timeout(inst->ua, g_ms->bg_timer); g_ms->bg_timer = NULL; } - - struct ms_session* s = g_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(g_ms); g_ms = NULL; +const uint8_t* member_sync_get_hash(struct UTUN_INSTANCE* inst, const char* ch_id, + uint8_t level, uint64_t prefix64) { + return merkle_sync_get_hash(inst, ch_id, level, prefix64); } diff --git a/tools/chatgui/transport/member_sync.h b/tools/chatgui/transport/member_sync.h index 6d489eaf..e9671125 100644 --- a/tools/chatgui/transport/member_sync.h +++ b/tools/chatgui/transport/member_sync.h @@ -1,6 +1,51 @@ +/* + * ── member_sync — синхронизация участников чата ── + * + * Тонкая прослойка над merkle_sync, адаптированная под мемберов каналов. + * Данные хранятся в таблицах: nodes, node_addresses, peers_. + * Хеш мембера = SHA256(node_id || x25519 || ed25519 || join_sig || addrs || online). + * + * ── Использование ── + * + * // разово: инициализация (из chat_sync_init) + * member_sync_init(inst); + * + * // асинхронный запуск синхронизации (из cs_on_conn_up) + * member_sync_start(inst, peer, ch_id, on_done, ch); + * + * // коллбэк: синхронизация завершена + * static void on_done(uint64_t peer, const char* ch_id, int result, void* arg) { + * struct channel_cache* ch = (struct channel_cache*)arg; + * if (result == MT_OK) ch->synced = CS_SYNC_DONE; + * } + * + * // добавление/обновление мембера (из chat_core_create_channel, cs_handle_welcome) + * member_sync_put(inst, ch_id, node_id, x25519, ed25519, join_sig, addrs, ac); + * // дерево автоматически пересчитано + * + * // онлайн-статус (из cs_on_conn_up / cs_on_conn_down) + * member_sync_set_online(inst, node_id, 1); // online + * member_sync_set_online(inst, node_id, 0); // offline + * // пересчитывает дерево для всех каналов, где состоит node_id + * + * // удаление мембера + * member_sync_del(inst, ch_id, node_id); + * + * // отмена синхронизации (коллбэк НЕ вызывается) + * member_sync_cancel(inst, peer, ch_id); + * + * // количество мемберов в канале + * int n = member_sync_count(inst, ch_id); + * + * // разово: завершение (из chat_sync_destroy) + * member_sync_destroy(inst); + */ + #ifndef MEMBER_SYNC_H #define MEMBER_SYNC_H +#include "merkle_sync.h" + #ifdef __cplusplus extern "C" { #endif @@ -10,52 +55,66 @@ extern "C" { struct UTUN_INSTANCE; -#define ETCP_RT_ID_MEMBER_SYNC 0x31 -#define MS_MAX_BATCH 32 -#define MS_MAX_LEVEL 5 -#define MS_BUCKETS 32 -#define MS_HASH_SIZE 32 - -/* Wire message types (etcp_router service 0x31) */ -#define MS_MSG_HASHES 0x01 -#define MS_MSG_REQUEST 0x02 -#define MS_MSG_BATCH 0x03 - -/* === Initialization === */ +/* + * Инициализировать модуль: создаёт merkle_sync с коллбэками для мемберов + * (ETCP сервис 0x31), регистрирует _on_node_updated для пересчёта дерева + * при обновлении node info. + */ int member_sync_init(struct UTUN_INSTANCE* inst); + +/* Завершить модуль. */ void member_sync_destroy(struct UTUN_INSTANCE* inst); -/* === CRUD (DB write + tree recompute) === */ +/* + * Запустить синхронизацию мемберов канала ch_id с пиром peer. + * Делегирует в merkle_sync_start(). + */ +int member_sync_start(struct UTUN_INSTANCE* inst, uint64_t peer, + const char* ch_id, merkle_sync_done_cb done_cb, void* arg); + +/* + * Отменить синхронизацию. Делегирует в merkle_sync_cancel(). + */ +void member_sync_cancel(struct UTUN_INSTANCE* inst, uint64_t peer, const char* ch_id); + +/* + * Добавить/обновить мембера и пересчитать дерево. + * + * ch_id — ID канала + * member_id — node_id участника + * x25519 — X25519 публичный ключ (32 байта) + * ed25519 — Ed25519 публичный ключ (32 байта) + * join_sig — Ed25519 подпись join-сообщения (64 байта) + * сообщение: ch_id || '\0' || node_id(8LE) || x25519(32) || name || '\0' + * addrs_data — адреса в wire-формате: [family:1][ip:4|16][port:2]... + * addr_count — количество адресов (0 если нет) + * + * Пишет в таблицы: nodes (INSERT OR REPLACE), node_addresses, + * peers_ (INSERT OR REPLACE). + * Затем вызывает merkle_sync_recompute_path() для пересчёта дерева. + * + * Возвращает 0 при успехе, -1 при ошибке. + */ int member_sync_put(struct UTUN_INSTANCE* inst, const char* ch_id, - uint64_t member_id, const uint8_t* x25519, - const uint8_t* ed25519, const uint8_t* join_sig, - const uint8_t* addrs_data, int addr_count); + uint64_t member_id, const uint8_t* x25519, + const uint8_t* ed25519, const uint8_t* join_sig, + const uint8_t* addrs_data, int addr_count); + +/* Удалить мембера из канала и пересчитать дерево. */ int member_sync_del(struct UTUN_INSTANCE* inst, const char* ch_id, uint64_t member_id); + +/* Количество мемберов в канале (SELECT COUNT из peers_). */ int member_sync_count(struct UTUN_INSTANCE* inst, const char* ch_id); -/* === Online status (updates nodes.online + recomputes tree path) === */ +/* + * Установить онлайн-статус узла (nodes.online = 0/1). + * Пересчитывает дерево для всех каналов, где состоит node_id. + */ void member_sync_set_online(struct UTUN_INSTANCE* inst, uint64_t node_id, int online); -/* === Protocol primitives === */ -int member_sync_get_level_hashes(struct UTUN_INSTANCE* inst, - const char* ch_id, uint8_t level, - uint64_t prefix, uint8_t prefix_bytes, - uint32_t* bitmap, uint8_t hashes[MS_BUCKETS][MS_HASH_SIZE]); -int member_sync_get_bucket_members(struct UTUN_INSTANCE* inst, - const char* ch_id, uint8_t level, - uint64_t prefix, uint8_t prefix_bytes, - uint8_t* buf, size_t* len); - -/* === Background consistency check === */ -int member_sync_bg_check(struct UTUN_INSTANCE* inst, const char* ch_id); - -/* === Sync sessions === */ -void member_sync_start(struct UTUN_INSTANCE* inst, uint64_t peer, const char* ch_id); -void member_sync_cancel_peer(struct UTUN_INSTANCE* inst, uint64_t peer); - -/* === Test helper === */ +/* Получить хеш бакета из merkle_tree_hash (для тестов). */ const uint8_t* member_sync_get_hash(struct UTUN_INSTANCE* inst, - const char* ch_id, uint8_t level, uint64_t prefix64); + const char* ch_id, uint8_t level, uint64_t prefix64); #ifdef __cplusplus } diff --git a/tools/chatgui/transport/merkle_sync.c b/tools/chatgui/transport/merkle_sync.c new file mode 100644 index 00000000..30f00af2 --- /dev/null +++ b/tools/chatgui/transport/merkle_sync.c @@ -0,0 +1,618 @@ +#include "merkle_sync.h" + +#include "../../../src/utun_instance.h" +#include "../../../src/etcp_router.h" +#include "../../../src/etcp_api.h" +#include "../../../src/etcp.h" +#include "../../../src/topo_group.h" +#include "../../../lib/debug_config.h" +#include "../../../lib/mem.h" +#include "../../../lib/u_async.h" + +#include +#include +#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; + +struct bucket_entry { uint8_t level; uint8_t prefix_bytes; uint64_t prefix; }; + +struct ms_session { + struct ms_session* next; + char ns[64]; + uint64_t peer; + uint8_t active; + void* timer; + uint8_t retries; + merkle_sync_done_cb done_cb; + void* cb_arg; +}; + +struct merkle_sync { + struct UTUN_INSTANCE* inst; + uint8_t svc_id; + const struct merkle_sync_data_ops* ops; + void* data_ctx; + struct ms_session* sessions; + uint8_t initialized; + void* bg_timer; +}; + +/* ── DB access ── */ + +static sqlite3* _db(struct UTUN_INSTANCE* inst) { + return inst && inst->topo_groups ? inst->topo_groups->topo_sqlite_db : NULL; +} + +/* ── Prefix arithmetic ── */ + +uint64_t merkle_sync_level_prefix(uint64_t key, uint8_t level) { + int shift = 63 - (int)level * 5; if (shift < 0) shift = 0; + return (key >> shift) << shift; +} + +uint8_t merkle_sync_prefix_bytes(uint8_t level) { int bits = level * 5; return (uint8_t)((bits + 7) / 8); } + +static void _prefix_write(uint8_t* out, uint64_t prefix, uint8_t pb) { + for (int i = (int)pb - 1; i >= 0; i--) out[pb - 1 - i] = (uint8_t)(prefix >> (i * 8)); +} + +static uint64_t _prefix_read(const uint8_t* in, uint8_t pb) { + uint64_t v = 0; + for (uint8_t i = 0; i < pb && i < 8; i++) v = (v << 8) | in[i]; + return v; +} + +static int _popcount_u32(uint32_t v) { return __builtin_popcount(v); } + +/* ── Ensure merkle_tree_hash table ── */ + +static void _ensure_table(struct merkle_sync* ms) { + sqlite3* db = _db(ms->inst); if (!db) return; + sqlite3_exec(db, + "CREATE TABLE IF NOT EXISTS merkle_tree_hash (" + " namespace TEXT NOT NULL," + " level INTEGER NOT NULL CHECK(level BETWEEN 1 AND 5)," + " prefix64 INTEGER NOT NULL," + " hash BLOB NOT NULL," + " member_count INTEGER NOT NULL," + " PRIMARY KEY (namespace, level, prefix64))", + NULL, NULL, NULL); +} + +/* ── Tree operations ── */ + +static int _recompute_bucket(struct merkle_sync* ms, const char* ns, + uint8_t level, uint64_t prefix64) { + sqlite3* db = _db(ms->inst); if (!db) return -1; + + EVP_MD_CTX* ctx = EVP_MD_CTX_new(); + EVP_DigestInit_ex(ctx, EVP_sha256(), NULL); + int count = ms->ops->update_bucket_hash(ms->data_ctx, ns, level, prefix64, ctx); + + if (count < 0) { EVP_MD_CTX_free(ctx); return -1; } + if (count == 0) { + EVP_MD_CTX_free(ctx); + sqlite3_stmt* ds = NULL; + sqlite3_prepare_v2(db, + "DELETE FROM merkle_tree_hash WHERE namespace=? AND level=? AND prefix64=?", + -1, &ds, NULL); + if (ds) { sqlite3_bind_text(ds, 1, ns, -1, SQLITE_STATIC); sqlite3_bind_int(ds, 2, level); + sqlite3_bind_int64(ds, 3, (sqlite3_int64)prefix64); sqlite3_step(ds); sqlite3_finalize(ds); } + return 0; + } + + uint8_t hash[MT_HASH_SIZE]; + EVP_DigestFinal_ex(ctx, hash, NULL); + EVP_MD_CTX_free(ctx); + + sqlite3_stmt* is = NULL; + sqlite3_prepare_v2(db, + "INSERT OR REPLACE INTO merkle_tree_hash(namespace, level, prefix64, hash, member_count)" + " VALUES(?,?,?,?,?)", -1, &is, NULL); + if (is) { + sqlite3_bind_text(is, 1, ns, -1, SQLITE_STATIC); + sqlite3_bind_int(is, 2, level); + sqlite3_bind_int64(is, 3, (sqlite3_int64)prefix64); + sqlite3_bind_blob(is, 4, hash, MT_HASH_SIZE, SQLITE_STATIC); + sqlite3_bind_int(is, 5, count); + sqlite3_step(is); sqlite3_finalize(is); + } + return 0; +} + +void merkle_sync_recompute_path(struct UTUN_INSTANCE* inst, const char* ns, uint64_t key) { + struct merkle_sync* ms = g_merkle; + if (!ms || !ms->initialized) return; + for (uint8_t level = 1; level <= MT_MAX_LEVEL; level++) + _recompute_bucket(ms, ns, level, merkle_sync_level_prefix(key, level)); +} + +const uint8_t* merkle_sync_get_hash(struct UTUN_INSTANCE* inst, const char* ns, + uint8_t level, uint64_t prefix64) { + static uint8_t zero[MT_HASH_SIZE]; + sqlite3* db = _db(inst); if (!db) { memset(zero, 0, MT_HASH_SIZE); return zero; } + sqlite3_stmt* stmt = NULL; + sqlite3_prepare_v2(db, + "SELECT hash FROM merkle_tree_hash WHERE namespace=? AND level=? AND prefix64=?", + -1, &stmt, NULL); + if (!stmt) { memset(zero, 0, MT_HASH_SIZE); return zero; } + sqlite3_bind_text(stmt, 1, ns, -1, SQLITE_STATIC); + sqlite3_bind_int(stmt, 2, level); + sqlite3_bind_int64(stmt, 3, (sqlite3_int64)prefix64); + memset(zero, 0, MT_HASH_SIZE); + if (sqlite3_step(stmt) == SQLITE_ROW) { + const void* h = sqlite3_column_blob(stmt, 0); + if (h) memcpy(zero, h, MT_HASH_SIZE); + } + sqlite3_finalize(stmt); + return zero; +} + +static int _get_level_hashes(struct merkle_sync* ms, const char* ns, + uint8_t level, uint64_t prefix, uint8_t prefix_bytes, + uint32_t* bitmap, uint8_t hashes[MT_BUCKETS][MT_HASH_SIZE]) { + (void)prefix_bytes; + memset(hashes, 0, sizeof(uint8_t) * MT_BUCKETS * MT_HASH_SIZE); + *bitmap = 0; + if (level >= MT_MAX_LEVEL) return 0; + int next_shift = 63 - ((int)level + 1) * 5; + + for (int i = 0; i < MT_BUCKETS; i++) { + uint64_t child_prefix = prefix | ((uint64_t)i << (next_shift > 0 ? next_shift : 0)); + const uint8_t* h = merkle_sync_get_hash(ms->inst, ns, (uint8_t)(level + 1), child_prefix); + int empty = 1; + for (int j = 0; j < MT_HASH_SIZE; j++) { if (h[j] != 0) { empty = 0; break; } } + if (!empty) { *bitmap |= (1u << i); memcpy(hashes[i], h, MT_HASH_SIZE); } + } + return 0; +} + +/* ── Send helpers ── */ + +static int _send_msg(struct merkle_sync* ms, uint64_t peer, + const uint8_t* payload, size_t len) { + if (!ms->inst || len < 1) return -1; + uint8_t* buf = u_malloc(1 + len); + if (!buf) return -1; + buf[0] = ms->svc_id; + memcpy(buf + 1, payload, len); + struct ll_entry* entry = queue_entry_new(0); + if (!entry) { u_free(buf); return -1; } + entry->dgram = buf; entry->len = 1 + len; + int r = etcp_route_send(ms->inst, peer, entry, 0); + if (r != 0) { u_free(buf); queue_entry_free(entry); } + 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) { + uint8_t ch_len = (uint8_t)strlen(ns); + size_t max_sz = 1 + 1 + ch_len + 1 + 1 + prefix_bytes + 1 + 4 + MT_BUCKETS * MT_HASH_SIZE + 2 + 65536; + uint8_t* buf = u_malloc(max_sz); + if (!buf) return -1; + uint8_t* p = buf; + *p++ = 0x01; /* MS_MSG_HASHES */ + *p++ = ch_len; memcpy(p, ns, ch_len); p += ch_len; + *p++ = level; + *p++ = prefix_bytes; _prefix_write(p, prefix, prefix_bytes); p += prefix_bytes; + *p++ = (uint8_t)(is_data ? 1 : 0); + + if (is_data) { + size_t mlen = 65536; uint8_t* mbuf = u_malloc(mlen); + if (mbuf) { + if (ms->ops->get_items(ms->data_ctx, ns, level, prefix, prefix_bytes, mbuf, &mlen) == 0) + { memcpy(p, mbuf, mlen); p += mlen; } + else { uint16_t zero = 0; memcpy(p, &zero, 2); p += 2; } + u_free(mbuf); + } + } else { + uint32_t bitmap; uint8_t hashes[MT_BUCKETS][MT_HASH_SIZE]; + _get_level_hashes(ms, ns, level, prefix, prefix_bytes, &bitmap, hashes); + memcpy(p, &bitmap, 4); p += 4; + for (int i = 0; i < MT_BUCKETS; i++) + if (bitmap & (1u << i)) { memcpy(p, hashes[i], MT_HASH_SIZE); p += MT_HASH_SIZE; } + } + int r = _send_msg(ms, peer, buf, (size_t)(p - buf)); + u_free(buf); + return r; +} + +static int _send_batch(struct merkle_sync* ms, uint64_t peer, const char* ns, + struct bucket_entry* buckets, int count) { + uint8_t ch_len = (uint8_t)strlen(ns); + size_t max_sz = 1 + 1 + ch_len + 1 + (size_t)count * (1 + 1 + 8 + 1 + 4 + MT_BUCKETS * MT_HASH_SIZE + 65536); + uint8_t* buf = u_malloc(max_sz); + if (!buf) return -1; + uint8_t* p = buf; + *p++ = 0x03; /* MS_MSG_BATCH */ + *p++ = ch_len; memcpy(p, ns, ch_len); p += ch_len; + *p++ = (uint8_t)count; + + for (int i = 0; i < count; i++) { + *p++ = buckets[i].level; + *p++ = buckets[i].prefix_bytes; + _prefix_write(p, buckets[i].prefix, buckets[i].prefix_bytes); p += buckets[i].prefix_bytes; + + EVP_MD_CTX* tctx = EVP_MD_CTX_new(); + EVP_DigestInit_ex(tctx, EVP_sha256(), NULL); + int bcount = ms->ops->update_bucket_hash(ms->data_ctx, ns, + buckets[i].level, buckets[i].prefix, tctx); + if (bcount < 0) bcount = 0; + EVP_MD_CTX_free(tctx); + + int is_terminal = (buckets[i].level >= MT_MAX_LEVEL || bcount < 8); + *p++ = (uint8_t)(is_terminal ? 1 : 0); + + if (is_terminal) { + size_t mlen = 65536; uint8_t* mbuf = u_malloc(mlen); + if (mbuf) { + if (ms->ops->get_items(ms->data_ctx, ns, buckets[i].level, + buckets[i].prefix, buckets[i].prefix_bytes, mbuf, &mlen) == 0) + { memcpy(p, mbuf, mlen); p += mlen; } + else { uint16_t z = 0; memcpy(p, &z, 2); p += 2; } + u_free(mbuf); + } + } else { + uint32_t bm; uint8_t hs[MT_BUCKETS][MT_HASH_SIZE]; + _get_level_hashes(ms, ns, buckets[i].level, + buckets[i].prefix, buckets[i].prefix_bytes, &bm, hs); + memcpy(p, &bm, 4); p += 4; + for (int j = 0; j < MT_BUCKETS; j++) + if (bm & (1u << j)) { memcpy(p, hs[j], MT_HASH_SIZE); p += MT_HASH_SIZE; } + } + } + int r = _send_msg(ms, peer, buf, (size_t)(p - buf)); + u_free(buf); + return r; +} + +/* ── 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; + s->timer = NULL; + s->retries++; + if (s->retries > 3) { + DEBUG_WARN(DEBUG_CATEGORY_CONNECTIVITY, "%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_CONNECTIVITY, "%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; + return NULL; +} + +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 (s->done_cb) s->done_cb(s->peer, s->ns, result, s->cb_arg); +} + +/* ── Recv handlers ── */ + +static void _handle_hashes(struct merkle_sync* ms, uint64_t peer, const char* ns, + const uint8_t* pl, size_t plen) { + if (plen < 3) return; + uint8_t level = pl[0]; uint8_t pb = pl[1]; + uint64_t prefix = _prefix_read(pl + 2, pb); + uint8_t is_data = pl[2 + pb]; + const uint8_t* payload = pl + 3 + pb; + size_t paylen = plen - 3 - pb; + + 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; + ms->ops->apply_items(ms->data_ctx, ns, payload, paylen); + if (s && s->active) _session_done(s, MT_OK); + return; + } + + if (!s) { + 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; + } + s->active = 1; + + if (paylen < 4) return; + uint32_t remote_bm; memcpy(&remote_bm, payload, 4); + const uint8_t* rh = payload + 4; + + uint32_t local_bm; uint8_t lh[MT_BUCKETS][MT_HASH_SIZE]; + _get_level_hashes(ms, ns, level, prefix, pb, &local_bm, lh); + + uint32_t differs = remote_bm ^ local_bm; + for (int i = 0; i < MT_BUCKETS; i++) { + if (!(local_bm & (1u << i)) && !(remote_bm & (1u << i))) continue; + if ((local_bm & (1u << i)) && (remote_bm & (1u << i))) { + /* find offset of remote hash for slot i in the received blob */ + int rh_idx = 0; + for (int j = 0; j < i; j++) if (remote_bm & (1u << j)) rh_idx++; + if (memcmp(rh + (size_t)rh_idx * MT_HASH_SIZE, lh[i], MT_HASH_SIZE) == 0) + differs &= ~(1u << i); + } + } + + if (differs == 0) { _session_done(s, MT_OK); return; } + + int next_shift = 63 - ((int)level + 1) * 5; + struct bucket_entry requests[MT_BUCKETS]; int rcount = 0; + for (int i = 0; i < MT_BUCKETS && rcount < 32; i++) { + if (!(differs & (1u << i))) continue; + uint8_t nl = (uint8_t)(level + 1); if (nl > MT_MAX_LEVEL) nl = MT_MAX_LEVEL; + requests[rcount].level = nl; + requests[rcount].prefix_bytes = merkle_sync_prefix_bytes(nl); + requests[rcount].prefix = prefix | ((uint64_t)i << (next_shift > 0 ? next_shift : 0)); + rcount++; + } + + if (rcount > 0) { + uint8_t ch_len = (uint8_t)strlen(ns); + size_t rs = 1 + 1 + ch_len + 1; + for (int i = 0; i < rcount; i++) rs += 1 + 1 + requests[i].prefix_bytes + 1; + uint8_t* rbuf = u_malloc(rs); + if (rbuf) { + uint8_t* wr = rbuf; + *wr++ = 0x02; /* MS_MSG_REQUEST */ *wr++ = ch_len; + memcpy(wr, ns, ch_len); wr += ch_len; + *wr++ = (uint8_t)rcount; + for (int i = 0; i < rcount; i++) { + *wr++ = requests[i].level; + *wr++ = requests[i].prefix_bytes; + _prefix_write(wr, requests[i].prefix, requests[i].prefix_bytes); + wr += requests[i].prefix_bytes; + EVP_MD_CTX* tctx = EVP_MD_CTX_new(); + EVP_DigestInit_ex(tctx, EVP_sha256(), NULL); + int cnt = ms->ops->update_bucket_hash(ms->data_ctx, ns, + requests[i].level, requests[i].prefix, tctx); + if (cnt < 0) cnt = 0; + EVP_MD_CTX_free(tctx); + *wr++ = (uint8_t)(cnt < 8 || requests[i].level >= MT_MAX_LEVEL ? 1 : 0); + } + _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); + } +} + +static void _handle_request(struct merkle_sync* ms, uint64_t peer, const char* ns, + const uint8_t* pl, size_t plen) { + if (plen < 1) return; + uint8_t count = pl[0]; const uint8_t* bp = pl + 1; size_t off = 0; + + struct bucket_entry buckets[32]; int bc = 0; + for (uint8_t i = 0; i < count && bc < 32; i++) { + if (off + 2 > plen - 1) break; + uint8_t lvl = bp[off++]; uint8_t pb_i = bp[off++]; + if (off + pb_i > plen - 1) break; + uint64_t pr = _prefix_read(bp + off, pb_i); off += pb_i; + buckets[bc].level = lvl; buckets[bc].prefix_bytes = pb_i; buckets[bc].prefix = pr; bc++; + off++; /* skip is_terminal byte */ + } + if (bc > 0) _send_batch(ms, peer, ns, buckets, bc); +} + +static void _handle_batch(struct merkle_sync* ms, uint64_t peer, const char* ns, + const uint8_t* pl, size_t plen) { + if (plen < 1) return; + uint8_t count = pl[0]; const uint8_t* bp = pl + 1; size_t rem = plen - 1; + + for (uint8_t i = 0; i < count && rem >= 3; i++) { + uint8_t lvl = bp[0]; uint8_t pb_i = bp[1]; rem -= 2; bp += 2; + if (rem < pb_i + 1) break; + uint64_t pr = _prefix_read(bp, pb_i); bp += pb_i; rem -= pb_i; + uint8_t is_data = *bp++; rem--; + + if (is_data && rem >= 2) { + uint16_t mc; memcpy(&mc, bp, 2); + size_t item_len = rem; /* pass entire remainder to apply_items */ + ms->ops->apply_items(ms->data_ctx, ns, bp, item_len); + bp += rem; rem = 0; + } else if (!is_data && rem >= 4) { + 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; + _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; + memcpy(sub_pl + sub_len, &bm, 4); sub_len += 4; + int nh = _popcount_u32(bm); + if (rem >= (size_t)nh * MT_HASH_SIZE) { + memcpy(sub_pl + sub_len, bp, (size_t)nh * MT_HASH_SIZE); + sub_len += (size_t)nh * MT_HASH_SIZE; + bp += (size_t)nh * MT_HASH_SIZE; rem -= (size_t)nh * MT_HASH_SIZE; + } + _handle_hashes(ms, peer, ns, sub_pl, sub_len); + } + } +} + +static void _recv_cb(struct ETCP_CONN* conn, struct ll_entry* entry) { + if (!entry || entry->len < 4) { + if (entry) { if (entry->dgram) u_free(entry->dgram); queue_entry_free(entry); } + return; + } + if (!g_merkle || !g_merkle->initialized) { u_free(entry->dgram); queue_entry_free(entry); return; } + uint64_t peer = conn ? conn->peer_node_id : 0; + const uint8_t* d = entry->dgram; + size_t dlen = entry->len; + + uint8_t ch_len = d[1]; + if (dlen < (size_t)(2 + ch_len + 1)) { u_free(entry->dgram); queue_entry_free(entry); return; } + char ns[64]; memcpy(ns, d + 2, ch_len); ns[ch_len] = '\0'; + uint8_t type = d[2 + ch_len]; + const uint8_t* pl = d + 3 + ch_len; + size_t plen = dlen - 3 - ch_len; + + switch (type) { + case 0x01: _handle_hashes(g_merkle, peer, ns, pl, plen); break; + case 0x02: _handle_request(g_merkle, peer, ns, pl, plen); break; + case 0x03: _handle_batch(g_merkle, peer, ns, pl, plen); break; + } + + u_free(entry->dgram); queue_entry_free(entry); +} + +/* ── Background consistency check ── */ + +int merkle_sync_bg_check(struct UTUN_INSTANCE* inst, const char* ns) { + struct merkle_sync* ms = g_merkle; + if (!ms || !ms->initialized) return -1; + sqlite3* db = _db(inst); if (!db || !ns) return -1; + + sqlite3_stmt* stmt = NULL; + if (sqlite3_prepare_v2(db, + "SELECT prefix64 FROM merkle_tree_hash WHERE namespace=? AND level=? LIMIT 1", + -1, &stmt, NULL) != SQLITE_OK) return -1; + sqlite3_bind_text(stmt, 1, ns, -1, SQLITE_STATIC); + sqlite3_bind_int(stmt, 2, MT_MAX_LEVEL); + + if (sqlite3_step(stmt) != SQLITE_ROW) { sqlite3_finalize(stmt); return -1; } + uint64_t pref = (uint64_t)sqlite3_column_int64(stmt, 0); + sqlite3_finalize(stmt); + + uint8_t old_hash[MT_HASH_SIZE]; + memcpy(old_hash, merkle_sync_get_hash(inst, ns, MT_MAX_LEVEL, pref), MT_HASH_SIZE); + + _recompute_bucket(ms, ns, MT_MAX_LEVEL, pref); + + const uint8_t* new_hash = merkle_sync_get_hash(inst, ns, MT_MAX_LEVEL, pref); + if (memcmp(old_hash, new_hash, MT_HASH_SIZE) != 0) { + for (uint8_t lv = (uint8_t)(MT_MAX_LEVEL - 1); lv >= 1; lv--) + _recompute_bucket(ms, ns, lv, merkle_sync_level_prefix(pref, lv)); + return 1; + } + return 0; +} + +static void _bg_timer_cb(void* arg) { + struct merkle_sync* ms = (struct merkle_sync*)arg; + if (!ms || !ms->initialized || !ms->inst) return; + sqlite3* db = _db(ms->inst); if (!db) return; + + sqlite3_stmt* cs = NULL; + sqlite3_prepare_v2(db, "SELECT channel_id FROM channels", -1, &cs, NULL); + if (!cs) return; + while (sqlite3_step(cs) == SQLITE_ROW) { + const char* ch_id = (const char*)sqlite3_column_text(cs, 0); + if (!ch_id) continue; + int rc = merkle_sync_bg_check(ms->inst, ch_id); + if (rc < 0) continue; + break; + } + sqlite3_finalize(cs); + + ms->bg_timer = uasync_set_timeout(ms->inst->ua, (uint32_t)(MS_BG_INTERVAL_MS * 10), + ms, _bg_timer_cb, "ms_bg"); +} + +/* ── Public lifecycle ── */ + +int merkle_sync_init(struct UTUN_INSTANCE* inst, uint8_t svc_id, + const struct merkle_sync_data_ops* ops, void* data_ctx) { + if (!inst || !ops) return -1; + struct merkle_sync* ms = u_calloc(1, sizeof(*ms)); + if (!ms) return -1; + ms->inst = inst; ms->svc_id = svc_id; ms->ops = ops; ms->data_ctx = data_ctx; + ms->initialized = 1; + g_merkle = ms; + + _ensure_table(ms); + + etcp_router_bind(inst, svc_id, _recv_cb); + + ms->bg_timer = uasync_set_timeout(inst->ua, (uint32_t)(MS_BG_INTERVAL_MS * 10), + ms, _bg_timer_cb, "ms_bg"); + + DEBUG_INFO(DEBUG_CATEGORY_CONNECTIVITY, "%s: initialized svc=%02x", MS_ID, svc_id); + return 0; +} + +void merkle_sync_destroy(struct UTUN_INSTANCE* inst) { + struct merkle_sync* ms = g_merkle; + if (!ms || !inst) return; + ms->initialized = 0; g_merkle = NULL; + + etcp_router_bind(inst, ms->svc_id, NULL); + + if (ms->bg_timer) { uasync_cancel_timeout(inst->ua, ms->bg_timer); ms->bg_timer = NULL; } + + 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); +} + +/* ── Public async sync ── */ + +int merkle_sync_start(struct UTUN_INSTANCE* inst, uint64_t peer, + const char* ns, merkle_sync_done_cb done_cb, void* arg) { + struct merkle_sync* ms = g_merkle; + if (!inst || !ns || !ms || !ms->initialized) return -1; + _ensure_table(ms); + + struct ms_session* s = _session_find(ms, peer, ns); + if (!s) { + s = u_calloc(1, sizeof(*s)); + if (!s) return -1; + snprintf(s->ns, sizeof(s->ns), "%s", ns); + s->peer = peer; + s->next = ms->sessions; + ms->sessions = s; + } else { + /* replace callback if session already exists */ + if (s->timer) { uasync_cancel_timeout(inst->ua, s->timer); s->timer = NULL; } + } + s->active = 1; s->retries = 0; + s->done_cb = done_cb; s->cb_arg = arg; + + _send_hashes(ms, peer, ns, 1, 0, 1, 0); + _session_start_timer(s); + return 0; +} + +void merkle_sync_cancel(struct UTUN_INSTANCE* inst, uint64_t peer, const char* ns) { + struct merkle_sync* ms = g_merkle; + if (!ms || !ns) return; + struct ms_session** p = &ms->sessions; + 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; + } + p = &(*p)->next; + } +} diff --git a/tools/chatgui/transport/merkle_sync.h b/tools/chatgui/transport/merkle_sync.h new file mode 100644 index 00000000..972f24c0 --- /dev/null +++ b/tools/chatgui/transport/merkle_sync.h @@ -0,0 +1,250 @@ +#ifndef MERKLE_SYNC_H +#define MERKLE_SYNC_H + +#include +#include +#include + +struct UTUN_INSTANCE; + +#define MT_MAX_LEVEL 5 +#define MT_BUCKETS 32 +#define MT_HASH_SIZE 32 + +/* + * ── Архитектура ── + * + * merkle_sync — универсальный протокол синхронизации на Merkle-деревьях. + * Группирует элементы (items) в префиксное дерево: 5 уровней × 32 бакета, + * SHA256-хеши. Два пира обмениваются хешами уровней, находят различающиеся + * бакеты и передают только их содержимое. + * + * merkle_sync не знает, что такое "элемент" — потребитель (member_sync) + * предоставляет три коллбэка, описывающих модель данных: + * + * update_bucket_hash — перечислить элементы в диапазоне ключей, + * вычислить хеш каждого и подать в SHA256-контекст + * get_items — сериализовать элементы бакета в wire-формат + * apply_items — десериализовать и сохранить полученные элементы + * + * Потребитель (chat_sync/chat_core) работает только с member_sync + * и не видит merkle_sync напрямую: + * + * // ── разово ── + * member_sync_init(inst); + * + * // ── запустил синхронизацию — забыл ── + * member_sync_start(inst, peer, ch_id, on_done, my_ctx); + * + * // ── получил результат ── + * static void on_done(uint64_t peer, const char* ch_id, int result, void* arg) { + * if (result == MT_OK) printf("sync ok\n"); + * if (result == MT_ERR_TIMEOUT) printf("timeout\n"); + * } + * + * // ── изменил данные — дерево само пересчиталось ── + * member_sync_put(inst, ch_id, node_id, x25519, ed25519, join_sig, addrs, ac); + * member_sync_set_online(inst, node_id, 1); + * + * // ── отменил — коллбэк не вызовется ── + * member_sync_cancel(inst, peer, ch_id); + * + * // ── разово ── + * member_sync_destroy(inst); + * + * ── Формат дерева ── + * + * 5-уровневое префиксное дерево над 64-битными ключами (node_id). + * Каждый уровень берёт 5 старших бит ключа: уровень 1 — биты 59-63, + * уровень 5 — биты 39-63. Каждый узел дерева разбивается на 32 бакета + * (по 5 бит = 32 комбинации). + * + * level 1 [0..31] корень: 32 бакета по 5 бит + * level 2 [0..31]...[0..31] каждый — ещё 32 бакета + * ... + * level 5 [0..31].........[0..31] листья: 32^5 = 33M бакетов макс + * + * Хеш бакета = SHA256(хеш_элемента_1 || ... || хеш_элемента_N). + * Бакет считается терминальным (leaf) на уровне 5 или если в нём < 8 элементов. + * + * ── Wire-протокол (сервис ETCP, id задаётся при init) ── + * + * Каждое сообщение: [svc_id:1][ns_len:1][ns:var][type:1][payload:var] + * + * MSG_HASHES (0x01): level,prefix,is_data, [bitmap+hashes | member_data] + * MSG_REQUEST (0x02): count, [level,prefix,is_terminal]* + * MSG_BATCH (0x03): count, [level,prefix,is_terminal,[member_data|bitmap+hashes]]* + * + * Алгоритм: + * A → MSG_HASHES(level=1, bitmap+hashes всех 32 бакетов уровня 2) → B + * B сравнивает со своим деревом, находит различающиеся бакеты + * B → MSG_REQUEST(level=2, prefix=X, is_terminal) → A + * A → MSG_BATCH(данные бакета) → B + * B сохраняет, пересчитывает хеши + * Рекурсивно для подбакетов, пока хеши не совпадут. + * + * ── Сессии ── + * + * Таймаут 10 секунд, 3 ретрая. При таймауте сессия перезапускает + * MSG_HASHES с level=1. После 3 ретраев — done_cb с MT_ERR_TIMEOUT. + */ + +/* ── Data model callbacks ── */ + +struct merkle_sync_data_ops { + /* + * Подать хеши всех элементов в префиксном диапазоне в sha_ctx. + * + * ctx — data_ctx, переданный в merkle_sync_init + * ns — namespace (например channel_id) + * level — уровень дерева (1..5) + * prefix64 — префикс ключа (старшие level*5 бит, остальные нули) + * sha_ctx — уже инициализирован (EVP_DigestInit_ex), потребитель + * делает только EVP_DigestUpdate(sha_ctx, item_hash, 32) + * для каждого элемента + * + * Возвращает количество элементов (0 = бакет пуст, -1 = ошибка). + * + * Пример реализации для мемберов: + * SELECT ... FROM peers_ JOIN nodes + * WHERE (node_id & mask) == prefix64 ORDER BY node_id + * для каждой строки: _compute_member_hash() → EVP_DigestUpdate() + */ + int (*update_bucket_hash)(void* ctx, const char* ns, uint8_t level, + uint64_t prefix64, EVP_MD_CTX* sha_ctx); + + /* + * Сериализовать элементы бакета в wire-формат. + * + * buf — буфер для записи (выделяет merkle_sync, мин. 64K) + * *len — [in] размер буфера, [out] записанный размер + * + * Wire-формат для мемберов: + * [count:2][node_id:8][x25519:32][ed25519:32][join_sig:64][online:1][addr_cnt:1][addrs:var]... + * + * Возвращает 0 при успехе, <0 при ошибке, -2 если буфер мал. + */ + int (*get_items)(void* ctx, const char* ns, uint8_t level, + uint64_t prefix, uint8_t pbytes, uint8_t* buf, size_t* len); + + /* + * Десериализовать и сохранить элементы из wire-формата. + * + * data, len — данные в том же формате, что выдаёт get_items + * + * Вызывается при приёме MSG_HASHES(is_data=1) или MSG_BATCH(is_terminal=1). + * Должна сохранить элементы в БД и вызвать merkle_sync_recompute_path + * для каждого изменённого ключа (через member_sync_put). + * + * Возвращает 0 при успехе, <0 при ошибке. + */ + int (*apply_items)(void* ctx, const char* ns, + const uint8_t* data, size_t len); +}; + +/* + * Коллбэк завершения синхронизации. + * + * peer — node_id пира + * ns — namespace (channel_id) + * result — MT_OK (0) = успех, MT_ERR_TIMEOUT (-1) = исчерпаны ретраи + * arg — пользовательский контекст + * + * При отмене через merkle_sync_cancel() коллбэк НЕ вызывается. + */ + +#define MT_OK 0 +#define MT_ERR_TIMEOUT -1 + +typedef void (*merkle_sync_done_cb)(uint64_t peer, const char* ns, int result, void* arg); + +/* ── Lifecycle ── */ + +/* + * Инициализировать модуль. Создаёт таблицу merkle_tree_hash, биндит + * ETCP-сервис svc_id (напр. 0x31), запускает фоновую проверку (bg_timer). + * Вызывается один раз, обычно из member_sync_init(). + * + * inst — экземпляр uTun (нужен для etcp_route_send, uasync, sqlite3) + * svc_id — ID сервиса на ETCP-роутере + * ops — коллбэки модели данных (update_bucket_hash, get_items, apply_items) + * data_ctx — прозрачный контекст, передаваемый в коллбэки первым аргументом + * + * Возвращает 0 при успехе, -1 при ошибке. + */ +int merkle_sync_init(struct UTUN_INSTANCE* inst, uint8_t svc_id, + const struct merkle_sync_data_ops* ops, void* data_ctx); + +/* + * Завершить модуль. Отвязывает ETCP-сервис, останавливает bg_timer, + * удаляет все активные сессии (коллбэки НЕ вызываются). + */ +void merkle_sync_destroy(struct UTUN_INSTANCE* inst); + +/* ── Async sync ── */ + +/* + * Запустить синхронизацию namespace ns с пиром peer. + * Отправляет MSG_HASHES(level=1) и запускает таймаут 10s. + * + * Если сессия для (peer, ns) уже существует — перезапускает её + * с новым коллбэком (старый коллбэк теряется без вызова). + * + * Возвращает 0 при успехе, -1 при ошибке. + */ +int merkle_sync_start(struct UTUN_INSTANCE* inst, uint64_t peer, + const char* ns, merkle_sync_done_cb done_cb, void* arg); + +/* + * Отменить синхронизацию. Удаляет сессию, done_cb НЕ вызывается. + * Безопасно вызывать, если сессия не существует. + */ +void merkle_sync_cancel(struct UTUN_INSTANCE* inst, uint64_t peer, const char* ns); + +/* ── Recompute tree path after data change ── */ + +/* + * Пересчитать Merkle-дерево для ключа key — все уровни от 1 до 5. + * Вызывается после изменения данных: merkle_sync сам не знает, + * когда данные изменились — потребитель должен вызвать явно. + * member_sync делает это внутри put/del/set_online. + */ +void merkle_sync_recompute_path(struct UTUN_INSTANCE* inst, const char* ns, uint64_t key); + +/* ── Background consistency check ── */ + +/* + * Проверить консистентность одного случайного level-5 бакета в ns. + * Пересчитывает хеш из данных, сравнивает с сохранённым в merkle_tree_hash. + * При расхождении — пересчитывает весь путь вверх до корня. + * Возвращает 1 если были изменения, 0 если совпало, -1 ошибка. + */ +int merkle_sync_bg_check(struct UTUN_INSTANCE* inst, const char* ns); + +/* ── Test helper ── */ + +/* + * Получить сохранённый хеш бакета из merkle_tree_hash. + * Возвращает указатель на статический буфер (32 байта), перезаписывается + * при следующем вызове. Если бакет не найден — возвращает нулевой хеш. + */ +const uint8_t* merkle_sync_get_hash(struct UTUN_INSTANCE* inst, + const char* ns, uint8_t level, uint64_t prefix64); + +/* ── Prefix arithmetic (pure, exported for convenience) ── */ + +/* + * Выделить level*5 старших бит из 64-битного ключа. + * Уровень 1: биты 59-63 (сдвиг 59) + * Уровень 5: биты 39-63 (сдвиг 39) + * Пример: merkle_sync_level_prefix(0x1234567890ABCDEF, 1) = 0x1000000000000000 + */ +uint64_t merkle_sync_level_prefix(uint64_t key, uint8_t level); + +/* + * Количество байт, необходимое для хранения префикса уровня level. + * Уровень 1: 1 байт (5 бит), уровень 4: 3 байта (20 бит), уровень 5: 4 байта (25 бит). + */ +uint8_t merkle_sync_prefix_bytes(uint8_t level); + +#endif /* MERKLE_SYNC_H */