You can not select more than 25 topics Topics must start with a letter or number, can include dashes ('-') and can be up to 35 characters long.
 
 
 
 
 
 

790 lines
33 KiB

#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 <string.h>
#include <stdlib.h>
#include <sqlite3.h>
#define MS_ID "merkle_sync"
#define MS_BG_INTERVAL_MS 100
#define MS_PENDING_MAX 128
struct bucket_entry { uint8_t level; uint8_t prefix_bytes; uint64_t prefix; };
enum { SESS_SYNCING, SESS_SYNCED };
enum { MS_PEND_WAITING = 0, MS_PEND_RESOLVED = 1 };
struct ms_pending {
uint8_t level;
uint64_t prefix;
int8_t parent;
int sub_count;
uint8_t state;
};
struct ms_session {
struct ms_session* next;
char ns[64];
uint64_t peer;
uint8_t sess_state;
merkle_sync_done_cb done_cb;
void* cb_arg;
struct ms_pending pending[MS_PENDING_MAX];
int pending_count;
};
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;
}
/* ── Prefix arithmetic ── */
uint64_t merkle_sync_level_prefix(uint64_t key, uint8_t level) {
if (level == 0) return 0;
int shift = 64 - (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 = 0; i < (int)pb; i++)
out[i] = (uint8_t)(prefix >> (8 * (7 - i)));
}
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 << (8 * (8 - (int)pb));
}
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 = inst->msync;
if (!ms || !ms->initialized) { DEBUG_WARN(DEBUG_CATEGORY_DB_SYNC, "%s: recompute_path called but msync not initialized", MS_ID); 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 next_shift = 64 - ((int)level + 1) * 5; if (next_shift < 0) next_shift = 0;
sqlite3_stmt* stmt = NULL;
int query_level = (int)(level + 1);
const char* sql = (level == 0)
? "SELECT prefix64, hash FROM merkle_tree_hash"
" WHERE namespace=? AND level=? ORDER BY prefix64"
: "SELECT prefix64, hash FROM merkle_tree_hash"
" WHERE namespace=? AND level=? AND prefix64>=? AND prefix64<=?"
" ORDER BY prefix64";
if (sqlite3_prepare_v2(db, sql, -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, query_level);
if (level > 0) {
uint64_t range_end = prefix | (0x1FULL << (next_shift > 0 ? next_shift : 0));
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_broadcast_data(struct merkle_sync* ms, uint64_t peer, const char* ns,
const uint8_t* item_data, size_t item_len) {
uint8_t ch_len = (uint8_t)strlen(ns);
size_t sz = 1 + 1 + ch_len + 1 + 1 + 1 + 2 + item_len;
uint8_t* buf = u_malloc(sz); if (!buf) return -1;
uint8_t* p = buf;
*p++ = ch_len; memcpy(p, ns, ch_len); p += ch_len;
*p++ = 0x01; /* MSG_HASHES */
*p++ = 0; /* level=0 */
*p++ = 0; /* prefix_bytes=0 */
*p++ = 1; /* is_data=1 */
uint16_t dlen = (uint16_t)item_len; memcpy(p, &dlen, 2); p += 2;
memcpy(p, item_data, item_len); p += item_len;
int r = _send_msg(ms, peer, buf, (size_t)(p - buf));
u_free(buf);
return r;
}
static int _send_hashes(struct merkle_sync* ms, uint64_t peer, const char* ns,
uint8_t level, uint64_t prefix, uint8_t prefix_bytes, int is_data) {
DEBUG_TRACE(DEBUG_CATEGORY_DB_SYNC, "%s: send_hashes peer=%016llx ns=%s L%d P%016llx is_data=%d", MS_ID, (unsigned long long)peer, ns, level, (unsigned long long)prefix, is_data);
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) {
uint16_t dlen = (uint16_t)mlen; memcpy(p, &dlen, 2); p += 2;
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) {
uint16_t dlen = (uint16_t)mlen; memcpy(p, &dlen, 2); p += 2;
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 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 int _find_pending(struct ms_session* s, uint8_t level, uint64_t prefix) {
for (int i = 0; i < s->pending_count; i++)
if (s->pending[i].level == level && s->pending[i].prefix == prefix)
return i;
return -1;
}
static void _session_done(struct ms_session* s, int result) {
if (s->sess_state == SESS_SYNCED) { DEBUG_WARN(DEBUG_CATEGORY_DB_SYNC, "%s: session ALREADY DONE peer=%016llx ns=%s", MS_ID, (unsigned long long)s->peer, s->ns); return; }
merkle_sync_done_cb cb = s->done_cb; void* arg = s->cb_arg;
s->done_cb = NULL; s->cb_arg = NULL;
if (result == MT_OK) { s->sess_state = SESS_SYNCED; DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "%s: session SYNCED peer=%016llx ns=%s", MS_ID, (unsigned long long)s->peer, s->ns); }
else { s->sess_state = SESS_SYNCING; DEBUG_WARN(DEBUG_CATEGORY_DB_SYNC, "%s: session FAILED peer=%016llx ns=%s result=%d", MS_ID, (unsigned long long)s->peer, s->ns, result); }
if (cb) cb(s->peer, s->ns, result, arg);
}
/* ── Recv handlers ── */
static void _handle_hashes(struct merkle_sync* ms, uint64_t peer, const char* ns,
const uint8_t* pl, size_t plen, int parent_idx) {
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 (is_data) {
if (paylen < 2) return;
uint16_t dlen; memcpy(&dlen, payload, 2);
const uint8_t* pd = payload + 2;
if (paylen - 2 < dlen) return;
ms->ops->apply_items(ms->data_ctx, ns, peer, pd, dlen);
if (s && s->sess_state == SESS_SYNCING) _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;
} else {
s->pending_count = 0;
}
s->sess_state = SESS_SYNCING;
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 && level == 0) {
_send_hashes(ms, peer, ns, 0, 0, 0, 1);
_session_done(s, MT_OK);
return;
}
if (differs == 0) return;
int next_shift = 64 - ((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) {
if (s) {
for (int i = 0; i < rcount && s->pending_count < MS_PENDING_MAX; i++) {
int idx = s->pending_count++;
s->pending[idx].level = requests[i].level;
s->pending[idx].prefix = requests[i].prefix;
s->pending[idx].parent = (int8_t)parent_idx;
s->pending[idx].sub_count = 0;
s->pending[idx].state = MS_PEND_WAITING;
}
if (parent_idx >= 0 && parent_idx < MS_PENDING_MAX)
s->pending[parent_idx].sub_count++;
}
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);
}
} 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);
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;
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);
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 >= 4) {
uint16_t dlen; memcpy(&dlen, bp, 2); bp += 2; rem -= 2;
if (rem < dlen) break;
ms->ops->apply_items(ms->data_ctx, ns, peer, bp, dlen);
bp += dlen; rem -= dlen;
struct ms_session* sb = _session_find(ms, peer, ns);
if (sb) { int pidx = _find_pending(sb, lvl, pr); if (pidx >= 0) sb->pending[pidx].state = MS_PEND_RESOLVED; }
} else if (!is_data && rem >= 4) {
struct ms_session* sc = _session_find(ms, peer, ns);
int pidx = sc ? _find_pending(sc, lvl, pr) : -1;
if (pidx >= 0) sc->pending[pidx].state = MS_PEND_RESOLVED;
uint8_t sub_pl[4096]; size_t sub_len = 0;
sub_pl[sub_len++] = lvl;
sub_pl[sub_len++] = merkle_sync_prefix_bytes(lvl);
_prefix_write(sub_pl + sub_len, pr, merkle_sync_prefix_bytes(lvl));
sub_len += merkle_sync_prefix_bytes(lvl);
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, pidx);
}
}
{
struct ms_session* s = _session_find(ms, peer, ns);
if (!s) return;
int all_done = (s->pending_count > 0);
for (int i = 0; i < s->pending_count; i++) {
if (s->pending[i].state != MS_PEND_RESOLVED) { all_done = 0; break; }
}
if (all_done && s->sess_state == SESS_SYNCING) {
_send_hashes(ms, peer, ns, 0, 0, 0, 1);
_session_done(s, MT_OK);
}
}
}
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;
}
struct UTUN_INSTANCE* inst = conn ? conn->instance : NULL;
struct merkle_sync* ms = inst ? inst->msync : NULL;
if (!ms || !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 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(ms, peer, ns, pl, plen, -1); break;
case 0x02: _handle_request(ms, peer, ns, pl, plen); break;
case 0x03: _handle_batch(ms, peer, ns, pl, plen); break;
case 0x04:
if (plen >= 9 && ms->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);
ms->ops->apply_update(ms->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 = inst ? inst->msync : NULL;
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->sess_state != SESS_SYNCED || strcmp(s->ns, ns) != 0) continue;
struct ll_entry* entry = queue_entry_new(0);
if (!entry) continue;
uint8_t* dcopy = u_malloc(pkt);
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);
}
/* ── Broadcast to synced peers ── */
void merkle_sync_broadcast(struct UTUN_INSTANCE* inst, const char* ns,
uint64_t from_peer, const uint8_t* data, size_t len) {
struct merkle_sync* ms = inst ? inst->msync : NULL;
if (!ms || !ms->initialized || !ns || !data || len < 2) return;
DEBUG_DEBUG(DEBUG_CATEGORY_DB_SYNC, "%s: broadcast ns=%s from=%016llx len=%zu",
MS_ID, ns, (unsigned long long)from_peer, len);
for (struct ms_session* s = ms->sessions; s; s = s->next) {
if (s->sess_state != SESS_SYNCED || s->peer == from_peer || strcmp(s->ns, ns) != 0) continue;
_send_broadcast_data(ms, s->peer, ns, data, len);
}
}
/* ── Background consistency check ── */
int merkle_sync_bg_check(struct UTUN_INSTANCE* inst, const char* ns) {
struct merkle_sync* ms = inst ? inst->msync : NULL;
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;
inst->msync = 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) {
if (!inst) return;
struct merkle_sync* ms = inst->msync;
if (!ms) return;
ms->initialized = 0; inst->msync = 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;
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 = inst ? inst->msync : NULL;
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 {
s->pending_count = 0;
DEBUG_DEBUG(DEBUG_CATEGORY_DB_SYNC, "%s: start OVERWRITE peer=%016llx ns=%s", MS_ID, (unsigned long long)peer, ns);
}
s->sess_state = SESS_SYNCING;
s->done_cb = done_cb; s->cb_arg = arg;
_send_hashes(ms, peer, ns, 0, 0, 0, 0);
return 0;
}
void merkle_sync_cancel(struct UTUN_INSTANCE* inst, uint64_t peer, const char* ns) {
struct merkle_sync* ms = inst ? inst->msync : NULL;
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) {
*p = s->next; u_free(s);
return;
}
p = &(*p)->next;
}
}