Browse Source

Define join contract and require committed CHAT membership for BGP

proxy
evgeny 4 days ago
parent
commit
9745a987ac
  1. 51
      src/chat/chat_join.h
  2. 110
      src/chat/chat_sync.c
  3. 9
      src/chat/chat_sync.h
  4. 37
      src/routing_layer/topo_group.c
  5. 10
      src/routing_layer/topo_group.h
  6. 44
      tests/test_chat_join_e2e.c
  7. 4
      tests/test_group_ownership.c

51
src/chat/chat_join.h

@ -11,12 +11,59 @@
* │ │◀─── JOIN_INFO_REQ(key) ──────┤ (J подключился к C по ссылке)
* │ │── JOIN_INFO_RESP ───────────▶│ (ch_x25519, ch_ed25519, ch_sig, C-мембер)
* │ │◀─── JOIN_REQUEST ────────────┤ (key + полный набор мембера J)
* │ │ (валиден ключ → BGP new_conn, pending status=sent)
* │ │ валиден ключ → только pending, БЕЗ BGP new_conn
* │◀── REQUEST_FWD ───────────┤ │
* │ verif join_sig, node_id=derive(x25519) │
* │ member_sync_put(J, signed_by=A, signature=sign(A, J_x25519||"group_join"))
* │ ── merkle broadcast ──▶ (рекорд J дошёл по merkle) │
* │ │── JOIN_READY(ch) ────────────▶│ → topo_group_new_conn + ensure_channel_ready
* │ │── JOIN_READY(ch) ────────────▶│ добавление подтверждено у C
*
* КОНТРАКТ ЧЛЕНСТВА (обязателен для всех путей подключения):
* - A и C уже участвуют в канале; A может совпадать с C.
* - NCD даёт транспорт J–C для протокола добавления. Наличие транспорта,
* invite-ключа или локального TOPO_GROUP не даёт участия J в BGP-группе.
* Инициатор удерживает отдельный NCD handle до результата/таймаута добавления;
* закрывать его на транспортном UP, рассчитывая на владельца BGP, запрещено.
* - После проверки JOIN_REQUEST инвайтер сохраняет мембера через member_sync.
* Merkle распространяет запись между существующими участниками.
* - Каждый узел разрешает групповую сессию J после принятия и сохранения
* записи J в СВОЕЙ таблице участников. Подтверждение всех узлов не требуется.
* - Проверка invite-ключа разрешает только обработку запроса добавления.
* Временного членства и раннего BGP new_conn до сохранения мембера нет.
* - JOIN_READY означает наличие записи J у C. Это не подтверждение доставки
* всем узлам и не окончание обмена топологией. BGP JOIN_GROUP / TABLE_COMPLETE
* из topo_group.h относятся к отдельной групповой сессии.
*
* ПОРЯДОК И ПРОВЕРКИ:
* 1. A регистрирует (канал, join_key) -> A на C на JOIN_KEY_TTL_SECONDS.
* Ссылку можно публиковать только после успешного KEY_REGISTER_ACK.
* При A==C регистрация и обработка запроса локальные.
* 2. J открывает NCD к C; C проверяет ключ из JOIN_INFO_REQ и возвращает
* подписанное описание канала и свою запись участника. J проверяет ответ.
* Создание локальной инфраструктуры канала ещё не добавляет J на других узлах.
* 3. C проверяет JOIN_REQUEST: ключ и привязку идентичности J к транспорту.
* Запрос идёт к A по существующему маршруту канала; маршрут через J не нужен.
* 4. A проверяет запись, подписывает приглашение и сохраняет мембера.
* После локального commit или принятия записи по Merkle C отправляет JOIN_READY.
* 5. JOIN_GROUP, пришедший раньше записи мембера, не создаёт сессию у получателя.
* После commit member_sync повторно инициирует групповое согласование
* с этим пиром, если транспорт уже установлен. Новый ETCP-коннект не требуется.
*
* ОШИБКИ, ПОВТОРЫ И ОТМЕНА:
* - Неверный/истёкший ключ, несовпадение идентичности или неверная подпись
* не создают мембера. Ошибка/таймаут не дают временных прав.
* - Повторная запись проходит обычную проверку подписей и версий member_sync;
* повтор не создаёт второго мембера или второго владельца групповой сессии.
* - Ответ относится к ожидаемым каналу и пиру; ответ другой попытки не должен
* завершать текущее ожидание. UP другого пира не меняет адресата приглашения.
* - Отмена ожидания не откатывает уже сохранённую/реплицированную запись.
* Удаление локального канала отменяет его ожидания и подключения.
* - Recovery восстанавливает сессии существующих участников и не использует
* invite-ключ как замену локальной таблицы участников.
*
* J <-> C: ETCP_RT_ID_CHAT_SYNC, заголовок [service:1][group_id:8][type:1], payload
* описан в chat_sync.h. A <-> C: ETCP_RT_ID_JOIN через etcp_router с подписью
* и шифрованием; subcommands и payload описаны ниже.
*
* Модуль отвечает только за маршрутизируемую (etcp_router) часть:
* - регистрация ключа (инвайтер → connection),

110
src/chat/chat_sync.c

@ -17,6 +17,7 @@
#include "../utun_instance.h"
#include "../transport_layer/etcp_api.h"
#include "../transport_layer/etcp.h"
#include "../transport_layer/node_conn_direct.h"
#include "../transport_layer/secure_channel.h"
#include "../ntp_time.h"
#include "../../lib/u_async.h"
@ -73,6 +74,8 @@ struct chat_sync {
uint8_t pending_invite_is_inviter; /* 1 = inviter sends CHANNEL_INVITE on conn_up, 0 = joiner sends JOIN_INFO_REQ */
char pending_password[128]; /* пароль для JOIN_INFO_REQ */
uint64_t pending_join_key; /* join-ключ из invite-ссылки (джойнер) */
struct NODE_CONN_DIRECT* join_handle; /* транспорт протокола добавления, не членство */
void* join_transport_timer;
struct join_pending* join_pending; /* connection-узел: ожидающие ready джойнеры */
/* Pending incoming invites (для ask/deny) */
uint64_t pending_inv_in_ch_id[4];
@ -87,6 +90,47 @@ struct chat_sync {
#define CS_INVITE_SEND_TIMEOUT_MS 3000 /* inviter waits for joiner to respond */
#define CS_SYNC_INTERVAL_MS 1000
static void cs_info_req_timeout_cb(void* arg);
static void cs_join_connect_timeout(void* arg) {
struct chat_sync* cs = arg;
cs->join_transport_timer = NULL;
cs_info_req_timeout_cb(cs);
}
static void cs_join_release_transport(struct chat_sync* cs) {
if (cs->join_transport_timer) {
uasync_cancel_timeout(cs->inst->ua, cs->join_transport_timer); cs->join_transport_timer = NULL;
}
struct NODE_CONN_DIRECT* handle = cs->join_handle;
cs->join_handle = NULL;
if (handle) {
DEBUG_DEBUG(DEBUG_CATEGORY_MEMBER_SYNC, "join transport release peer=%016llx",
(unsigned long long)node_conn_direct_node_id(handle));
node_conn_direct_close(handle);
}
}
static int cs_join_hold_transport(struct chat_sync* cs, uint64_t node_id, struct TOPO_NODE* ni) {
struct NODE_CONN_DIRECT* handle = NULL;
int rc = ni ? node_conn_direct_open_node(cs->inst, node_id, NULL, NULL, &handle, ni, NULL)
: node_conn_direct_open(cs->inst, node_id, NULL, NULL, &handle, NULL);
if (rc == NCD_ERR) {
DEBUG_ERROR(DEBUG_CATEGORY_MEMBER_SYNC, "cannot acquire join transport peer=%016llx", (unsigned long long)node_id);
return -1;
}
cs_join_release_transport(cs);
cs->join_handle = handle;
cs->join_transport_timer = uasync_set_timeout(cs->inst->ua, CS_INFO_REQ_TIMEOUT_MS * 10,
cs, cs_join_connect_timeout, "join_transport_timeout");
if (!cs->join_transport_timer) {
DEBUG_ERROR(DEBUG_CATEGORY_MEMBER_SYNC, "cannot arm join transport timeout");
cs_join_release_transport(cs); return -1;
}
DEBUG_INFO(DEBUG_CATEGORY_MEMBER_SYNC, "join owns transport peer=%016llx", (unsigned long long)node_id);
return 0;
}
static const char* cs_msg_name(uint8_t type) {
switch (type) {
case CS_MSG_INIT_SYNC: return "INIT_SYNC";
@ -229,6 +273,13 @@ static void chat_sync_recv_cb(struct ETCP_CONN* conn, struct ll_entry* entry) {
DEBUG_INFO(DEBUG_CATEGORY_MEMBER_SYNC, "%s: RECV %s from=%016llx ch=%s len=%zu",
CS_ID, cs_msg_name(type), (unsigned long long)peer, ch_id, plen);
if ((type == CS_MSG_JOIN_INFO_RESP || type == CS_MSG_JOIN_READY || type == CS_MSG_ERROR) &&
(cs->pending_invite_node_id != peer || cs->pending_invite_ch_id != group_id)) {
DEBUG_WARN(DEBUG_CATEGORY_MEMBER_SYNC, "ignore stale join response peer=%016llx group=%016llx type=%u",
(unsigned long long)peer, (unsigned long long)group_id, type);
u_free(entry->dgram); queue_entry_free(entry); return;
}
switch (type) {
case CS_MSG_JOIN_INFO_REQ: DEBUG_TRACE(DEBUG_CATEGORY_MEMBER_SYNC, "%s: → handle JOIN_INFO_REQ", CS_ID); cs_handle_join_info_req(cs, peer, ch_id, pl, plen); break;
case CS_MSG_JOIN_INFO_RESP: DEBUG_TRACE(DEBUG_CATEGORY_MEMBER_SYNC, "%s: → handle JOIN_INFO_RESP", CS_ID); cs_handle_join_info_resp(cs, peer, ch_id, pl, plen); break;
@ -236,10 +287,11 @@ static void chat_sync_recv_cb(struct ETCP_CONN* conn, struct ll_entry* entry) {
case CS_MSG_JOIN_READY: DEBUG_TRACE(DEBUG_CATEGORY_MEMBER_SYNC, "%s: → handle JOIN_READY", CS_ID); cs_handle_join_ready(cs, peer, ch_id, pl, plen); break;
case CS_MSG_ERROR: {
DEBUG_WARN(DEBUG_CATEGORY_MEMBER_SYNC, "%s: RECV ERROR from=%016llx ch=%s", CS_ID, (unsigned long long)peer, ch_id);
if (cs->info_req_timer) { uasync_cancel_timeout(cs->inst->ua, cs->info_req_timer); cs->info_req_timer = NULL; }
cs_cancel_proto_timers(cs);
uint8_t err[20]; memcpy(err, &peer, 8); int r = -1; memcpy(err + 8, &r, 4);
memcpy(err + 12, &cs->pending_invite_ch_id, 8); chat_event_post(cs->inst, CHAT_EVT_CONNECT_RESULT, err, 20);
cs->pending_invite_node_id = 0; cs->pending_invite_ch_id = 0;
cs_join_release_transport(cs);
break;
}
case CS_MSG_CHANNEL_INVITE: cs_handle_channel_invite(cs, peer, ch_id, pl, plen); break;
@ -267,6 +319,7 @@ static void cs_info_req_timeout_cb(void* arg) {
chat_event_post(cs->inst, CHAT_EVT_CONNECT_RESULT, err, 20);
cs->pending_invite_node_id = 0;
cs->pending_invite_ch_id = 0;
cs_join_release_transport(cs);
}
static void cs_join_timeout_cb(void* arg) {
@ -283,6 +336,7 @@ static void cs_join_timeout_cb(void* arg) {
cs->pending_invite_node_id = 0;
cs->pending_invite_ch_id = 0;
cs->pending_invite_is_inviter = 0;
cs_join_release_transport(cs);
}
static void cs_invite_send_timeout_cb(void* arg) {
@ -299,6 +353,7 @@ static void cs_invite_send_timeout_cb(void* arg) {
cs->pending_invite_node_id = 0;
cs->pending_invite_ch_id = 0;
cs->pending_invite_is_inviter = 0;
cs_join_release_transport(cs);
}
/* ── helper: find active ETCP_CONN for node ── */
@ -518,7 +573,10 @@ static void _on_invite_sync_done(uint64_t peer, const char* ns, int result, void
/* ── send JOIN_INFO_REQ to start join handshake ── */
static void cs_start_channel_join(struct chat_sync* cs, uint64_t peer) {
if (!cs || cs->info_req_timer) return;
if (!cs || cs->info_req_timer || cs->join_timer) return;
if (cs->join_transport_timer) {
uasync_cancel_timeout(cs->inst->ua, cs->join_transport_timer); cs->join_transport_timer = NULL;
}
char ch_id_str[64];
snprintf(ch_id_str, sizeof(ch_id_str), "%llu", (unsigned long long)cs->pending_invite_ch_id);
/* payload: [join_key:8][pass_len:1][password] */
@ -548,11 +606,10 @@ static void cs_on_conn_up(struct ETCP_CONN* conn, int event, void* arg) { (void)
DEBUG_DEBUG(DEBUG_CATEGORY_MEMBER_SYNC, "%s: conn_up peer=%016llx pending_invite=%llu links_up=%d", CS_ID, (unsigned long long)peer, (unsigned long long)cs->pending_invite_ch_id, conn->links_up);
if (cs->pending_invite_ch_id != 0) {
if (cs->pending_invite_node_id != 0 && cs->pending_invite_node_id != peer)
DEBUG_WARN(DEBUG_CATEGORY_MEMBER_SYNC, "%s: invite node_id MISMATCH: invite=0x%016llx ETCP_peer=0x%016llx — invite is STALE!",
CS_ID, (unsigned long long)cs->pending_invite_node_id, (unsigned long long)peer);
cs->pending_invite_node_id = peer;
if (cs->pending_invite_ch_id != 0 && cs->pending_invite_node_id == peer) {
if (cs->join_transport_timer) {
uasync_cancel_timeout(cs->inst->ua, cs->join_transport_timer); cs->join_transport_timer = NULL;
}
if (cs->pending_invite_is_inviter) {
/* we are the inviter — send CHANNEL_INVITE to the joiner */
char ch_id[64]; snprintf(ch_id, sizeof(ch_id), "%llu",
@ -584,8 +641,11 @@ static void cs_on_conn_down(struct ETCP_CONN* conn, int event, void* arg) { (voi
cs_on_peer_status_changed(cs, peer, 0);
cs_cancel_proto_timers(cs);
if (cs->pending_invite_node_id == peer) {
cs_cancel_proto_timers(cs);
if (cs->join_transport_timer) uasync_cancel_timeout(cs->inst->ua, cs->join_transport_timer);
cs->join_transport_timer = uasync_set_timeout(cs->inst->ua, CS_INFO_REQ_TIMEOUT_MS * 10,
cs, cs_join_connect_timeout, "join_reconnect_timeout");
DEBUG_WARN(DEBUG_CATEGORY_MEMBER_SYNC, "%s: conn down while waiting invite resp peer=%016llx, timers cancelled, state kept for retry on reconnect",
CS_ID, (unsigned long long)peer);
/* invite state survives connection flaps: only timeout or explicit protocol end clears it */
@ -733,6 +793,7 @@ void chat_sync_destroy(struct UTUN_INSTANCE* inst) {
if (cs->refresh_wait) { standby_wait_cancel(cs->refresh_wait); cs->refresh_wait = NULL; }
#endif
cs_cancel_proto_timers(cs);
cs_join_release_transport(cs);
for (int i = 0; i < cs->channel_count; i++)
if (cs->channels[i].peer_ids) u_free(cs->channels[i].peer_ids);
@ -762,6 +823,11 @@ void chat_sync_connect_node(struct UTUN_INSTANCE* inst, uint64_t node_id) {
void chat_sync_on_channel_deleted(struct UTUN_INSTANCE* inst, const char* ch_id) {
struct chat_sync* cs = cs_of(inst);
if (!cs || !ch_id || !ch_id[0]) return;
if (cs->pending_invite_ch_id == strtoull(ch_id, NULL, 10)) {
cs_cancel_proto_timers(cs);
cs->pending_invite_ch_id = 0; cs->pending_invite_node_id = 0;
cs_join_release_transport(cs);
}
for (int i = 0; i < cs->channel_count; i++) {
if (strcmp(cs->channels[i].channel_id, ch_id) != 0) continue;
if (cs->channels[i].peer_ids) u_free(cs->channels[i].peer_ids);
@ -896,13 +962,13 @@ static void cm_invite_trampoline(void* arg) {
DEBUG_INFO(DEBUG_CATEGORY_GENERAL,
"channel invite: link to 0x%016llx already up — skipping connect, join handshake directly",
(unsigned long long)node_id);
cs_start_channel_join(cs, node_id);
if (cs_join_hold_transport(cs, node_id, NULL) == 0) cs_start_channel_join(cs, node_id);
} else {
/* topo_group_invite_to_channel: прямое ncd-подключение к connection-узлу БЕЗ INVITE_INFO.
Реальный join-протокол (JOIN_INFO_REQ → … → JOIN_READY) запускается по cs_on_conn_up. */
/* Транспорт принадлежит join-протоколу до его результата, а не будущей BGP-группе.
JOIN_INFO_REQ → … → JOIN_READY запускается по cs_on_conn_up. */
DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "channel invite: connecting to 0x%016llx via %d addresses",
(unsigned long long)node_id, inv->addr_count);
int r = topo_group_invite_to_channel(inst, channel_id, ni, node_id);
int r = cs_join_hold_transport(cs, node_id, ni);
if (r < 0)
DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "channel invite: FAILED — topo_group_invite_to_channel error=%d", r);
}
@ -989,6 +1055,7 @@ void chat_sync_join_channel(struct UTUN_INSTANCE* inst, const char* ch_id, uint6
cs->pending_invite_ch_id = channel_id;
cs->pending_invite_node_id = target_node_id;
if (cs_join_hold_transport(cs, target_node_id, NULL) < 0) return;
cs_start_channel_join(cs, target_node_id);
cs_on_peer_status_changed(cs, target_node_id, 1);
}
@ -1328,13 +1395,8 @@ static void cs_handle_join_request(struct chat_sync* cs, uint64_t peer,
return;
}
/* валидный ключ → сразу добавляем узел в группу для BGP sync */
{
uint64_t gid = strtoull(ch_id, NULL, 10);
struct TOPO_GROUP* g = topo_groups_find(cs->inst->topo_groups, gid);
if (g && g->group_type == TOPO_GROUP_TYPE_CHAT)
topo_group_new_conn(g, jconn);
}
/* Ключ разрешает протокол добавления. BGP начнётся после сохранения мембера,
локально либо через Merkle; существующий транспорт обслуживает JOIN_REQUEST. */
/* добавить в pending (ждут JOIN_READY по приходу рекорда через merkle) */
cs_pending_add(cs, m.node_id, ch_id, jconn);
@ -1353,6 +1415,7 @@ static void cs_handle_join_request(struct chat_sync* cs, uint64_t peer,
cs->pending_invite_ch_id = 0;
cs->pending_invite_node_id = 0;
cs->pending_invite_is_inviter = 0;
cs_join_release_transport(cs);
}
} else {
chat_join_forward_request(cs->inst, channel_id, inviter, join_key, &m);
@ -1394,6 +1457,7 @@ static void cs_handle_join_ready(struct chat_sync* cs, uint64_t peer,
chat_event_post(cs->inst, CHAT_EVT_CONNECT_RESULT, cevt, 20);
DEBUG_INFO(DEBUG_CATEGORY_MEMBER_SYNC, "%s: JOIN_READY invite success ch=%s peer=%016llx", CS_ID, ch_id, (unsigned long long)peer);
cs->pending_invite_node_id = 0; cs->pending_invite_ch_id = 0; cs->pending_join_key = 0;
cs_join_release_transport(cs);
}
}
@ -1602,7 +1666,7 @@ void chat_sync_invite_to_channel(struct UTUN_INSTANCE* inst, const char* ch_id,
}
}
/* Адреса грузятся из node_addresses в БД внутри node_conn_direct_open */
int r = topo_group_invite_to_channel(inst, channel_id, NULL, target_node_id);
int r = cs_join_hold_transport(cs, target_node_id, NULL);
if (r < 0) {
DEBUG_ERROR(DEBUG_CATEGORY_MEMBER_SYNC, "%s: invite_to_channel topo_group_invite_to_channel failed rc=%d", CS_ID, r);
}
@ -1618,6 +1682,8 @@ void chat_sync_invite_to_channel(struct UTUN_INSTANCE* inst, const char* ch_id,
}
/* Already connected — send invite now, wait for joiner to respond */
if (cs_join_hold_transport(cs, target_node_id, NULL) < 0) return;
uasync_cancel_timeout(inst->ua, cs->join_transport_timer); cs->join_transport_timer = NULL;
cs_send_channel_invite(cs, inst, ch_id, target_node_id, ch_name, owner, ch_x25519, ch_ed, ch_sig);
cs->pending_invite_ch_id = channel_id;
cs->pending_invite_node_id = target_node_id;
@ -1661,6 +1727,8 @@ static void cs_invite_trampoline(void* arg) {
char ch_name[128]; uint8_t ch_x2[32], ch_ed2[32], ch_sig[64]; uint64_t ch_owner;
if (topo_node_sqlite_channel_get(inst->topo_sqlite_db,
ch_id, ch_name, (int)sizeof(ch_name), &ch_owner, ch_x2, ch_ed2, ch_sig) == 0) {
if (cs_join_hold_transport(cs, target_node_id, NULL) < 0) { u_free(w); return; }
uasync_cancel_timeout(inst->ua, cs->join_transport_timer); cs->join_transport_timer = NULL;
cs_send_channel_invite(cs, inst, ch_id, target_node_id, ch_name, ch_owner, ch_x2, ch_ed2, ch_sig);
cs->pending_invite_ch_id = channel_id;
cs->pending_invite_node_id = target_node_id;
@ -1717,7 +1785,7 @@ static void cs_invite_trampoline(void* arg) {
cs->pending_invite_node_id = target_node_id;
cs->pending_invite_is_inviter = 1;
int r = topo_group_invite_to_channel(inst, channel_id, ni, target_node_id);
int r = cs_join_hold_transport(cs, target_node_id, ni);
if (r < 0) {
DEBUG_ERROR(DEBUG_CATEGORY_MEMBER_SYNC, "%s: invite_with_addrs topo_group_invite_to_channel failed rc=%d", CS_ID, r);
cs->pending_invite_node_id = 0; cs->pending_invite_ch_id = 0;

9
src/chat/chat_sync.h

@ -12,7 +12,7 @@
*
* Обрабатывает invite-ссылки (chat_sync_connect_from_invite):
* - сохраняет pubkey, join_key и адреса invite-узла в БД (nodes, node_addresses)
* - создаёт TOPO_GROUP, запускает conn_mgr для подключения
* - создаёт инфраструктуру канала, использует NCD для прямого подключения
* - после поднятия ETCP-соединения запускает join-протокол
*
* Использование:
@ -95,7 +95,7 @@ void chat_sync_invite_to_channel_with_addrs(
Отправляет CS_MSG_CHANNEL_INVITE с полной информацией о канале и своих адресах */
void chat_sync_invite_to_channel(struct UTUN_INSTANCE* inst, const char* ch_id, uint64_t target_node_id);
/* Принять входящее приглашение — продолжаем join-протокол (CHANNEL_JOIN → WELCOME) */
/* Принять входящее приглашение — продолжить протокол из chat_join.h. */
void chat_sync_accept_invite(struct UTUN_INSTANCE* inst, uint64_t channel_id, uint64_t inviter_node_id);
/* Отклонить входящее приглашение */
@ -103,8 +103,9 @@ void chat_sync_deny_invite(struct UTUN_INSTANCE* inst, uint64_t channel_id, uint
/* Присоединиться к каналу через уже подключённый узел (uasync-поток).
Соединение с target_node_id должно быть установлено (links_up && initialized).
Протокол: CHANNEL_INFO_REQ → CHANNEL_INFO_RESP → CHANNEL_JOIN → WELCOME.
Канал (DB + TOPO_GROUP) создаётся локально при получении CHANNEL_INFO_RESP. */
Протокол: JOIN_INFO_REQ → JOIN_INFO_RESP → JOIN_REQUEST → JOIN_READY (chat_join.h).
Канал (DB + TOPO_GROUP) создаётся локально при получении JOIN_INFO_RESP;
это само по себе не даёт членства на других узлах. */
void chat_sync_join_channel(struct UTUN_INSTANCE* inst, const char* ch_id, uint64_t target_node_id);
/* Перезапустить групповые подключения для каналов без активных соединений.

37
src/routing_layer/topo_group.c

@ -225,6 +225,19 @@ static void topo_group_broadcast_withdraw(struct TOPO_GROUP* group, uint64_t nod
// Приём пакетов
// ============================================================================
static int topo_group_has_sender(struct TOPO_GROUP* group, struct ETCP_CONN* conn) {
for (struct ll_entry* e = group->senders_list ? group->senders_list->head : NULL; e; e = e->next)
if (((struct TOPO_GROUP_CONN_ITEM*)e->data)->conn == conn) return 1;
return 0;
}
/* Транспорт не даёт членства: CHAT сверяет локальную таблицу участников. */
static int topo_group_peer_allowed(struct TOPO_GROUP* group, uint64_t node_id) {
if (group->group_type != TOPO_GROUP_TYPE_CHAT) return 1;
sqlite3* db = group->instance->topo_sqlite_db;
return db && group->channel_id[0] && topo_node_sqlite_member_in_channel(db, group->channel_id, node_id);
}
/* Приёмник BGP-пакетов от ETCP: диспетчер по subcmd (NODEINFO/WITHDRAW/...). */
static void topo_group_receive_cbk(struct ETCP_CONN* from_conn, struct ll_entry* entry) {
if (!from_conn || !entry || entry->len < 2) { if (entry) { queue_dgram_free(entry); queue_entry_free(entry); } return; }
@ -265,6 +278,12 @@ static void topo_group_receive_cbk(struct ETCP_CONN* from_conn, struct ll_entry*
DEBUG_DEBUG(DEBUG_CATEGORY_BGP,"BGP recv %s from=%s group=%016llx len=%u",group_subcmd_name(subcmd),
from_conn->log_name,(unsigned long long)pkt_group_id,entry->len);
if (!topo_group_peer_allowed(group, from_conn->peer_node_id)) {
DEBUG_INFO(DEBUG_CATEGORY_BGP, "BGP deferred: member unknown peer=%016llx group=%016llx cmd=%u",
(unsigned long long)from_conn->peer_node_id, (unsigned long long)group->group_id, subcmd);
queue_dgram_free(entry); queue_entry_free(entry); return;
}
if (subcmd == TOPO_SUBCMD_NODEINFO) { nodeinfo_dump_log(data, entry->len); topo_group_process_nodeinfo(group, from_conn, data, entry->len); }
else if (subcmd == TOPO_SUBCMD_WITHDRAW) topo_group_process_withdraw(group, from_conn, data, entry->len);
else if (subcmd == TOPO_SUBCMD_REQUEST_TABLE) topo_group_handle_request_table(group, from_conn);
@ -701,6 +720,11 @@ int topo_group_new_conn(struct TOPO_GROUP* group, struct ETCP_CONN* conn) {
if (!group || !conn) { DEBUG_ERROR(DEBUG_CATEGORY_BGP, "invalid args"); return -1; }
if (!conn->instance) { DEBUG_ERROR(DEBUG_CATEGORY_BGP, "conn->instance is NULL"); return -1; }
if (!conn->instance->rt) { DEBUG_ERROR(DEBUG_CATEGORY_BGP, "conn->instance->rt is NULL"); return -1; }
if (!topo_group_peer_allowed(group, conn->peer_node_id)) {
DEBUG_INFO(DEBUG_CATEGORY_BGP, "group connect deferred: member unknown peer=%016llx group=%016llx",
(unsigned long long)conn->peer_node_id, (unsigned long long)group->group_id);
return -1;
}
/* дедуп: тот же conn стреляет ETCP_CONN_STATUS_UP дважды (UDP-линк, затем TCP-линк).
* Если conn уже в senders_list — повторно не обрабатываем, иначе active_conn_count
@ -731,6 +755,7 @@ int topo_group_new_conn(struct TOPO_GROUP* group, struct ETCP_CONN* conn) {
DEBUG_DEBUG(DEBUG_CATEGORY_BGP, "topo_group_new_conn: peer=%016llx group=%016llx type=%d ch=%s", (unsigned long long)conn->peer_node_id, (unsigned long long)group->group_id, group->group_type, group->channel_id);
topo_group_send_join_group(group, conn);
topo_group_send_table_request(group, conn);
topo_group_connect_on_up(group, conn);
return 0;
@ -1211,6 +1236,11 @@ static void topo_group_send_full_table(struct TOPO_GROUP* group, struct ETCP_CON
/* Обрабатывает REQUEST_TABLE: шлёт свой nodeinfo + полную таблицу + TABLE_COMPLETE. */
static void topo_group_handle_request_table(struct TOPO_GROUP* group, struct ETCP_CONN* conn) {
if (!group || !conn) return;
if (!topo_group_peer_allowed(group, conn->peer_node_id)) {
DEBUG_INFO(DEBUG_CATEGORY_BGP, "table request deferred: member unknown peer=%016llx group=%016llx",
(unsigned long long)conn->peer_node_id, (unsigned long long)group->group_id);
return;
}
topo_group_send_nodeinfo(group, group->local_node, conn, 0);
topo_group_send_full_table(group, conn);
/* senders_list заполняется только через topo_group_new_conn (add + BGP);
@ -1222,9 +1252,10 @@ static void topo_group_handle_request_table(struct TOPO_GROUP* group, struct ETC
static void topo_group_handle_join_group(struct TOPO_GROUP* group, struct ETCP_CONN* conn) {
if (!group || !conn) return;
DEBUG_INFO(DEBUG_CATEGORY_BGP, "handle_join_group: from %s grp=%016llx type=%d", conn->log_name, (unsigned long long)group->group_id, group->group_type);
/* входящий запросил членство: для non-CHAT (UTUN) — добавить запросившего узла + инициировать BGP обратно */
if (group->group_type != TOPO_GROUP_TYPE_CHAT)
topo_group_new_conn(group, conn);
/* Проверка таблицы участников общая для входящих и локальных запросов. */
int existing = topo_group_has_sender(group, conn);
if (topo_group_new_conn(group, conn) == 0 && existing)
topo_group_send_table_request(group, conn);
}
/* Обрабатывает RESYNC: если мы VPN-клиент пира — заново шлём JOIN_GROUP по группам. */

10
src/routing_layer/topo_group.h

@ -90,6 +90,16 @@ void topo_fire_nodeinfo_cbk(struct UTUN_INSTANCE* instance, struct TOPO_GROUP* g
// ETCP ID для пакетов топологии: ETCP_ID_TOPO_ENTRY (см. etcp_api.h)
/* JOIN_GROUP восстанавливает BGP-связь с группой, не добавляет мембера.
* Для CHAT источник членства — локальная таблица участников канала.
* Запись приходит по протоколу из chat_join.h / member_sync; до её сохранения
* транспорт может обслуживать приглашение, но не даёт доступа к топологии канала.
* Поздний commit мембера повторно инициирует согласование на готовом NCD.
* JOIN_READY из chat_sync.h — результат добавления участника, а TABLE_COMPLETE —
* результат обмена топологией одной группы; ни один не означает готовности других групп.
* Конфиг, auto-connect и recovery обязаны использовать общую проверку членства.
*/
// Sub-команды
#define TOPO_SUBCMD_NODEINFO 0x04 // полная информация об узле + подсети
#define TOPO_SUBCMD_REQUEST_TABLE 0x05 // запрос полной таблицы

44
tests/test_chat_join_e2e.c

@ -75,6 +75,22 @@ struct e2e_shared {
};
static struct e2e_shared g_sh;
static int early_bgp_membership;
/* Проверяем границу доступа на каждом шаге реального event loop. */
static void poll_checked(struct UTUN_INSTANCE* inst) {
uasync_poll(inst->ua, POLL_MS);
if (inst->node_id == g_sh.nid[IDX_J]) return;
if (topo_node_sqlite_member_in_channel(inst->topo_sqlite_db, g_sh.ch_id, g_sh.nid[IDX_J])) return;
struct TOPO_GROUP* group = topo_groups_find(inst->topo_groups, strtoull(g_sh.ch_id, NULL, 10));
for (struct ll_entry* e = group && group->senders_list ? group->senders_list->head : NULL; e; e = e->next) {
struct TOPO_GROUP_CONN_ITEM* item = (struct TOPO_GROUP_CONN_ITEM*)e->data;
if (item->conn && item->conn->peer_node_id == g_sh.nid[IDX_J]) {
DEBUG_ERROR(DEBUG_CATEGORY_BGP, "test: J entered BGP before member commit group=%s", g_sh.ch_id);
early_bgp_membership = 1;
}
}
}
struct invite_result { int calls, result; char link[1024]; };
static void invite_ready(void* arg, int result, const char* link) {
@ -193,7 +209,7 @@ static int wait_conn(struct UTUN_INSTANCE* inst, uint64_t nid, int max_iter) {
for (int a = 0; a < max_iter; a++) {
struct ETCP_CONN* c = instance_find_conn(inst, nid);
if (c && c->initialized && c->links_up) return 1;
uasync_poll(inst->ua, POLL_MS);
poll_checked(inst);
}
{
struct ETCP_CONN* c = instance_find_conn(inst, nid);
@ -212,7 +228,7 @@ static int wait_bgp(struct UTUN_INSTANCE* inst, uint64_t nid, int max_iter) {
struct TOPO_GROUP_NODE* nq = topo_node_find_by_id(g, nid);
if (nq && nq->paths && nq->paths->head) return 1;
}
uasync_poll(inst->ua, POLL_MS);
poll_checked(inst);
}
{
struct TOPO_GROUP* g = topo_groups_find(inst->topo_groups, gid);
@ -231,7 +247,7 @@ static int wait_signed(struct UTUN_INSTANCE* inst, uint64_t nid, uint64_t want_s
uint64_t sb = 0; uint8_t sig[64];
if (topo_node_sqlite_member_get_sign(inst->topo_sqlite_db, g_sh.ch_id, nid, &sb, sig) == 0
&& sb == want_sb) return 1;
uasync_poll(inst->ua, POLL_MS);
poll_checked(inst);
}
{
uint64_t sb = 0; uint8_t sig[64];
@ -245,7 +261,7 @@ static int wait_signed(struct UTUN_INSTANCE* inst, uint64_t nid, uint64_t want_s
static int wait_not_member(struct UTUN_INSTANCE* inst, uint64_t nid, int max_iter) {
for (int a = 0; a < max_iter; a++) {
if (topo_node_sqlite_member_in_channel(inst->topo_sqlite_db, g_sh.ch_id, nid)) return 0;
uasync_poll(inst->ua, POLL_MS);
poll_checked(inst);
}
return 1;
}
@ -253,7 +269,7 @@ static int wait_not_member(struct UTUN_INSTANCE* inst, uint64_t nid, int max_ite
static int wait_file(struct UTUN_INSTANCE* inst, const char* path, int max_iter) {
for (int a = 0; a < max_iter; a++) {
if (access(path, F_OK) == 0) return 1;
uasync_poll(inst->ua, POLL_MS);
poll_checked(inst);
}
return 0;
}
@ -272,7 +288,7 @@ static void join_event_handler(struct UTUN_INSTANCE* inst, int type, const uint8
}
static int wait_join_result(struct UTUN_INSTANCE* inst, int max_iter) {
for (int a = 0; a < max_iter && !g_j_join_ok; a++) uasync_poll(inst->ua, POLL_MS);
for (int a = 0; a < max_iter && !g_j_join_ok; a++) poll_checked(inst);
return g_j_join_ok;
}
@ -356,7 +372,7 @@ static int check_gui_superseded(struct UTUN_INSTANCE* inst) {
}
int deferred = gui_results == 0;
uint64_t deadline = get_time_tb() + 10000;
while (!gui_results && get_time_tb() < deadline) uasync_poll(inst->ua, POLL_MS);
while (!gui_results && get_time_tb() < deadline) poll_checked(inst);
chat_event_set_handler(inst, join_event_handler);
return deferred && gui_results == 1 && gui_result_valid && !inst->invite_gui_request;
}
@ -381,7 +397,7 @@ static int check_headless_invites(struct UTUN_INSTANCE* inst) {
uint64_t deadline = get_time_tb() + 10000;
int used = 0, lines = 0;
while (lines < 2 && get_time_tb() < deadline) {
uasync_poll(inst->ua, POLL_MS);
poll_checked(inst);
ssize_t r = recv(fd, reply + used, sizeof(reply) - 1 - used, 0);
if (r > 0) {
for (int i = 0; i < r; i++) if (reply[used + i] == '\n') lines++;
@ -408,7 +424,7 @@ static int check_headless_invites(struct UTUN_INSTANCE* inst) {
if (send(fd, command, n, 0) != n) goto out;
used = 0; reply[0] = 0; deadline = get_time_tb() + 10000;
while (get_time_tb() < deadline) {
uasync_poll(inst->ua, POLL_MS);
poll_checked(inst);
ssize_t r = recv(fd, reply + used, sizeof(reply) - 1 - used, 0);
if (r == 0) break;
if (r > 0) { used += (int)r; reply[used] = 0; }
@ -466,10 +482,10 @@ static int child_main(const char* role, const char* dir, int invalid_key) {
/* conn_presence может отставать от paths; повторяем только ошибку запуска. */
for (int a = 0; a < 1000 && !pending; a++) {
pending = chat_invite_build_link(inst, ch_num, target, NULL, invite_ready, &result);
if (!pending) uasync_poll(inst->ua, POLL_MS);
if (!pending) poll_checked(inst);
}
uint64_t deadline = get_time_tb() + JOIN_REGISTER_TIMEOUT_TB + 10000;
while (pending && !result.calls && get_time_tb() < deadline) uasync_poll(inst->ua, POLL_MS);
while (pending && !result.calls && get_time_tb() < deadline) poll_checked(inst);
if (!result.calls && pending) chat_invite_build_cancel(pending);
if (result.calls != 1 || result.result != CHAT_JOIN_OK || !result.link[0]) {
fprintf(stderr, "A: invite confirmation failed calls=%d result=%d\n", result.calls, result.result); goto out;
@ -556,6 +572,10 @@ static int child_main(const char* role, const char* dir, int invalid_key) {
if (!wait_join_result(inst, 8000)) {
fprintf(stderr, "J: no JOIN_READY\n"); goto out;
}
uint64_t connection_node = g_sh.nid[g_sh.degraded ? IDX_A : IDX_C];
if (!wait_bgp(inst, connection_node, 8000)) {
fprintf(stderr, "J: JOIN_READY received but BGP did not recover\n"); goto out;
}
char ready[512]; snprintf(ready, sizeof(ready), "%s/join_ready", dir);
wf(ready, "ready");
rc = 0;
@ -569,7 +589,7 @@ static int child_main(const char* role, const char* dir, int invalid_key) {
}
out:
if (early_link || delayed_registration) rc = 1;
if (early_link || delayed_registration || early_bgp_membership) rc = 1;
fprintf(stderr, "%s: %s\n", role, rc == 0 ? "OK" : "FAIL");
fflush(stdout); fflush(stderr);
_exit(rc);

4
tests/test_group_ownership.c

@ -17,8 +17,8 @@ int main(void) {
"my_public_key=b3193173def895bd0fcea6f86af077c7d77216f10395275f627ac18242ec0f01\n"
"[server: udp]\naddr=127.0.0.1:0\ntype=public\n");
assert(inst && utun_instance_init(inst) == 0);
struct TOPO_GROUP* a = topo_groups_create_group(inst->topo_groups, 42, TOPO_GROUP_TYPE_CHAT, NULL);
struct TOPO_GROUP* b = topo_groups_create_group(inst->topo_groups, 43, TOPO_GROUP_TYPE_CHAT, NULL);
struct TOPO_GROUP* a = topo_groups_create_group(inst->topo_groups, 42, TOPO_GROUP_TYPE_UTUN, NULL);
struct TOPO_GROUP* b = topo_groups_create_group(inst->topo_groups, 43, TOPO_GROUP_TYPE_UTUN, NULL);
assert(a && b);
struct SC_MYKEYS keys; assert(sc_generate_keypair(&keys) == SC_OK);
struct TOPO_ADDR4 addr = { .addr = {127, 0, 0, 1}, .port = 9, .protocol = TOPO_PROTO_UDP };

Loading…
Cancel
Save