Browse Source

member_sync: apply_record compare+update per-block, put/del/set_online, send_to

topo_upd
evgeny 2 months ago
parent
commit
d0784a3ed7
  1. 470
      src/chat/member_sync.c
  2. 50
      src/chat/member_sync.h

470
src/chat/member_sync.c

@ -293,6 +293,219 @@ void member_sync_remove_props_cbk(node_props_changed_fn fn, void* arg) {
}
}
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;
}
/* Записывает адреса узла в node_addresses (DELETE + INSERT). Только локальная запись. */
static void _member_write_addrs(sqlite3* db, uint64_t member_id,
const uint8_t* addrs_data, int addr_count) {
if (!addrs_data || addr_count <= 0) return;
DEBUG_INFO(DEBUG_CATEGORY_MEMBER_SYNC, "%s: [ADDR_SYNC] write_addrs node=0x%016llx addr_count=%d",
MS_ID, (unsigned long long)member_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=?", -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] write_addrs DELETE 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(?,?,?,?,?,?,?)", -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, addr_type_from_ip(fam, p - ip_sz - 2));
sqlite3_bind_int(as, 7, (int)sid);
sqlite3_step(as); sqlite3_reset(as);
if (fam == 4) {
const uint8_t* ip = p - 6;
DEBUG_INFO(DEBUG_CATEGORY_MEMBER_SYNC, "%s: [ADDR_SYNC] write_addrs 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] write_addrs 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] write_addrs DONE node=0x%016llx: %d addrs written",
MS_ID, (unsigned long long)member_id, written);
}
sqlite3_exec(db, "COMMIT", NULL, NULL, NULL);
}
int member_sync_apply_record(struct UTUN_INSTANCE* inst, const char* ch_id,
uint64_t from_peer, const struct ms_member_rec* m) {
(void)from_peer;
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; }
uint64_t self = inst->node_id;
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] = "";
{
sqlite3_stmt* st = NULL;
char sql[512]; snprintf(sql, sizeof(sql),
"SELECT join_sig, join_ts, update_ts, 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)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);
}
sqlite3_finalize(st);
}
}
/* ── Блок A (мембер): verify update_sig + сравнить ver (update_ts) ── */
int block_a_ok = 1; /* без update_sig (join-only) — валиден как identity */
if (m->update_sig) {
const uint8_t* jsig = m->join_sig ? m->join_sig : (has_local ? local_join_sig : NULL);
if (!jsig) { block_a_ok = 0; }
else {
uint8_t vmsg[512]; size_t vlen = 0;
memcpy(vmsg + vlen, jsig, 64); vlen += 64;
memcpy(vmsg + vlen, &m->update_ts, 8); vlen += 8;
size_t nl2 = m->userinfo ? strlen(m->userinfo) : 0;
memcpy(vmsg + vlen, m->userinfo ? m->userinfo : "", nl2 + 1); vlen += nl2 + 1;
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, vlen) == 1);
if (!ok) { 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); block_a_ok = 0; }
EVP_MD_CTX_free(ver);
}
EVP_PKEY_free(pkey);
}
}
}
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->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_a_changed) {
const uint8_t* jsig = m->join_sig ? m->join_sig : (has_local ? local_join_sig : 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) != 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;
}
/* адреса — только для чужих узлов, при обновлении блока A */
if (m->node_id != self)
_member_write_addrs(db, m->node_id, m->addrs_data, m->addr_count);
changed = 1;
}
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;
/* ── 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 (данные идентичны) */
}
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);
}
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;
@ -303,6 +516,8 @@ static int _member_apply_items(void* ctx, const char* ns, uint64_t 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 >= 148; i++) {
uint64_t nid; memcpy(&nid, mp, 8); mp += 8; mrem -= 8;
const uint8_t* x25 = mp; mp += 32; mrem -= 32;
@ -337,106 +552,38 @@ static int _member_apply_items(void* ctx, const char* ns, uint64_t from_peer,
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) {
char ver_str[32]; adm_ver = json_flat_get(atags, "ver", ver_str, sizeof(ver_str)) == 0 ? atoi(ver_str) : 0;
adm_storage = (json_flat_get(atags, "storage", ver_str, sizeof(ver_str)) == 0 && strcmp(ver_str, "yes") == 0) ? 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 and resolve conflicts */
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, adm_tags_sig 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);
const uint8_t* local_sig = (const uint8_t*)sqlite3_column_blob(av, 1);
if (local_tags) {
int local_ver = 0;
char lv[32];
if (json_flat_get(local_tags, "ver", lv, sizeof(lv)) == 0) local_ver = atoi(lv);
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;
} else if (adm_ver == local_ver && atsig && local_sig && sqlite3_column_bytes(av, 1) >= 64) {
int tie = memcmp(atsig, local_sig, 64);
if (tie >= 0) {
DEBUG_WARN(DEBUG_CATEGORY_MEMBER_SYNC, "%s: equal ver=%d sig %s local — keeping local node=0x%016llx ns=%s",
MS_ID, adm_ver, tie > 0 ? ">" : "==", (unsigned long long)nid, ns);
adm_ver = -1;
}
}
}
}
sqlite3_finalize(av);
}
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.addrs_data = nid != inst->node_id ? addrs : NULL;
m.addr_count = nid != inst->node_id ? (int)ac : 0;
int r = member_sync_apply_record(inst, ns, from_peer, &m);
DEBUG_INFO(DEBUG_CATEGORY_MEMBER_SYNC, "%s: [ADDR_SYNC] member_sync recv ns=%s nid=0x%016llx ac=%d r=%d", MS_ID, ns, (unsigned long long)nid, ac, r);
if (r < 0) { mp += consumed; mrem -= (size_t)consumed; continue; }
if (r & MS_APPLY_CHANGED) {
changed_count++;
/* 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; }
}
}
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; }
}
if (r & MS_APPLY_STALE) {
/* у нас версия новее — отправить наш полный рекорд автору */
member_sync_send_to(inst, ns, nid, from_peer);
}
mp += consumed; mrem -= (size_t)consumed;
}
if (count > 0)
if (changed_count > 0)
merkle_sync_broadcast(inst, ns, from_peer, data, len);
return 0;
}
@ -511,9 +658,12 @@ void member_sync_cancel(struct UTUN_INSTANCE* inst, uint64_t peer, const char* c
merkle_sync_cancel(inst, peer, ch_id);
}
void member_sync_broadcast_one(struct UTUN_INSTANCE* inst, const char* ch_id, uint64_t member_id) {
if (!inst || !ch_id) return;
sqlite3* db = _db(inst); if (!db) return;
/* Сериализует одного мембера в 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;
@ -522,11 +672,11 @@ void member_sync_broadcast_one(struct UTUN_INSTANCE* inst, const char* ch_id, ui
" join_sig, join_ts, update_sig, update_ts, userinfo, adm_tags, adm_tags_sig"
" 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: broadcast_one — query failed ch=%s nid=0x%016llx", MS_ID, ch_id, (unsigned long long)member_id);
return;
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; }
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);
@ -538,7 +688,7 @@ void member_sync_broadcast_one(struct UTUN_INSTANCE* inst, const char* ch_id, ui
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) { sqlite3_finalize(stmt); return; }
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;
@ -569,7 +719,10 @@ void member_sync_broadcast_one(struct UTUN_INSTANCE* inst, const char* ch_id, ui
}
uint8_t flags = (sig && jts) ? PEERS_FLAG_HAS_JOIN : 0;
uint8_t buf[8192]; size_t off = 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 + 1 + (size_t)addr_off;
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;
@ -587,11 +740,25 @@ void member_sync_broadcast_one(struct UTUN_INSTANCE* inst, const char* ch_id, ui
memcpy(buf + off, addrs, (size_t)addr_off); off += (size_t)addr_off;
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,
@ -602,76 +769,18 @@ int member_sync_put(struct UTUN_INSTANCE* inst, const char* ch_id,
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=?", -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 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(?,?,?,?,?,?,?)", -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, addr_type_from_ip(fam, p - ip_sz - 2));
sqlite3_bind_int(as, 7, (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;
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.addrs_data = addrs_data; m.addr_count = addr_count;
return member_sync_apply_record(inst, ch_id, inst->node_id, &m);
}
int member_sync_del(struct UTUN_INSTANCE* inst, const char* ch_id, uint64_t member_id) {
@ -679,8 +788,7 @@ int member_sync_del(struct UTUN_INSTANCE* inst, const char* ch_id, uint64_t memb
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;
return merkle_sync_recompute_path(inst, ch_id, member_id);
}
int member_sync_count(struct UTUN_INSTANCE* inst, const char* ch_id) {
@ -700,11 +808,17 @@ 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;
int member_sync_set_online(struct UTUN_INSTANCE* inst, uint64_t node_id, int online) {
/* ВАЖНО: онлайн-статус НЕ участвует в merkle sync — не хешируется и не
* версионируется. Только локальная запись в nodes.online. Рассылка online
* идёт отдельным лёгким push_update (MSG_ITEM_UPDATE), вне дерева. */
if (!inst) return 0;
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);
sqlite3* db = _db(inst); if (!db) return 0;
int val = online ? 1 : 0;
if (topo_node_sqlite_node_get_online(db, node_id) == val) return 0;
topo_node_sqlite_node_set_online(db, node_id, val);
return 1;
}
const uint8_t* member_sync_get_hash(struct UTUN_INSTANCE* inst, const char* ch_id,

50
src/chat/member_sync.h

@ -54,6 +54,27 @@ extern "C" {
struct UTUN_INSTANCE;
/* Результат применения рекорда мембера (compare+update по версиям двух блоков). */
#define MS_APPLY_CHANGED 0x01 /* блок стал новее — обновлён, нужен recompute + relay остальным */
#define MS_APPLY_STALE 0x02 /* у нас версия новее — отправить наш полный рекорд автору */
/* Распарсенный рекорд мембера (два подписанных блока + адреса). */
struct ms_member_rec {
uint64_t node_id;
const uint8_t* x25519; /* 32 байта, обязательно */
const uint8_t* ed25519; /* 32 байта, обязательно */
const uint8_t* join_sig; /* 64 байта, может быть NULL */
uint64_t join_ts;
const uint8_t* update_sig; /* 64 байта, может быть NULL */
uint64_t update_ts; /* ver блока мембера */
const char* userinfo; /* JSON {"name":...}, может быть NULL */
const char* adm_tags; /* JSON владельца с "ver", может быть NULL */
const uint8_t* adm_tags_sig; /* 64 байта, может быть NULL */
int storage; /* 0/1, производная от adm_tags.storage */
const uint8_t* addrs_data; /* wire-адреса, может быть NULL */
int addr_count;
};
/*
* Инициализировать модуль: создаёт merkle_sync с коллбэками для мемберов
* (ETCP сервис 0x31), регистрирует _on_node_updated для пересчёта дерева
@ -105,6 +126,26 @@ int member_sync_put(struct UTUN_INSTANCE* inst, const char* ch_id,
const uint8_t* addrs_data, int addr_count,
const char* adm_tags, const uint8_t* adm_tags_sig, int storage);
/*
* Применить рекорд мембера с по-блочным сравнением версий:
* - блок A (мембер): ver = update_ts, подписан ключом мембера (update_sig);
* - блок B (владелец): ver = adm_tags.ver, подписан канальным ключом (adm_tags_sig).
* Каждый блок обновляется независимо, только если его ver увеличился.
* Если ver пришедшего блока меньше локального — блок игнорируется и выставляется
* MS_APPLY_STALE (нужно отправить нашу свежую запись автору from_peer).
* Адреса пишутся только при обновлении блока A и только для чужих узлов.
* Возвращает битовую маску MS_APPLY_*, <0 — ошибка.
*/
int member_sync_apply_record(struct UTUN_INSTANCE* inst, const char* ch_id,
uint64_t from_peer, const struct ms_member_rec* m);
/*
* Отправить наш полный рекорд мембера конкретному пиру (send-back при stale).
* Возвращает 0 при успехе, -1 если нет сессии/соединения.
*/
int member_sync_send_to(struct UTUN_INSTANCE* inst, const char* ch_id,
uint64_t member_id, uint64_t target_peer);
/* Удалить мембера из канала и пересчитать дерево. */
int member_sync_del(struct UTUN_INSTANCE* inst, const char* ch_id, uint64_t member_id);
@ -117,8 +158,15 @@ void member_sync_broadcast_one(struct UTUN_INSTANCE* inst, const char* ch_id, ui
/*
* Установить онлайн-статус узла (nodes.online = 0/1).
*
* ВАЖНО: онлайн-статус НЕ участвует в merkle sync — он не хешируется и не
* версионируется (в _compute_member_hash поля online нет). Это только локальная
* запись в БД. Рассылка online идёт отдельным лёгким push_update (MSG_ITEM_UPDATE),
* вне дерева, и только если значение реально изменилось (возвращает 1).
*
* Возвращает 1 если значение изменилось, 0 если совпало.
*/
void member_sync_set_online(struct UTUN_INSTANCE* inst, uint64_t node_id, int online);
int member_sync_set_online(struct UTUN_INSTANCE* inst, uint64_t node_id, int online);
/* Получить хеш бакета из merkle_tree_hash (для тестов). */
const uint8_t* member_sync_get_hash(struct UTUN_INSTANCE* inst,

Loading…
Cancel
Save