From d0784a3ed747390c131bce9227b356a3dd9d044f Mon Sep 17 00:00:00 2001 From: evgeny Date: Fri, 14 Aug 2026 16:23:03 +0300 Subject: [PATCH] member_sync: apply_record compare+update per-block, put/del/set_online, send_to --- src/chat/member_sync.c | 470 +++++++++++++++++++++++++---------------- src/chat/member_sync.h | 50 ++++- 2 files changed, 341 insertions(+), 179 deletions(-) diff --git a/src/chat/member_sync.c b/src/chat/member_sync.c index f3ca53fa..db84076c 100644 --- a/src/chat/member_sync.c +++ b/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, diff --git a/src/chat/member_sync.h b/src/chat/member_sync.h index a2509386..0f307837 100644 --- a/src/chat/member_sync.h +++ b/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,