diff --git a/src/chat/chat_join.h b/src/chat/chat_join.h index 0b94819a..473d2a22 100644 --- a/src/chat/chat_join.h +++ b/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), diff --git a/src/chat/chat_sync.c b/src/chat/chat_sync.c index 38183730..6f729fd6 100644 --- a/src/chat/chat_sync.c +++ b/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; diff --git a/src/chat/chat_sync.h b/src/chat/chat_sync.h index 3851340f..6c4fbee3 100644 --- a/src/chat/chat_sync.h +++ b/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); /* Перезапустить групповые подключения для каналов без активных соединений. diff --git a/src/routing_layer/topo_group.c b/src/routing_layer/topo_group.c index 3425cbf3..9a3ac2fb 100644 --- a/src/routing_layer/topo_group.c +++ b/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 по группам. */ diff --git a/src/routing_layer/topo_group.h b/src/routing_layer/topo_group.h index 3019cd90..5f7da44d 100644 --- a/src/routing_layer/topo_group.h +++ b/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 // запрос полной таблицы diff --git a/tests/test_chat_join_e2e.c b/tests/test_chat_join_e2e.c index be134c84..97da5d0f 100644 --- a/tests/test_chat_join_e2e.c +++ b/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); diff --git a/tests/test_group_ownership.c b/tests/test_group_ownership.c index 4ea1226c..b9d0a5d2 100644 --- a/tests/test_group_ownership.c +++ b/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 };