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.
 
 
 
 
 
 

621 lines
30 KiB

#include "member_sync.h"
#include "../routing_layer/topo_node_sqlite.h"
#include "chat_core.h"
#include "../utun_instance.h"
#include "../../lib/debug_config.h"
#include "../../lib/mem.h"
#include <string.h>
#include <stdlib.h>
#include <sqlite3.h>
#include <openssl/evp.h>
#define MS_ID "member_sync"
struct addr_item { uint8_t family; uint8_t addr[16]; uint16_t port; };
/* ── DB access ── */
static sqlite3* _db(struct UTUN_INSTANCE* inst) {
return inst ? inst->topo_sqlite_db : NULL;
}
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);
}
/* ── 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;
const struct addr_item* ib = (const struct addr_item*)b;
if (ia->family != ib->family) return (int)ia->family - (int)ib->family;
int r = memcmp(ia->addr, ib->addr, (size_t)(ia->family == 4 ? 4 : 16));
if (r) return r;
return (int)ia->port - (int)ib->port;
}
static void _compute_member_hash(uint64_t node_id, const uint8_t* x25519,
const uint8_t* ed25519,
const uint8_t* join_sig, uint64_t join_ts,
const uint8_t* update_sig, uint64_t update_ts,
const uint8_t* addrs_data, int addr_count,
const char* adm_tags, const uint8_t* adm_tags_sig,
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);
EVP_DigestUpdate(ctx, x25519, 32);
EVP_DigestUpdate(ctx, ed25519, 32);
EVP_DigestUpdate(ctx, join_sig, 64);
EVP_DigestUpdate(ctx, &join_ts, 8);
if (update_sig) EVP_DigestUpdate(ctx, update_sig, 64); else { static const uint8_t z[64]; EVP_DigestUpdate(ctx, z, 64); }
EVP_DigestUpdate(ctx, &update_ts, 8);
uint8_t ac = (uint8_t)addr_count;
EVP_DigestUpdate(ctx, &ac, 1);
if (addr_count > 0 && addrs_data) {
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;
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;
n++;
}
qsort(items, (size_t)n, sizeof(struct addr_item), _addr_cmp);
for (int i = 0; i < n; i++) {
EVP_DigestUpdate(ctx, &items[i].family, 1);
EVP_DigestUpdate(ctx, items[i].addr, (size_t)(items[i].family == 4 ? 4 : 16));
uint8_t port_be[2] = { (uint8_t)(items[i].port >> 8), (uint8_t)(items[i].port & 0xFF) };
EVP_DigestUpdate(ctx, port_be, 2);
}
}
uint8_t atl = adm_tags ? (uint8_t)strnlen(adm_tags, 255) : 0;
EVP_DigestUpdate(ctx, &atl, 1);
if (atl) EVP_DigestUpdate(ctx, adm_tags, atl);
if (adm_tags_sig) EVP_DigestUpdate(ctx, adm_tags_sig, 64); else { static const uint8_t z64[64]; EVP_DigestUpdate(ctx, z64, 64); }
EVP_DigestFinal_ex(ctx, hash_out, NULL);
EVP_MD_CTX_free(ctx);
}
/* ── merkle_sync_data_ops implementation ── */
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) { DEBUG_ERROR(DEBUG_CATEGORY_MEMBER_SYNC, "%s: bucket_hash — db is NULL", MS_ID); return -1; }
if (level == 0) return 0;
DEBUG_TRACE(DEBUG_CATEGORY_MEMBER_SYNC, "%s: bucket_hash ns=%s L%d/P%016llx", MS_ID, ns, level, (unsigned long long)prefix64);
char peers_tbl[128]; _peers_table(ns, peers_tbl, sizeof(peers_tbl));
int mask_shift = 64 - (int)level * 5;
uint64_t mask = (mask_shift >= 0 && mask_shift < 64) ? (~0ULL << mask_shift) : UINT64_MAX;
char sql[512]; snprintf(sql, sizeof(sql),
"SELECT node_id, x25519_pubkey, ed25519_pubkey,"
" join_sig, join_ts, update_sig, update_ts, adm_tags, adm_tags_sig"
" FROM \"%s\""
" WHERE (node_id & %lld) == %lld ORDER BY node_id ASC",
peers_tbl, (long long)mask, (long long)prefix64);
sqlite3_stmt* stmt = NULL;
if (sqlite3_prepare_v2(db, sql, -1, &stmt, NULL) != SQLITE_OK) { DEBUG_ERROR(DEBUG_CATEGORY_MEMBER_SYNC, "%s: bucket_hash prepare failed ns=%s", MS_ID, ns); return -1; }
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);
const uint8_t* ed = (const uint8_t*)sqlite3_column_blob(stmt, 2);
const uint8_t* sig = (const uint8_t*)sqlite3_column_blob(stmt, 3);
uint64_t jts = (uint64_t)sqlite3_column_int64(stmt, 4);
const uint8_t* usig = (const uint8_t*)sqlite3_column_blob(stmt, 5);
uint64_t uts = (uint64_t)sqlite3_column_int64(stmt, 6);
const char* atags = (const char*)sqlite3_column_text(stmt, 7);
const uint8_t* atsig = (const uint8_t*)sqlite3_column_blob(stmt, 8);
if (!x25 || !ed) continue;
uint8_t addrs[2048]; int addr_off = 0; int addr_count = 0;
sqlite3_stmt* as = NULL;
sqlite3_prepare_v2(db,
"SELECT family, socket_id, protocol, address, port FROM node_addresses WHERE node_id=? AND addr_type=0"
" ORDER BY family, address, port", -1, &as, NULL);
if (as) {
sqlite3_bind_int64(as, 1, (sqlite3_int64)nid);
while (sqlite3_step(as) == SQLITE_ROW && addr_off < (int)sizeof(addrs) - 9) {
int fam = sqlite3_column_int(as, 0);
int sid = sqlite3_column_int(as, 1);
int proto = sqlite3_column_int(as, 2);
int ip_sz = fam == 4 ? 4 : 16;
if (addr_off + 3 + ip_sz + 2 > (int)sizeof(addrs)) break;
addrs[addr_off++] = (uint8_t)fam;
addrs[addr_off++] = (uint8_t)sid;
addrs[addr_off++] = (uint8_t)proto;
memcpy(addrs + addr_off, sqlite3_column_blob(as, 3), (size_t)ip_sz);
addr_off += ip_sz;
uint16_t p = (uint16_t)sqlite3_column_int(as, 4);
addrs[addr_off++] = (uint8_t)(p >> 8);
addrs[addr_off++] = (uint8_t)(p & 0xFF);
addr_count++;
}
sqlite3_finalize(as);
}
uint8_t mh[MT_HASH_SIZE];
_compute_member_hash(nid, x25, ed, sig, jts, usig, uts, addrs, addr_count, atags, atsig, mh);
EVP_DigestUpdate(sha_ctx, mh, MT_HASH_SIZE);
count++;
}
sqlite3_finalize(stmt);
return count;
}
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) { DEBUG_ERROR(DEBUG_CATEGORY_MEMBER_SYNC, "%s: get_items — invalid args db=%p buf=%p len=%p", MS_ID, (void*)db, (void*)buf, (void*)len); return -1; }
if (level == 0) { *len = 0; return 0; }
DEBUG_TRACE(DEBUG_CATEGORY_MEMBER_SYNC, "%s: get_items ns=%s L%d/P%016llx", MS_ID, ns, level, (unsigned long long)prefix);
char peers_tbl[128]; _peers_table(ns, peers_tbl, sizeof(peers_tbl));
int mask_shift = 64 - (int)level * 5;
uint64_t mask = (mask_shift >= 0 && mask_shift < 64) ? (~0ULL << mask_shift) : UINT64_MAX;
char sql[512]; snprintf(sql, sizeof(sql),
"SELECT node_id, x25519_pubkey, ed25519_pubkey,"
" join_sig, join_ts, update_sig, update_ts, userinfo, adm_tags, adm_tags_sig"
" FROM \"%s\""
" WHERE (node_id & %lld) == %lld ORDER BY node_id ASC",
peers_tbl, (long long)mask, (long long)prefix);
sqlite3_stmt* stmt = NULL;
if (sqlite3_prepare_v2(db, sql, -1, &stmt, NULL) != SQLITE_OK) { DEBUG_ERROR(DEBUG_CATEGORY_MEMBER_SYNC, "%s: get_items prepare failed ns=%s", MS_ID, ns); return -1; }
size_t off = 0;
if (off + 2 > *len) { sqlite3_finalize(stmt); return -2; }
uint16_t* cnt = (uint16_t*)(buf + off); off += 2; *cnt = 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);
const uint8_t* ed = (const uint8_t*)sqlite3_column_blob(stmt, 2);
const uint8_t* sig = (const uint8_t*)sqlite3_column_blob(stmt, 3);
uint64_t jts = (uint64_t)sqlite3_column_int64(stmt, 4);
const uint8_t* usig = (const uint8_t*)sqlite3_column_blob(stmt, 5);
uint64_t uts = (uint64_t)sqlite3_column_int64(stmt, 6);
const char* nm = (const char*)sqlite3_column_text(stmt, 7);
const char* atags = (const char*)sqlite3_column_text(stmt, 8);
const uint8_t* atsig = (const uint8_t*)sqlite3_column_blob(stmt, 9);
if (!x25 || !ed) continue;
uint8_t nl = nm ? (uint8_t)strnlen(nm, 255) : 0;
uint8_t atl = atags ? (uint8_t)strnlen(atags, 255) : 0;
sqlite3_stmt* as = NULL;
sqlite3_prepare_v2(db,
"SELECT family, socket_id, protocol, address, port FROM node_addresses WHERE node_id=? AND addr_type=0"
" ORDER BY family, address, port", -1, &as, NULL);
uint8_t addrs[2048]; int addr_off = 0; int addr_count = 0;
if (as) {
sqlite3_bind_int64(as, 1, (sqlite3_int64)nid);
while (sqlite3_step(as) == SQLITE_ROW && addr_off < (int)sizeof(addrs) - 9) {
int fam = sqlite3_column_int(as, 0);
int sid = sqlite3_column_int(as, 1);
int proto = sqlite3_column_int(as, 2);
int ip_sz = fam == 4 ? 4 : 16;
if (addr_off + 3 + ip_sz + 2 > (int)sizeof(addrs)) break;
addrs[addr_off++] = (uint8_t)fam;
addrs[addr_off++] = (uint8_t)sid;
addrs[addr_off++] = (uint8_t)proto;
memcpy(addrs + addr_off, sqlite3_column_blob(as, 3), (size_t)ip_sz);
addr_off += ip_sz;
uint16_t p = (uint16_t)sqlite3_column_int(as, 4);
addrs[addr_off++] = (uint8_t)(p >> 8);
addrs[addr_off++] = (uint8_t)(p & 0xFF);
addr_count++;
}
sqlite3_finalize(as);
}
uint8_t flags = (sig && jts) ? PEERS_FLAG_HAS_JOIN : 0;
size_t need = 8 + 32 + 32 + 1 + (flags ? 72ULL : 0ULL) + 64 + 8 + 1 + (size_t)nl + 1 + (size_t)atl + 64 + 1 + (size_t)addr_off;
if (off + need > *len) { sqlite3_finalize(stmt); return -2; }
memcpy(buf + off, &nid, 8); off += 8;
memcpy(buf + off, x25, 32); off += 32;
memcpy(buf + off, ed, 32); off += 32;
buf[off++] = flags;
if (flags & PEERS_FLAG_HAS_JOIN) {
memcpy(buf + off, sig, 64); off += 64;
memcpy(buf + off, &jts, 8); off += 8;
}
if (usig && uts) {
memcpy(buf + off, usig, 64); off += 64;
memcpy(buf + off, &uts, 8); off += 8;
} else {
memset(buf + off, 0, 64); off += 64;
uint64_t z = 0; memcpy(buf + off, &z, 8); off += 8;
}
buf[off++] = nl;
if (nl) { memcpy(buf + off, nm, nl); off += nl; }
buf[off++] = atl;
if (atl) { memcpy(buf + off, atags, atl); off += atl; }
if (atsig && sqlite3_column_bytes(stmt, 9) >= 64) { memcpy(buf + off, atsig, 64); off += 64; }
else { memset(buf + off, 0, 64); off += 64; }
buf[off++] = (uint8_t)addr_count;
memcpy(buf + off, addrs, (size_t)addr_off); off += (size_t)addr_off;
(*cnt)++;
}
sqlite3_finalize(stmt);
*len = off;
DEBUG_DEBUG(DEBUG_CATEGORY_MEMBER_SYNC, "%s: get_items ns=%s L%d/P%016llx mask=%016llx count=%u off=%zu",
MS_ID, ns, level, (unsigned long long)prefix, (unsigned long long)mask, *cnt, off);
return 0;
}
/* ── node_props_changed (adm_tags) callbacks ── */
struct ms_props_cbk { node_props_changed_fn fn; void* arg; struct ms_props_cbk* next; };
static struct ms_props_cbk* g_props_cbks = NULL;
void member_sync_add_props_cbk(node_props_changed_fn fn, void* arg) {
if (!fn) return;
struct ms_props_cbk* e = u_malloc(sizeof(*e));
if (!e) return;
e->fn = fn; e->arg = arg;
e->next = g_props_cbks;
g_props_cbks = e;
}
void member_sync_remove_props_cbk(node_props_changed_fn fn, void* arg) {
if (!fn) return;
struct ms_props_cbk** p = &g_props_cbks;
while (*p) {
if ((*p)->fn == fn && (*p)->arg == arg) {
struct ms_props_cbk* r = *p;
*p = r->next; u_free(r);
return;
}
p = &(*p)->next;
}
}
static int _member_apply_items(void* ctx, const char* ns, uint64_t from_peer,
const uint8_t* data, size_t len) {
struct UTUN_INSTANCE* inst = (struct UTUN_INSTANCE*)ctx;
if (len < 2) { DEBUG_DEBUG(DEBUG_CATEGORY_MEMBER_SYNC, "%s: apply_items — no items from peer=%016llx ns=%.*s (empty, sync complete)", MS_ID, (unsigned long long)from_peer, (int)strnlen(ns, 32), ns); return 0; }
uint16_t count; memcpy(&count, data, 2);
DEBUG_TRACE(DEBUG_CATEGORY_MEMBER_SYNC, "%s: apply_items ns=%s count=%u from=%016llx", MS_ID, ns, count, (unsigned long long)from_peer);
const uint8_t* mp = data + 2; size_t mrem = len - 2;
sqlite3* db = _db(inst);
for (uint16_t i = 0; i < count && mrem >= 148; 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;
uint8_t flags = *mp++; mrem--;
const uint8_t* jsig = NULL; uint64_t jts = 0;
if (flags & PEERS_FLAG_HAS_JOIN) {
if (mrem < 72) break; jsig = mp; mp += 64; mrem -= 64; memcpy(&jts, mp, 8); mp += 8; mrem -= 8;
} else {
/* lookup join from local DB */
uint8_t buf_jsig[64]; uint64_t buf_jts;
if (db && topo_node_sqlite_member_get_join(db, ns, nid, buf_jsig, &buf_jts) == 0) { jsig = buf_jsig; jts = buf_jts; }
}
if (mrem < 72) break;
const uint8_t* usig = mp; mp += 64; mrem -= 64; uint64_t uts; memcpy(&uts, mp, 8); mp += 8; mrem -= 8;
uint8_t nl = *mp++; mrem--;
char nm[256] = "";
if (nl && mrem >= nl) { memcpy(nm, mp, nl); nm[nl] = '\0'; mp += nl; mrem -= nl; }
if (mrem < 1) break;
uint8_t atl = *mp++; mrem--;
char atags[256] = "";
const uint8_t* atsig = NULL;
if (atl && mrem >= atl) { memcpy(atags, mp, atl); atags[atl] = '\0'; mp += atl; mrem -= atl; }
if (mrem >= 64) { atsig = mp; mp += 64; mrem -= 64; }
uint8_t ac = *mp++; mrem--;
const uint8_t* addrs = mp;
int consumed = 0;
for (int a = 0; a < (int)ac && mrem >= (size_t)(2 + consumed); a++) {
uint8_t fam = mp[consumed]; consumed++;
uint8_t sid = mp[consumed]; consumed++;
uint8_t proto = mp[consumed]; consumed++;
int sz = fam == 4 ? 4 : 16;
consumed += sz + 2;
}
/* verify update_sig if available */
if (usig && uts && jsig && ed) {
uint8_t vmsg[256]; size_t vlen = 0;
memcpy(vmsg + vlen, jsig, 64); vlen += 64;
memcpy(vmsg + vlen, &uts, 8); vlen += 8;
size_t nl2 = nm ? strlen(nm) : 0; memcpy(vmsg + vlen, nm ? nm : "", nl2 + 1); vlen += nl2 + 1;
{ EVP_PKEY* pkey = EVP_PKEY_new_raw_public_key(EVP_PKEY_ED25519, NULL, ed, 32);
if (pkey) {
EVP_MD_CTX* ver = EVP_MD_CTX_new();
if (ver) {
int ok = (EVP_DigestVerifyInit(ver, NULL, NULL, NULL, pkey) == 1)
&& (EVP_DigestVerify(ver, usig, 64, vmsg, vlen) == 1);
if (!ok)
DEBUG_ERROR(DEBUG_CATEGORY_MEMBER_SYNC, "%s: member_sync invalid update_sig node=0x%016llx ns=%s", MS_ID,
(unsigned long long)nid, ns);
EVP_MD_CTX_free(ver);
}
EVP_PKEY_free(pkey);
}
}
}
/* verify adm_tags_sig and version check */
int adm_ver = 0, adm_storage = 0;
if (atags[0] && atsig) {
const char* ver_s = strstr(atags, "ver=");
if (ver_s) adm_ver = atoi(ver_s + 4);
adm_storage = strstr(atags, "storage=yes") ? 1 : 0;
uint8_t ch_ed_pub[32] = {0};
if (db && topo_node_sqlite_channel_get(db, ns, NULL, 0, NULL, NULL, ch_ed_pub, NULL) == 0) {
uint8_t amsg[264]; size_t aoff = 0;
size_t atl2 = strlen(atags);
if (atl2 > 191) atl2 = 191;
memcpy(amsg + aoff, atags, atl2); aoff += atl2;
memcpy(amsg + aoff, &nid, 8); aoff += 8;
EVP_PKEY* pkey = EVP_PKEY_new_raw_public_key(EVP_PKEY_ED25519, NULL, ch_ed_pub, 32);
if (pkey) {
EVP_MD_CTX* vctx = EVP_MD_CTX_new();
if (vctx) {
int ok = (EVP_DigestVerifyInit(vctx, NULL, NULL, NULL, pkey) == 1)
&& (EVP_DigestVerify(vctx, atsig, 64, amsg, aoff) == 1);
if (!ok) {
DEBUG_ERROR(DEBUG_CATEGORY_MEMBER_SYNC, "%s: member_sync invalid adm_tags_sig node=0x%016llx ns=%s tags=%s",
MS_ID, (unsigned long long)nid, ns, atags);
adm_ver = -1; /* skip this record */
}
EVP_MD_CTX_free(vctx);
}
EVP_PKEY_free(pkey);
}
} else {
DEBUG_WARN(DEBUG_CATEGORY_MEMBER_SYNC, "%s: cannot verify adm_tags_sig — no channel key for ns=%s", MS_ID, ns);
adm_ver = -1;
}
/* check version is strictly greater than local */
if (adm_ver > 0 && db) {
char peers_tbl[128]; _peers_table(ns, peers_tbl, sizeof(peers_tbl));
sqlite3_stmt* av = NULL;
char asql[256]; snprintf(asql, sizeof(asql), "SELECT adm_tags FROM \"%s\" WHERE node_id=?", peers_tbl);
if (sqlite3_prepare_v2(db, asql, -1, &av, NULL) == SQLITE_OK) {
sqlite3_bind_int64(av, 1, (sqlite3_int64)nid);
if (sqlite3_step(av) == SQLITE_ROW) {
const char* local_tags = (const char*)sqlite3_column_text(av, 0);
if (local_tags) {
int local_ver = 0; const char* lv = strstr(local_tags, "ver=");
if (lv) local_ver = atoi(lv + 4);
if (adm_ver <= local_ver) {
DEBUG_WARN(DEBUG_CATEGORY_MEMBER_SYNC, "%s: stale adm_tags node=0x%016llx ns=%s ver=%d <= local=%d — skipping",
MS_ID, (unsigned long long)nid, ns, adm_ver, local_ver);
adm_ver = -1;
}
}
}
sqlite3_finalize(av);
}
}
}
if (adm_ver < 0) { mp += consumed; mrem -= (size_t)consumed; continue; }
member_sync_put(inst, ns, nid, x25, ed, jsig, jts, usig, uts, nm,
nid != inst->node_id ? addrs : NULL,
nid != inst->node_id ? (int)ac : 0,
atags[0] ? atags : NULL, atsig, adm_storage);
DEBUG_INFO(DEBUG_CATEGORY_MEMBER_SYNC, "%s: [ADDR_SYNC] member_sync recv ns=%s nid=0x%016llx ac=%d adm_ver=%d storage=%d", MS_ID, ns, (unsigned long long)nid, ac, adm_ver, adm_storage);
/* fire node_props_changed callbacks */
{ struct ms_props_cbk* pc = g_props_cbks;
if (pc) {
const char* tags_final = atags[0] ? atags : NULL;
while (pc) { pc->fn(nid, tags_final, pc->arg); pc = pc->next; }
}
}
mp += consumed; mrem -= (size_t)consumed;
}
if (count > 0)
merkle_sync_broadcast(inst, ns, from_peer, data, len);
return 0;
}
static member_sync_node_updated_fn g_node_updated_cb = NULL;
static int _ms_apply_update(void* ctx, const char* ns, uint64_t key,
uint8_t type, const uint8_t* data, size_t len) {
(void)ns;
struct UTUN_INSTANCE* inst = (struct UTUN_INSTANCE*)ctx;
if (type == 0x01 && len >= 1) {
int online = data[0];
member_sync_set_online(inst, key, online);
if (g_node_updated_cb) g_node_updated_cb(key, online);
}
return 0;
}
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,
.apply_update = _ms_apply_update,
};
/* ── Public API ── */
static void _rebuild_all_trees(struct UTUN_INSTANCE* inst) {
sqlite3* db = _db(inst); if (!db) return;
sqlite3_stmt* st = NULL;
if (sqlite3_prepare_v2(db,
"SELECT name FROM sqlite_master WHERE type='table' AND name LIKE 'peers_%'",
-1, &st, NULL) != SQLITE_OK) return;
while (sqlite3_step(st) == SQLITE_ROW) {
const char* tbl = (const char*)sqlite3_column_text(st, 0);
if (!tbl || strncmp(tbl, "peers_", 6) != 0) continue;
const char* ch_id = tbl + 6;
sqlite3_stmt* ns = NULL;
char sql[256]; snprintf(sql, sizeof(sql), "SELECT node_id FROM \"%s\"", tbl);
if (sqlite3_prepare_v2(db, sql, -1, &ns, NULL) == SQLITE_OK) {
while (sqlite3_step(ns) == SQLITE_ROW) {
uint64_t nid = (uint64_t)sqlite3_column_int64(ns, 0);
merkle_sync_recompute_path(inst, ch_id, nid);
DEBUG_DEBUG(DEBUG_CATEGORY_MEMBER_SYNC, "%s: rebuild tree ns=%s nid=%016llx", MS_ID, ch_id, (unsigned long long)nid);
}
sqlite3_finalize(ns);
}
}
sqlite3_finalize(st);
}
int member_sync_init(struct UTUN_INSTANCE* inst) {
if (!inst) { DEBUG_ERROR(DEBUG_CATEGORY_MEMBER_SYNC, "%s: init — inst is NULL", MS_ID); return -1; }
int rc = merkle_sync_init(inst, 0x31, &g_member_ops, inst);
if (rc != 0) return rc;
_rebuild_all_trees(inst);
DEBUG_INFO(DEBUG_CATEGORY_MEMBER_SYNC, "%s: initialized, merkle_rc=%d", MS_ID, rc);
return 0;
}
void member_sync_destroy(struct UTUN_INSTANCE* inst) {
if (!inst) return;
merkle_sync_destroy(inst);
}
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(struct UTUN_INSTANCE* inst, uint64_t peer, const char* ch_id) {
merkle_sync_cancel(inst, peer, ch_id);
}
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, uint64_t join_ts,
const uint8_t* update_sig, uint64_t update_ts,
const char* userinfo,
const uint8_t* addrs_data, int addr_count,
const char* adm_tags, const uint8_t* adm_tags_sig, int storage) {
if (!inst || !ch_id) { DEBUG_ERROR(DEBUG_CATEGORY_MEMBER_SYNC, "%s: put — inst=%p ch_id=%s", MS_ID, (void*)inst, ch_id ? ch_id : "(null)"); return -1; }
DEBUG_TRACE(DEBUG_CATEGORY_MEMBER_SYNC, "%s: put ch=%s nid=%016llx userinfo=%s ac=%d adm_tags=%s storage=%d", MS_ID, ch_id, (unsigned long long)member_id, userinfo ? userinfo : "", addr_count, adm_tags ? adm_tags : "", storage);
sqlite3* db = _db(inst); if (!db) { DEBUG_ERROR(DEBUG_CATEGORY_MEMBER_SYNC, "%s: put — db is NULL ch=%s", MS_ID, ch_id); return -1; }
int rc = topo_node_sqlite_member_put(db, ch_id, member_id,
join_sig, join_ts, update_sig, update_ts,
x25519, ed25519, userinfo, NULL,
adm_tags, adm_tags_sig, storage);
if (rc != 0) {
DEBUG_ERROR(DEBUG_CATEGORY_MEMBER_SYNC, "%s: member_put FAILED ch=%s id=0x%016llx userinfo=%s rc=%d",
MS_ID, ch_id, (unsigned long long)member_id, userinfo ? userinfo : "", rc);
return -1;
}
if (addrs_data && addr_count > 0) {
DEBUG_INFO(DEBUG_CATEGORY_MEMBER_SYNC, "%s: [ADDR_SYNC] member_sync_put node=0x%016llx ch=%s addr_count=%d",
MS_ID, (unsigned long long)member_id, ch_id, addr_count);
sqlite3_exec(db, "BEGIN", NULL, NULL, NULL);
sqlite3_stmt* ds = NULL;
sqlite3_prepare_v2(db, "DELETE FROM node_addresses WHERE node_id=? AND addr_type=0", -1, &ds, NULL);
if (ds) { sqlite3_bind_int64(ds, 1, (sqlite3_int64)member_id); sqlite3_step(ds);
int deleted = sqlite3_changes(db);
DEBUG_INFO(DEBUG_CATEGORY_MEMBER_SYNC, "%s: [ADDR_SYNC] member_sync_put DELETE addr_type=0 node=0x%016llx: %d rows deleted",
MS_ID, (unsigned long long)member_id, deleted);
sqlite3_finalize(ds); }
sqlite3_stmt* as = NULL;
sqlite3_prepare_v2(db,
"INSERT INTO node_addresses(node_id,family,protocol,address,port,addr_type,socket_id)"
" VALUES(?,?,?,?,?,0,?)", -1, &as, NULL);
if (as) {
const uint8_t* p = addrs_data;
int written = 0;
for (int i = 0; i < addr_count; i++) {
uint8_t fam = *p++;
uint8_t sid = *p++;
uint8_t proto = *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_int(as, 3, (int)proto);
sqlite3_bind_blob(as, 4, p, ip_sz, SQLITE_STATIC);
p += ip_sz;
uint16_t port = ((uint16_t)p[0] << 8) | p[1]; p += 2;
sqlite3_bind_int(as, 5, port);
sqlite3_bind_int(as, 6, (int)sid);
sqlite3_step(as); sqlite3_reset(as);
if (fam == 4) {
const uint8_t* ip = p - 6; /* p advanced by ip_sz(4) + port(2) - back to start of ip */
DEBUG_INFO(DEBUG_CATEGORY_MEMBER_SYNC, "%s: [ADDR_SYNC] member_sync_put INSERT node=0x%016llx proto=%d sock=%d %d.%d.%d.%d:%d",
MS_ID, (unsigned long long)member_id, proto, sid, ip[0], ip[1], ip[2], ip[3], port);
} else {
DEBUG_INFO(DEBUG_CATEGORY_MEMBER_SYNC, "%s: [ADDR_SYNC] member_sync_put INSERT v6 node=0x%016llx sock=%d port=%d",
MS_ID, (unsigned long long)member_id, sid, port);
}
written++;
}
sqlite3_finalize(as);
DEBUG_INFO(DEBUG_CATEGORY_MEMBER_SYNC, "%s: [ADDR_SYNC] member_sync_put DONE node=0x%016llx: %d addrs written",
MS_ID, (unsigned long long)member_id, written);
}
sqlite3_exec(db, "COMMIT", NULL, NULL, NULL);
} else {
if (!addrs_data)
DEBUG_INFO(DEBUG_CATEGORY_MEMBER_SYNC, "%s: [ADDR_SYNC] member_sync_put node=0x%016llx ch=%s addrs_data=NULL — NO addresses (skipping)", MS_ID, (unsigned long long)member_id, ch_id);
else
DEBUG_INFO(DEBUG_CATEGORY_MEMBER_SYNC, "%s: [ADDR_SYNC] member_sync_put node=0x%016llx ch=%s addr_count=%d — empty, NO DELETE (skipping)", MS_ID, (unsigned long long)member_id, ch_id, addr_count);
}
merkle_sync_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) { DEBUG_ERROR(DEBUG_CATEGORY_MEMBER_SYNC, "%s: del — inst=%p ch_id=%s", MS_ID, (void*)inst, ch_id ? ch_id : "(null)"); return -1; }
DEBUG_TRACE(DEBUG_CATEGORY_MEMBER_SYNC, "%s: del ch=%s nid=%016llx", MS_ID, ch_id, (unsigned long long)member_id);
sqlite3* db = _db(inst); if (!db) { DEBUG_ERROR(DEBUG_CATEGORY_MEMBER_SYNC, "%s: del — db is NULL ch=%s", MS_ID, ch_id); 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_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;
}
void member_sync_set_node_updated_cb(member_sync_node_updated_fn cb) {
g_node_updated_cb = cb;
}
void member_sync_set_online(struct UTUN_INSTANCE* inst, uint64_t node_id, int online) {
if (!inst) return;
DEBUG_TRACE(DEBUG_CATEGORY_MEMBER_SYNC, "%s: set_online nid=%016llx online=%d", MS_ID, (unsigned long long)node_id, online);
sqlite3* db = _db(inst); if (!db) return;
topo_node_sqlite_node_set_online(db, node_id, online);
}
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);
}