diff --git a/src/routing_layer/topo_group.c b/src/routing_layer/topo_group.c index 3394a098..f11d018c 100644 --- a/src/routing_layer/topo_group.c +++ b/src/routing_layer/topo_group.c @@ -55,7 +55,7 @@ static struct TOPO_GROUP_CONN_ITEM* topo_group_peer(const struct TOPO_GROUP* gro int topo_group_peer_ready(const struct TOPO_GROUP* group, uint64_t peer_id) { const struct TOPO_GROUP_CONN_ITEM* peer = topo_group_peer(group, peer_id); - return peer && peer->conn && peer->conn->links_up && !peer->conn->close_requested && peer->table_received && peer->table_sent; + return peer && peer->conn && peer->conn->links_up && !peer->conn->close_requested && peer->accepted && peer->table_received && peer->table_sent; } static void topo_group_peer_free(struct TOPO_GROUP_CONN_ITEM* peer) { @@ -154,126 +154,82 @@ static void topo_group_log_exchange(struct TOPO_GROUP* group, struct TOPO_GROUP_ topo_group_peer_progressed(group, peer->node_id); DEBUG_INFO(DEBUG_CATEGORY_BGP, "group exchange: group=%016llx peer=%016llx exchange=%016llx sent=%u received=%u ready=%d", (unsigned long long)group->group_id, (unsigned long long)peer->conn->peer_node_id, - (unsigned long long)peer->exchange_id, peer->table_sent, peer->table_received, + (unsigned long long)peer->local_epoch, peer->table_sent, peer->table_received, topo_group_peer_ready(group, peer->conn->peer_node_id)); } -/* Шлёт пиру запрос полной таблицы узлов группы (REQUEST_TABLE). */ +static struct TOPOMSG_HEADER topo_header(struct TOPO_GROUP_CONN_ITEM* peer, uint8_t cmd) { + return (struct TOPOMSG_HEADER){ .cmd = ETCP_ID_TOPO_ENTRY, .subcmd = cmd, .group_id = peer->group->group_id, + .src_epoch = peer->local_epoch, .dst_epoch = peer->peer_epoch }; +} + +static int topo_send_control(struct ETCP_CONN* conn, const void* data, size_t size) { + struct ll_entry* e = ll_alloc_lldgram(size); + if (!e) { DEBUG_ERROR(DEBUG_CATEGORY_BGP, "group control allocation failed size=%zu", size); return -1; } + memcpy(e->dgram, data, size); e->len = size; + if (etcp_send(conn, e) == 0) return 0; + DEBUG_ERROR(DEBUG_CATEGORY_BGP, "group control send failed peer=%016llx cmd=%u", + (unsigned long long)conn->peer_node_id, ((const uint8_t*)data)[1]); + queue_dgram_free(e); queue_entry_free(e); return -1; +} + +static int topo_group_reset_exchange(struct TOPO_GROUP_CONN_ITEM* peer) { + peer->accepted = 0; peer->table_received = 0; peer->table_sent = 0; peer->peer_epoch = 0; + uint64_t epoch; + if (random_bytes((uint8_t*)&epoch, sizeof(epoch)) != 0 || !epoch) { + peer->local_epoch = 0; + DEBUG_ERROR(DEBUG_CATEGORY_BGP, "cannot generate group epoch group=%016llx", (unsigned long long)peer->group->group_id); + return -1; + } + peer->local_epoch = epoch; + topo_group_peer_progressed(peer->group, peer->node_id); + DEBUG_INFO(DEBUG_CATEGORY_BGP, "group session started: group=%016llx peer=%016llx local=%016llx", + (unsigned long long)peer->group->group_id, (unsigned long long)peer->node_id, (unsigned long long)epoch); + return 0; +} + static void topo_group_send_table_request(struct TOPO_GROUP* group, struct ETCP_CONN* conn) { - if (!group || !conn) return; struct TOPO_GROUP_CONN_ITEM* peer = topo_group_peer(group, conn->peer_node_id); - if (!peer || peer->conn != conn) { DEBUG_ERROR(DEBUG_CATEGORY_BGP, "table request without group peer"); return; } - if (!peer->exchange_id && (random_bytes((uint8_t*)&peer->exchange_id, sizeof(peer->exchange_id)) != 0 || !peer->exchange_id)) { - peer->exchange_id = 0; - DEBUG_ERROR(DEBUG_CATEGORY_BGP, "cannot generate group exchange id group=%016llx", (unsigned long long)group->group_id); - return; - } - DEBUG_DEBUG(DEBUG_CATEGORY_BGP, "Sending table request to %s grp=%016llx", conn->log_name, (unsigned long long)group->group_id); - struct TOPOMSG_TABLE_REQ* req = u_calloc(1, sizeof(struct TOPOMSG_TABLE_REQ)); - if (!req) { DEBUG_ERROR(DEBUG_CATEGORY_BGP, "table request allocation failed"); return; } - req->cmd = ETCP_ID_TOPO_ENTRY; - req->subcmd = TOPO_SUBCMD_REQUEST_TABLE; - req->group_id = group->group_id; - req->exchange_id = peer->exchange_id; - struct ll_entry* e = queue_entry_new(0); - if (!e) { DEBUG_ERROR(DEBUG_CATEGORY_BGP, "table request entry allocation failed"); u_free(req); return; } - e->dgram = (uint8_t*)req; e->len = sizeof(struct TOPOMSG_TABLE_REQ); - if (etcp_send(conn, e) != 0) { DEBUG_ERROR(DEBUG_CATEGORY_BGP, "cannot send table request"); u_free(req); queue_entry_free(e); } + if (!peer || !peer->accepted) return; + struct TOPOMSG_HEADER msg = topo_header(peer, TOPO_SUBCMD_REQUEST_TABLE); + topo_send_control(conn, &msg, sizeof(msg)); } -/* Запрашивает у пира членство в группе (JOIN_GROUP). */ static void topo_group_send_join_group(struct TOPO_GROUP* group, struct ETCP_CONN* conn) { - if (!group || !conn) return; - DEBUG_DEBUG(DEBUG_CATEGORY_BGP, "Sending join group request to %s grp=%016llx", conn->log_name, (unsigned long long)group->group_id); - struct TOPOMSG_JOIN_GROUP* req = u_calloc(1, sizeof(struct TOPOMSG_JOIN_GROUP)); - if (!req) return; - req->cmd = ETCP_ID_TOPO_ENTRY; - req->subcmd = TOPO_SUBCMD_JOIN_GROUP; - req->group_id = group->group_id; - struct ll_entry* e = queue_entry_new(0); - if (!e) { u_free(req); return; } - e->dgram = (uint8_t*)req; e->len = sizeof(struct TOPOMSG_JOIN_GROUP); - if (etcp_send(conn, e) != 0) { u_free(req); queue_entry_free(e); } + struct TOPO_GROUP_CONN_ITEM* peer = topo_group_peer(group, conn->peer_node_id); + if (!peer || peer->conn != conn || !peer->local_epoch) return; + struct TOPOMSG_JOIN_GROUP msg = { .h = topo_header(peer, TOPO_SUBCMD_JOIN_GROUP), .group_type = group->group_type }; + msg.h.dst_epoch = 0; + topo_send_control(conn, &msg, sizeof(msg)); } -/* Шлёт пиру RESYNC — просьбу заново анонсировать свои группы. */ static void topo_group_send_resync(struct ETCP_CONN* conn) { - if (!conn) return; - DEBUG_INFO(DEBUG_CATEGORY_BGP, "Sending RESYNC to %s", conn->log_name); - struct TOPOMSG_RESYNC* req = u_calloc(1, sizeof(struct TOPOMSG_RESYNC)); - if (!req) return; - req->cmd = ETCP_ID_TOPO_ENTRY; - req->subcmd = TOPO_SUBCMD_RESYNC; - struct ll_entry* e = queue_entry_new(0); - if (!e) { u_free(req); return; } - e->dgram = (uint8_t*)req; e->len = sizeof(struct TOPOMSG_RESYNC); - if (etcp_send(conn, e) != 0) { u_free(req); queue_entry_free(e); } -} - -/* Сигнализирует пиру об окончании начальной синхронизации таблицы. */ -static int topo_group_send_table_complete(struct TOPO_GROUP* group, struct ETCP_CONN* conn, uint64_t exchange_id) { - struct TOPOMSG_TABLE_REQ* req = u_calloc(1, sizeof(struct TOPOMSG_TABLE_REQ)); - if (!req) { DEBUG_ERROR(DEBUG_CATEGORY_BGP, "table complete allocation failed"); return -1; } - req->cmd = ETCP_ID_TOPO_ENTRY; - req->subcmd = TOPO_SUBCMD_TABLE_COMPLETE; - req->group_id = group->group_id; - req->exchange_id = exchange_id; - struct ll_entry* e = queue_entry_new(0); - if (!e) { DEBUG_ERROR(DEBUG_CATEGORY_BGP, "table complete entry allocation failed"); u_free(req); return -1; } - e->dgram = (uint8_t*)req; e->len = sizeof(struct TOPOMSG_TABLE_REQ); - if (etcp_send(conn, e) != 0) { - DEBUG_ERROR(DEBUG_CATEGORY_BGP, "cannot send table complete"); u_free(req); queue_entry_free(e); return -1; - } - return 0; + struct TOPOMSG_RESYNC msg = { .cmd = ETCP_ID_TOPO_ENTRY, .subcmd = TOPO_SUBCMD_RESYNC }; + topo_send_control(conn, &msg, sizeof(msg)); +} + +static int topo_group_send_table_complete(struct TOPO_GROUP* group, struct ETCP_CONN* conn) { + struct TOPO_GROUP_CONN_ITEM* peer = topo_group_peer(group, conn->peer_node_id); + if (!peer || !peer->accepted) return -1; + struct TOPOMSG_HEADER msg = topo_header(peer, TOPO_SUBCMD_TABLE_COMPLETE); + return topo_send_control(conn, &msg, sizeof(msg)); +} + +static void topo_group_reject(struct ETCP_CONN* conn, const struct TOPOMSG_HEADER* h, enum topo_join_reject reason) { + struct TOPOMSG_JOIN_GROUP msg = { .h = { .cmd = ETCP_ID_TOPO_ENTRY, .subcmd = TOPO_SUBCMD_JOIN_REJECT, + .group_id = h->group_id, .dst_epoch = h->src_epoch }, .group_type = reason }; + DEBUG_INFO(DEBUG_CATEGORY_BGP, "group JOIN rejected: group=%016llx peer=%016llx requested=%016llx reason=%u", + (unsigned long long)h->group_id, (unsigned long long)conn->peer_node_id, (unsigned long long)h->src_epoch, reason); + topo_send_control(conn, &msg, sizeof(msg)); } static int topo_group_add_to_senders(struct TOPO_GROUP* group, struct ETCP_CONN* conn); static bool topo_group_should_send_to(const struct TOPO_GROUP_NODE* nq, uint64_t target_id); static int topo_group_send_full_table(struct TOPO_GROUP* group, struct ETCP_CONN* conn); -static void topo_group_handle_request_table(struct TOPO_GROUP* group, struct ETCP_CONN* conn, uint64_t exchange_id); -static void topo_group_handle_join_group(struct TOPO_GROUP* group, struct ETCP_CONN* conn); -static void topo_group_send_resync(struct ETCP_CONN* conn); +static void topo_group_handle_request_table(struct TOPO_GROUP* group, struct ETCP_CONN* conn); +static void topo_group_handle_join_group(struct TOPO_GROUP* group, struct ETCP_CONN* conn, const struct TOPOMSG_JOIN_GROUP* msg); static void topo_group_handle_resync(struct UTUN_INSTANCE* instance, struct ETCP_CONN* conn); -/* Запоминает REQUEST_TABLE для группы, которой у нас ещё нет (CHAT-канал не загружен). */ -static void topo_group_remember_table_req(struct TOPO_GROUPS* g, uint64_t node_id, uint64_t group_id, uint64_t exchange_id) { - if (!g) return; - for (int i = 0; i < g->pending_table_req_count; i++) - if (g->pending_table_reqs[i].node_id == node_id && g->pending_table_reqs[i].group_id == group_id) { - g->pending_table_reqs[i].exchange_id = exchange_id; return; - } - if (g->pending_table_req_count >= TOPO_MAX_PENDING_TABLE_REQS) { - memmove(g->pending_table_reqs, g->pending_table_reqs + 1, - (TOPO_MAX_PENDING_TABLE_REQS - 1) * sizeof(g->pending_table_reqs[0])); - g->pending_table_req_count = TOPO_MAX_PENDING_TABLE_REQS - 1; - } - g->pending_table_reqs[g->pending_table_req_count].node_id = node_id; - g->pending_table_reqs[g->pending_table_req_count].group_id = group_id; - g->pending_table_reqs[g->pending_table_req_count].exchange_id = exchange_id; - g->pending_table_req_count++; - DEBUG_INFO(DEBUG_CATEGORY_BGP, "remembered table request node=%016llx grp=%016llx (group not created yet)", - (unsigned long long)node_id, (unsigned long long)group_id); -} - -/* Когда группа создана — отвечаем тем, кто просил её таблицу раньше. */ -static void topo_group_fulfill_table_reqs(struct TOPO_GROUPS* g, uint64_t group_id) { - if (!g) return; - struct TOPO_GROUP* group = topo_groups_find(g, group_id); - if (!group) return; - int w = 0; - for (int i = 0; i < g->pending_table_req_count; i++) { - struct topo_pending_table_req* req = &g->pending_table_reqs[i]; - if (req->group_id != group_id) { g->pending_table_reqs[w++] = *req; continue; } - struct ETCP_CONN* conn = instance_find_conn(g->instance, req->node_id); - if (conn && conn->links_up) { - DEBUG_INFO(DEBUG_CATEGORY_BGP, "fulfilling remembered table request node=%016llx grp=%016llx", - (unsigned long long)req->node_id, (unsigned long long)group_id); - topo_group_handle_request_table(group, conn, req->exchange_id); - } - } - g->pending_table_req_count = w; -} - /* Краткий DEBUG-дамп принятого NODEINFO (id/ver/счётчики/hops). */ static void nodeinfo_dump_log(const uint8_t* data, size_t len) { if (!data || len < sizeof(struct TOPOMSG_NODEINFO_PKT)) return; @@ -326,6 +282,8 @@ static const char* group_subcmd_name(uint8_t subcmd) { case TOPO_SUBCMD_TABLE_COMPLETE: return "TABLE_COMPLETE"; case TOPO_SUBCMD_ERR_GROUP_MISMATCH: return "ERR_GROUP_MISMATCH"; case TOPO_SUBCMD_JOIN_GROUP: return "JOIN_GROUP"; + case TOPO_SUBCMD_JOIN_ACCEPT: return "JOIN_ACCEPT"; + case TOPO_SUBCMD_JOIN_REJECT: return "JOIN_REJECT"; case TOPO_SUBCMD_RESYNC: return "RESYNC"; case TOPO_SUBCMD_LEAVE_GROUP: return "LEAVE_GROUP"; default: return "?"; @@ -340,20 +298,12 @@ static const char* group_subcmd_name(uint8_t subcmd) { /* Рассылает WITHDRAW всем BGP-пирам (кроме exclude) об удалении узла. */ static void topo_group_broadcast_withdraw(struct TOPO_GROUP* group, uint64_t node_id, uint64_t wd_source, struct ETCP_CONN* exclude) { if (!group) { DEBUG_ERROR(DEBUG_CATEGORY_BGP, "group is NULL"); return; } - struct TOPOMSG_WITHDRAW_PKT* pkt = u_calloc(1, sizeof(struct TOPOMSG_WITHDRAW_PKT)); - if (!pkt) return; - pkt->cmd = ETCP_ID_TOPO_ENTRY; pkt->subcmd = TOPO_SUBCMD_WITHDRAW; - pkt->group_id = group->group_id; pkt->node_id = node_id; pkt->wd_source = wd_source; - struct ll_entry* e = group->senders_list ? group->senders_list->head : NULL; - while (e) { - struct TOPO_GROUP_CONN_ITEM* item = (struct TOPO_GROUP_CONN_ITEM*)e->data; - if (item && item->conn && item->conn != exclude) { - struct ll_entry* copy = queue_entry_new(0); - if (copy) { copy->dgram = u_malloc(sizeof(*pkt)); if (copy->dgram) { memcpy(copy->dgram, pkt, sizeof(*pkt)); copy->len = sizeof(*pkt); if (etcp_send(item->conn, copy) != 0) { u_free(copy->dgram); queue_entry_free(copy); } } else { queue_entry_free(copy); } } - } - e = e->next; + for (struct ll_entry* e = group->senders_list ? group->senders_list->head : NULL; e; e = e->next) { + struct TOPO_GROUP_CONN_ITEM* peer = (struct TOPO_GROUP_CONN_ITEM*)e->data; + if (!peer->conn || !peer->accepted || peer->conn == exclude) continue; + struct TOPOMSG_WITHDRAW_PKT msg = { .h = topo_header(peer, TOPO_SUBCMD_WITHDRAW), .node_id = node_id, .wd_source = wd_source }; + topo_send_control(peer->conn, &msg, sizeof(msg)); } - u_free(pkt); } @@ -382,96 +332,97 @@ static int topo_group_peer_allowed(struct TOPO_GROUP* group, uint64_t 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; } + if (!entry) return; + if (!from_conn || !entry->dgram || entry->len < 2) { + DEBUG_WARN(DEBUG_CATEGORY_BGP, "invalid BGP packet length=%u", entry->len); goto done; + } struct UTUN_INSTANCE* instance = from_conn->instance; - if (!instance || !instance->topo_groups) { DEBUG_ERROR(DEBUG_CATEGORY_BGP, "invalid instance/topo_groups"); queue_dgram_free(entry); queue_entry_free(entry); return; } - uint8_t* data = entry->dgram; uint8_t cmd = data[0]; uint8_t subcmd = data[1]; - if (cmd != ETCP_ID_TOPO_ENTRY) { queue_dgram_free(entry); queue_entry_free(entry); return; } - - /* RESYNC — общий сигнал без group_id, обрабатывается до поиска группы */ - if (subcmd == TOPO_SUBCMD_RESYNC) { - topo_group_handle_resync(instance, from_conn); - queue_dgram_free(entry); queue_entry_free(entry); return; - } - - if ((subcmd == TOPO_SUBCMD_REQUEST_TABLE || subcmd == TOPO_SUBCMD_TABLE_COMPLETE) && - (entry->len != sizeof(struct TOPOMSG_TABLE_REQ) || !((struct TOPOMSG_TABLE_REQ*)data)->exchange_id)) { - DEBUG_WARN(DEBUG_CATEGORY_BGP, "invalid table control peer=%016llx cmd=%u length=%u", - (unsigned long long)from_conn->peer_node_id, subcmd, entry->len); - queue_dgram_free(entry); queue_entry_free(entry); return; - } - - uint64_t pkt_group_id = 0; - if (subcmd == TOPO_SUBCMD_NODEINFO && entry->len >= sizeof(struct TOPOMSG_NODEINFO_PKT)) { - pkt_group_id = ((struct TOPOMSG_NODEINFO_PKT*)data)->node.group_id; - } else if (subcmd == TOPO_SUBCMD_WITHDRAW && entry->len >= sizeof(struct TOPOMSG_WITHDRAW_PKT)) { - pkt_group_id = ((struct TOPOMSG_WITHDRAW_PKT*)data)->group_id; - } else if (subcmd == TOPO_SUBCMD_REQUEST_TABLE && entry->len >= sizeof(struct TOPOMSG_TABLE_REQ)) { - pkt_group_id = ((struct TOPOMSG_TABLE_REQ*)data)->group_id; - } else if (subcmd == TOPO_SUBCMD_TABLE_COMPLETE && entry->len >= sizeof(struct TOPOMSG_TABLE_REQ)) { - pkt_group_id = ((struct TOPOMSG_TABLE_REQ*)data)->group_id; - } else if (subcmd == TOPO_SUBCMD_ERR_GROUP_MISMATCH && entry->len >= sizeof(struct TOPOMSG_ERR_GROUP_MISMATCH)) { - pkt_group_id = ((struct TOPOMSG_ERR_GROUP_MISMATCH*)data)->group_id; - } else if ((subcmd == TOPO_SUBCMD_JOIN_GROUP || subcmd == TOPO_SUBCMD_LEAVE_GROUP) && - entry->len == sizeof(struct TOPOMSG_JOIN_GROUP)) { - pkt_group_id = ((struct TOPOMSG_JOIN_GROUP*)data)->group_id; - } - - struct TOPO_GROUP* group = topo_groups_find(instance->topo_groups, pkt_group_id); - if (!group) { - if (subcmd == TOPO_SUBCMD_REQUEST_TABLE) - topo_group_remember_table_req(instance->topo_groups, from_conn->peer_node_id, pkt_group_id, - ((struct TOPOMSG_TABLE_REQ*)data)->exchange_id); - DEBUG_DEBUG(DEBUG_CATEGORY_BGP, "BGP recv %s from %s: group %016llx not found, dropping", group_subcmd_name(subcmd), from_conn->log_name, (unsigned long long)pkt_group_id); - queue_dgram_free(entry); queue_entry_free(entry); return; - } - - 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 (!instance || !instance->topo_groups) { DEBUG_ERROR(DEBUG_CATEGORY_BGP, "invalid instance/topo_groups"); goto done; } + const uint8_t* data = entry->dgram; + uint8_t subcmd = data[1]; + if (data[0] != ETCP_ID_TOPO_ENTRY) { DEBUG_WARN(DEBUG_CATEGORY_BGP, "invalid BGP command=%u", data[0]); goto done; } + if (subcmd == TOPO_SUBCMD_RESYNC && entry->len == sizeof(struct TOPOMSG_RESYNC)) { + topo_group_handle_resync(instance, from_conn); goto done; + } + size_t size = sizeof(struct TOPOMSG_HEADER); + if (subcmd == TOPO_SUBCMD_JOIN_GROUP || subcmd == TOPO_SUBCMD_JOIN_ACCEPT || subcmd == TOPO_SUBCMD_JOIN_REJECT) + size = sizeof(struct TOPOMSG_JOIN_GROUP); + else if (subcmd == TOPO_SUBCMD_NODEINFO) size = sizeof(struct TOPOMSG_NODEINFO_PKT); + else if (subcmd == TOPO_SUBCMD_WITHDRAW) size = sizeof(struct TOPOMSG_WITHDRAW_PKT); + else if (subcmd == TOPO_SUBCMD_ERR_GROUP_MISMATCH) size = sizeof(struct TOPOMSG_ERR_GROUP_MISMATCH); + else if (subcmd != TOPO_SUBCMD_REQUEST_TABLE && subcmd != TOPO_SUBCMD_TABLE_COMPLETE && subcmd != TOPO_SUBCMD_LEAVE_GROUP) { + DEBUG_WARN(DEBUG_CATEGORY_BGP, "unknown BGP command=%u", subcmd); goto done; + } + if (entry->len < size || (subcmd != TOPO_SUBCMD_NODEINFO && entry->len != size)) { + DEBUG_WARN(DEBUG_CATEGORY_BGP, "invalid BGP length=%u cmd=%u expected=%zu", entry->len, subcmd, size); goto done; + } + const struct TOPOMSG_HEADER* h = (const struct TOPOMSG_HEADER*)data; + struct TOPO_GROUP* group = topo_groups_find(instance->topo_groups, h->group_id); + if (subcmd == TOPO_SUBCMD_JOIN_GROUP) { + if (!h->src_epoch || h->dst_epoch) { DEBUG_WARN(DEBUG_CATEGORY_BGP, "invalid JOIN epochs"); goto done; } + if (!group || group->stopping) topo_group_reject(from_conn, h, TOPO_JOIN_NO_GROUP); + else if (!topo_group_peer_allowed(group, from_conn->peer_node_id)) topo_group_reject(from_conn, h, TOPO_JOIN_NOT_MEMBER); + else if (((const struct TOPOMSG_JOIN_GROUP*)data)->group_type != group->group_type) + topo_group_reject(from_conn, h, TOPO_JOIN_WRONG_TYPE); + else topo_group_handle_join_group(group, from_conn, (const struct TOPOMSG_JOIN_GROUP*)data); + goto done; + } + struct TOPO_GROUP_CONN_ITEM* peer = topo_group_peer(group, from_conn->peer_node_id); + /* До получения JOIN ответа отменяем только ещё не согласованное предложение + * с тем же src_epoch. У READY-сессии LEAVE всегда проверяет оба поколения. */ + if (subcmd == TOPO_SUBCMD_LEAVE_GROUP && peer && peer->conn == from_conn && !peer->accepted && + !h->dst_epoch && h->src_epoch && h->src_epoch == peer->peer_epoch) { + topo_group_remove_conn(group, from_conn, TOPO_REMOVE_REMOTE_LEAVE); goto done; + } + if (!peer || peer->conn != from_conn || !h->dst_epoch || h->dst_epoch != peer->local_epoch) goto stale; + if (subcmd == TOPO_SUBCMD_JOIN_REJECT) { + DEBUG_WARN(DEBUG_CATEGORY_BGP, "group JOIN refused: group=%016llx peer=%016llx local=%016llx reason=%u", + (unsigned long long)h->group_id, (unsigned long long)peer->node_id, (unsigned long long)peer->local_epoch, + ((const struct TOPOMSG_JOIN_GROUP*)data)->group_type); + topo_group_remove_conn(group, from_conn, TOPO_REMOVE_REJECTED); goto done; + } + if (!h->src_epoch || (peer->peer_epoch && peer->peer_epoch != h->src_epoch)) goto stale; 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; + DEBUG_WARN(DEBUG_CATEGORY_BGP, "group peer no longer a member group=%016llx peer=%016llx", + (unsigned long long)h->group_id, (unsigned long long)peer->node_id); goto done; } - - if ((subcmd == TOPO_SUBCMD_LEAVE_GROUP || (group->group_type == TOPO_GROUP_TYPE_UTUN && - (subcmd == TOPO_SUBCMD_NODEINFO || subcmd == TOPO_SUBCMD_WITHDRAW))) && - !topo_group_has_sender(group, from_conn)) { - DEBUG_DEBUG(DEBUG_CATEGORY_BGP, "ignore %s from detached peer=%016llx group=%016llx", - group_subcmd_name(subcmd), (unsigned long long)from_conn->peer_node_id, (unsigned long long)pkt_group_id); - 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, ((struct TOPOMSG_TABLE_REQ*)data)->exchange_id); - else if (subcmd == TOPO_SUBCMD_JOIN_GROUP) topo_group_handle_join_group(group, from_conn); - else if (subcmd == TOPO_SUBCMD_LEAVE_GROUP) { - DEBUG_INFO(DEBUG_CATEGORY_BGP, "peer left group: peer=%016llx group=%016llx", - (unsigned long long)from_conn->peer_node_id, (unsigned long long)group->group_id); - topo_group_remove_conn(group, from_conn, TOPO_REMOVE_REMOTE_LEAVE); - } - else if (subcmd == TOPO_SUBCMD_TABLE_COMPLETE) { - struct TOPO_GROUP_CONN_ITEM* peer = topo_group_peer(group, from_conn->peer_node_id); - uint64_t id = ((struct TOPOMSG_TABLE_REQ*)data)->exchange_id; - if (!peer || peer->conn != from_conn || peer->exchange_id != id) { - DEBUG_DEBUG(DEBUG_CATEGORY_BGP, "ignore stale TABLE_COMPLETE group=%016llx peer=%016llx exchange=%016llx", - (unsigned long long)group->group_id, (unsigned long long)from_conn->peer_node_id, (unsigned long long)id); - } else if (!peer->table_received) { - peer->table_received = 1; - topo_group_log_exchange(group, peer); + if (subcmd == TOPO_SUBCMD_JOIN_ACCEPT) { + if (((const struct TOPOMSG_JOIN_GROUP*)data)->group_type != group->group_type) { + DEBUG_WARN(DEBUG_CATEGORY_BGP, "JOIN_ACCEPT type mismatch"); goto done; } - } - else if (subcmd == TOPO_SUBCMD_ERR_GROUP_MISMATCH) { - if (entry->len >= sizeof(struct TOPOMSG_ERR_GROUP_MISMATCH)) { - struct TOPOMSG_ERR_GROUP_MISMATCH* err = (struct TOPOMSG_ERR_GROUP_MISMATCH*)data; - DEBUG_ERROR(DEBUG_CATEGORY_BGP, "Group mismatch from %s: expected=%u received=0x%02x group=%016llx", - from_conn->log_name, err->expected_type, err->received_flags, (unsigned long long)err->group_id); + if (!peer->accepted) { + peer->peer_epoch = h->src_epoch; peer->accepted = 1; + topo_group_peer_progressed(group, peer->node_id); + DEBUG_INFO(DEBUG_CATEGORY_BGP, "group JOIN accepted: group=%016llx peer=%016llx local=%016llx remote=%016llx", + (unsigned long long)h->group_id, (unsigned long long)peer->node_id, + (unsigned long long)peer->local_epoch, (unsigned long long)peer->peer_epoch); + topo_group_send_table_request(group, from_conn); } + goto done; } - + if (subcmd == TOPO_SUBCMD_LEAVE_GROUP && peer->peer_epoch) { + topo_group_remove_conn(group, from_conn, TOPO_REMOVE_REMOTE_LEAVE); goto done; + } + if (!peer->accepted || !peer->peer_epoch) goto stale; + DEBUG_DEBUG(DEBUG_CATEGORY_BGP, "BGP recv %s peer=%016llx group=%016llx local=%016llx remote=%016llx len=%u", + group_subcmd_name(subcmd), (unsigned long long)peer->node_id, (unsigned long long)h->group_id, + (unsigned long long)peer->local_epoch, (unsigned long long)peer->peer_epoch, entry->len); + 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); + else if (subcmd == TOPO_SUBCMD_TABLE_COMPLETE && !peer->table_received) { + peer->table_received = 1; topo_group_log_exchange(group, peer); + } else if (subcmd == TOPO_SUBCMD_ERR_GROUP_MISMATCH) { + const struct TOPOMSG_ERR_GROUP_MISMATCH* err = (const struct TOPOMSG_ERR_GROUP_MISMATCH*)data; + DEBUG_ERROR(DEBUG_CATEGORY_BGP, "group mismatch peer=%016llx expected=%u received=%u", + (unsigned long long)peer->node_id, err->expected_type, err->received_flags); + } + goto done; +stale: + DEBUG_DEBUG(DEBUG_CATEGORY_BGP, "ignore unmatched group message cmd=%u group=%016llx peer=%016llx src=%016llx dst=%016llx", + subcmd, (unsigned long long)h->group_id, (unsigned long long)from_conn->peer_node_id, + (unsigned long long)h->src_epoch, (unsigned long long)h->dst_epoch); +done: queue_dgram_free(entry); queue_entry_free(entry); } @@ -492,11 +443,10 @@ static void topo_group_conn_status(struct ETCP_CONN* conn, int status, void* arg for (struct ll_entry* se = g->senders_list->head; se; se = se->next) { struct TOPO_GROUP_CONN_ITEM* peer = (struct TOPO_GROUP_CONN_ITEM*)se->data; if (peer->conn != conn) continue; - peer->exchange_id = 0; peer->table_received = 0; peer->table_sent = 0; - topo_group_peer_progressed(g, peer->node_id); + if (topo_group_reset_exchange(peer) < 0) continue; DEBUG_INFO(DEBUG_CATEGORY_BGP,"BGP session resync: peer=%016llx group=%016llx", (unsigned long long)conn->peer_node_id,(unsigned long long)g->group_id); - topo_group_send_join_group(g,conn); topo_group_send_table_request(g,conn); + topo_group_send_join_group(g, conn); break; } } @@ -823,7 +773,6 @@ struct TOPO_GROUP* topo_groups_create_group(struct TOPO_GROUPS* g, uint64_t grou DEBUG_INFO(DEBUG_CATEGORY_BGP, "Group created: group_id=%016llx type=%u ch_id=%s", (unsigned long long)group_id, group_type, group->channel_id); if (group_type == TOPO_GROUP_TYPE_CHAT && chat_setting_get_int(group->instance, "group_autoconnect", 1)) topo_group_connect_init(group); - topo_group_fulfill_table_reqs(g, group_id); return group; } @@ -923,8 +872,9 @@ static int topo_group_begin_conn(struct TOPO_GROUP* group, struct ETCP_CONN* con 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); + peer = topo_group_peer(group, conn->peer_node_id); + if (topo_group_reset_exchange(peer) < 0) return -1; topo_group_send_join_group(group, conn); - topo_group_send_table_request(group, conn); topo_group_connect_on_up(group, conn); return 0; } @@ -1100,15 +1050,11 @@ int topo_group_remove_path(struct TOPO_GROUP_NODE* nq, struct ETCP_CONN* conn) { /* Шлёт пиру ошибку несоответствия типа группы (ERR_GROUP_MISMATCH). */ static void topo_group_send_err_group_mismatch(struct TOPO_GROUP* group, struct ETCP_CONN* conn, uint8_t expected_type, uint8_t received_flags) { if (!group || !conn) return; - struct TOPOMSG_ERR_GROUP_MISMATCH* pkt = u_calloc(1, sizeof(struct TOPOMSG_ERR_GROUP_MISMATCH)); - if (!pkt) return; - pkt->cmd = ETCP_ID_TOPO_ENTRY; pkt->subcmd = TOPO_SUBCMD_ERR_GROUP_MISMATCH; - pkt->group_id = group->group_id; - pkt->expected_type = expected_type; pkt->received_flags = received_flags; - struct ll_entry* e = queue_entry_new(0); - if (!e) { u_free(pkt); return; } - e->dgram = (uint8_t*)pkt; e->len = sizeof(struct TOPOMSG_ERR_GROUP_MISMATCH); - if (etcp_send(conn, e) != 0) { DEBUG_ERROR(DEBUG_CATEGORY_BGP, "etcp_send ERR_GROUP_MISMATCH failed"); u_free(pkt); queue_entry_free(e); } + struct TOPO_GROUP_CONN_ITEM* peer = topo_group_peer(group, conn->peer_node_id); + if (!peer || !peer->accepted) return; + struct TOPOMSG_ERR_GROUP_MISMATCH msg = { .h = topo_header(peer, TOPO_SUBCMD_ERR_GROUP_MISMATCH), + .expected_type = expected_type, .received_flags = received_flags }; + topo_send_control(conn, &msg, sizeof(msg)); } /* Обрабатывает NODEINFO: проверка подписи/версии, обновление узла, пути, роутинг, форвард. */ @@ -1158,8 +1104,8 @@ int topo_group_process_nodeinfo(struct TOPO_GROUP* group, struct ETCP_CONN* from paths = nodeinfo1->paths; } else { is_new_node = 1; } - const uint8_t* ser_data = data + 2; - size_t ser_len = len - 2; + const uint8_t* ser_data = data + sizeof(struct TOPOMSG_HEADER); + size_t ser_len = len - sizeof(struct TOPOMSG_HEADER); struct TOPO_NODE* new_ni = NULL; struct TOPO_NODESUBNETS* new_subnets = NULL; uint64_t* new_hop_list = NULL; uint8_t new_hop_count = 0; @@ -1350,7 +1296,10 @@ int topo_group_send_nodeinfo(struct TOPO_GROUP* group, struct TOPO_GROUP_NODE* n uint8_t* p = u_malloc(max_sz); if (!p) { DEBUG_ERROR(DEBUG_CATEGORY_BGP, "send_nodeinfo: allocation failed"); return -1; } - p[0] = ETCP_ID_TOPO_ENTRY; p[1] = TOPO_SUBCMD_NODEINFO; + struct TOPO_GROUP_CONN_ITEM* peer = topo_group_peer(group, conn->peer_node_id); + if (!peer || !peer->accepted) { u_free(p); return 0; } /* initial exchange sends the current state after ACCEPT */ + struct TOPOMSG_HEADER h = topo_header(peer, TOPO_SUBCMD_NODEINFO); + memcpy(p, &h, sizeof(h)); struct TOPO_NODE* sni = topo_node_registry_find(group->instance->topo_groups, node->node_id); if (!sni) { @@ -1367,7 +1316,7 @@ int topo_group_send_nodeinfo(struct TOPO_GROUP* group, struct TOPO_GROUP_NODE* n } DEBUG_DEBUG(DEBUG_CATEGORY_BGP, "send_nodeinfo: node %016llx ver=%d grp=%016llx to conn=%p name='%s'", (unsigned long long)node->node_id, sni->ver, (unsigned long long)group->group_id, (void*)conn, conn->log_name); uint8_t sflags = (group->group_type == TOPO_GROUP_TYPE_CHAT) ? 0 : TOPO_FLAG_SEND_SUBNETS; - int ser_len = topo_node_serialize(sni, node, group->group_id, sflags, p + 2, max_sz - 2, cumulative_rtt); + int ser_len = topo_node_serialize(sni, node, group->group_id, sflags, p + sizeof(h), max_sz - sizeof(h), cumulative_rtt); if (ser_len < 0) { DEBUG_ERROR(DEBUG_CATEGORY_BGP, "send_nodeinfo: serialize failed for node %016llx", (unsigned long long)node->node_id); u_free(p); return -1; @@ -1375,7 +1324,7 @@ int topo_group_send_nodeinfo(struct TOPO_GROUP* group, struct TOPO_GROUP_NODE* n struct ll_entry* e = queue_entry_new(0); if (!e) { DEBUG_ERROR(DEBUG_CATEGORY_BGP, "send_nodeinfo: entry allocation failed"); u_free(p); return -1; } - e->dgram = p; e->len = (size_t)ser_len + 2; + e->dgram = p; e->len = (size_t)ser_len + sizeof(h); if (etcp_send(conn, e) != 0) { DEBUG_ERROR(DEBUG_CATEGORY_BGP, "send_nodeinfo: etcp_send FAILED for node %016llx to %s", @@ -1437,6 +1386,7 @@ int topo_group_join_peer(struct TOPO_GROUP* group, struct NODE_CONN_DIRECT* hand } struct topo_leave { + struct TOPOMSG_HEADER message; struct UTUN_INSTANCE* instance; uint64_t group_id; struct NODE_CONN_DIRECT* handle; @@ -1470,15 +1420,7 @@ static void topo_leave_ready(struct ll_queue* queue, void* arg) { (unsigned long long)leave->group_id); topo_leave_finish(leave); return; } - struct ll_entry* packet = ll_alloc_lldgram(sizeof(struct TOPOMSG_JOIN_GROUP)); - int result = -1; - if (packet) { - struct TOPOMSG_JOIN_GROUP msg = { .cmd = ETCP_ID_TOPO_ENTRY, .subcmd = TOPO_SUBCMD_LEAVE_GROUP, - .group_id = leave->group_id }; - memcpy(packet->dgram, &msg, sizeof(msg)); packet->len = sizeof(msg); - result = etcp_send(conn, packet); - if (result != 0) { queue_dgram_free(packet); queue_entry_free(packet); } - } + int result = topo_send_control(conn, &leave->message, sizeof(leave->message)); if (result != 0) DEBUG_ERROR(DEBUG_CATEGORY_BGP, "cannot send LEAVE_GROUP group=%016llx", (unsigned long long)leave->group_id); else DEBUG_INFO(DEBUG_CATEGORY_BGP, "LEAVE_GROUP sent: peer=%016llx group=%016llx", (unsigned long long)node_conn_direct_node_id(leave->handle), (unsigned long long)leave->group_id); @@ -1497,6 +1439,7 @@ int topo_group_leave_peer(struct TOPO_GROUP* group, struct NODE_CONN_DIRECT* han DEBUG_ERROR(DEBUG_CATEGORY_BGP, "LEAVE_GROUP allocation failed"); topo_group_remove_conn(group, conn, TOPO_REMOVE_LOCAL_LEAVE); return -1; } + leave->message = topo_header(topo_group_peer(group, conn->peer_node_id), TOPO_SUBCMD_LEAVE_GROUP); leave->instance = group->instance; leave->group_id = group->group_id; leave->queue = conn->send_input_q; if (node_conn_direct_open(group->instance, conn->peer_node_id, topo_leave_event, leave, &leave->handle, NULL) == NCD_ERR) { DEBUG_ERROR(DEBUG_CATEGORY_BGP, "cannot retain NCD handle for LEAVE_GROUP"); u_free(leave); @@ -1526,7 +1469,7 @@ static int topo_group_send_full_table(struct TOPO_GROUP* group, struct ETCP_CONN } /* Обрабатывает REQUEST_TABLE: шлёт свой nodeinfo + полную таблицу + TABLE_COMPLETE. */ -static void topo_group_handle_request_table(struct TOPO_GROUP* group, struct ETCP_CONN* conn, uint64_t exchange_id) { +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", @@ -1534,25 +1477,28 @@ static void topo_group_handle_request_table(struct TOPO_GROUP* group, struct ETC return; } if (topo_group_send_nodeinfo(group, group->local_node, conn, 0) < 0 || topo_group_send_full_table(group, conn) < 0) { - DEBUG_ERROR(DEBUG_CATEGORY_BGP, "table exchange incomplete group=%016llx peer=%016llx exchange=%016llx", - (unsigned long long)group->group_id, (unsigned long long)conn->peer_node_id, (unsigned long long)exchange_id); + DEBUG_ERROR(DEBUG_CATEGORY_BGP, "table exchange incomplete group=%016llx peer=%016llx", + (unsigned long long)group->group_id, (unsigned long long)conn->peer_node_id); return; } /* senders_list заполняется только через topo_group_new_conn (add + BGP); здесь не добавляем, иначе new_conn (на JOIN_GROUP) упирается в дедуп и не шлёт TABLE_REQ обратно. */ - if (topo_group_send_table_complete(group, conn, exchange_id) < 0) return; + if (topo_group_send_table_complete(group, conn) < 0) return; struct TOPO_GROUP_CONN_ITEM* peer = topo_group_peer(group, conn->peer_node_id); if (peer && peer->conn == conn && !peer->table_sent) { peer->table_sent = 1; topo_group_log_exchange(group, peer); } } /* Обрабатывает JOIN_GROUP: для UTUN — добавить запросившего и инициировать BGP обратно. */ -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); - /* Проверка таблицы участников общая для входящих и локальных запросов. */ - int existing = topo_group_has_sender(group, conn); - if (topo_group_begin_conn(group, conn, 0) == 0 && existing) - topo_group_send_table_request(group, conn); +static void topo_group_handle_join_group(struct TOPO_GROUP* group, struct ETCP_CONN* conn, const struct TOPOMSG_JOIN_GROUP* msg) { + if (topo_group_begin_conn(group, conn, 0) < 0) { topo_group_reject(conn, &msg->h, TOPO_JOIN_UNAVAILABLE); return; } + struct TOPO_GROUP_CONN_ITEM* peer = topo_group_peer(group, conn->peer_node_id); + if (peer->peer_epoch && peer->peer_epoch != msg->h.src_epoch && peer->accepted) { + if (topo_group_reset_exchange(peer) < 0) return; + topo_group_send_join_group(group, conn); + } + peer->peer_epoch = msg->h.src_epoch; + struct TOPOMSG_JOIN_GROUP reply = { .h = topo_header(peer, TOPO_SUBCMD_JOIN_ACCEPT), .group_type = group->group_type }; + topo_send_control(conn, &reply, sizeof(reply)); } /* Обрабатывает RESYNC: если мы VPN-клиент пира — заново шлём JOIN_GROUP по группам. */ diff --git a/src/routing_layer/topo_group.h b/src/routing_layer/topo_group.h index 9bef7a5d..740bb8aa 100644 --- a/src/routing_layer/topo_group.h +++ b/src/routing_layer/topo_group.h @@ -91,20 +91,27 @@ 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 обязаны использовать общую проверку членства. +/* Групповая сессия принадлежит паре (group, peer), транспорт разделяется через NCD. + * JOIN_GROUP не добавляет мембера: CHAT проверяет локальную таблицу участников + * (chat_join.h / member_sync). Поздний commit мембера повторяет JOIN на готовом NCD. * - * REQUEST_TABLE содержит ненулевой exchange_id, сгенерированный запросившей - * стороной для текущего обмена (группа, пир). TABLE_COMPLETE возвращает тот же - * id только после успешной постановки всех NODEINFO в транспорт. Дубликаты - * запроса используют прежний id; новое присоединение и REINIT создают новый. - * Ответ с другим id не завершает текущий обмен. Формат изменён без совместимости. + * Каждая сторона создаёт случайный ненулевой local_epoch. JOIN несёт src_epoch, + * dst_epoch=0 и тип группы. Получатель проверяет группу/членство/тип, посылает свой + * JOIN и JOIN_ACCEPT(src=local_epoch, dst=полученный src). JOIN_REJECT возвращает + * dst исходного запроса и причину; чужой/старый ответ не меняет текущую сессию. + * Одновременные JOIN допустимы; повтор того же JOIN не сбрасывает состояние. + * LEAVE до получения поколения пира имеет dst=0 и отменяет только предложение + * с совпадающим src у ещё не согласованного получателя. После ACCEPT нужны оба id. + * После ACCEPT стороны запрашивают таблицы. Все остальные сообщения содержат + * оба поколения и принимаются только в согласованной сессии. TABLE_COMPLETE + * отправляется после передачи транспорту всей таблицы отправителя. + * READY = живой транспорт + ACCEPT + table_sent + table_received. + * + * Новое присоединение и REINIT создают новое поколение; REINIT сохраняет пути. + * Поколения защищают сообщения, а доступность существующего пути определяется + * его соединением. Recovery дополнительно ожидает READY группового пира. + * JOIN_READY чата означает добавление мембера и не заменяет этот обмен. + * Формат протокола изменяется без обратной совместимости. */ // Sub-команды @@ -113,71 +120,47 @@ void topo_fire_nodeinfo_cbk(struct UTUN_INSTANCE* instance, struct TOPO_GROUP* g #define TOPO_SUBCMD_WITHDRAW 0x06 // узел стал недоступен #define TOPO_SUBCMD_TABLE_COMPLETE 0x0B // завершение начальной синхронизации таблицы #define TOPO_SUBCMD_ERR_GROUP_MISMATCH 0x0C // ошибка несоответствия типа группы -#define TOPO_SUBCMD_JOIN_GROUP 0x0D // запрос членства в группе +#define TOPO_SUBCMD_JOIN_GROUP 0x0D // согласование участия в группе #define TOPO_SUBCMD_RESYNC 0x0E // общий запрос: «я переподключился, заново анонсируй свои группы» #define TOPO_SUBCMD_LEAVE_GROUP 0x0F // завершить участие, сохранив общий транспорт +#define TOPO_SUBCMD_JOIN_ACCEPT 0x10 +#define TOPO_SUBCMD_JOIN_REJECT 0x11 + +enum topo_join_reject { TOPO_JOIN_NO_GROUP = 1, TOPO_JOIN_NOT_MEMBER, TOPO_JOIN_WRONG_TYPE, TOPO_JOIN_UNAVAILABLE }; + #define MAX_HOPS 16 #define BGP_NODES_HASH_SIZE 256 -/** - * @brief Пакет с информацией об узле (NODEINFO) - */ +/* Заголовок всех групповых сообщений. RESYNC — единственное исключение. */ +struct TOPOMSG_HEADER { + uint8_t cmd, subcmd; + uint64_t group_id; + uint64_t src_epoch, dst_epoch; +} __attribute__((packed)); + struct TOPOMSG_NODEINFO_PKT { - uint8_t cmd; // ETCP_ID_TOPO_ENTRY - uint8_t subcmd; // TOPO_SUBCMD_NODEINFO + struct TOPOMSG_HEADER h; struct TOPOMSG_NODE node; } __attribute__((packed)); -/** - * @brief Пакет WITHDRAW (фиксированный) - */ struct TOPOMSG_WITHDRAW_PKT { - uint8_t cmd; - uint8_t subcmd; - uint64_t group_id; // идентификатор группы - uint64_t node_id; // удаляемый узел (который стал недоступен) - uint64_t wd_source; // узел который инициировал withdraw (при удалении он должен быть в hoplist или = current node_id) -} __attribute__((packed)); - -/** - * @brief Пакет запроса полной таблицы / завершения синхронизации - */ -struct TOPOMSG_TABLE_REQ { - uint8_t cmd; - uint8_t subcmd; - uint64_t group_id; // идентификатор группы - uint64_t exchange_id; // ненулевой id запроса; TABLE_COMPLETE возвращает его без изменений + struct TOPOMSG_HEADER h; + uint64_t node_id, wd_source; } __attribute__((packed)); -/** - * @brief Пакет запроса членства в группе (отправитель просит добавить его в группу) - */ struct TOPOMSG_JOIN_GROUP { - uint8_t cmd; - uint8_t subcmd; - uint64_t group_id; // идентификатор группы + struct TOPOMSG_HEADER h; + uint8_t group_type; /* JOIN/ACCEPT: тип группы; REJECT: topo_join_reject */ } __attribute__((packed)); -/** - * @brief Общий запрос ресинхронизации: отправитель (пассивная сторона) переподключился - * и просит пира заново анонсировать свои группы (без group_id — группа решается - * принимающей стороной по её конфигу). - */ struct TOPOMSG_RESYNC { - uint8_t cmd; // ETCP_ID_TOPO_ENTRY - uint8_t subcmd; // TOPO_SUBCMD_RESYNC + uint8_t cmd, subcmd; } __attribute__((packed)); -/** - * @brief Пакет ошибки несоответствия типа группы - */ struct TOPOMSG_ERR_GROUP_MISMATCH { - uint8_t cmd; - uint8_t subcmd; - uint64_t group_id; // идентификатор группы - uint8_t expected_type; - uint8_t received_flags; + struct TOPOMSG_HEADER h; + uint8_t expected_type, received_flags; } __attribute__((packed)); struct TOPO_GROUP_CONN_ITEM { @@ -188,8 +171,9 @@ struct TOPO_GROUP_CONN_ITEM { struct TOPO_PEER_REQUEST* requests; uint64_t progress; // время последнего прогресса обмена (timebase) uint8_t retained; // участие запрошено постоянным владельцем или достигло READY - uint64_t exchange_id; // текущий исходящий REQUEST_TABLE - uint8_t table_received; // TABLE_COMPLETE для exchange_id принят + uint64_t local_epoch, peer_epoch; // поколения сторон текущего присоединения + uint8_t accepted; // JOIN_ACCEPT подтвердил local_epoch + uint8_t table_received; // TABLE_COMPLETE текущей сессии принят uint8_t table_sent; // ответ на запрос пира целиком передан транспорту }; @@ -224,16 +208,6 @@ struct TOPO_GROUP { uint8_t stopping; }; -/* Отложенный REQUEST_TABLE: пир запросил таблицу группы, которой у нас ещё нет. - * Запоминаем (node_id, group_id) и отвечаем, когда группа будет создана. */ -#define TOPO_MAX_PENDING_TABLE_REQS 32 - -struct topo_pending_table_req { - uint64_t node_id; /* кто запросил таблицу */ - uint64_t group_id; /* какую группу */ - uint64_t exchange_id; -}; - /** * @brief Контейнер всех групп топологии экземпляра */ @@ -248,8 +222,6 @@ struct TOPO_GROUPS { struct memory_pool* v4_subnet_pool; struct memory_pool* v6_subnet_pool; topo_node_updated_fn node_updated_cb; /* optional — set by chatgui's member_sync */ - struct topo_pending_table_req pending_table_reqs[TOPO_MAX_PENDING_TABLE_REQS]; - int pending_table_req_count; }; /** @@ -314,7 +286,7 @@ int topo_group_new_conn(struct TOPO_GROUP* group, struct ETCP_CONN* conn); /* Готовность относится только к паре (group, peer): собственный снимок отправлен, * снимок пира принят до TABLE_COMPLETE текущего REQUEST_TABLE. Транспортный UP * этого не гарантирует. Повторное new_conn не сбрасывает готовность; REINIT - * начинает новый обмен с новым exchange_id. Старый TABLE_COMPLETE игнорируется. + * начинает новый обмен с новым поколением. Старый TABLE_COMPLETE игнорируется. * Это состояние обмена таблицами, а не добавление мембера (см. chat_join.h). */ int topo_group_peer_ready(const struct TOPO_GROUP* group, uint64_t peer_id); @@ -341,7 +313,8 @@ enum topo_group_remove_reason { TOPO_REMOVE_TRANSPORT_DOWN, TOPO_REMOVE_LOCAL_LEAVE, TOPO_REMOVE_REMOTE_LEAVE, - TOPO_REMOVE_MEMBER_INVALID + TOPO_REMOVE_MEMBER_INVALID, + TOPO_REMOVE_REJECTED }; /* Только потеря транспорта запускает recovery каскадно потерянных маршрутов. */ void topo_group_remove_conn(struct TOPO_GROUP* group, struct ETCP_CONN* conn, enum topo_group_remove_reason reason); diff --git a/tests/test_group_exchange.c b/tests/test_group_exchange.c index 3299a9cf..72901109 100644 --- a/tests/test_group_exchange.c +++ b/tests/test_group_exchange.c @@ -1,5 +1,6 @@ #include #include +#include #include "utun_instance.h" #include "routing_layer/topo_group.h" #include "transport_layer/etcp.h" @@ -13,21 +14,53 @@ static struct TOPO_GROUP_CONN_ITEM* peer(struct TOPO_GROUP* group) { return (struct TOPO_GROUP_CONN_ITEM*)group->senders_list->head->data; } -static void receive(struct ETCP_CONN* conn, uint64_t group_id, uint8_t cmd, uint64_t exchange_id, size_t length) { - struct TOPOMSG_TABLE_REQ msg = { .cmd = ETCP_ID_TOPO_ENTRY, .subcmd = cmd, - .group_id = group_id, .exchange_id = exchange_id }; +static void packet(struct ETCP_CONN* conn, const void* data, size_t size) { + struct ll_entry* e = ll_alloc_lldgram(size); assert(e); + memcpy(e->dgram, data, size); e->len = size; + conn->instance->api_bindings.callbacks[ETCP_ID_TOPO_ENTRY](conn, e); +} + +static void nodeinfo(struct TOPO_GROUP* group, struct ETCP_CONN* conn, struct TOPO_NODE* node, + struct SC_MYKEYS* keys, uint64_t generation, uint64_t timestamp) { + node->timestamp = timestamp; + EVP_PKEY* key = EVP_PKEY_new_raw_private_key(EVP_PKEY_ED25519, NULL, keys->private_key, 32); + size_t key_len = 32; + assert(key && EVP_PKEY_get_raw_public_key(key, node->ed25519_public_key, &key_len) == 1); + EVP_PKEY_free(key); + uint8_t signed_data[TOPO_SIG_MSG_MAX_SIZE]; + int length = topo_node_build_sig_msg(node, signed_data, sizeof(signed_data)); assert(length > 0); + assert(sc_ed25519_sign(keys->private_key, signed_data, length, node->x25519_self_sig) == SC_OK); + struct TOPOMSG_HEADER h = { .cmd = ETCP_ID_TOPO_ENTRY, .subcmd = TOPO_SUBCMD_NODEINFO, + .group_id = group->group_id, .src_epoch = group->group_id + 1000, .dst_epoch = generation }; + uint8_t wire[TOPO_NODE_WIRE_MAX_SIZE + sizeof(struct TOPOMSG_NODEINFO_PKT)]; + memcpy(wire, &h, sizeof(h)); + struct TOPO_GROUP_NODE nq = {0}; + length = topo_node_serialize(node, &nq, group->group_id, TOPO_FLAG_SEND_SUBNETS, + wire + sizeof(h), sizeof(wire) - sizeof(h), 0); assert(length > 0); + packet(conn, wire, sizeof(h) + length); +} + +static void receive(struct ETCP_CONN* conn, uint64_t group_id, uint8_t cmd, uint64_t local_epoch, size_t length) { + struct TOPOMSG_JOIN_GROUP msg = { .h = { .cmd = ETCP_ID_TOPO_ENTRY, .subcmd = cmd, .group_id = group_id, + .src_epoch = group_id + 1000, .dst_epoch = local_epoch }, .group_type = TOPO_GROUP_TYPE_UTUN }; assert(length <= sizeof(msg)); struct ll_entry* packet = ll_alloc_lldgram(sizeof(msg)); assert(packet); memcpy(packet->dgram, &msg, sizeof(msg)); packet->len = length; conn->instance->api_bindings.callbacks[ETCP_ID_TOPO_ENTRY](conn, packet); } -static void complete(struct TOPO_GROUP* group, struct ETCP_CONN* conn, uint64_t exchange_id) { - receive(conn, group->group_id, TOPO_SUBCMD_TABLE_COMPLETE, exchange_id, sizeof(struct TOPOMSG_TABLE_REQ)); +static void accept_join(struct TOPO_GROUP* group, struct ETCP_CONN* conn) { + receive(conn, group->group_id, TOPO_SUBCMD_JOIN_GROUP, 0, sizeof(struct TOPOMSG_JOIN_GROUP)); + receive(conn, group->group_id, TOPO_SUBCMD_JOIN_ACCEPT, peer(group)->local_epoch, sizeof(struct TOPOMSG_JOIN_GROUP)); + assert(peer(group)->accepted); +} + +static void complete(struct TOPO_GROUP* group, struct ETCP_CONN* conn, uint64_t local_epoch) { + receive(conn, group->group_id, TOPO_SUBCMD_TABLE_COMPLETE, local_epoch, sizeof(struct TOPOMSG_HEADER)); } static void request(struct TOPO_GROUP* group, struct ETCP_CONN* conn) { - receive(conn, group->group_id, TOPO_SUBCMD_REQUEST_TABLE, 123, sizeof(struct TOPOMSG_TABLE_REQ)); + receive(conn, group->group_id, TOPO_SUBCMD_REQUEST_TABLE, peer(group)->local_epoch, sizeof(struct TOPOMSG_HEADER)); } int main(void) { @@ -55,9 +88,12 @@ int main(void) { * Event loop не запускаем: транспортный handshake проверяют интеграционные тесты. */ conn->links_up = 1; assert(topo_group_new_conn(a, conn) == 0 && topo_group_new_conn(b, conn) == 0); - uint64_t a_id = peer(a)->exchange_id, b_id = peer(b)->exchange_id; + uint64_t a_id = peer(a)->local_epoch, b_id = peer(b)->local_epoch; assert(a_id && b_id && a_id != b_id); assert(!topo_group_peer_ready(a, node.node_id) && !topo_group_peer_ready(b, node.node_id)); + complete(a, conn, a_id); + assert(!peer(a)->table_received); /* no topology before ACCEPT */ + accept_join(a, conn); accept_join(b, conn); complete(a, conn, b_id); receive(conn, a->group_id, TOPO_SUBCMD_TABLE_COMPLETE, a_id, 2); assert(!peer(a)->table_received); @@ -65,9 +101,9 @@ int main(void) { assert(peer(a)->table_received && !topo_group_peer_ready(a, node.node_id)); request(a, conn); assert(topo_group_peer_ready(a, node.node_id) && !topo_group_peer_ready(b, node.node_id)); - assert(topo_group_new_conn(a, conn) == 0 && peer(a)->exchange_id == a_id); + assert(topo_group_new_conn(a, conn) == 0 && peer(a)->local_epoch == a_id); receive(conn, a->group_id, TOPO_SUBCMD_JOIN_GROUP, 0, sizeof(struct TOPOMSG_JOIN_GROUP)); - assert(peer(a)->exchange_id == a_id); + assert(peer(a)->local_epoch == a_id); assert(queue_entry_count(a->senders_list) == 1 && topo_group_peer_ready(a, node.node_id)); /* Другой порядок доставки: сначала наш снимок, потом подтверждение чужого. */ request(b, conn); @@ -75,29 +111,58 @@ int main(void) { complete(b, conn, b_id); assert(topo_group_peer_ready(b, node.node_id)); + nodeinfo(a, conn, &node, &keys, a_id, 1); + assert(topo_node_find_by_id(a, node.node_id)->last_timestamp == 1); etcp_fire_conn_status(conn, ETCP_CONN_STATUS_REINIT); - assert(peer(a)->exchange_id != a_id && peer(b)->exchange_id != b_id); + assert(topo_group_find_conn_for_node(a, node.node_id) == conn); /* REINIT preserves known paths */ + assert(peer(a)->local_epoch != a_id && peer(b)->local_epoch != b_id); assert(!topo_group_peer_ready(a, node.node_id) && !topo_group_peer_ready(b, node.node_id)); complete(a, conn, a_id); complete(b, conn, b_id); assert(!peer(a)->table_received && !peer(b)->table_received); + accept_join(a, conn); /* Ошибка отправки снимка не разрешает TABLE_COMPLETE и READY. */ int saved_limit = conn->send_input_q->size_limit; conn->send_input_q->size_limit = 0; request(a, conn); assert(!peer(a)->table_sent); conn->send_input_q->size_limit = saved_limit; - complete(a, conn, peer(a)->exchange_id); request(a, conn); + complete(a, conn, peer(a)->local_epoch); request(a, conn); assert(topo_group_peer_ready(a, node.node_id)); - a_id = peer(a)->exchange_id; + nodeinfo(a, conn, &node, &keys, a_id, 2); + assert(topo_node_find_by_id(a, node.node_id)->last_timestamp == 1); + nodeinfo(a, conn, &node, &keys, peer(a)->local_epoch, 2); + assert(topo_node_find_by_id(a, node.node_id)->last_timestamp == 2); + struct TOPOMSG_WITHDRAW_PKT wd = { .h = { .cmd = ETCP_ID_TOPO_ENTRY, .subcmd = TOPO_SUBCMD_WITHDRAW, + .group_id = a->group_id, .src_epoch = a->group_id + 1000, .dst_epoch = a_id }, + .node_id = node.node_id, .wd_source = node.node_id }; + packet(conn, &wd, sizeof(wd)); + assert(topo_node_find_by_id(a, node.node_id)); + wd.h.dst_epoch = peer(a)->local_epoch; packet(conn, &wd, sizeof(wd)); + assert(!topo_node_find_by_id(a, node.node_id)); + a_id = peer(a)->local_epoch; topo_group_remove_conn(a, conn, TOPO_REMOVE_LOCAL_LEAVE); complete(a, conn, a_id); assert(!topo_group_peer_ready(a, node.node_id)); - assert(topo_group_new_conn(a, conn) == 0 && peer(a)->exchange_id != a_id); + assert(topo_group_new_conn(a, conn) == 0 && peer(a)->local_epoch != a_id); + receive(conn, a->group_id, TOPO_SUBCMD_JOIN_ACCEPT, a_id, sizeof(struct TOPOMSG_JOIN_GROUP)); + assert(!peer(a)->accepted); + accept_join(a, conn); + receive(conn, a->group_id, TOPO_SUBCMD_LEAVE_GROUP, a_id, sizeof(struct TOPOMSG_HEADER)); + assert(peer(a)->accepted); /* old LEAVE cannot detach the new session */ complete(a, conn, a_id); request(a, conn); assert(!topo_group_peer_ready(a, node.node_id)); - complete(a, conn, peer(a)->exchange_id); + complete(a, conn, peer(a)->local_epoch); assert(topo_group_peer_ready(a, node.node_id) && !topo_group_peer_ready(b, node.node_id)); + receive(conn, b->group_id, TOPO_SUBCMD_JOIN_REJECT, b_id, sizeof(struct TOPOMSG_JOIN_GROUP)); + assert(b->senders_list->head); + receive(conn, b->group_id, TOPO_SUBCMD_JOIN_REJECT, peer(b)->local_epoch, sizeof(struct TOPOMSG_JOIN_GROUP)); + assert(!b->senders_list->head && !conn->close_requested); + /* Cancellation while only JOIN has been received must not leave a half-open peer. */ + receive(conn, b->group_id, TOPO_SUBCMD_JOIN_GROUP, 0, sizeof(struct TOPOMSG_JOIN_GROUP)); + assert(peer(b)->peer_epoch && !peer(b)->accepted); + receive(conn, b->group_id, TOPO_SUBCMD_LEAVE_GROUP, 0, sizeof(struct TOPOMSG_HEADER)); + assert(!b->senders_list->head); conn->links_up = 0; assert(!topo_group_peer_ready(a, node.node_id)); node_conn_direct_close(owner); diff --git a/tests/test_group_recovery.c b/tests/test_group_recovery.c index adab6bf2..7ccc3df1 100644 --- a/tests/test_group_recovery.c +++ b/tests/test_group_recovery.c @@ -75,16 +75,27 @@ static void route(struct fixture* f, int target, int via) { } static void table(struct fixture* f, int i, uint8_t subcmd, uint64_t id) { - struct TOPOMSG_TABLE_REQ msg = { .cmd = ETCP_ID_TOPO_ENTRY, .subcmd = subcmd, .group_id = f->group->group_id, .exchange_id = id }; + struct TOPOMSG_HEADER msg = { .cmd = ETCP_ID_TOPO_ENTRY, .subcmd = subcmd, .group_id = f->group->group_id, + .src_epoch = 7, .dst_epoch = id }; struct ll_entry* entry = ll_alloc_lldgram(sizeof(msg)); assert(entry); memcpy(entry->dgram, &msg, sizeof(msg)); entry->len = sizeof(msg); f->inst->api_bindings.callbacks[ETCP_ID_TOPO_ENTRY](f->conns[i], entry); } +static void accept_join(struct fixture* f, int i) { + struct TOPOMSG_JOIN_GROUP msg = { .h = { .cmd = ETCP_ID_TOPO_ENTRY, .subcmd = TOPO_SUBCMD_JOIN_ACCEPT, + .group_id = f->group->group_id, .src_epoch = 7, .dst_epoch = peer(f, i)->local_epoch }, .group_type = TOPO_GROUP_TYPE_UTUN }; + struct ll_entry* e = ll_alloc_lldgram(sizeof(msg)); assert(e); + memcpy(e->dgram, &msg, sizeof(msg)); e->len = sizeof(msg); + f->inst->api_bindings.callbacks[ETCP_ID_TOPO_ENTRY](f->conns[i], e); + assert(peer(f, i)->accepted); +} + static void ready(struct fixture* f, int i) { - assert(peer(f, i) && peer(f, i)->exchange_id); - table(f, i, TOPO_SUBCMD_REQUEST_TABLE, 7); - table(f, i, TOPO_SUBCMD_TABLE_COMPLETE, peer(f, i)->exchange_id); + assert(peer(f, i) && peer(f, i)->local_epoch); + accept_join(f, i); + table(f, i, TOPO_SUBCMD_REQUEST_TABLE, peer(f, i)->local_epoch); + table(f, i, TOPO_SUBCMD_TABLE_COMPLETE, peer(f, i)->local_epoch); assert(topo_group_peer_ready(f->group, f->ids[i])); } @@ -110,7 +121,7 @@ static void partial_and_external_routes(void) { assert(peer(&f, 0)->handle == ownership && topo_group_peer_ready(f.group, f.ids[0])); route(&f, 1, 0); poll_events(&f); /* another READY peer restores the remaining subtree */ up(&f, 2); route(&f, 2, 2); - table(&f, 2, TOPO_SUBCMD_TABLE_COMPLETE, peer(&f, 0)->exchange_id); poll_events(&f); + table(&f, 2, TOPO_SUBCMD_TABLE_COMPLETE, peer(&f, 0)->local_epoch); poll_events(&f); assert(f.group->recovery && !topo_group_peer_ready(f.group, f.ids[2])); ready(&f, 2); poll_events(&f); assert(!f.group->recovery && !peer(&f, 1)); @@ -124,10 +135,11 @@ static void stalled_and_exhausted(void) { topo_recovery_add_node(f.group, f.ids[1], f.ids[0], 20); topo_recovery_start(f.group); poll_events(&f); up(&f, 0); + accept_join(&f, 0); uint64_t old_progress = get_time_tb() - TOPO_RECOVERY_SYNC_TIMEOUT_MS * 10ULL - 1; peer(&f, 0)->progress = old_progress; /* Реальный шаг обмена обновляет срок, пустая повторная проверка — нет. */ - table(&f, 0, TOPO_SUBCMD_REQUEST_TABLE, 8); poll_events(&f); + table(&f, 0, TOPO_SUBCMD_REQUEST_TABLE, peer(&f, 0)->local_epoch); poll_events(&f); assert(peer(&f, 0) && !peer(&f, 1) && peer(&f, 0)->progress > old_progress); peer(&f, 0)->progress = old_progress; topo_recovery_changed(f.group); poll_events(&f); diff --git a/tests/test_node_conn_direct.c b/tests/test_node_conn_direct.c index 79bc728e..a3d38972 100644 --- a/tests/test_node_conn_direct.c +++ b/tests/test_node_conn_direct.c @@ -104,9 +104,9 @@ static int update_addresses(int up) { struct TOPO_GROUP_NODE nq = {0}; struct TOPO_GROUP group = { .instance = g_a, .group_type = TOPO_GROUP_TYPE_CHAT, .group_id = 42 }; uint8_t wire[4096]; - int len = topo_node_serialize(source, &nq, 42, 0, wire + 2, sizeof(wire) - 2, 0); + int len = topo_node_serialize(source, &nq, 42, 0, wire + sizeof(struct TOPOMSG_HEADER), sizeof(wire) - sizeof(struct TOPOMSG_HEADER), 0); struct TOPO_NODE* copy = NULL; - if (len < 0 || topo_node_deserialize(&group, wire + 2, (size_t)len, ©, NULL, NULL, NULL, NULL) < 0) return -1; + if (len < 0 || topo_node_deserialize(&group, wire + sizeof(struct TOPOMSG_HEADER), (size_t)len, ©, NULL, NULL, NULL, NULL) < 0) return -1; struct TOPO_ADDR4* addr = memory_pool_alloc(g_a->topo_groups->v4_addr_pool); if (!addr) { topo_node_destroy(g_a->topo_groups, copy); return -1; } *addr = (struct TOPO_ADDR4){ .next = copy->v4_addrs, .addr = {127,0,0,2}, .port = (uint16_t)(pb + up), @@ -116,9 +116,9 @@ static int update_addresses(int up) { uint8_t message[TOPO_SIG_MSG_MAX_SIZE]; len = topo_node_build_sig_msg(copy, message, sizeof(message)); if (len < 0 || sc_ed25519_sign(g_b->my_ed25519_privkey, message, (size_t)len, copy->x25519_self_sig) != SC_OK) return -1; - len = topo_node_serialize(copy, &nq, 42, 0, wire + 2, sizeof(wire) - 2, 0); + len = topo_node_serialize(copy, &nq, 42, 0, wire + sizeof(struct TOPOMSG_HEADER), sizeof(wire) - sizeof(struct TOPOMSG_HEADER), 0); group.nodes = queue_new(ua, 16, 0, 8, "ncd_address_test"); - int result = topo_group_process_nodeinfo(&group, conn, wire, (size_t)len + 2); + int result = topo_group_process_nodeinfo(&group, conn, wire, (size_t)len + sizeof(struct TOPOMSG_HEADER)); int after = 0; for (struct ETCP_LINK* l = conn->links; l; l = l->next) after++; if (result != 0 || after != before + 1 || node_conn_direct_update_node(g_a, nid_b) != 0) result = -1; diff --git a/tests/test_node_snapshot.c b/tests/test_node_snapshot.c index c7760258..773bdcf8 100644 --- a/tests/test_node_snapshot.c +++ b/tests/test_node_snapshot.c @@ -43,22 +43,22 @@ static void test_bgp(struct UTUN_INSTANCE* signer) { struct TOPO_NODE* ni = record(signer, INT64_MAX - 7, "BGP", 1); struct TOPO_GROUP_NODE nq = {0}; uint8_t packet[4096] = {0}; - int len = topo_node_serialize(ni, &nq, group.group_id, 0, packet + 2, sizeof(packet) - 2, 0); + int len = topo_node_serialize(ni, &nq, group.group_id, 0, packet + sizeof(struct TOPOMSG_HEADER), sizeof(packet) - sizeof(struct TOPOMSG_HEADER), 0); assert(len > 0); - assert(topo_group_process_nodeinfo(&group, &from, packet, (size_t)len + 2) == 0); + assert(topo_group_process_nodeinfo(&group, &from, packet, (size_t)len + sizeof(struct TOPOMSG_HEADER)) == 0); struct TOPO_GROUP_NODE* peer = topo_node_find_by_id(&group, ni->node_id); assert(peer && peer->last_timestamp == ni->timestamp && queue_entry_count(peer->paths) == 1); - assert(topo_group_process_nodeinfo(&group, &from, packet, (size_t)len + 2) == 0); + assert(topo_group_process_nodeinfo(&group, &from, packet, (size_t)len + sizeof(struct TOPOMSG_HEADER)) == 0); assert(queue_entry_count(peer->paths) == 1); struct TOPOMSG_NODEINFO_PKT* wire = (struct TOPOMSG_NODEINFO_PKT*)packet; wire->node.timestamp++; /* подделка не удаляет узел и ранее принятый маршрут */ - assert(topo_group_process_nodeinfo(&group, &from, packet, (size_t)len + 2) < 0); + assert(topo_group_process_nodeinfo(&group, &from, packet, (size_t)len + sizeof(struct TOPOMSG_HEADER)) < 0); assert(topo_node_find_by_id(&group, ni->node_id) == peer && queue_entry_count(peer->paths) == 1); wire->node.timestamp--; ni->timestamp--; sign_record(signer, ni); - len = topo_node_serialize(ni, &nq, group.group_id, 0, packet + 2, sizeof(packet) - 2, 0); - assert(topo_group_process_nodeinfo(&group, &from, packet, (size_t)len + 2) == 0); + len = topo_node_serialize(ni, &nq, group.group_id, 0, packet + sizeof(struct TOPOMSG_HEADER), sizeof(packet) - sizeof(struct TOPOMSG_HEADER), 0); + assert(topo_group_process_nodeinfo(&group, &from, packet, (size_t)len + sizeof(struct TOPOMSG_HEADER)) == 0); assert(peer->last_timestamp == INT64_MAX - 7 && queue_entry_count(peer->paths) == 1); topo_group_remove_path(peer, &from); topo_nodeq_remove_node(&group, peer);