#include "merkle_sync.h" #include "../../../src/utun_instance.h" #include "etcp_api.h" #include "etcp.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; uint8_t synced; 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_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) { 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) { DEBUG_TRACE(DEBUG_CATEGORY_DB_SYNC, "%s: ns=%s L%d/P%016llx", MS_ID, ns, level, (unsigned long long)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) { DEBUG_TRACE(DEBUG_CATEGORY_DB_SYNC, "%s: → error", MS_ID); EVP_MD_CTX_free(ctx); return -1; } if (count == 0) { DEBUG_TRACE(DEBUG_CATEGORY_DB_SYNC, "%s: → empty, deleted", MS_ID); 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); DEBUG_TRACE(DEBUG_CATEGORY_DB_SYNC, "%s: → saved, count=%d", MS_ID, count); 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; DEBUG_TRACE(DEBUG_CATEGORY_DB_SYNC, "%s: ns=%s key=%016llx", MS_ID, ns, (unsigned long long)key); 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; 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)); sqlite3_stmt* stmt = NULL; if (sqlite3_prepare_v2(db, "SELECT prefix64, hash FROM merkle_tree_hash" " WHERE namespace=? AND level=? AND prefix64>=? AND prefix64<=?" " ORDER BY prefix64", -1, &stmt, NULL) != SQLITE_OK || !stmt) { DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "%s: _get_level_hashes SQL error ns=%s L%d — %s", MS_ID, ns, level, sqlite3_errmsg(db)); return -1; } sqlite3_bind_text(stmt, 1, ns, -1, SQLITE_STATIC); sqlite3_bind_int(stmt, 2, (int)(level + 1)); sqlite3_bind_int64(stmt, 3, (sqlite3_int64)prefix); sqlite3_bind_int64(stmt, 4, (sqlite3_int64)range_end); while (sqlite3_step(stmt) == SQLITE_ROW) { sqlite3_int64 row_pfx = sqlite3_column_int64(stmt, 0); const void* h = sqlite3_column_blob(stmt, 1); if (!h || sqlite3_column_bytes(stmt, 1) < MT_HASH_SIZE) continue; int bucket = (int)(((uint64_t)row_pfx >> next_shift) & 0x1F); *bitmap |= (1u << bucket); memcpy(hashes[bucket], h, MT_HASH_SIZE); } sqlite3_finalize(stmt); DEBUG_DEBUG(DEBUG_CATEGORY_DB_SYNC, "%s: _get_level_hashes ns=%s L%d/P%016llx bitmap=%08x", MS_ID, ns, level, (unsigned long long)prefix, *bitmap); return 0; } /* ── Send helpers ── */ static struct ETCP_CONN* ms_find_conn_for_node(struct UTUN_INSTANCE* inst, uint64_t node_id) { if (!inst->connections) { DEBUG_DEBUG(DEBUG_CATEGORY_DB_SYNC, "%s: find_conn — inst->connections is NULL", MS_ID); return NULL; } struct ll_entry* e = queue_find_data_by_index(inst->connections, (const uint8_t*)&node_id); if (!e) { DEBUG_DEBUG(DEBUG_CATEGORY_DB_SYNC, "%s: find_conn — node=%016llx NOT FOUND in connections", MS_ID, (unsigned long long)node_id); return NULL; } struct conn_queue_entry* ce = (struct conn_queue_entry*)e->data; DEBUG_DEBUG(DEBUG_CATEGORY_DB_SYNC, "%s: find_conn — node=%016llx found=%d initialized=%d links_up=%d", MS_ID, (unsigned long long)node_id, 1, ce->conn->initialized, ce->conn->links_up); if (ce->conn->initialized && ce->conn->links_up) return ce->conn; return NULL; } 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; struct ETCP_CONN* conn = ms_find_conn_for_node(ms->inst, peer); if (!conn) { DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "%s: no conn for node %016llx", MS_ID, (unsigned long long)peer); u_free(buf); queue_entry_free(entry); return -1; } int r = etcp_send(conn, entry); if (r != 0) { u_free(buf); queue_entry_free(entry); } DEBUG_DEBUG(DEBUG_CATEGORY_DB_SYNC, "%s: send OK peer=%016llx len=%zu rc=%d", MS_ID, (unsigned long long)peer, len, r); 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); 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++ = ch_len; memcpy(p, ns, ch_len); p += ch_len; *p++ = 0x01; /* MS_MSG_HASHES */ *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); DEBUG_DEBUG(DEBUG_CATEGORY_DB_SYNC, "%s: send_hashes peer=%016llx ns=%s level=%d prefix=%016llx bitmap=%08x", MS_ID, (unsigned long long)peer, ns, level, (unsigned long long)prefix, bitmap); 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) { DEBUG_TRACE(DEBUG_CATEGORY_DB_SYNC, "%s: send_batch peer=%016llx ns=%s count=%d", MS_ID, (unsigned long long)peer, ns, 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++ = ch_len; memcpy(p, ns, ch_len); p += ch_len; *p++ = 0x03; /* MS_MSG_BATCH */ *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; 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; 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 (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); } /* ── 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; 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; 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))) { 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); } } DEBUG_DEBUG(DEBUG_CATEGORY_DB_SYNC, "%s: handle_hashes peer=%016llx ns=%s level=%d remote_bm=%08x local_bm=%08x differs=%08x is_data=%d", MS_ID, (unsigned long long)peer, ns, level, remote_bm, local_bm, differs, is_data); 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++ = ch_len; memcpy(wr, ns, ch_len); wr += ch_len; *wr++ = 0x02; /* MS_MSG_REQUEST */ *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; DEBUG_TRACE(DEBUG_CATEGORY_DB_SYNC, "%s: handle_request peer=%016llx ns=%s count=%d", MS_ID, (unsigned long long)peer, ns, count); DEBUG_DEBUG(DEBUG_CATEGORY_DB_SYNC, "%s: handle_request peer=%016llx ns=%s count=%d", MS_ID, (unsigned long long)peer, ns, count); 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); 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); } } } 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; DEBUG_TRACE(DEBUG_CATEGORY_DB_SYNC, "%s: handle_batch peer=%016llx ns=%s count=%d len=%zu", MS_ID, (unsigned long long)peer, ns, count, plen - 1); DEBUG_DEBUG(DEBUG_CATEGORY_DB_SYNC, "%s: handle_batch peer=%016llx ns=%s count=%d len=%zu", MS_ID, (unsigned long long)peer, ns, count, 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]; DEBUG_TRACE(DEBUG_CATEGORY_DB_SYNC, "%s: recv type=%02x from=%016llx ns=%s len=%zu", MS_ID, type, (unsigned long long)peer, ns, dlen); 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; case 0x04: if (plen >= 9 && g_merkle->ops->apply_update) { uint64_t key; memcpy(&key, pl, 8); uint8_t utype = pl[8]; DEBUG_DEBUG(DEBUG_CATEGORY_DB_SYNC, "%s: recv ITEM_UPDATE key=%016llx type=%02x from=%016llx ns=%s", MS_ID, (unsigned long long)key, utype, (unsigned long long)peer, ns); g_merkle->ops->apply_update(g_merkle->data_ctx, ns, key, utype, pl + 9, plen - 9); } break; } u_free(entry->dgram); queue_entry_free(entry); } /* ── Push lightweight update to synced peers ── */ 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) { struct merkle_sync* ms = g_merkle; if (!ms || !ms->initialized || !ns || !data) return; DEBUG_DEBUG(DEBUG_CATEGORY_DB_SYNC, "%s: push_update ns=%s key=%016llx type=%02x len=%zu", MS_ID, ns, (unsigned long long)key, type, len); uint8_t ch_len = (uint8_t)strlen(ns); if (ch_len > 63) return; size_t pkt = 1 + 1 + ch_len + 1 + 8 + 1 + len; uint8_t* buf = u_malloc(pkt); if (!buf) return; uint8_t* p = buf; *p++ = ms->svc_id; *p++ = ch_len; memcpy(p, ns, ch_len); p += ch_len; *p++ = 0x04; /* MSG_ITEM_UPDATE */ memcpy(p, &key, 8); p += 8; *p++ = type; memcpy(p, data, len); for (struct ms_session* s = ms->sessions; s; s = s->next) { if (!s->synced || strcmp(s->ns, ns) != 0) continue; struct ll_entry* entry = queue_entry_new(0); if (!entry) continue; uint8_t* dcopy = u_malloc(pkt); if (!dcopy) { queue_entry_free(entry); continue; } memcpy(dcopy, buf, pkt); entry->dgram = dcopy; entry->len = pkt; struct ETCP_CONN* conn = ms_find_conn_for_node(inst, s->peer); if (!conn) { u_free(dcopy); queue_entry_free(entry); continue; } int r = etcp_send(conn, entry); DEBUG_DEBUG(DEBUG_CATEGORY_DB_SYNC, "%s: push_update sent type=%02x to=%016llx rc=%d", MS_ID, type, (unsigned long long)s->peer, r); if (r != 0) { u_free(dcopy); queue_entry_free(entry); } } u_free(buf); } /* ── 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; DEBUG_TRACE(DEBUG_CATEGORY_DB_SYNC, "%s: bg_check ns=%s → %s", MS_ID, ch_id, rc == 0 ? "consistent" : "RECALCULATED"); 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; DEBUG_TRACE(DEBUG_CATEGORY_DB_SYNC, "%s: init svc=%02x", MS_ID, svc_id); 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_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_DB_SYNC, "%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_unbind(inst, ms->svc_id); 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); DEBUG_TRACE(DEBUG_CATEGORY_DB_SYNC, "%s: start peer=%016llx ns=%s new=%d", MS_ID, (unsigned long long)peer, ns, s ? 0 : 1); 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->synced = 0; } 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; DEBUG_TRACE(DEBUG_CATEGORY_DB_SYNC, "%s: cancel peer=%016llx ns=%s", MS_ID, (unsigned long long)peer, ns); 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; } }