Browse Source

chat: BGP→member_sync→consumers callback, BGP не удаляет member, member_sync не лезет в BGP-RAM

topo_upd
evgeny 2 months ago
parent
commit
7f160666d9
  1. 6
      src/chat/chat_channel.c
  2. 72
      src/chat/member_sync.c
  3. 5
      src/chat/member_sync.h
  4. 4
      src/routing_layer/topo_group.c

6
src/chat/chat_channel.c

@ -63,8 +63,10 @@ void chat_core_ensure_channel_ready(const char* ch_id) {
struct DB_SYNC_INSTANCE* si = db_sync_instance_add(g_cc.inst, tbl_msg, ch_hash, 1);
uint64_t gid = strtoull(ch_id, NULL, 10);
if (gid != 0 && g_cc.inst->topo_groups && !topo_groups_find(g_cc.inst->topo_groups, gid))
topo_groups_create_group(g_cc.inst->topo_groups, gid, TOPO_GROUP_TYPE_CHAT, ch_id);
if (gid != 0 && g_cc.inst->topo_groups && !topo_groups_find(g_cc.inst->topo_groups, gid)) {
struct TOPO_GROUP* g = topo_groups_create_group(g_cc.inst->topo_groups, gid, TOPO_GROUP_TYPE_CHAT, ch_id);
if (g) member_sync_subscribe_group(g_cc.inst, g);
}
if (si) {
si_register(si, ch_id);
db_sync_set_insert_cb(si, on_msg_inserted, u_strdup(ch_id));

72
src/chat/member_sync.c

@ -1,6 +1,7 @@
#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"
@ -229,6 +230,11 @@ static int _member_get_items(void* ctx, const char* ns, uint8_t level,
struct ms_props_cbk { node_props_changed_fn fn; void* arg; struct ms_props_cbk* next; };
static struct ms_props_cbk* g_props_cbks = NULL;
static void _fire_props_changed(uint64_t node_id, const char* adm_tags) {
struct ms_props_cbk* pc = g_props_cbks;
while (pc) { pc->fn(node_id, adm_tags, pc->arg); pc = pc->next; }
}
void member_sync_add_props_cbk(node_props_changed_fn fn, void* arg) {
if (!fn) return;
struct ms_props_cbk* e = u_malloc(sizeof(*e));
@ -251,6 +257,35 @@ void member_sync_remove_props_cbk(node_props_changed_fn fn, void* arg) {
}
}
/* ── 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(node_id, tags[0] ? tags : NULL);
}
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;
@ -440,13 +475,6 @@ int member_sync_apply_record(struct UTUN_INSTANCE* inst, const char* ch_id,
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));
/* если BGP уже знает адреса узла (пришли раньше мембера) — персистим их в node_addresses */
struct TOPO_NODE* ni = inst->topo_groups
? topo_node_registry_find(inst->topo_groups, m->node_id) : NULL;
if (ni && (ni->v4_addrs || ni->v6_addrs)) {
topo_node_sqlite_addrs_put(db, m->node_id, ni);
topo_node_sqlite_nodeinfo_updated(db, m->node_id);
}
}
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) {
@ -458,6 +486,10 @@ int member_sync_apply_record(struct UTUN_INSTANCE* inst, const char* ch_id,
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);
@ -586,17 +618,43 @@ static void _rebuild_all_trees(struct UTUN_INSTANCE* inst) {
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, 0x31, &g_member_ops, inst);
if (rc != 0) return rc;
_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;
}
}
merkle_sync_destroy(inst);
}

5
src/chat/member_sync.h

@ -50,6 +50,7 @@ extern "C" {
#include <stddef.h>
struct UTUN_INSTANCE;
struct TOPO_GROUP;
/* Результат применения рекорда мембера (compare+update по версиям двух блоков). */
#define MS_APPLY_CHANGED 0x01 /* блок стал новее — обновлён, нужен recompute + relay остальным */
@ -80,6 +81,10 @@ int member_sync_init(struct UTUN_INSTANCE* inst);
/* Завершить модуль. */
void member_sync_destroy(struct UTUN_INSTANCE* inst);
/* Подписаться на BGP node-события chat-группы (BGP → member_sync → node_props_changed).
Вызывается при создании канала (chat_core_ensure_channel_ready). */
void member_sync_subscribe_group(struct UTUN_INSTANCE* inst, struct TOPO_GROUP* group);
/*
* Запустить синхронизацию мемберов канала ch_id с пиром peer.
* Делегирует в merkle_sync_start().

4
src/routing_layer/topo_group.c

@ -851,8 +851,8 @@ int topo_group_process_withdraw(struct TOPO_GROUP* group, struct ETCP_CONN* send
struct ll_entry* entry = queue_find_data_by_index(group->nodes, &node_id);
if (entry) { queue_remove_data(group->nodes, entry); queue_entry_free(entry); }
topo_group_broadcast_withdraw(group, node_id, wd_source, sender);
if (group->group_type == TOPO_GROUP_TYPE_CHAT && group->instance->topo_sqlite_db)
topo_node_sqlite_member_del(group->instance->topo_sqlite_db, group->channel_id, node_id);
/* BGP не удаляет member-запись: он работает только со своей RAM-таблицей.
Удаление из peers_* — только битая подпись при верификации (member_sync). */
{ /* fire BGP REMOVE callback */
struct topo_node_cbk_entry* c = group->node_cbks;
while (c) { c->fn(group, node_id, TOPO_NODE_EVENT_REMOVE, c->arg); c = c->next; }

Loading…
Cancel
Save