Browse Source

Negotiate group sessions and scope topology messages to peer epochs

proxy
evgeny 4 days ago
parent
commit
792776a1f6
  1. 422
      src/routing_layer/topo_group.c
  2. 121
      src/routing_layer/topo_group.h
  3. 93
      tests/test_group_exchange.c
  4. 24
      tests/test_group_recovery.c
  5. 8
      tests/test_node_conn_direct.c
  6. 12
      tests/test_node_snapshot.c

422
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); } }
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));
}
e = e->next;
}
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; }
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 (!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 (!entry) return;
if (!from_conn || !entry->dgram || entry->len < 2) {
DEBUG_WARN(DEBUG_CATEGORY_BGP, "invalid BGP packet length=%u", entry->len); 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;
struct UTUN_INSTANCE* instance = from_conn->instance;
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;
}
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);
/* До получения 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_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_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 по группам. */

121
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);

93
tests/test_group_exchange.c

@ -1,5 +1,6 @@
#include <assert.h>
#include <string.h>
#include <openssl/evp.h>
#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);

24
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);

8
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, &copy, NULL, NULL, NULL, NULL) < 0) return -1;
if (len < 0 || topo_node_deserialize(&group, wire + sizeof(struct TOPOMSG_HEADER), (size_t)len, &copy, 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;

12
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);

Loading…
Cancel
Save