Browse Source

chat: убрать PEER_REMOVE (conn_down ≠ удаление мембера), удалить мёртвый member_sync_del

topo_upd
evgeny 2 months ago
parent
commit
b69dcb4eab
  1. 98
      src/chat/chat_sync.c
  2. 1
      src/chat/chat_sync.h
  3. 8
      src/chat/member_sync.c
  4. 6
      src/chat/member_sync.h

98
src/chat/chat_sync.c

@ -6,7 +6,7 @@
#include "chat_event.h" #include "chat_event.h"
#include "../routing_layer/topo_node_sqlite.h" #include "../routing_layer/topo_node_sqlite.h"
#include "../routing_layer/topo_group.h" #include "../routing_layer/topo_group.h"
#include "../routing_layer/conn_mgr.h" #include "../routing_layer/topo_group_invite.h"
#include "../routing_layer/topo_group_connect.h" #include "../routing_layer/topo_group_connect.h"
#include "../../lib/json_flat.h" #include "../../lib/json_flat.h"
#include "member_sync.h" #include "member_sync.h"
@ -86,7 +86,6 @@ static const char* cs_msg_name(uint8_t type) {
case CS_MSG_WELCOME: return "WELCOME"; case CS_MSG_WELCOME: return "WELCOME";
case CS_MSG_PEER_UPSERT: return "PEER_UPSERT"; case CS_MSG_PEER_UPSERT: return "PEER_UPSERT";
case CS_MSG_ERROR: return "ERROR"; case CS_MSG_ERROR: return "ERROR";
case CS_MSG_PEER_REMOVE: return "PEER_REMOVE";
case CS_MSG_CHANNEL_INVITE: return "CHANNEL_INVITE"; case CS_MSG_CHANNEL_INVITE: return "CHANNEL_INVITE";
default: return "???"; default: return "???";
} }
@ -151,8 +150,6 @@ static void cs_handle_welcome(struct chat_sync* cs, uint64_t peer,
const char* ch_id, const uint8_t* pl, size_t len); const char* ch_id, const uint8_t* pl, size_t len);
static void cs_handle_peer_upsert(struct chat_sync* cs, uint64_t peer, static void cs_handle_peer_upsert(struct chat_sync* cs, uint64_t peer,
const char* ch_id, const uint8_t* pl, size_t len); const char* ch_id, const uint8_t* pl, size_t len);
static void cs_handle_peer_remove(struct chat_sync* cs, uint64_t peer,
const char* ch_id, const uint8_t* pl, size_t len);
static void cs_handle_channel_invite(struct chat_sync* cs, uint64_t peer, static void cs_handle_channel_invite(struct chat_sync* cs, uint64_t peer,
const char* ch_id, const uint8_t* pl, size_t len); const char* ch_id, const uint8_t* pl, size_t len);
static int cs_send_channel_invite(struct chat_sync* cs, struct UTUN_INSTANCE* inst, static int cs_send_channel_invite(struct chat_sync* cs, struct UTUN_INSTANCE* inst,
@ -202,7 +199,6 @@ static void chat_sync_recv_cb(struct ETCP_CONN* conn, struct ll_entry* entry) {
g_cs->pending_invite_node_id = 0; g_cs->pending_invite_ch_id = 0; g_cs->pending_invite_node_id = 0; g_cs->pending_invite_ch_id = 0;
break; break;
} }
case CS_MSG_PEER_REMOVE: cs_handle_peer_remove(g_cs, peer, ch_id, pl, plen); break;
case CS_MSG_CHANNEL_INVITE: cs_handle_channel_invite(g_cs, peer, ch_id, pl, plen); break; case CS_MSG_CHANNEL_INVITE: cs_handle_channel_invite(g_cs, peer, ch_id, pl, plen); break;
default: DEBUG_WARN(DEBUG_CATEGORY_MEMBER_SYNC, "%s: UNKNOWN msg type=%02x from=%016llx", CS_ID, type, (unsigned long long)peer); break; default: DEBUG_WARN(DEBUG_CATEGORY_MEMBER_SYNC, "%s: UNKNOWN msg type=%02x from=%016llx", CS_ID, type, (unsigned long long)peer); break;
} }
@ -219,7 +215,7 @@ static void cs_info_req_timeout_cb(void* arg) {
CS_ID, (unsigned long long)cs->pending_invite_node_id, CS_ID, (unsigned long long)cs->pending_invite_node_id,
(unsigned long long)cs->pending_invite_ch_id); (unsigned long long)cs->pending_invite_ch_id);
uint8_t err[20]; memcpy(err, &cs->pending_invite_node_id, 8); uint8_t err[20]; memcpy(err, &cs->pending_invite_node_id, 8);
int r = CONN_EVENT_TIMEOUT; memcpy(err + 8, &r, 4); int r = 2; memcpy(err + 8, &r, 4); /* timeout */
memcpy(err + 12, &cs->pending_invite_ch_id, 8); memcpy(err + 12, &cs->pending_invite_ch_id, 8);
chat_event_post(CHAT_EVT_CONNECT_RESULT, err, 20); chat_event_post(CHAT_EVT_CONNECT_RESULT, err, 20);
cs->pending_invite_node_id = 0; cs->pending_invite_node_id = 0;
@ -234,7 +230,7 @@ static void cs_join_timeout_cb(void* arg) {
DEBUG_WARN(DEBUG_CATEGORY_MEMBER_SYNC, "%s: CHANNEL_JOIN timeout peer=%016llx ch=%llu", DEBUG_WARN(DEBUG_CATEGORY_MEMBER_SYNC, "%s: CHANNEL_JOIN timeout peer=%016llx ch=%llu",
CS_ID, (unsigned long long)peer, (unsigned long long)cs->pending_invite_ch_id); CS_ID, (unsigned long long)peer, (unsigned long long)cs->pending_invite_ch_id);
uint8_t err[20]; memcpy(err, &peer, 8); uint8_t err[20]; memcpy(err, &peer, 8);
int r = CONN_EVENT_TIMEOUT; memcpy(err + 8, &r, 4); int r = 2; memcpy(err + 8, &r, 4); /* timeout */
memcpy(err + 12, &cs->pending_invite_ch_id, 8); memcpy(err + 12, &cs->pending_invite_ch_id, 8);
chat_event_post(CHAT_EVT_CONNECT_RESULT, err, 20); chat_event_post(CHAT_EVT_CONNECT_RESULT, err, 20);
cs->pending_invite_node_id = 0; cs->pending_invite_node_id = 0;
@ -250,7 +246,7 @@ static void cs_invite_send_timeout_cb(void* arg) {
DEBUG_WARN(DEBUG_CATEGORY_MEMBER_SYNC, "%s: CHANNEL_INVITE timeout — joiner did not respond peer=%016llx ch=%llu", DEBUG_WARN(DEBUG_CATEGORY_MEMBER_SYNC, "%s: CHANNEL_INVITE timeout — joiner did not respond peer=%016llx ch=%llu",
CS_ID, (unsigned long long)peer, (unsigned long long)cs->pending_invite_ch_id); CS_ID, (unsigned long long)peer, (unsigned long long)cs->pending_invite_ch_id);
uint8_t err[20]; memcpy(err, &peer, 8); uint8_t err[20]; memcpy(err, &peer, 8);
int r = CONN_EVENT_TIMEOUT; memcpy(err + 8, &r, 4); int r = 2; memcpy(err + 8, &r, 4); /* timeout */
memcpy(err + 12, &cs->pending_invite_ch_id, 8); memcpy(err + 12, &cs->pending_invite_ch_id, 8);
chat_event_post(CHAT_EVT_CONNECT_RESULT, err, 20); chat_event_post(CHAT_EVT_CONNECT_RESULT, err, 20);
cs->pending_invite_node_id = 0; cs->pending_invite_node_id = 0;
@ -527,12 +523,8 @@ static void cs_on_conn_down(struct ETCP_CONN* conn, int event, void* arg) { (voi
int found = 0; int found = 0;
for (int j = 0; j < g_cs->channels[i].peer_count; j++) for (int j = 0; j < g_cs->channels[i].peer_count; j++)
if (g_cs->channels[i].peer_ids[j] == peer) { found = 1; break; } if (g_cs->channels[i].peer_ids[j] == peer) { found = 1; break; }
if (found) { if (found)
g_cs->channels[i].synced = CS_SYNC_NONE; g_cs->channels[i].synced = CS_SYNC_NONE;
uint8_t rem[9]; rem[0] = CS_MSG_PEER_REMOVE;
memcpy(rem + 1, &peer, 8);
cs_propagate(g_cs, g_cs->channels[i].channel_id, peer, rem, 9);
}
} }
} }
@ -697,35 +689,17 @@ struct chat_invite {
int addr_count; int addr_count;
}; };
/* коллбэк invite-подключения: conn_mgr сам всё делает, здесь только лог результата. /* коллбэк invite-подключения (topo_group_invite): nq->handle уже установлен модулем,
реальный join-протокол запускается через cs_on_conn_up (глобальный callback), здесь только лог результата. Реальный join-протокол запускается через cs_on_conn_up
сохранение в БД — через cs_handle_channel_info_resp после верификации */ (глобальный callback), сохранение в БД — через cs_handle_channel_info_resp. */
static void cs_invite_conn_cb(struct CONN_MGR_HANDLE* h, static void cs_tgi_cb(uint64_t node_id, uint64_t group_id, int event, void* arg) {
uint64_t node_id, uint64_t group_id, (void)arg; (void)group_id;
enum conn_mgr_event event, void* arg) { if (event == TGI_EVENT_JOIN) {
(void)arg;
if (event == CONN_EVENT_JOIN) {
DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "channel invite: JOINED — connected to 0x%016llx, online now", DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "channel invite: JOINED — connected to 0x%016llx, online now",
(unsigned long long)node_id); (unsigned long long)node_id);
struct UTUN_INSTANCE* inst = chat_core_get_inst(); } else {
if (inst && inst->topo_groups) {
struct TOPO_GROUP* g = topo_groups_find(inst->topo_groups, group_id);
if (g) {
struct TOPO_GROUP_NODE* nq = topo_node_find_by_id(g, node_id);
if (nq) nq->handle = h;
}
}
} else if (event == CONN_EVENT_TIMEOUT) {
DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "channel invite: TIMEOUT — could not connect to 0x%016llx", DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "channel invite: TIMEOUT — could not connect to 0x%016llx",
(unsigned long long)node_id); (unsigned long long)node_id);
struct UTUN_INSTANCE* inst = chat_core_get_inst();
if (inst && inst->topo_groups) {
struct TOPO_GROUP* g = topo_groups_find(inst->topo_groups, group_id);
if (g) {
struct TOPO_GROUP_NODE* nq = topo_node_find_by_id(g, node_id);
if (nq) nq->handle = NULL;
}
}
} }
} }
@ -788,13 +762,13 @@ static void cm_invite_trampoline(void* arg) {
u_free(ni); u_free(w); return; u_free(ni); u_free(w); return;
} }
/* conn_mgr_open_invite сам создаст TOPO_GROUP, найдёт/создаст соединение, запустит INVITE_INFO. /* topo_group_invite_join сам создаст TOPO_GROUP, найдёт/создаст соединение (ncd),
DB channel + sync создаются позже в cs_handle_channel_info_resp после верификации */ запустит INVITE_INFO. DB channel + sync создаются позже в cs_handle_channel_info_resp после верификации */
DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "channel invite: connecting to 0x%016llx via %d addresses", DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "channel invite: connecting to 0x%016llx via %d addresses",
(unsigned long long)node_id, inv->addr_count); (unsigned long long)node_id, inv->addr_count);
int r = conn_mgr_open_invite(inst, channel_id, ni, node_id, cs_invite_conn_cb, NULL, NULL); int r = topo_group_invite_join(inst, channel_id, ni, node_id, cs_tgi_cb, NULL);
if (r < 0) { if (r < 0) {
DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "channel invite: FAILED — conn_mgr_open_invite error=%d", r); DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "channel invite: FAILED — topo_group_invite_join error=%d", r);
} }
/* free caller-owned ni */ /* free caller-owned ni */
@ -1412,28 +1386,6 @@ static void cs_handle_peer_upsert(struct chat_sync* cs, uint64_t peer,
CS_ID, (unsigned long long)node_id, ch_id); CS_ID, (unsigned long long)node_id, ch_id);
} }
/* ─── PEER_REMOVE (0x0E) ─── */
static void cs_handle_peer_remove(struct chat_sync* cs, uint64_t peer,
const char* ch_id, const uint8_t* pl, size_t len) {
if (len < 8) { DEBUG_WARN(DEBUG_CATEGORY_MEMBER_SYNC, "%s: PEER_REMOVE too short len=%zu", CS_ID, len); return; }
uint64_t node_id; memcpy(&node_id, pl, 8);
sqlite3* db = cs->inst->topo_sqlite_db;
topo_node_sqlite_member_del(db, ch_id, node_id);
cs_refresh_channels(cs);
{ uint8_t mevt[1 + 64 + 8]; uint8_t ml = (uint8_t)strlen(ch_id);
mevt[0] = ml; memcpy(mevt + 1, ch_id, ml);
memcpy(mevt + 1 + ml, &node_id, 8);
chat_event_post(CHAT_EVT_MEMBER_REMOVED, mevt, 1 + ml + 8); }
cs_propagate(cs, ch_id, peer, pl, len);
DEBUG_INFO(DEBUG_CATEGORY_MEMBER_SYNC, "%s: PEER_REMOVE node=0x%016llx ch=%s",
CS_ID, (unsigned long long)node_id, ch_id);
}
/* ─── CHANNEL_INVITE (0x0F): inviter → joiner (push-приглашение в канал) ─── */ /* ─── CHANNEL_INVITE (0x0F): inviter → joiner (push-приглашение в канал) ─── */
static void cs_handle_channel_invite(struct chat_sync* cs, uint64_t peer, static void cs_handle_channel_invite(struct chat_sync* cs, uint64_t peer,
@ -1593,7 +1545,7 @@ void chat_sync_invite_to_channel(struct UTUN_INSTANCE* inst, const char* ch_id,
/* Check if already connected */ /* Check if already connected */
struct ETCP_CONN* conn = cs_find_conn_for_node(inst, target_node_id); struct ETCP_CONN* conn = cs_find_conn_for_node(inst, target_node_id);
if (!conn || !conn->links_up || !conn->initialized) { if (!conn || !conn->links_up || !conn->initialized) {
/* Try to connect via conn_mgr */ /* Try to connect via direct (ncd) */
if (inst->topo_groups) { if (inst->topo_groups) {
struct TOPO_GROUP* g = topo_groups_find(inst->topo_groups, channel_id); struct TOPO_GROUP* g = topo_groups_find(inst->topo_groups, channel_id);
if (g) { if (g) {
@ -1607,16 +1559,10 @@ void chat_sync_invite_to_channel(struct UTUN_INSTANCE* inst, const char* ch_id,
return; return;
} }
} }
/* Use conn_mgr_open_invite with channel group context */ /* Адреса грузятся из node_addresses в БД внутри node_conn_direct_open */
struct TOPO_NODE* ni = u_calloc(1, sizeof(*ni)); int r = topo_group_invite_to_channel(inst, channel_id, NULL, target_node_id);
if (ni) {
ni->node_id = target_node_id;
/* Addresses loaded from node_addresses in DB by conn_mgr */
int r = conn_mgr_open_invite(inst, channel_id, ni, target_node_id, NULL, NULL, NULL);
if (r < 0) { if (r < 0) {
DEBUG_ERROR(DEBUG_CATEGORY_MEMBER_SYNC, "%s: invite_to_channel conn_mgr_open_invite failed rc=%d", CS_ID, r); DEBUG_ERROR(DEBUG_CATEGORY_MEMBER_SYNC, "%s: invite_to_channel topo_group_invite_to_channel failed rc=%d", CS_ID, r);
}
u_free(ni);
} }
} else { } else {
DEBUG_WARN(DEBUG_CATEGORY_MEMBER_SYNC, "%s: invite_to_channel — no topo_groups, cannot connect", CS_ID); DEBUG_WARN(DEBUG_CATEGORY_MEMBER_SYNC, "%s: invite_to_channel — no topo_groups, cannot connect", CS_ID);
@ -1725,9 +1671,9 @@ static void cs_invite_trampoline(void* arg) {
g_cs->pending_invite_node_id = target_node_id; g_cs->pending_invite_node_id = target_node_id;
g_cs->pending_invite_is_inviter = 1; g_cs->pending_invite_is_inviter = 1;
int r = conn_mgr_open_invite(inst, channel_id, ni, target_node_id, cs_invite_conn_cb, NULL, NULL); int r = topo_group_invite_to_channel(inst, channel_id, ni, target_node_id);
if (r < 0) { if (r < 0) {
DEBUG_ERROR(DEBUG_CATEGORY_MEMBER_SYNC, "%s: invite_with_addrs conn_mgr_open_invite failed rc=%d", CS_ID, r); DEBUG_ERROR(DEBUG_CATEGORY_MEMBER_SYNC, "%s: invite_with_addrs topo_group_invite_to_channel failed rc=%d", CS_ID, r);
g_cs->pending_invite_node_id = 0; g_cs->pending_invite_ch_id = 0; g_cs->pending_invite_node_id = 0; g_cs->pending_invite_ch_id = 0;
g_cs->pending_invite_is_inviter = 0; g_cs->pending_invite_is_inviter = 0;
} }

1
src/chat/chat_sync.h

@ -52,7 +52,6 @@ struct UASYNC;
#define CS_MSG_CHANNEL_JOIN 0x0B #define CS_MSG_CHANNEL_JOIN 0x0B
#define CS_MSG_WELCOME 0x0C #define CS_MSG_WELCOME 0x0C
#define CS_MSG_PEER_UPSERT 0x0D #define CS_MSG_PEER_UPSERT 0x0D
#define CS_MSG_PEER_REMOVE 0x0E
#define CS_MSG_CHANNEL_INVITE 0x0F /* push-invite: inviter → joiner */ #define CS_MSG_CHANNEL_INVITE 0x0F /* push-invite: inviter → joiner */
/* Protocol constants */ /* Protocol constants */

8
src/chat/member_sync.c

@ -702,14 +702,6 @@ int member_sync_put(struct UTUN_INSTANCE* inst, const char* ch_id,
return member_sync_apply_record(inst, ch_id, inst->node_id, &m); 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) {
if (!inst || !ch_id) { DEBUG_ERROR(DEBUG_CATEGORY_MEMBER_SYNC, "%s: del — inst=%p ch_id=%s", MS_ID, (void*)inst, ch_id ? ch_id : "(null)"); return -1; }
DEBUG_TRACE(DEBUG_CATEGORY_MEMBER_SYNC, "%s: del ch=%s nid=%016llx", MS_ID, ch_id, (unsigned long long)member_id);
sqlite3* db = _db(inst); if (!db) { DEBUG_ERROR(DEBUG_CATEGORY_MEMBER_SYNC, "%s: del — db is NULL ch=%s", MS_ID, ch_id); return -1; }
topo_node_sqlite_member_del(db, ch_id, member_id);
return merkle_sync_recompute_path(inst, ch_id, member_id);
}
int member_sync_count(struct UTUN_INSTANCE* inst, const char* ch_id) { int member_sync_count(struct UTUN_INSTANCE* inst, const char* ch_id) {
sqlite3* db = _db(inst); sqlite3* db = _db(inst);
if (!db || !ch_id) return 0; if (!db || !ch_id) return 0;

6
src/chat/member_sync.h

@ -27,9 +27,6 @@
* member_sync_set_online(inst, node_id, 1); // online * member_sync_set_online(inst, node_id, 1); // online
* member_sync_set_online(inst, node_id, 0); // offline * member_sync_set_online(inst, node_id, 0); // offline
* *
* // удаление мембера
* member_sync_del(inst, ch_id, node_id);
*
* // отмена синхронизации (коллбэк НЕ вызывается) * // отмена синхронизации (коллбэк НЕ вызывается)
* member_sync_cancel(inst, peer, ch_id); * member_sync_cancel(inst, peer, ch_id);
* *
@ -160,9 +157,6 @@ int member_sync_apply_record(struct UTUN_INSTANCE* inst, const char* ch_id,
int member_sync_send_to(struct UTUN_INSTANCE* inst, const char* ch_id, int member_sync_send_to(struct UTUN_INSTANCE* inst, const char* ch_id,
uint64_t member_id, uint64_t target_peer); uint64_t member_id, uint64_t target_peer);
/* Удалить мембера из канала и пересчитать дерево. */
int member_sync_del(struct UTUN_INSTANCE* inst, const char* ch_id, uint64_t member_id);
/* Количество мемберов в канале (SELECT COUNT из peers_<ch_id>). */ /* Количество мемберов в канале (SELECT COUNT из peers_<ch_id>). */
int member_sync_count(struct UTUN_INSTANCE* inst, const char* ch_id); int member_sync_count(struct UTUN_INSTANCE* inst, const char* ch_id);

Loading…
Cancel
Save