Browse Source

Track topology exchange readiness per group and peer

proxy
evgeny 4 days ago
parent
commit
66c59501fc
  1. 4
      src/chat/chat_status.c
  2. 130
      src/routing_layer/topo_group.c
  3. 23
      src/routing_layer/topo_group.h
  4. 2
      src/transport_layer/etcp.h
  5. 6
      src/transport_layer/etcp_api.c
  6. 15
      src/transport_layer/etcp_api.h
  7. 5
      src/transport_layer/etcp_api_doc.md
  8. 36
      src/transport_layer/etcp_connect.c
  9. 12
      src/transport_layer/etcp_connect.h
  10. 12
      src/transport_layer/etcp_connect_doc.md
  11. 4
      src/transport_layer/etcp_dump.c
  12. 5
      tests/Makefile.am
  13. 4
      tests/test_etcp_router_reconnect.c
  14. 108
      tests/test_group_exchange.c

4
src/chat/chat_status.c

@ -218,7 +218,7 @@ static void collect_conn_metrics(struct UTUN_INSTANCE* inst, uint64_t peer_node_
"=== PEER: %04llX\u2192%04llX [%s] ===\n"
"Role: %s\n"
"Status: %s Links: %d/%d Initialized: %s MTU: %d\n"
"Reinit: %u Reset: %u Tx_state: %d Routing_exchange: %d\n"
"Reinit: %u Reset: %u Tx_state: %d\n"
"\n--- ETCP ---\n"
"RTT last/avg10: %.1f/%.1f ms Jitter: %.1f ms\n"
"Bytes sent: %u Retrans: %u ACKs: %u\n"
@ -230,7 +230,7 @@ static void collect_conn_metrics(struct UTUN_INSTANCE* inst, uint64_t peer_node_
role_str,
conn->links_up ? "UP" : "DOWN", links_up_count, link_count,
conn->initialized ? "yes" : "no", conn->mtu,
conn->reinit_count, conn->reset_count, conn->tx_state, conn->routing_exchange_active,
conn->reinit_count, conn->reset_count, conn->tx_state,
conn->rtt_last / 10.0f, conn->rtt_avg_10 / 10.0f, conn->jitter / 10.0f,
conn->bytes_sent_total, conn->retransmissions_count, conn->ack_packets_count,
conn->unacked_bytes, conn->max_inflight,

130
src/routing_layer/topo_group.c

@ -37,19 +37,47 @@
// Вспомогательные функции
// ============================================================================
static struct TOPO_GROUP_CONN_ITEM* topo_group_peer(const struct TOPO_GROUP* group, uint64_t peer_id) {
for (struct ll_entry* e = group && group->senders_list ? group->senders_list->head : NULL; e; e = e->next) {
struct TOPO_GROUP_CONN_ITEM* item = (struct TOPO_GROUP_CONN_ITEM*)e->data;
if (item->conn && item->conn->peer_node_id == peer_id) return item;
}
return NULL;
}
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->links_up && peer->table_received && peer->table_sent;
}
static void topo_group_log_exchange(struct TOPO_GROUP* group, struct TOPO_GROUP_CONN_ITEM* peer) {
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,
topo_group_peer_ready(group, peer->conn->peer_node_id));
}
/* Шлёт пиру запрос полной таблицы узлов группы (REQUEST_TABLE). */
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) return;
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) { u_free(req); return; }
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) { u_free(req); queue_entry_free(e); }
if (etcp_send(conn, e) != 0) { DEBUG_ERROR(DEBUG_CATEGORY_BGP, "cannot send table request"); u_free(req); queue_entry_free(e); }
}
/* Запрашивает у пира членство в группе (JOIN_GROUP). */
@ -82,32 +110,37 @@ static void topo_group_send_resync(struct ETCP_CONN* conn) {
}
/* Сигнализирует пиру об окончании начальной синхронизации таблицы. */
static void topo_group_send_table_complete(struct TOPO_GROUP* group, struct ETCP_CONN* conn) {
if (!group || !conn) return;
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) return;
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) { u_free(req); return; }
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) { u_free(req); queue_entry_free(e); }
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;
}
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 void 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);
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_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) {
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) return;
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]));
@ -115,6 +148,7 @@ static void topo_group_remember_table_req(struct TOPO_GROUPS* g, uint64_t node_i
}
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);
@ -133,7 +167,7 @@ static void topo_group_fulfill_table_reqs(struct TOPO_GROUPS* g, uint64_t group_
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);
topo_group_handle_request_table(group, conn, req->exchange_id);
}
}
g->pending_table_req_count = w;
@ -259,6 +293,13 @@ static void topo_group_receive_cbk(struct ETCP_CONN* from_conn, struct ll_entry*
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;
@ -278,7 +319,8 @@ static void topo_group_receive_cbk(struct ETCP_CONN* from_conn, struct ll_entry*
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);
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;
}
@ -302,14 +344,25 @@ static void topo_group_receive_cbk(struct ETCP_CONN* from_conn, struct ll_entry*
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_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) { etcp_set_routing_exchange_state(from_conn, 3); DEBUG_INFO(DEBUG_CATEGORY_BGP, "BGP sync complete with %s: %d nodes grp=%016llx", from_conn->log_name, queue_entry_count(group->nodes), (unsigned long long)group->group_id); }
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);
}
}
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;
@ -336,7 +389,9 @@ static void topo_group_conn_status(struct ETCP_CONN* conn, int status, void* arg
for (struct ll_entry* e = groups->group_list->head; e; e = e->next) {
struct TOPO_GROUP* g = (struct TOPO_GROUP*)e;
for (struct ll_entry* se = g->senders_list->head; se; se = se->next) {
if (((struct TOPO_GROUP_CONN_ITEM*)se->data)->conn != conn) continue;
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;
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);
@ -1163,12 +1218,12 @@ int topo_group_process_withdraw(struct TOPO_GROUP* group, struct ETCP_CONN* send
}
/* Сериализует и шлёт NODEINFO узла конкретному пиру. */
void topo_group_send_nodeinfo(struct TOPO_GROUP* group, struct TOPO_GROUP_NODE* node, struct ETCP_CONN* conn, uint16_t cumulative_rtt) {
if (!group || !node || !conn) return;
int topo_group_send_nodeinfo(struct TOPO_GROUP* group, struct TOPO_GROUP_NODE* node, struct ETCP_CONN* conn, uint16_t cumulative_rtt) {
if (!group || !node || !conn) { DEBUG_ERROR(DEBUG_CATEGORY_BGP, "send_nodeinfo: invalid arguments"); return -1; }
size_t max_sz = sizeof(struct TOPOMSG_NODEINFO_PKT) + TOPO_NODE_WIRE_MAX_SIZE;
uint8_t* p = u_malloc(max_sz);
if (!p) return;
if (!p) { DEBUG_ERROR(DEBUG_CATEGORY_BGP, "send_nodeinfo: allocation failed"); return -1; }
p[0] = ETCP_ID_TOPO_ENTRY; p[1] = TOPO_SUBCMD_NODEINFO;
@ -1183,18 +1238,26 @@ void topo_group_send_nodeinfo(struct TOPO_GROUP* group, struct TOPO_GROUP_NODE*
node->paths ? queue_entry_count(node->paths) : -1,
group->nodes ? queue_entry_count(group->nodes) : -1);
u_free(p);
return;
return -1;
}
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);
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; }
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;
}
struct ll_entry* e = queue_entry_new(0);
if (!e) { u_free(p); return; }
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;
if (etcp_send(conn, e) != 0) { DEBUG_ERROR(DEBUG_CATEGORY_BGP, "send_nodeinfo: etcp_send FAILED for node %016llx to %s", (unsigned long long)node->node_id, conn->log_name); u_free(p); queue_entry_free(e); }
if (etcp_send(conn, e) != 0) {
DEBUG_ERROR(DEBUG_CATEGORY_BGP, "send_nodeinfo: etcp_send FAILED for node %016llx to %s",
(unsigned long long)node->node_id, conn->log_name);
u_free(p); queue_entry_free(e); return -1;
}
return 0;
}
/* Добавляет conn в senders_list (дедуп), если его там ещё нет. */
@ -1315,32 +1378,37 @@ int topo_group_leave_peer(struct TOPO_GROUP* group, struct NODE_CONN_DIRECT* han
}
/* Шлёт пиру полную таблицу узлов группы (все узлы, кроме достижимых через него). */
static void topo_group_send_full_table(struct TOPO_GROUP* group, struct ETCP_CONN* conn) {
if (!group || !conn) return;
static int topo_group_send_full_table(struct TOPO_GROUP* group, struct ETCP_CONN* conn) {
uint64_t target = conn->peer_node_id;
struct ll_entry* e = group->nodes ? group->nodes->head : NULL;
while (e) {
struct TOPO_GROUP_NODE* nq = (struct TOPO_GROUP_NODE*)e;
if (nq->node_id == group->instance->node_id) { e = e->next; continue; }
if (!nq->paths) { DEBUG_WARN(DEBUG_CATEGORY_BGP, "node has no paths"); e = e->next; continue; }
if (topo_group_should_send_to(nq, target)) topo_group_send_nodeinfo(group, nq, conn, topo_get_chain_rtt(nq));
if (topo_group_should_send_to(nq, target) && topo_group_send_nodeinfo(group, nq, conn, topo_get_chain_rtt(nq)) < 0) return -1;
e = e->next;
}
return 0;
}
/* Обрабатывает REQUEST_TABLE: шлёт свой nodeinfo + полную таблицу + TABLE_COMPLETE. */
static void topo_group_handle_request_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) {
if (!group || !conn) return;
if (!topo_group_peer_allowed(group, conn->peer_node_id)) {
DEBUG_INFO(DEBUG_CATEGORY_BGP, "table request deferred: member unknown peer=%016llx group=%016llx",
(unsigned long long)conn->peer_node_id, (unsigned long long)group->group_id);
return;
}
topo_group_send_nodeinfo(group, group->local_node, conn, 0);
topo_group_send_full_table(group, conn);
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);
return;
}
/* senders_list заполняется только через topo_group_new_conn (add + BGP);
здесь не добавляем, иначе new_conn (на JOIN_GROUP) упирается в дедуп и не шлёт TABLE_REQ обратно. */
topo_group_send_table_complete(group, conn);
if (topo_group_send_table_complete(group, conn, exchange_id) < 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 обратно. */

23
src/routing_layer/topo_group.h

@ -98,6 +98,12 @@ void topo_fire_nodeinfo_cbk(struct UTUN_INSTANCE* instance, struct TOPO_GROUP* g
* JOIN_READY из chat_sync.h — результат добавления участника, а TABLE_COMPLETE —
* результат обмена топологией одной группы; ни один не означает готовности других групп.
* Конфиг, auto-connect и recovery обязаны использовать общую проверку членства.
*
* REQUEST_TABLE содержит ненулевой exchange_id, сгенерированный запросившей
* стороной для текущего обмена (группа, пир). TABLE_COMPLETE возвращает тот же
* id только после успешной постановки всех NODEINFO в транспорт. Дубликаты
* запроса используют прежний id; новое присоединение и REINIT создают новый.
* Ответ с другим id не завершает текущий обмен. Формат изменён без совместимости.
*/
// Sub-команды
@ -140,6 +146,7 @@ struct TOPOMSG_TABLE_REQ {
uint8_t cmd;
uint8_t subcmd;
uint64_t group_id; // идентификатор группы
uint64_t exchange_id; // ненулевой id запроса; TABLE_COMPLETE возвращает его без изменений
} __attribute__((packed));
/**
@ -173,9 +180,11 @@ struct TOPOMSG_ERR_GROUP_MISMATCH {
} __attribute__((packed));
struct TOPO_GROUP_CONN_ITEM {
struct ll_entry ll;
struct ETCP_CONN* conn;
struct NODE_CONN_DIRECT* handle; // владение BGP этим conn (NCD-handle)
uint64_t exchange_id; // текущий исходящий REQUEST_TABLE
uint8_t table_received; // TABLE_COMPLETE для exchange_id принят
uint8_t table_sent; // ответ на запрос пира целиком передан транспорту
};
struct ETCP_LINK;
@ -194,7 +203,7 @@ struct TOPO_GROUP {
uint64_t group_id; // уникальный идентификатор группы
uint8_t group_type; // TOPO_GROUP_TYPE_*
struct UTUN_INSTANCE* instance;
struct ll_queue* senders_list; // TOPO_GROUP_CONN_ITEM{ll_entry,ETCP_CONN*} — активные BGP-пиры, без хеш-индекса
struct ll_queue* senders_list; // ll_entry.data = TOPO_GROUP_CONN_ITEM, без хеш-индекса
struct ll_queue* nodes; // TOPO_GROUP_NODE{ll_entry,node_id,paths,...} — per-group узлы, хеш-индекс по node_id(8B)
struct TOPO_GROUP_NODE* local_node; // свой узел в этой группе
uint8_t ed25519_public_key[SC_PUBKEY_SIZE];
@ -215,6 +224,7 @@ struct TOPO_GROUP {
struct topo_pending_table_req {
uint64_t node_id; /* кто запросил таблицу */
uint64_t group_id; /* какую группу */
uint64_t exchange_id;
};
/**
@ -294,6 +304,13 @@ void topo_groups_remove_group(struct TOPO_GROUPS* g, uint64_t group_id);
*/
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 игнорируется.
* Это состояние обмена таблицами, а не добавление мембера (см. chat_join.h). */
int topo_group_peer_ready(const struct TOPO_GROUP* group, uint64_t peer_id);
/**
* @brief Удаляет conn из senders_list, очищает paths во всех nodes, отправляет withdraw если node unreachable.
*
@ -332,7 +349,7 @@ int topo_group_process_withdraw(struct TOPO_GROUP* group, struct ETCP_CONN* send
/**
* @brief Отправляет NODEINFO пакет одному conn (всегда local_node).
*/
void topo_group_send_nodeinfo(struct TOPO_GROUP* group, struct TOPO_GROUP_NODE* node, struct ETCP_CONN* conn, uint16_t cumulative_rtt);
int topo_group_send_nodeinfo(struct TOPO_GROUP* group, struct TOPO_GROUP_NODE* node, struct ETCP_CONN* conn, uint16_t cumulative_rtt);
/**
* @brief Отправляет WITHDRAW для node_id (вызывает broadcast_withdraw).

2
src/transport_layer/etcp.h

@ -194,7 +194,6 @@ struct ETCP_CONN {
// uint32_t total_packets_sent; // Total packets sent counter - Not used
// Flags
uint8_t routing_exchange_active; // 0-не активен, 1-надо инициировать (клиент), 2-обмен идёт, 3-завершён, 4-пропущен (нет BGP)
uint8_t got_initial_pkt; //
uint8_t initialized; // 0 - только созданный ETCP, 1 - хотя бы один линк проинициалзирован (обмен ключами произведен)
uint64_t reset_id; // локальная эпоха потока (меняется только при локальном reset)
@ -216,7 +215,6 @@ struct ETCP_CONN {
// Unified callback chain with event mask (init/reinit/up/down/node_changed)
struct etcp_cbk_entry* cbks;
uint8_t reinit_pending; // 1 = reinit в процессе, ждём завершения
void (*bgp_ready_cbk)(struct ETCP_CONN* conn); // вызывается когда BGP готов (завершён или пропущен)
uint8_t fin_wait : 1; // close-pending: ожидаем подтверждение закрытия от remote
void (*fin_wait_clear_cb)(struct ETCP_CONN* conn, void* arg); // вызывается при сбросе fin_wait

6
src/transport_layer/etcp_api.c

@ -155,12 +155,6 @@ void etcp_socket_cbk_fire(struct ETCP_SOCKET* sock, int event) {
struct etcp_socket_cbk_entry* cbe = sock->instance->socket_cbks;
while (cbe) { struct etcp_socket_cbk_entry* n = cbe->next; if (cbe->event_mask & event) cbe->fn(sock, event, cbe->arg); cbe = n; }
}
void etcp_set_routing_exchange_state(struct ETCP_CONN* conn, uint8_t new_state) {
if (!conn || conn->state == 2) return;
conn->routing_exchange_active = new_state;
if (new_state >= 3 && conn->bgp_ready_cbk)
conn->bgp_ready_cbk(conn);
}
int etcp_bind(struct UTUN_INSTANCE* inst, uint8_t id, etcp_recv_fn callback) {
if (!inst || !callback || (unsigned)id >= ETCP_MAX_BINDINGS) return -1;

15
src/transport_layer/etcp_api.h

@ -265,7 +265,6 @@ void etcp_fire_link_status_cbk(struct ETCP_LINK* link, int old_state, int old_st
// ---- Background connection initialization ----
#define ETCP_CONNECT_EARLY 1
#define ETCP_CONNECT_LATE 2
#define ETCP_CONNECT_BGP_READY 4 // BGP-синхронизация завершена
typedef void (*etcp_connect_callback_t)(void* arg, struct ETCP_CONN* conn, int type);
@ -276,11 +275,11 @@ typedef void (*etcp_connect_callback_t)(void* arg, struct ETCP_CONN* conn, int t
* @param node Узел из node registry (TOPO_GROUP_NODE)
* @param cb Коллбэк завершения (сигнатура etcp_connect_callback_t)
* @param arg Пользовательский аргумент
* @param flags Битовая маска этапов (ETCP_CONNECT_EARLY / LATE / BGP_READY)
* @param flags Битовая маска этапов (ETCP_CONNECT_EARLY / LATE)
* @return 0 при успехе, -1 при ошибке
*
* @note Коллбэк вызывается на каждом запрошенном этапе с параметром
* type = ETCP_CONNECT_EARLY / ETCP_CONNECT_LATE / ETCP_CONNECT_BGP_READY
* type = ETCP_CONNECT_EARLY / ETCP_CONNECT_LATE
* (или 0 при таймауте/неудаче). При уже готовом соединении EARLY+LATE
* доставляются немедленно.
*/
@ -397,16 +396,6 @@ void etcp_remove_socket_cbk(struct UTUN_INSTANCE* inst, etcp_socket_cbk_fn fn, v
*/
void etcp_socket_cbk_fire(struct ETCP_SOCKET* sock, int event);
/**
* @brief Установить состояние обмена маршрут-таблицами (BGP)
*
* @param conn ETCP соединение
* @param new_state 0-не активен, 1-надо инициировать (клиент), 2-обмен идёт,
* 3-завершён, 4-пропущен (нет BGP)
*
* @note При new_state >= 3 и наличии bgp_ready_cbk вызывается этот коллбэк
*/
void etcp_set_routing_exchange_state(struct ETCP_CONN* conn, uint8_t new_state);
/**
* @brief Внутренняя функция etcp: Коллбэк для очередей output ll_queue normalizer

5
src/transport_layer/etcp_api_doc.md

@ -70,10 +70,10 @@ etcp_set_new_conn_cbk(instance, on_new_conn, my_data);
```c
etcp_connect(instance, node, on_connect, arg,
ETCP_CONNECT_EARLY | ETCP_CONNECT_BGP_READY);
ETCP_CONNECT_EARLY | ETCP_CONNECT_LATE);
```
Флаги: `ETCP_CONNECT_EARLY` — раннее подключение, `ETCP_CONNECT_LATE` — отложенное, `ETCP_CONNECT_BGP_READY` — после BGP-синхронизации.
Флаги: `ETCP_CONNECT_EARLY` — раннее подключение, `ETCP_CONNECT_LATE` — отложенное.
## 3. API
@ -107,7 +107,6 @@ etcp_connect(instance, node, on_connect, arg,
| `etcp_conn_add_down_cbk/remove_down_cbk` | Аналогично для down |
| `etcp_set_new_conn_cbk(inst, fn, arg)` | Установить коллбэк на входящее соединение |
| `etcp_add_new_conn_cbk/remove_new_conn_cbk` | Добавить/убрать из цепочки new_conn |
| `etcp_set_routing_exchange_state(conn, state)` | Установить состояние обмена маршрутами; при state≥3 вызывает `bgp_ready_cbk` |
| `etcp_connect(inst, node, cb, arg, flags)` | Фоновое подключение к узлу |
### Командные ID (сырые, первый байт кодограммы)

36
src/transport_layer/etcp_connect.c

@ -138,13 +138,6 @@ static void connect_create_links_v6(struct ETCP_CONNECT* ctx, struct TOPO_GROUP_
addr_count, sock_count, link_count, addr_count == 0 ? " (v6_addrs EMPTY)" : "", (unsigned long long)node->node_id);
}
static void connect_bgp_ready_cb(struct ETCP_CONN* conn) {
struct ETCP_CONNECT* ctx = connect_find(conn->instance, conn->peer_node_id);
DEBUG_DEBUG(DEBUG_CATEGORY_BGP, "[etcp_connect] bgp_ready_cb: node=%016llx ctx=%p ctx_done=%d ra=%d",
(unsigned long long)conn->peer_node_id, (void*)ctx, ctx ? ctx->done : -1, conn->routing_exchange_active);
if (!ctx || ctx->done) return;
connect_deliver(ctx, ETCP_CONNECT_BGP_READY, 1);
}
static void connect_initial_timeout_cb(void* arg) {
struct ETCP_CONNECT* ctx = (struct ETCP_CONNECT*)arg;
@ -229,7 +222,9 @@ static void connect_init_cb(struct ETCP_CONN* conn, int event, void* arg) { (voi
int etcp_connect(struct UTUN_INSTANCE* inst, struct TOPO_GROUP_NODE* node,
etcp_connect_callback_t cb, void* arg, uint8_t flags) {
if (!inst || !node || !cb) return -1;
if (!inst || !node || !cb || !flags || (flags & ~(ETCP_CONNECT_EARLY | ETCP_CONNECT_LATE))) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP_CONNECT, "invalid connect arguments/flags=%u", flags); return -1;
}
struct TOPO_NODE* ni = topo_node_registry_find(inst->topo_groups, node->node_id);
if (!ni) { DEBUG_ERROR(DEBUG_CATEGORY_ETCP_CONNECT, "[etcp_connect] node not in registry"); return -1; }
@ -294,7 +289,6 @@ int etcp_connect(struct UTUN_INSTANCE* inst, struct TOPO_GROUP_NODE* node,
if (!cn) { u_free(ctx); return -1; }
cn->cb = cb; cn->arg = arg; cn->flags = flags; ctx->cb_list = cn;
etcp_conn_add_cbk(conn, connect_init_cb, ctx, ETCP_CBK_EVENT_INIT);
conn->bgp_ready_cbk = connect_bgp_ready_cb;
ctx->initial_timer = uasync_set_timeout(inst->ua, (int)inst->etcp_connect_timeout_tb, ctx,
connect_initial_timeout_cb, "etcp_connect_init");
ctx->next = inst->pending_connects; inst->pending_connects = ctx;
@ -318,25 +312,13 @@ int etcp_connect(struct UTUN_INSTANCE* inst, struct TOPO_GROUP_NODE* node,
DEBUG_INFO(DEBUG_CATEGORY_ETCP_CONNECT, "[etcp_connect] node 0x%016llx: existing ready conn, delivering EARLY+LATE",
(unsigned long long)node_id);
struct etcp_connect_cb_node* n = ctx->cb_list;
struct etcp_connect_cb_node* keep = NULL;
while (n) {
struct etcp_connect_cb_node* next = n->next;
if (n->flags & ETCP_CONNECT_EARLY) n->cb(n->arg, ctx->conn, ETCP_CONNECT_EARLY);
if (n->flags & ETCP_CONNECT_LATE) n->cb(n->arg, ctx->conn, ETCP_CONNECT_LATE);
if (n->flags & ETCP_CONNECT_BGP_READY) {
if (conn->routing_exchange_active >= 3) n->cb(n->arg, ctx->conn, ETCP_CONNECT_BGP_READY);
else { n->next = keep; keep = n; n = next; continue; }
}
u_free(n);
n = next;
}
ctx->cb_list = keep;
if (keep) {
conn->bgp_ready_cbk = connect_bgp_ready_cb;
ctx->next = inst->pending_connects; inst->pending_connects = ctx;
} else {
u_free(ctx);
if (n->flags & ETCP_CONNECT_EARLY) n->cb(n->arg, ctx->conn, ETCP_CONNECT_EARLY);
if (n->flags & ETCP_CONNECT_LATE) n->cb(n->arg, ctx->conn, ETCP_CONNECT_LATE);
u_free(n); n = next;
}
u_free(ctx);
return 0;
}
@ -379,7 +361,6 @@ int etcp_connect(struct UTUN_INSTANCE* inst, struct TOPO_GROUP_NODE* node,
}
etcp_conn_add_cbk(conn, connect_init_cb, ctx, ETCP_CBK_EVENT_INIT);
conn->bgp_ready_cbk = connect_bgp_ready_cb;
connect_create_links_v4(ctx, node);
connect_create_links_v6(ctx, node);
@ -415,7 +396,7 @@ void etcp_connect_cancel_for_conn(struct UTUN_INSTANCE* inst, struct ETCP_CONN*
}
}
/* Полная отмена pending-коннекта по node_id: снимает connect_init_cb/bgp_ready_cbk
/* Полная отмена pending-коннекта по node_id: снимает connect_init_cb
* с conn (conn остаётся жив, например, переиспользуется ncd), отменяет таймеры,
* вынимает ctx из pending_connects и освобождает. Безопасно вызывать, если ctx нет. */
void etcp_connect_cancel_by_node(struct UTUN_INSTANCE* inst, uint64_t node_id) {
@ -427,7 +408,6 @@ void etcp_connect_cancel_by_node(struct UTUN_INSTANCE* inst, uint64_t node_id) {
*pp = ctx->next;
if (ctx->conn) {
etcp_conn_remove_cbk(ctx->conn, connect_init_cb, ctx);
ctx->conn->bgp_ready_cbk = NULL;
}
if (ctx->initial_timer) { uasync_cancel_timeout(inst->ua, ctx->initial_timer); ctx->initial_timer = NULL; }
if (ctx->settle_timer) { uasync_cancel_timeout(inst->ua, ctx->settle_timer); ctx->settle_timer = NULL; }

12
src/transport_layer/etcp_connect.h

@ -25,7 +25,7 @@ typedef void (*etcp_connect_callback_t)(void* arg, struct ETCP_CONN* conn, int t
* @param cb коллбэк (etcp_connect_callback_t)
* @param arg пользовательский аргумент, передаваемый в cb
* @param flags битовая маска интересующих фаз: ETCP_CONNECT_EARLY(1) |
* ETCP_CONNECT_LATE(2) | ETCP_CONNECT_BGP_READY(4)
* ETCP_CONNECT_LATE(2)
* @return 0 при успехе, -1 при ошибке (до вызова коллбэка)
*
*
@ -78,11 +78,7 @@ typedef void (*etcp_connect_callback_t)(void* arg, struct ETCP_CONN* conn, int t
* │
* └─ conn->state == 1 (ready, уже проинициализирован):
* → connect_deliver(ctx, ETCP_CONNECT_LATE) сразу
* → если routing_exchange_active >= 3:
* connect_deliver(ctx, ETCP_CONNECT_BGP_READY)
* → если в cb_list остались неудовлетворённые коллбэки
* (ждут BGP_READY, который ещё не готов):
* ctx остаётся в pending_connects с conn->bgp_ready_cbk
* → контекст освобождается; готовность топологии принадлежит группе
*
* 4. !conn && !ctx — новое подключение (основной путь)
* → etcp_connection_create(inst, NULL)
@ -118,10 +114,6 @@ typedef void (*etcp_connect_callback_t)(void* arg, struct ETCP_CONN* conn, int t
* ├─ Убирает connect_init_cb из conn->cbks
* └─ connect_deliver(ctx, ETCP_CONNECT_LATE) → ctx освобождается
*
* connect_bgp_ready_cb (вызывается при routing_exchange_active >= 3):
* └─ connect_deliver(ctx, ETCP_CONNECT_BGP_READY)
* (может произойти в любой момент после установки соединения)
*
* connect_initial_timeout_cb (initial_timer):
* ├─ ctx->done = 1
* ├─ Закрывает TCP-линк

12
src/transport_layer/etcp_connect_doc.md

@ -1,7 +1,7 @@
# ETCP Connect (etcp_connect)
## 1. Назначение
Управление исходящими ETCP-соединениями. Предоставляет асинхронный API `etcp_connect()` для установки соединения с удалённым узлом: создаёт ETCP_CONN, настраивает криптографию, поднимает UDP-линки и опционально TCP/STCP-транспорт, отслеживает прогресс установки с таймерами. Доставляет коллбэки о фазах готовности (EARLY, LATE, BGP_READY).
Управление исходящими ETCP-соединениями. Предоставляет асинхронный API `etcp_connect()` для установки соединения с удалённым узлом: создаёт ETCP_CONN, настраивает криптографию, поднимает UDP-линки и опционально TCP/STCP-транспорт, отслеживает прогресс установки с таймерами. Доставляет коллбэки о фазах готовности (EARLY, LATE).
## 2. Как пользоваться
@ -21,8 +21,6 @@ void my_connect_cb(void* arg, struct ETCP_CONN* conn, int type) {
// первый линк готов (можно начинать обмен)
if (type & ETCP_CONNECT_LATE)
// settle-фаза завершена (все линки проверены, неудачные удалены)
if (type & ETCP_CONNECT_BGP_READY)
// BGP-синхронизация завершена (routing_exchange_active >= 3)
}
```
@ -30,8 +28,7 @@ void my_connect_cb(void* arg, struct ETCP_CONN* conn, int type) {
1. `etcp_connect()` — создаёт ETCP_CONN, crypto_ctx, UDP-линки и TCP/STCP
2. `connect_ready_cb()` — первый линк стал ready → EARLY-коллбэк → запуск settle-таймера
3. `connect_settle_timeout_cb()` — settle (min_rtt × 8, clamp 500..2000ms) → удаление неудачных линков → LATE-коллбэк
4. `connect_bgp_ready_cb()` — BGP завершён → BGP_READY-коллбэк
5. При таймауте (`connect_initial_timeout_cb`) — type=0 (conn=NULL), все ресурсы освобождаются
4. При таймауте (`connect_initial_timeout_cb`) — type=0 (conn=NULL), все ресурсы освобождаются
**Нюансы:**
- Если соединение уже установлено — коллбэки вызываются немедленно
@ -57,7 +54,6 @@ void my_connect_cb(void* arg, struct ETCP_CONN* conn, int type) {
|---|---|
| `ETCP_CONNECT_EARLY` (1) | Первый линк готов |
| `ETCP_CONNECT_LATE` (2) | Settle-фаза завершена |
| `ETCP_CONNECT_BGP_READY` (4) | BGP-синхронизация завершена |
### Функции
| Функция | Описание |
@ -69,3 +65,7 @@ void my_connect_cb(void* arg, struct ETCP_CONN* conn, int type) {
1. **Таймаут установки** (`connect_initial_timeout_cb`): stcp_link_close → etcp_connection_close → connect_cancel → type=0
2. **Settle-таймаут** (`connect_settle_timeout_cb`): удаление неудачных линков/stcp → LATE → connect_cancel
3. **Штатное закрытие** (`utun_instance_destroy`): фаза 1 (detach) → фаза 2 (deferred: `etcp_connect_cancel_for_conn`)
Готовность BGP проверяется для конкретной группы через `topo_group_peer_ready(group, peer_id)`.
Общий транспорт не имеет состояния завершения группового обмена. Для владения
соединениями используйте NCD; создание транспорта само по себе не вступает в группу.

4
src/transport_layer/etcp_dump.c

@ -56,9 +56,9 @@ void etcp_dump_conn_state(struct ETCP_CONN* conn) {
DHDR("=== ETCP CONN [%s] ===", conn->log_name);
DLOG("GENERAL: peer=0x%016llx state=%d init=%d l_up=%d tx=%d mtu=%d rout_ex=%d",
DLOG("GENERAL: peer=0x%016llx state=%d init=%d l_up=%d tx=%d mtu=%d",
(unsigned long long)conn->peer_node_id, conn->state, conn->initialized, conn->links_up,
conn->tx_state, conn->mtu, conn->routing_exchange_active);
conn->tx_state, conn->mtu);
DLOG("IDS: next_tx=%u last_rx=%u last_del=%u rx_ack_till=%u got_init=%d reset_done=%d",
conn->next_tx_id, conn->last_rx_id, conn->last_delivered_id, conn->rx_ack_till,

5
tests/Makefile.am

@ -65,6 +65,7 @@ check_PROGRAMS = \
test_etcp_connect \
test_node_conn_direct \
test_group_ownership \
test_group_exchange \
test_node_snapshot \
test_ncd_config \
test_db_sync \
@ -394,6 +395,10 @@ test_group_ownership_SOURCES = test_group_ownership.c
test_group_ownership_CFLAGS = -I$(top_srcdir)/src -I$(top_srcdir)/lib
test_group_ownership_LDADD = $(LIBUTUN) $(CRYPTO_LIBS) $(COMMON_LIBS)
test_group_exchange_SOURCES = test_group_exchange.c
test_group_exchange_CFLAGS = -I$(top_srcdir)/src -I$(top_srcdir)/lib
test_group_exchange_LDADD = $(LIBUTUN) $(CRYPTO_LIBS) $(COMMON_LIBS)
test_node_snapshot_SOURCES = test_node_snapshot.c
test_node_snapshot_CFLAGS = -I$(top_srcdir)/src -I$(top_srcdir)/lib
test_node_snapshot_LDADD = $(LIBUTUN) $(CRYPTO_LIBS) $(COMMON_LIBS)

4
tests/test_etcp_router_reconnect.c

@ -140,7 +140,8 @@ static int send_one_pkt(void) {
static int transit_routes_ready(void) {
struct TOPO_GROUP* group = g_b ? topo_groups_get_default(g_b->topo_groups) : NULL;
return group && topo_group_find_conn_for_node(group, node_a) && topo_group_find_conn_for_node(group, node_c);
return group && topo_group_peer_ready(group, node_a) && topo_group_peer_ready(group, node_c) &&
topo_group_find_conn_for_node(group, node_a) && topo_group_find_conn_for_node(group, node_c);
}
static const char* state_name(int s) {
@ -190,6 +191,7 @@ static void state_step(void) {
{
struct ETCP_CONN* direct = instance_find_conn(g_a, node_c);
if (direct && direct->links_up && direct->initialized &&
topo_group_peer_ready(topo_groups_get_default(g_a->topo_groups), node_c) &&
topo_group_find_conn_for_node(topo_groups_get_default(g_a->topo_groups), node_c) == direct) {
fprintf(stderr,"RECOVERY: direct a-c path established without test reconnect\n");
g_phase_sent = 0; g_state = ST_PHASE3_SEND;

108
tests/test_group_exchange.c

@ -0,0 +1,108 @@
#include <assert.h>
#include <string.h>
#include "utun_instance.h"
#include "routing_layer/topo_group.h"
#include "transport_layer/etcp.h"
#include "transport_layer/etcp_api.h"
#include "transport_layer/node_conn_direct.h"
#include "../lib/u_async.h"
#include "../lib/debug_config.h"
static struct TOPO_GROUP_CONN_ITEM* peer(struct TOPO_GROUP* group) {
assert(group->senders_list->head && !group->senders_list->head->next);
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 };
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 request(struct TOPO_GROUP* group, struct ETCP_CONN* conn) {
receive(conn, group->group_id, TOPO_SUBCMD_REQUEST_TABLE, 123, sizeof(struct TOPOMSG_TABLE_REQ));
}
int main(void) {
debug_config_init(); debug_set_level(DEBUG_LEVEL_INFO);
utun_instance_set_tun_init_enabled(0);
struct UASYNC* ua = uasync_create(); assert(ua);
struct UTUN_INSTANCE* inst = utun_instance_create_from_str(ua,
"[global]\n"
"my_private_key=704f2e012c8fa8768130cb0f988a997dccb628372bc5ceccacc78dcbfec5916f\n"
"my_public_key=b3193173def895bd0fcea6f86af077c7d77216f10395275f627ac18242ec0f01\n"
"[server: udp]\naddr=127.0.0.1:0\ntype=public\n");
assert(inst && utun_instance_init(inst) == 0);
struct TOPO_GROUP* a = topo_groups_create_group(inst->topo_groups, 42, TOPO_GROUP_TYPE_UTUN, NULL);
struct TOPO_GROUP* b = topo_groups_create_group(inst->topo_groups, 43, TOPO_GROUP_TYPE_UTUN, NULL);
assert(a && b);
struct SC_MYKEYS keys; assert(sc_generate_keypair(&keys) == SC_OK);
struct TOPO_ADDR4 addr = { .addr = {127, 0, 0, 1}, .port = 9, .protocol = TOPO_PROTO_UDP };
struct TOPO_NODE node = { .v4_addrs = &addr };
memcpy(node.public_key, keys.public_key, SC_PUBKEY_SIZE);
node.node_id = sc_derive_node_id_from_pubkey(node.public_key);
struct NODE_CONN_DIRECT* owner = NULL;
assert(node_conn_direct_open_node(inst, node.node_id, NULL, NULL, &owner, &node, NULL) == NCD_NEW);
struct ETCP_CONN* conn = node_conn_direct_get_conn(owner); assert(conn);
/* Детерминированная доставка управляющих пакетов через штатный BGP receiver.
* 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;
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, b_id);
receive(conn, a->group_id, TOPO_SUBCMD_TABLE_COMPLETE, a_id, 2);
assert(!peer(a)->table_received);
complete(a, conn, a_id);
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);
receive(conn, a->group_id, TOPO_SUBCMD_JOIN_GROUP, 0, sizeof(struct TOPOMSG_JOIN_GROUP));
assert(peer(a)->exchange_id == a_id);
assert(queue_entry_count(a->senders_list) == 1 && topo_group_peer_ready(a, node.node_id));
/* Другой порядок доставки: сначала наш снимок, потом подтверждение чужого. */
request(b, conn);
assert(!topo_group_peer_ready(b, node.node_id));
complete(b, conn, b_id);
assert(topo_group_peer_ready(b, node.node_id));
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_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);
/* Ошибка отправки снимка не разрешает 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);
assert(topo_group_peer_ready(a, node.node_id));
a_id = peer(a)->exchange_id;
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);
complete(a, conn, a_id); request(a, conn);
assert(!topo_group_peer_ready(a, node.node_id));
complete(a, conn, peer(a)->exchange_id);
assert(topo_group_peer_ready(a, node.node_id) && !topo_group_peer_ready(b, node.node_id));
conn->links_up = 0;
assert(!topo_group_peer_ready(a, node.node_id));
node_conn_direct_close(owner);
inst->running = 0; utun_instance_destroy(inst);
uasync_poll(ua, 0); uasync_destroy(ua, 0);
DEBUG_INFO(DEBUG_CATEGORY_BGP, "group exchange: isolation, duplicate attach, stale replies, send failure and REINIT passed");
return 0;
}
Loading…
Cancel
Save