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.
 
 
 
 
 
 

1095 lines
53 KiB

#include "member_sync.h"
#include "../routing_layer/topo_node_sqlite.h"
#include "../routing_layer/topo_group.h"
#include "chat_core.h"
#include "chat_core_priv.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>
#include "../../lib/json_flat.h"
#define MS_ID "member_sync"
/* ── 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 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 char* adm_tags, const uint8_t* adm_tags_sig,
uint64_t signed_by, const uint8_t* signature,
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 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_DigestUpdate(ctx, &signed_by, 8);
if (signature) EVP_DigestUpdate(ctx, signature, 64); else { static const uint8_t zsig[64]; EVP_DigestUpdate(ctx, zsig, 64); }
EVP_DigestFinal_ex(ctx, hash_out, NULL);
EVP_MD_CTX_free(ctx);
}
int member_sync_build_update_msg(const uint8_t* join_sig, uint64_t update_ts,
const char* userinfo,
uint8_t* out, int out_sz) {
if (!join_sig || !out || out_sz < 0) return -1;
int off = 0;
if (off + 64 > out_sz) return -1;
memcpy(out + off, join_sig, 64); off += 64;
if (off + 8 > out_sz) return -1;
memcpy(out + off, &update_ts, 8); off += 8;
size_t ul = userinfo ? strlen(userinfo) : 0;
if (off + (int)ul + 1 > out_sz) return -1;
memcpy(out + off, userinfo ? userinfo : "", ul + 1); off += (int)ul + 1;
return off;
}
int member_sync_build_join_msg(const uint8_t* ch_x25519, const uint8_t* ch_ed25519,
uint64_t node_id, const uint8_t* node_x25519, uint64_t join_ts,
uint8_t* out, int out_sz) {
if (!ch_x25519 || !ch_ed25519 || !node_x25519 || !out || out_sz < 0) return -1;
int off = 0;
if (off + 32 > out_sz) return -1;
memcpy(out + off, ch_x25519, 32); off += 32;
if (off + 32 > out_sz) return -1;
memcpy(out + off, ch_ed25519, 32); off += 32;
if (off + 8 > out_sz) return -1;
memcpy(out + off, &node_id, 8); off += 8;
if (off + 32 > out_sz) return -1;
memcpy(out + off, node_x25519, 32); off += 32;
if (off + 8 > out_sz) return -1;
memcpy(out + off, &join_ts, 8); off += 8;
return off;
}
/* ── 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, signed_by, signature, source"
" 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);
uint64_t sb = (uint64_t)sqlite3_column_int64(stmt, 9);
const uint8_t* sgn = (const uint8_t*)sqlite3_column_blob(stmt, 10);
int src = sqlite3_column_int(stmt, 11);
if (!x25 || !ed || src != 0) continue;
uint8_t mh[MT_HASH_SIZE];
_compute_member_hash(nid, x25, ed, sig, jts, usig, uts, atags, atsig, sb, sgn, 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,"
" signed_by, signature, source"
" 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);
uint64_t sb = (uint64_t)sqlite3_column_int64(stmt, 10);
const uint8_t* sgn = (const uint8_t*)sqlite3_column_blob(stmt, 11);
int src = sqlite3_column_int(stmt, 12);
if (!x25 || !ed || src != 0) continue;
uint8_t nl = nm ? (uint8_t)strnlen(nm, 255) : 0;
uint8_t atl = atags ? (uint8_t)strnlen(atags, 255) : 0;
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 + 8 + 64;
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; }
memcpy(buf + off, &sb, 8); off += 8;
if (sgn && sqlite3_column_bytes(stmt, 11) >= 64) { memcpy(buf + off, sgn, 64); off += 64; }
else { memset(buf + off, 0, 64); off += 64; }
(*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 void _fire_props_changed(struct UTUN_INSTANCE* inst, uint64_t node_id, const char* adm_tags, const char* channel_id) {
if (!inst || !CC(inst)) return;
struct ms_props_cbk* pc = CC(inst)->props_cbks;
while (pc) { pc->fn(node_id, adm_tags, channel_id, pc->arg); pc = pc->next; }
}
void member_sync_add_props_cbk(struct UTUN_INSTANCE* inst, node_props_changed_fn fn, void* arg) {
if (!inst || !fn || !CC(inst)) return;
struct ms_props_cbk* e = u_malloc(sizeof(*e));
if (!e) return;
e->fn = fn; e->arg = arg;
e->next = CC(inst)->props_cbks;
CC(inst)->props_cbks = e;
}
void member_sync_remove_props_cbk(struct UTUN_INSTANCE* inst, node_props_changed_fn fn, void* arg) {
if (!inst || !fn || !CC(inst)) return;
struct ms_props_cbk** p = &CC(inst)->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;
}
}
/* ── Коллбэк применения рекорда мембера (ready-хендлинг join на connection-узле) ── */
struct ms_apply_cbk { member_apply_cbk_fn fn; void* arg; struct ms_apply_cbk* next; };
static void _fire_apply_cbk(struct UTUN_INSTANCE* inst, const char* ch_id, uint64_t node_id,
const uint8_t* x25519, uint64_t signed_by,
const uint8_t* signature) {
if (!inst || !CC(inst)) return;
struct ms_apply_cbk* ac = CC(inst)->apply_cbks;
while (ac) { ac->fn(ch_id, node_id, x25519, signed_by, signature, ac->arg); ac = ac->next; }
}
void member_sync_add_apply_cbk(struct UTUN_INSTANCE* inst, member_apply_cbk_fn fn, void* arg) {
if (!inst || !fn || !CC(inst)) return;
struct ms_apply_cbk* e = u_malloc(sizeof(*e));
if (!e) return;
e->fn = fn; e->arg = arg;
e->next = CC(inst)->apply_cbks;
CC(inst)->apply_cbks = e;
}
void member_sync_remove_apply_cbk(struct UTUN_INSTANCE* inst, member_apply_cbk_fn fn, void* arg) {
if (!inst || !fn || !CC(inst)) return;
struct ms_apply_cbk** p = &CC(inst)->apply_cbks;
while (*p) {
if ((*p)->fn == fn && (*p)->arg == arg) {
struct ms_apply_cbk* r = *p;
*p = r->next; u_free(r);
return;
}
p = &(*p)->next;
}
}
/* ── BGP node event callback: BGP сообщает member_sync об обновлении адресов/классификации.
Классификацию node_type/storage пересчитывает сам BGP (nodeinfo_updated);
member_sync только рассылает node_props_changed потребителям. ── */
static void _ms_on_bgp_node(struct TOPO_GROUP* group, uint64_t node_id, int event, void* arg) {
struct UTUN_INSTANCE* inst = (struct UTUN_INSTANCE*)arg;
if (!inst || !group) return;
sqlite3* db = _db(inst); if (!db) return;
/* REMOVE не удаляет member (удаление — только битая подпись при верификации) */
char tags[256] = "";
char peers_tbl[128]; _peers_table(group->channel_id, peers_tbl, sizeof(peers_tbl));
sqlite3_stmt* st = NULL;
char sql[256]; snprintf(sql, sizeof(sql), "SELECT COALESCE(adm_tags,'') FROM \"%s\" WHERE node_id=?", peers_tbl);
if (sqlite3_prepare_v2(db, sql, -1, &st, NULL) == SQLITE_OK) {
sqlite3_bind_int64(st, 1, (sqlite3_int64)node_id);
if (sqlite3_step(st) == SQLITE_ROW) {
const char* t = (const char*)sqlite3_column_text(st, 0);
if (t) snprintf(tags, sizeof(tags), "%s", t);
}
sqlite3_finalize(st);
}
DEBUG_DEBUG(DEBUG_CATEGORY_MEMBER_SYNC, "%s: bgp_node ev=%d nid=0x%016llx tags=%s",
MS_ID, event, (unsigned long long)node_id, tags);
_fire_props_changed(inst, node_id, tags[0] ? tags : NULL, group->channel_id);
}
static int _sig_is_zero64(const uint8_t* sig) {
if (!sig) return 1;
for (int i = 0; i < 64; i++) if (sig[i]) return 0;
return 1;
}
/* Верифицирует join_sig: sign(ch_x25519||ch_ed25519||node_id||x25519||join_ts) ключом ed25519. */
static int _verify_join_sig(sqlite3* db, const char* ch_id, const struct ms_member_rec* m) {
if (!m->join_sig || _sig_is_zero64(m->join_sig)) return 0;
uint8_t ch_x[32], ch_ed[32];
if (topo_node_sqlite_channel_get(db, ch_id, NULL, 0, NULL, ch_x, ch_ed, NULL) != 0) {
DEBUG_WARN(DEBUG_CATEGORY_MEMBER_SYNC, "%s: verify_join_sig — no channel key ns=%s", MS_ID, ch_id);
return 0;
}
uint8_t vmsg[128];
int vlen = member_sync_build_join_msg(ch_x, ch_ed, m->node_id, m->x25519, m->join_ts,
vmsg, (int)sizeof(vmsg));
if (vlen < 0) return 0;
EVP_PKEY* pk = EVP_PKEY_new_raw_public_key(EVP_PKEY_ED25519, NULL, m->ed25519, 32);
if (!pk) return 0;
EVP_MD_CTX* vctx = EVP_MD_CTX_new();
int ok = vctx && (EVP_DigestVerifyInit(vctx, NULL, NULL, NULL, pk) == 1)
&& (EVP_DigestVerify(vctx, m->join_sig, 64, vmsg, (size_t)vlen) == 1);
if (vctx) EVP_MD_CTX_free(vctx);
EVP_PKEY_free(pk);
if (!ok) DEBUG_ERROR(DEBUG_CATEGORY_MEMBER_SYNC, "%s: apply_record invalid join_sig node=0x%016llx ns=%s", MS_ID, (unsigned long long)m->node_id, ch_id);
return ok ? 1 : 0;
}
/* ── Дерево приглашений: подпись pubkey узла ── */
int member_sync_build_sign_msg(const uint8_t* x25519, uint8_t* out, int out_sz) {
if (!x25519 || !out || out_sz < MEMBER_SIGN_MSG_SIZE) return -1;
memcpy(out, x25519, SC_PUBKEY_SIZE);
memcpy(out + SC_PUBKEY_SIZE, MEMBER_SIGN_SUFFIX, MEMBER_SIGN_SUFFIX_LEN);
return MEMBER_SIGN_MSG_SIZE;
}
int member_sync_sign_pubkey(const uint8_t* signer_ed_priv, const uint8_t* x25519, uint8_t sig[64]) {
if (!signer_ed_priv || !x25519 || !sig) return -1;
uint8_t msg[MEMBER_SIGN_MSG_SIZE];
if (member_sync_build_sign_msg(x25519, msg, (int)sizeof(msg)) < 0) return -1;
if (sc_ed25519_sign(signer_ed_priv, msg, sizeof(msg), sig) != SC_OK) {
DEBUG_ERROR(DEBUG_CATEGORY_MEMBER_SYNC, "%s: sign_pubkey — sc_ed25519_sign FAILED", MS_ID);
return -1;
}
return 0;
}
int member_sync_verify_signature(sqlite3* db, const char* ch_id, const struct ms_member_rec* m) {
if (!db || !ch_id || !m || !m->x25519 || !m->signature || _sig_is_zero64(m->signature)) return 0;
uint8_t msg[MEMBER_SIGN_MSG_SIZE];
if (member_sync_build_sign_msg(m->x25519, msg, (int)sizeof(msg)) < 0) return -1;
uint8_t signer_pub[32];
if (m->signed_by == 0) {
if (topo_node_sqlite_channel_get(db, ch_id, NULL, 0, NULL, NULL, signer_pub, NULL) != 0) {
DEBUG_DEBUG(DEBUG_CATEGORY_MEMBER_SYNC, "%s: verify_signature — no channel key ns=%s (defer)", MS_ID, ch_id);
return -1;
}
} else {
if (topo_node_sqlite_get_ed25519_pubkey(db, m->signed_by, signer_pub) != 0) {
DEBUG_DEBUG(DEBUG_CATEGORY_MEMBER_SYNC, "%s: verify_signature — no ed25519 pubkey for signed_by=0x%016llx ns=%s (defer)", MS_ID, (unsigned long long)m->signed_by, ch_id);
return -1;
}
}
int ok = sc_ed25519_verify(signer_pub, msg, sizeof(msg), m->signature) == SC_OK;
if (!ok)
DEBUG_ERROR(DEBUG_CATEGORY_MEMBER_SYNC, "%s: verify_signature FAIL nid=0x%016llx signed_by=0x%016llx ns=%s",
MS_ID, (unsigned long long)m->node_id, (unsigned long long)m->signed_by, ch_id);
return ok ? 1 : 0;
}
int member_sync_validate_pubkey(struct UTUN_INSTANCE* inst, const char* ch_id,
const uint8_t* x25519_pubkey) {
if (!inst || !ch_id || !x25519_pubkey) return 0;
sqlite3* db = _db(inst); if (!db) return 0;
uint64_t node_id = sc_derive_node_id_from_pubkey(x25519_pubkey);
if (node_id == 0) return 0;
return topo_node_sqlite_member_in_channel(db, ch_id, node_id);
}
int member_sync_apply_record(struct UTUN_INSTANCE* inst, const char* ch_id,
uint64_t from_peer, const struct ms_member_rec* m) {
if (!inst || !ch_id || !m || !m->x25519 || !m->ed25519) {
DEBUG_ERROR(DEBUG_CATEGORY_MEMBER_SYNC, "%s: apply_record — invalid args", MS_ID);
return -1;
}
sqlite3* db = _db(inst); if (!db) { DEBUG_ERROR(DEBUG_CATEGORY_MEMBER_SYNC, "%s: apply_record — db is NULL", MS_ID); return -1; }
int stale = 0, changed = 0;
DEBUG_TRACE(DEBUG_CATEGORY_MEMBER_SYNC, "%s: apply_record ch=%s nid=%016llx uts=%llu tags=%s",
MS_ID, ch_id, (unsigned long long)m->node_id, (unsigned long long)m->update_ts,
m->adm_tags ? m->adm_tags : "");
char peers_tbl[128]; _peers_table(ch_id, peers_tbl, sizeof(peers_tbl));
/* ── загрузить локальную запись ── */
int has_local = 0;
uint8_t local_join_sig[64] = {0}; uint64_t local_join_ts = 0;
uint64_t local_update_ts = 0;
char local_tags[256] = "";
uint64_t local_signed_by = 0;
uint8_t local_signature[64] = {0};
int local_src = 0;
{
sqlite3_stmt* st = NULL;
char sql[512]; snprintf(sql, sizeof(sql),
"SELECT join_sig, join_ts, update_ts, adm_tags, signed_by, signature, source FROM \"%s\" WHERE node_id=?", peers_tbl);
if (sqlite3_prepare_v2(db, sql, -1, &st, NULL) == SQLITE_OK) {
sqlite3_bind_int64(st, 1, (sqlite3_int64)m->node_id);
if (sqlite3_step(st) == SQLITE_ROW) {
has_local = 1;
const uint8_t* js = (const uint8_t*)sqlite3_column_blob(st, 0);
if (js) memcpy(local_join_sig, js, 64);
local_join_ts = (uint64_t)sqlite3_column_int64(st, 1);
local_update_ts = (uint64_t)sqlite3_column_int64(st, 2);
const char* tags = (const char*)sqlite3_column_text(st, 3);
if (tags) snprintf(local_tags, sizeof(local_tags), "%s", tags);
local_signed_by = (uint64_t)sqlite3_column_int64(st, 4);
const uint8_t* ls = (const uint8_t*)sqlite3_column_blob(st, 5);
if (ls && sqlite3_column_bytes(st, 5) >= 64) memcpy(local_signature, ls, 64);
local_src = sqlite3_column_int(st, 6);
}
sqlite3_finalize(st);
}
}
/* ── Блок A (мембер): verify update_sig (или join_sig) + сравнить ver ── */
int is_local = (from_peer == inst->node_id);
int identity_ok = 1;
if (!is_local) {
/* жёсткая привязка идентичности: node_id = derive(x25519), ed25519 — доверенный из nodes */
if (sc_derive_node_id_from_pubkey(m->x25519) != m->node_id) {
DEBUG_ERROR(DEBUG_CATEGORY_MEMBER_SYNC, "%s: apply_record node_id/x25519 mismatch (forgery) nid=0x%016llx ns=%s", MS_ID, (unsigned long long)m->node_id, ch_id);
identity_ok = 0;
} else {
uint8_t trusted_ed[32];
if (topo_node_sqlite_get_ed25519_pubkey(db, m->node_id, trusted_ed) == 0
&& memcmp(m->ed25519, trusted_ed, 32) != 0) {
DEBUG_ERROR(DEBUG_CATEGORY_MEMBER_SYNC, "%s: apply_record ed25519 mismatch trusted (forgery) nid=0x%016llx ns=%s", MS_ID, (unsigned long long)m->node_id, ch_id);
identity_ok = 0;
}
}
/* подпись дерева приглашений: -1 = родитель ещё не пришёл (defer до verify_local), 0 = явная подделка */
if (identity_ok && m->signature && !_sig_is_zero64(m->signature)) {
int sv = member_sync_verify_signature(db, ch_id, m);
if (sv == 0) {
DEBUG_ERROR(DEBUG_CATEGORY_MEMBER_SYNC, "%s: apply_record invalid tree-signature (forgery) nid=0x%016llx signed_by=0x%016llx ns=%s",
MS_ID, (unsigned long long)m->node_id, (unsigned long long)m->signed_by, ch_id);
identity_ok = 0;
}
}
}
int block_a_ok = 0;
if (is_local) {
/* локальная запись (member_sync_put) — данные доверенные, без верификации */
block_a_ok = 1;
} else if (!identity_ok) {
block_a_ok = 0;
} else if (m->update_sig && !_sig_is_zero64(m->update_sig)) {
const uint8_t* jsig = m->join_sig ? m->join_sig : (has_local ? local_join_sig : NULL);
if (jsig) {
uint8_t vmsg[8192];
int vlen = member_sync_build_update_msg(jsig, m->update_ts,
m->userinfo ? m->userinfo : "", vmsg, (int)sizeof(vmsg));
if (vlen >= 0) {
EVP_PKEY* pkey = EVP_PKEY_new_raw_public_key(EVP_PKEY_ED25519, NULL, m->ed25519, 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, m->update_sig, 64, vmsg, (size_t)vlen) == 1);
if (ok) block_a_ok = 1;
else DEBUG_ERROR(DEBUG_CATEGORY_MEMBER_SYNC, "%s: apply_record invalid update_sig node=0x%016llx ns=%s", MS_ID, (unsigned long long)m->node_id, ch_id);
EVP_MD_CTX_free(ver);
}
EVP_PKEY_free(pkey);
}
}
}
} else if (m->join_sig && !_sig_is_zero64(m->join_sig)) {
/* join-only запись — верифицируем join_sig напрямую */
block_a_ok = _verify_join_sig(db, ch_id, m);
}
int block_a_changed = 0, block_a_stale = 0;
if (block_a_ok) {
if (!has_local) {
block_a_changed = 1;
} else {
if (m->join_sig && m->join_ts > local_join_ts) block_a_changed = 1; /* re-join / key rotation */
else if (m->join_sig && m->join_ts < local_join_ts) block_a_stale = 1; /* старая эпоха */
else {
if (m->update_ts > local_update_ts) block_a_changed = 1;
else if (m->update_ts < local_update_ts) block_a_stale = 1;
}
}
}
/* ── Блок B (владелец): verify adm_tags_sig + сравнить ver (adm_tags.ver) ── */
int block_b_ok = 0, block_b_changed = 0, block_b_stale = 0;
int adm_ver = 0, adm_storage = 0;
if (m->adm_tags && m->adm_tags[0] && m->adm_tags_sig && !_sig_is_zero64(m->adm_tags_sig)) {
char ver_str[32];
adm_ver = json_flat_get(m->adm_tags, "ver", ver_str, sizeof(ver_str)) == 0 ? atoi(ver_str) : 0;
adm_storage = (json_flat_get(m->adm_tags, "storage", ver_str, sizeof(ver_str)) == 0 && strcmp(ver_str, "yes") == 0) ? 1 : 0;
uint8_t ch_ed_pub[32] = {0};
if (topo_node_sqlite_channel_get(db, ch_id, NULL, 0, NULL, NULL, ch_ed_pub, NULL) == 0) {
uint8_t amsg[264]; size_t aoff = 0;
size_t atl2 = strlen(m->adm_tags);
if (atl2 > 191) atl2 = 191;
memcpy(amsg + aoff, m->adm_tags, atl2); aoff += atl2;
memcpy(amsg + aoff, &m->node_id, 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, m->adm_tags_sig, 64, amsg, aoff) == 1);
if (!ok) DEBUG_ERROR(DEBUG_CATEGORY_MEMBER_SYNC, "%s: apply_record invalid adm_tags_sig node=0x%016llx ns=%s tags=%s", MS_ID, (unsigned long long)m->node_id, ch_id, m->adm_tags);
else block_b_ok = 1;
EVP_MD_CTX_free(vctx);
}
EVP_PKEY_free(pkey);
}
} else {
DEBUG_WARN(DEBUG_CATEGORY_MEMBER_SYNC, "%s: apply_record cannot verify adm_tags_sig — no channel key for ns=%s", MS_ID, ch_id);
}
}
if (block_b_ok) {
int local_ver = 0;
if (has_local && local_tags[0]) {
char lv[32];
if (json_flat_get(local_tags, "ver", lv, sizeof(lv)) == 0) local_ver = atoi(lv);
}
if (adm_ver > local_ver) block_b_changed = 1;
else if (adm_ver < local_ver) block_b_stale = 1;
}
if (block_b_changed && local_src != 0)
DEBUG_ERROR(DEBUG_CATEGORY_MEMBER_SYNC, "%s: adm_tags change on placeholder member (source=%d) — ignored by merkle, won't propagate nid=0x%016llx ns=%s",
MS_ID, local_src, (unsigned long long)m->node_id, ch_id);
/* ── запись изменённых блоков ── */
if (block_a_changed) {
const uint8_t* jsig = m->join_sig ? m->join_sig : (has_local ? local_join_sig : NULL);
/* дерево приглашений: если в обновлении подписи нет — сохраняем локальную (set один раз при join/create) */
uint64_t w_signed_by = m->signature ? m->signed_by : local_signed_by;
const uint8_t* w_signature = m->signature ? m->signature : (has_local ? local_signature : NULL);
if (topo_node_sqlite_member_block_put(db, ch_id, m->node_id, jsig, m->join_ts,
m->update_sig, m->update_ts, m->x25519, m->ed25519,
m->userinfo ? m->userinfo : "", NULL, w_signed_by, w_signature) != 0) {
DEBUG_ERROR(DEBUG_CATEGORY_MEMBER_SYNC, "%s: apply_record block_put FAILED ch=%s nid=0x%016llx", MS_ID, ch_id, (unsigned long long)m->node_id);
return -1;
}
changed = 1;
/* узел в nodes пишет только merkle-путь (verified identity) */
char node_name[128] = "";
if (m->userinfo) json_flat_get(m->userinfo, "name", node_name, sizeof(node_name));
topo_node_sqlite_node_update_verified(db, m->node_id, node_name, m->x25519, m->ed25519,
m->update_ts ? m->update_ts : m->join_ts,
ntp_time_get_seconds(inst));
}
if (block_b_changed) {
if (topo_node_sqlite_member_owner_put(db, ch_id, m->node_id, m->adm_tags, m->adm_tags_sig, adm_storage) != 0) {
DEBUG_ERROR(DEBUG_CATEGORY_MEMBER_SYNC, "%s: apply_record owner_put FAILED ch=%s nid=0x%016llx", MS_ID, ch_id, (unsigned long long)m->node_id);
return -1;
}
changed = 1;
}
if (block_a_stale || block_b_stale) stale = 1;
/* пересчёт классификации (node_type/storage) — DB-хелпер, читает node_addresses (адреса BGP) */
if (changed)
topo_node_sqlite_nodeinfo_updated(db, m->node_id);
/* ── recompute и подтвердить реальное изменение по хешу ── */
if (changed) {
int rc = merkle_sync_recompute_path(inst, ch_id, m->node_id);
if (rc < 0) return -1;
if (rc == 0) changed = 0; /* запись была no-op (данные идентичны) */
}
/* ── пассивное добавление: если узел уже подключён напрямую — добавить в CHAT-группу ── */
if (changed && inst->topo_groups) {
uint64_t gid = strtoull(ch_id, NULL, 10);
struct TOPO_GROUP* g = topo_groups_find(inst->topo_groups, gid);
if (g && g->group_type == TOPO_GROUP_TYPE_CHAT) {
struct ETCP_CONN* c = instance_find_conn(inst, m->node_id);
if (c) topo_group_new_conn(g, c);
}
}
DEBUG_DEBUG(DEBUG_CATEGORY_MEMBER_SYNC, "%s: apply_record ch=%s nid=0x%016llx → changed=%d stale=%d",
MS_ID, ch_id, (unsigned long long)m->node_id, changed, stale);
return (changed ? MS_APPLY_CHANGED : 0) | (stale ? MS_APPLY_STALE : 0);
}
/* ── Верификация локальной записи при чтении из БД ──
Единственный сценарий удаления мембера: подпись присутствует и невалидна,
либо node_id != derive(x25519). Отсутствие подписи = «не верифицировано», не удаляем. */
int member_sync_verify_local_record(struct sqlite3* db, const char* ch_id, const struct ms_member_rec* m) {
if (!db || !ch_id || !m || !m->x25519 || !m->ed25519) return 1; /* нет данных — не битая */
/* жёсткая привязка идентичности */
if (sc_derive_node_id_from_pubkey(m->x25519) != m->node_id) {
DEBUG_ERROR(DEBUG_CATEGORY_MEMBER_SYNC, "%s: verify_local — node_id/x25519 mismatch nid=0x%016llx ns=%s",
MS_ID, (unsigned long long)m->node_id, ch_id);
return 0;
}
/* update_sig присутствует — верифицируем */
if (m->update_sig && !_sig_is_zero64(m->update_sig)) {
uint8_t vmsg[8192];
int vlen = member_sync_build_update_msg(m->join_sig, m->update_ts,
m->userinfo ? m->userinfo : "", vmsg, (int)sizeof(vmsg));
if (vlen < 0 || !m->join_sig) return 0;
EVP_PKEY* pkey = EVP_PKEY_new_raw_public_key(EVP_PKEY_ED25519, NULL, m->ed25519, 32);
if (!pkey) return 0;
EVP_MD_CTX* ver = EVP_MD_CTX_new();
int ok = ver && (EVP_DigestVerifyInit(ver, NULL, NULL, NULL, pkey) == 1)
&& (EVP_DigestVerify(ver, m->update_sig, 64, vmsg, (size_t)vlen) == 1);
if (ver) EVP_MD_CTX_free(ver);
EVP_PKEY_free(pkey);
if (!ok) {
DEBUG_ERROR(DEBUG_CATEGORY_MEMBER_SYNC, "%s: verify_local — invalid update_sig nid=0x%016llx ns=%s",
MS_ID, (unsigned long long)m->node_id, ch_id);
return 0;
}
} else if (m->join_sig && !_sig_is_zero64(m->join_sig)) {
/* join-only — верифицируем join_sig напрямую */
if (!_verify_join_sig(db, ch_id, m)) return 0;
}
/* adm_tags присутствует — верифицируем adm_tags_sig */
if (m->adm_tags && m->adm_tags[0] && m->adm_tags_sig && !_sig_is_zero64(m->adm_tags_sig)) {
uint8_t ch_ed_pub[32] = {0};
if (topo_node_sqlite_channel_get(db, ch_id, NULL, 0, NULL, NULL, ch_ed_pub, NULL) != 0)
return 1; /* нет канального ключа — не можем проверить, не удаляем */
uint8_t amsg[264]; size_t aoff = 0;
size_t atl2 = strlen(m->adm_tags);
if (atl2 > 191) atl2 = 191;
memcpy(amsg + aoff, m->adm_tags, atl2); aoff += atl2;
memcpy(amsg + aoff, &m->node_id, 8); aoff += 8;
EVP_PKEY* pkey = EVP_PKEY_new_raw_public_key(EVP_PKEY_ED25519, NULL, ch_ed_pub, 32);
if (!pkey) return 1;
EVP_MD_CTX* vctx = EVP_MD_CTX_new();
int ok = vctx && (EVP_DigestVerifyInit(vctx, NULL, NULL, NULL, pkey) == 1)
&& (EVP_DigestVerify(vctx, m->adm_tags_sig, 64, amsg, aoff) == 1);
if (vctx) EVP_MD_CTX_free(vctx);
EVP_PKEY_free(pkey);
if (!ok) {
DEBUG_ERROR(DEBUG_CATEGORY_MEMBER_SYNC, "%s: verify_local — invalid adm_tags_sig nid=0x%016llx ns=%s tags=%s",
MS_ID, (unsigned long long)m->node_id, ch_id, m->adm_tags);
return 0;
}
}
/* дерево приглашений: подпись есть и невалидна (подписант резолвится) — битая запись.
-1 = родитель ещё не пришёл — не трогаем (не битая). */
if (m->signature && !_sig_is_zero64(m->signature)) {
int sv = member_sync_verify_signature(db, ch_id, m);
if (sv == 0) {
DEBUG_ERROR(DEBUG_CATEGORY_MEMBER_SYNC, "%s: verify_local — invalid tree-signature nid=0x%016llx signed_by=0x%016llx ns=%s",
MS_ID, (unsigned long long)m->node_id, (unsigned long long)m->signed_by, ch_id);
return 0;
}
}
return 1;
}
/* Прогоняет все записи канала: битая → удалить + пересчитать хеши. */
int member_sync_verify_and_purge(struct UTUN_INSTANCE* inst, const char* ch_id) {
if (!inst || !ch_id) return 0;
sqlite3* db = _db(inst); if (!db) return 0;
char peers_tbl[128]; _peers_table(ch_id, peers_tbl, sizeof(peers_tbl));
sqlite3_stmt* st = NULL;
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,"
" signed_by, signature, source"
" FROM \"%s\" ORDER BY node_id ASC", peers_tbl);
if (sqlite3_prepare_v2(db, sql, -1, &st, NULL) != SQLITE_OK) return 0;
int purged = 0;
while (sqlite3_step(st) == SQLITE_ROW) {
uint64_t nid = (uint64_t)sqlite3_column_int64(st, 0);
const uint8_t* x25 = (const uint8_t*)sqlite3_column_blob(st, 1);
const uint8_t* ed = (const uint8_t*)sqlite3_column_blob(st, 2);
const uint8_t* jsig = (const uint8_t*)sqlite3_column_blob(st, 3);
uint64_t jts = (uint64_t)sqlite3_column_int64(st, 4);
const uint8_t* usig = (const uint8_t*)sqlite3_column_blob(st, 5);
uint64_t uts = (uint64_t)sqlite3_column_int64(st, 6);
const char* nm = (const char*)sqlite3_column_text(st, 7);
const char* atags = (const char*)sqlite3_column_text(st, 8);
const uint8_t* atsig = (const uint8_t*)sqlite3_column_blob(st, 9);
uint64_t sb = (uint64_t)sqlite3_column_int64(st, 10);
const uint8_t* sgn = (const uint8_t*)sqlite3_column_blob(st, 11);
int src = sqlite3_column_int(st, 12);
if (!x25 || !ed) continue;
if (src != 0) continue; /* плейсхолдер topo_group — не верифицируем */
struct ms_member_rec m;
memset(&m, 0, sizeof(m));
m.node_id = nid;
m.x25519 = x25; m.ed25519 = ed;
m.join_sig = jsig; m.join_ts = jts;
m.update_sig = usig; m.update_ts = uts;
m.userinfo = nm;
m.adm_tags = atags; m.adm_tags_sig = atsig;
m.signed_by = sb; m.signature = sgn;
if (!member_sync_verify_local_record(db, ch_id, &m)) {
DEBUG_ERROR(DEBUG_CATEGORY_MEMBER_SYNC, "%s: verify_and_purge — PURGE broken record nid=0x%016llx ns=%s",
MS_ID, (unsigned long long)nid, ch_id);
topo_node_sqlite_member_del(db, ch_id, nid);
merkle_sync_recompute_path(inst, ch_id, nid);
/* пассивное удаление: если узел подключён напрямую — убрать из CHAT-группы */
if (inst->topo_groups) {
uint64_t gid = strtoull(ch_id, NULL, 10);
struct TOPO_GROUP* g = topo_groups_find(inst->topo_groups, gid);
if (g && g->group_type == TOPO_GROUP_TYPE_CHAT) {
struct ETCP_CONN* c = instance_find_conn(inst, nid);
if (c) topo_group_remove_conn(g, c);
}
}
purged++;
}
}
sqlite3_finalize(st);
if (purged > 0)
DEBUG_WARN(DEBUG_CATEGORY_MEMBER_SYNC, "%s: verify_and_purge — purged %d broken records ns=%s",
MS_ID, purged, ch_id);
return purged;
}
/* Прогнать очистку битых записей по всем каналам. */
void member_sync_verify_and_purge_all(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;
member_sync_verify_and_purge(inst, tbl + 6);
}
sqlite3_finalize(st);
}
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);
int changed_count = 0;
for (uint16_t i = 0; i < count && mrem >= 73; 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; }
uint64_t sb = 0; const uint8_t* sgn = NULL;
if (mrem >= 72) { memcpy(&sb, mp, 8); mp += 8; mrem -= 8; sgn = mp; mp += 64; mrem -= 64; }
struct ms_member_rec m;
memset(&m, 0, sizeof(m));
m.node_id = nid;
m.x25519 = x25; m.ed25519 = ed;
m.join_sig = jsig; m.join_ts = jts;
m.update_sig = usig; m.update_ts = uts;
m.userinfo = nm;
m.adm_tags = atags[0] ? atags : NULL;
m.adm_tags_sig = atsig;
m.signed_by = sb; m.signature = sgn;
int r = member_sync_apply_record(inst, ns, from_peer, &m);
DEBUG_DEBUG(DEBUG_CATEGORY_MEMBER_SYNC, "%s: member_sync recv ns=%s nid=0x%016llx r=%d", MS_ID, ns, (unsigned long long)nid, r);
if (r < 0) continue;
if (r & MS_APPLY_CHANGED) {
changed_count++;
/* fire node_props_changed callbacks */
_fire_props_changed(inst, nid, atags[0] ? atags : NULL, ns);
/* fire apply callbacks (ready-хендлинг join на connection-узле) */
_fire_apply_cbk(inst, ns, nid, x25, sb, sgn);
}
if (r & MS_APPLY_STALE) {
/* у нас версия новее — отправить наш полный рекорд автору */
member_sync_send_to(inst, ns, nid, from_peer);
}
}
if (changed_count > 0)
merkle_sync_broadcast(inst, ns, from_peer, data, len);
return 0;
}
static int _member_validate_peer(void* ctx, const char* ns, uint64_t peer) {
struct UTUN_INSTANCE* inst = (struct UTUN_INSTANCE*)ctx;
sqlite3* db = _db(inst);
if (!db) return 0;
return topo_node_sqlite_member_in_channel(db, ns, peer);
}
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 = NULL, /* online-статус больше не рассылается через merkle */
.validate_peer = _member_validate_peer,
};
/* ── 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;
/* полный пересчёт: снести всё дерево канала, затем собрать заново из peers */
sqlite3_stmt* del = NULL;
if (sqlite3_prepare_v2(db, "DELETE FROM merkle_tree_hash WHERE namespace=?", -1, &del, NULL) == SQLITE_OK) {
sqlite3_bind_text(del, 1, ch_id, -1, SQLITE_STATIC);
sqlite3_step(del);
sqlite3_finalize(del);
}
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);
}
/* Подписаться на BGP node-события chat-группы (BGP → member_sync → node_props_changed). */
void member_sync_subscribe_group(struct UTUN_INSTANCE* inst, struct TOPO_GROUP* group) {
if (!inst || !group || group->group_type != TOPO_GROUP_TYPE_CHAT) return;
topo_group_add_node_cbk(group, _ms_on_bgp_node, inst);
DEBUG_DEBUG(DEBUG_CATEGORY_MEMBER_SYNC, "%s: subscribe_group ch=%s", MS_ID, group->channel_id);
}
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, ETCP_RT_ID_MEMBER_SYNC, &g_member_ops, inst);
if (rc != 0) return rc;
/* очистка битых записей во всех каналах перед построением дерева */
member_sync_verify_and_purge_all(inst);
_rebuild_all_trees(inst);
/* подписка на BGP-события существующих chat-групп */
if (inst->topo_groups && inst->topo_groups->group_list) {
struct ll_entry* ge = inst->topo_groups->group_list->head;
while (ge) {
struct TOPO_GROUP* g = (struct TOPO_GROUP*)ge;
if (g->group_type == TOPO_GROUP_TYPE_CHAT && g->channel_id[0])
member_sync_subscribe_group(inst, g);
ge = ge->next;
}
}
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;
if (inst->topo_groups && inst->topo_groups->group_list) {
struct ll_entry* ge = inst->topo_groups->group_list->head;
while (ge) {
struct TOPO_GROUP* g = (struct TOPO_GROUP*)ge;
if (g->group_type == TOPO_GROUP_TYPE_CHAT && g->channel_id[0])
topo_group_remove_node_cbk(g, _ms_on_bgp_node, inst);
ge = ge->next;
}
}
/* освободить callback-списки (иначе утекут при chat_core_destroy) */
if (CC(inst)) {
struct ms_props_cbk* pc = CC(inst)->props_cbks;
while (pc) { struct ms_props_cbk* n = pc->next; u_free(pc); pc = n; }
CC(inst)->props_cbks = NULL;
struct ms_apply_cbk* ac = CC(inst)->apply_cbks;
while (ac) { struct ms_apply_cbk* n = ac->next; u_free(ac); ac = n; }
CC(inst)->apply_cbks = NULL;
}
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) {
member_sync_verify_and_purge(inst, ch_id);
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);
}
void member_sync_cancel_channel(struct UTUN_INSTANCE* inst, const char* ch_id) {
merkle_sync_cancel_ns(inst, ch_id);
}
/* Сериализует одного мембера в wire-формат [count:2][member...].
Возвращает 0 и *out_len, или -1 если мембер не найден/нет БД. */
static int _serialize_member(struct UTUN_INSTANCE* inst, const char* ch_id, uint64_t member_id,
uint8_t* buf, size_t buf_sz, size_t* out_len) {
if (!inst || !ch_id || !buf || !out_len) return -1;
sqlite3* db = _db(inst); if (!db) return -1;
char peers_tbl[128]; _peers_table(ch_id, peers_tbl, sizeof(peers_tbl));
sqlite3_stmt* stmt = NULL;
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,"
" signed_by, signature"
" FROM \"%s\" WHERE node_id=?", peers_tbl);
if (sqlite3_prepare_v2(db, sql, -1, &stmt, NULL) != SQLITE_OK) {
DEBUG_ERROR(DEBUG_CATEGORY_MEMBER_SYNC, "%s: serialize_member — query failed ch=%s nid=0x%016llx", MS_ID, ch_id, (unsigned long long)member_id);
return -1;
}
sqlite3_bind_int64(stmt, 1, (sqlite3_int64)member_id);
if (sqlite3_step(stmt) != SQLITE_ROW) { sqlite3_finalize(stmt); return -1; }
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);
uint64_t sb = (uint64_t)sqlite3_column_int64(stmt, 10);
const uint8_t* sgn = (const uint8_t*)sqlite3_column_blob(stmt, 11);
if (!x25 || !ed) { sqlite3_finalize(stmt); return -1; }
uint8_t nl = nm ? (uint8_t)strnlen(nm, 255) : 0;
uint8_t atl = atags ? (uint8_t)strnlen(atags, 255) : 0;
uint8_t flags = (sig && jts) ? PEERS_FLAG_HAS_JOIN : 0;
size_t need = 2 + 8 + 32 + 32 + 1 + (flags & PEERS_FLAG_HAS_JOIN ? 72 : 0)
+ 72 + 1 + (size_t)nl + 1 + (size_t)atl + 64 + 8 + 64;
if (need > buf_sz) { sqlite3_finalize(stmt); return -1; }
size_t off = 0;
uint16_t wcnt = 1; memcpy(buf, &wcnt, 2); off += 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, 72); off += 72; }
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); }
else { memset(buf + off, 0, 64); }
off += 64;
memcpy(buf + off, &sb, 8); off += 8;
if (sgn && sqlite3_column_bytes(stmt, 11) >= 64) { memcpy(buf + off, sgn, 64); }
else { memset(buf + off, 0, 64); }
off += 64;
sqlite3_finalize(stmt);
*out_len = off;
return 0;
}
void member_sync_broadcast_one(struct UTUN_INSTANCE* inst, const char* ch_id, uint64_t member_id) {
uint8_t buf[8192]; size_t off = 0;
if (_serialize_member(inst, ch_id, member_id, buf, sizeof(buf), &off) != 0) return;
merkle_sync_broadcast(inst, ch_id, inst->node_id, buf, off);
DEBUG_INFO(DEBUG_CATEGORY_MEMBER_SYNC, "%s: broadcast_one ch=%s nid=0x%016llx size=%zu",
MS_ID, ch_id, (unsigned long long)member_id, off);
}
int member_sync_send_to(struct UTUN_INSTANCE* inst, const char* ch_id,
uint64_t member_id, uint64_t target_peer) {
uint8_t buf[8192]; size_t off = 0;
if (_serialize_member(inst, ch_id, member_id, buf, sizeof(buf), &off) != 0) return -1;
return merkle_sync_send_to(inst, ch_id, target_peer, buf, off);
}
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 char* adm_tags, const uint8_t* adm_tags_sig, int storage,
uint64_t signed_by, const uint8_t* signature) {
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 adm_tags=%s storage=%d signed_by=%016llx", MS_ID, ch_id, (unsigned long long)member_id, userinfo ? userinfo : "", adm_tags ? adm_tags : "", storage, (unsigned long long)signed_by);
struct ms_member_rec m;
memset(&m, 0, sizeof(m));
m.node_id = member_id;
m.x25519 = x25519; m.ed25519 = ed25519;
m.join_sig = join_sig; m.join_ts = join_ts;
m.update_sig = update_sig; m.update_ts = update_ts;
m.userinfo = userinfo;
m.adm_tags = adm_tags; m.adm_tags_sig = adm_tags_sig; m.storage = storage;
m.signed_by = signed_by; m.signature = signature;
return member_sync_apply_record(inst, ch_id, inst->node_id, &m);
}
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;
}
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);
}