Browse Source

Fix ETCP INIT collision scope and backpressure regression tests

master
evgeny 6 days ago
parent
commit
ac45554439
  1. 25
      src/transport_layer/etcp_connections.c
  2. 7
      src/transport_layer/etcp_connections.h
  3. 179
      tests/test_etcp_link_stress.c
  4. 126
      tests/test_etcp_router.c

25
src/transport_layer/etcp_connections.c

@ -512,16 +512,25 @@ static void etcp_link_update_remote_addr(struct ETCP_LINK* link, struct ETCP_SOC
link->etcp->log_name, sockaddr_storage_to_str(new_addr).str);
}
/* Выбор линка для отправки collision INIT: outbound-линк (is_server==0), чей remote_addr
* совпадает с адресом пришедшего INIT (симметричный сокет). NULL — симметричного сокета
* нет (асимметрия/NAT) → коллизия не срабатывает, inbound принимается как есть. */
struct ETCP_LINK* etcp_select_collision_link(struct ETCP_CONN* conn, const struct sockaddr_storage* addr) {
if (!conn || !addr) return NULL;
/* Коллизия INIT возможна только на одной паре: принимающий сокет + адрес отправителя.
* Исходящий линк через другой локальный сокет — независимый путь к тому же пиру. */
struct ETCP_LINK* etcp_select_collision_link(struct ETCP_CONN* conn, struct ETCP_SOCKET* e_sock,
const struct sockaddr_storage* addr) {
if (!conn || !e_sock || !addr) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "collision lookup: invalid conn=%p socket=%p addr=%p", conn, e_sock, addr);
return NULL;
}
struct ETCP_LINK* l = conn->links;
while (l) {
if (l->is_server == 0 && sockaddr_equal(&l->remote_addr, addr)) return l;
if (l->is_server == 0 && l->conn == e_sock && sockaddr_equal(&l->remote_addr, addr)) {
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] collision lookup: recv_socket=%s remote=%s outbound_link=%u/%u",
conn->log_name, e_sock->name, sockaddr_storage_to_str(addr).str, l->local_link_id, l->remote_link_id);
return l;
}
l = l->next;
}
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] collision lookup: recv_socket=%s remote=%s no matching outbound link",
conn->log_name, e_sock->name, sockaddr_storage_to_str(addr).str);
return NULL;
}
@ -2427,7 +2436,7 @@ static void etcp_process_packet(struct ETCP_SOCKET* e_sock, uint8_t* data,
send_reset = 1;
} else {
// Check if WE have an outbound (master) link on this conn
struct ETCP_LINK* ml = etcp_select_collision_link(conn, &addr);
struct ETCP_LINK* ml = etcp_select_collision_link(conn, e_sock, &addr);
if (ml) {
if (conn->instance->node_id < peer_id) {
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] COLLISION: we=master node_id=%016llx < peer=%016llx — sending collision INIT",
@ -2477,7 +2486,7 @@ create_new_link:
send_reset = 1;
} else {
// Check if WE have an outbound (master) link
struct ETCP_LINK* ml = etcp_select_collision_link(conn, &addr);
struct ETCP_LINK* ml = etcp_select_collision_link(conn, e_sock, &addr);
if (ml && conn->instance->node_id < peer_id) {
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] COLLISION: we=master node_id=%016llx < peer=%016llx — sending collision INIT",
conn->log_name, (unsigned long long)conn->instance->node_id, (unsigned long long)peer_id);

7
src/transport_layer/etcp_connections.h

@ -384,9 +384,10 @@ int etcp_encrypt_send(struct ETCP_DGRAM* dgram);// зашифровывает и
// find link by address
struct ETCP_LINK* etcp_link_find_by_addr(struct ETCP_SOCKET* e_sock, struct sockaddr_storage* addr, int is_tcp);
// Выбор линка для collision INIT: outbound-линк (is_server==0), чей remote_addr совпадает
// с адресом пришедшего INIT (симметричный сокет). NULL — симметричного сокета нет.
struct ETCP_LINK* etcp_select_collision_link(struct ETCP_CONN* conn, const struct sockaddr_storage* addr);
// Выбор исходящего линка для collision INIT по принимающему сокету и адресу отправителя.
// Линк через другой локальный сокет не подходит. NULL — встречной попытки на этой паре нет.
struct ETCP_LINK* etcp_select_collision_link(struct ETCP_CONN* conn, struct ETCP_SOCKET* e_sock,
const struct sockaddr_storage* addr);
// find free local_link_id for connection
// scans all links in connection, marks used ids in bit array

179
tests/test_etcp_link_stress.c

@ -1,5 +1,5 @@
// test_etcp_link_stress.c - Multilink test: добавляет линки в 5 вариантах (V1/V2 one-way,
// V3a/b/c коллизия), проверяет per-link трафик и инвариант «нет реинита пока жив хотя бы один линк».
// test_etcp_link_stress.c — одностороннее добавление линков, настоящие и ложные коллизии INIT,
// уникальность пар сокет/адрес, per-link трафик и сохранение надёжного потока при failover.
#include <stdio.h>
#include <stdlib.h>
#include <string.h>
@ -174,15 +174,32 @@ static struct ETCP_LINK* first_link_on_socket(struct ETCP_SOCKET* sock) {
return lqe ? lqe->link : NULL;
}
static uint64_t client_link_acked(uint16_t srv_port) {
struct ETCP_LINK* l = find_link_by_remote_port(get_conn(client_instance), srv_port);
static uint64_t client_link_acked(struct ETCP_SOCKET* sock, uint16_t srv_port) {
struct sockaddr_storage addr; mk_addr(&addr, srv_port);
struct ETCP_LINK* l = etcp_link_find_by_addr(sock, &addr, 0);
return l ? l->acked_packets : 0;
}
static void dump_state(void) {
printf("--- STATE DUMP ---\n");
if (server_instance) etcp_dump_all_conns(server_instance);
if (client_instance) etcp_dump_all_conns(client_instance);
debug_level_t saved = g_debug_config.category_levels[DEBUG_CATEGORY_ETCP_DUMP];
debug_set_category_level(DEBUG_CATEGORY_ETCP_DUMP, DEBUG_LEVEL_DEBUG);
if (server_instance) etcp_dump_all(server_instance);
if (client_instance) etcp_dump_all(client_instance);
struct UTUN_INSTANCE* instances[] = { server_instance, client_instance };
for (int i = 0; i < 2; i++) {
struct ETCP_CONN* conn = get_conn(instances[i]);
if (!conn) continue;
DEBUG_ERROR(DEBUG_CATEGORY_DEBUG, "%s: links=%d up=%d reinit=%u",
i ? "client" : "server", count_links(conn), count_links_up(conn), conn->reinit_count);
for (struct ETCP_LINK* link = conn->links; link; link = link->next) {
DEBUG_ERROR(DEBUG_CATEGORY_DEBUG, "link=%u/%u socket=%s local=%s remote=%s server=%u up=%u init=%u state=%u",
link->local_link_id, link->remote_link_id, link->conn ? link->conn->name : "tcp",
sockaddr_storage_to_str(&link->local_bound_addr).str, sockaddr_storage_to_str(&link->remote_addr).str,
link->is_server, link->link_status, link->initialized, link->link_state);
}
}
debug_set_category_level(DEBUG_CATEGORY_ETCP_DUMP, saved);
printf("--- END DUMP ---\n");
fflush(stdout);
}
@ -194,6 +211,64 @@ static void dump_state(void) {
dump_state(); \
goto fail; } } while (0)
// Проверяет уникальность UDP-путей и соответствие списка линков индексу каждого сокета.
static int check_unique_links(struct ETCP_CONN* conn) {
for (struct ETCP_LINK* link = conn->links; link; link = link->next) {
if (link->is_tcp) continue;
if (etcp_link_find_by_addr(link->conn, &link->remote_addr, 0) != link) {
DEBUG_ERROR(DEBUG_CATEGORY_DEBUG, "link absent/replaced in socket index: conn=%s socket=%s link=%u remote=%s",
conn->log_name, link->conn ? link->conn->name : "-", link->local_link_id,
sockaddr_storage_to_str(&link->remote_addr).str);
return -1;
}
if (queue_check_consistency(link->conn->links_queue) != 0) return -1;
}
return 0;
}
// Проверяет выбор collision-линка независимо от порядка списка и состояния handshake.
static int test_collision_selection(void) {
struct ETCP_CONN conn = {0};
struct ETCP_SOCKET sockets[3] = {0};
struct ETCP_LINK links[4] = {0};
struct sockaddr_storage addr;
mk_addr(&addr, 21000);
for (int i = 0; i < 4; i++) {
links[i].conn = &sockets[i == 0 ? 1 : 0];
links[i].etcp = &conn;
links[i].remote_addr = addr;
links[i].local_link_id = i + 1;
links[i].next = i < 3 ? &links[i + 1] : NULL;
}
conn.links = links;
links[1].is_server = 1;
mk_addr(&links[2].remote_addr, 21001);
CHECK(etcp_select_collision_link(&conn, &sockets[0], &addr) == &links[3], "selected wrong socket/role/port");
CHECK(etcp_select_collision_link(&conn, &sockets[1], &addr) == &links[0], "second local socket not distinguished");
CHECK(!etcp_select_collision_link(&conn, &sockets[2], &addr), "unrelated socket reported collision");
links[3].initialized = links[3].link_status = 1;
CHECK(etcp_select_collision_link(&conn, &sockets[0], &addr) == &links[3], "established outbound excluded from recovery");
links[3].is_server = 1;
CHECK(!etcp_select_collision_link(&conn, &sockets[0], &addr), "inbound link reported collision");
links[3].is_server = 0;
((struct sockaddr_in*)&addr)->sin_addr.s_addr = htonl(0x7f000002);
CHECK(!etcp_select_collision_link(&conn, &sockets[0], &addr), "different peer IP reported collision");
memset(&addr, 0, sizeof(addr));
struct sockaddr_in6* addr6 = (struct sockaddr_in6*)&addr;
addr6->sin6_family = AF_INET6; addr6->sin6_addr.s6_addr[15] = 1; addr6->sin6_port = htons(21000);
CHECK(!etcp_select_collision_link(&conn, &sockets[0], &addr), "different address family reported collision");
links[0].remote_addr = links[3].remote_addr = addr;
CHECK(etcp_select_collision_link(&conn, &sockets[0], &addr) == &links[3], "IPv6 selected wrong local socket");
addr6->sin6_addr.s6_addr[15] = 2;
CHECK(!etcp_select_collision_link(&conn, &sockets[0], &addr), "different IPv6 peer reported collision");
conn.links = NULL;
CHECK(!etcp_select_collision_link(&conn, &sockets[0], &addr), "empty connection reported collision");
printf("collision selection: socket, role, port, IPv4/IPv6, pending/established — PASSED\n");
return 0;
fail:
return -1;
}
// ===== poll conditions =====
static int c_conn_up_both(void) { return conn_up(server_instance) && conn_up(client_instance); }
@ -201,6 +276,11 @@ static int c_link_up(void) { return g_link && g_link->initialized && g_link->lin
static int c_link_down(void) { return g_link && g_link->link_status == 0; }
static int c_socket_link_up(void) { struct ETCP_LINK* l = first_link_on_socket(g_sock); return l && l->initialized && l->link_status == 1; }
static int c_link_on_port_up(void) { struct ETCP_LINK* l = find_link_by_remote_port(g_conn, g_port); return l && l->initialized && l->link_status == 1; }
static int c_link_on_pair_up(void) {
struct sockaddr_storage addr; mk_addr(&addr, g_port);
struct ETCP_LINK* link = etcp_link_find_by_addr(g_sock, &addr, 0);
return link && link->etcp == g_conn && link->initialized && link->link_status == 1;
}
static int c_client_all_down(void) { struct ETCP_CONN* c = get_conn(client_instance); return !c || c->links_up == 0; }
static int c_server_all_down(void) { struct ETCP_CONN* c = get_conn(server_instance); return !c || c->links_up == 0; }
@ -258,14 +338,14 @@ static int run_traffic(int n, uint64_t timeout_ms) {
}
}
// Проверяет, что конкретный линк (по серверному порту) реально перенёс данные.
// Проверяет трафик конкретной пары: клиентский сокет + серверный порт.
// Линк может пропустить короткий burst из-за keepalive-стабилизации link_status,
// поэтому повторяем до нескольких keepalive-тиков. Возвращает 0 если линк нёс трафик.
static int check_link_traffic(uint16_t srv_port, const char* label) {
static int check_link_traffic(struct ETCP_SOCKET* sock, uint16_t srv_port, const char* label) {
for (int attempt = 0; attempt < 6; attempt++) {
uint64_t before = client_link_acked(srv_port);
uint64_t before = client_link_acked(sock, srv_port);
if (run_traffic(TRAFFIC_PACKETS, TRAFFIC_TIMEOUT_MS) != 0) return -1;
if (client_link_acked(srv_port) > before) return 0;
if (client_link_acked(sock, srv_port) > before) return 0;
uasync_poll(ua, 50);
}
printf(" [%s] link never carried traffic\n", label);
@ -327,7 +407,7 @@ static int session_a(void) {
// -- per-link трафик на исходном линке --
printf("[A] traffic on initial link...\n");
CHECK(check_link_traffic(server_port, "initial") == 0, "initial link carried no traffic");
CHECK(check_link_traffic(c0, server_port, "initial") == 0, "initial link carried no traffic");
CHECK(sc->reinit_count == srv_reinit0 && cc->reinit_count == cli_reinit0, "reinit during initial traffic");
// -- V1: клиент дозванивается до нового серверного сокета S1 --
@ -343,7 +423,7 @@ static int session_a(void) {
g_sock = s1;
CHECK(poll_until(c_socket_link_up, LINK_UP_TIMEOUT_MS, "V1 link up (server)") == 0, "V1 server link did not come up");
CHECK(sc->reinit_count == srv_reinit0 && cc->reinit_count == cli_reinit0, "reinit during V1");
CHECK(check_link_traffic(sp1, "V1") == 0, "V1 link carried no traffic");
CHECK(check_link_traffic(c0, sp1, "V1") == 0, "V1 link carried no traffic");
// -- V2: сервер дозванивается до клиента --
int sp2 = alloc_port();
@ -358,7 +438,7 @@ static int session_a(void) {
g_conn = cc; g_port = sp2;
CHECK(poll_until(c_link_on_port_up, LINK_UP_TIMEOUT_MS, "V2 link up (client)") == 0, "V2 client link did not come up");
CHECK(sc->reinit_count == srv_reinit0 && cc->reinit_count == cli_reinit0, "reinit during V2");
CHECK(check_link_traffic(sp2, "V2") == 0, "V2 link carried no traffic");
CHECK(check_link_traffic(c0, sp2, "V2") == 0, "V2 link carried no traffic");
printf("[A] after adds: server links_up=%d/%d client links_up=%d/%d\n",
count_links_up(sc), count_links(sc), count_links_up(cc), count_links(cc));
@ -455,15 +535,17 @@ static int session_collision(int which) {
if (which == COL_A) {
CHECK(run_traffic(TRAFFIC_PACKETS, TRAFFIC_TIMEOUT_MS) == 0, "traffic during collision failed");
g_conn = cc; g_port = sp3;
CHECK(poll_until(c_link_on_port_up, COLLISION_TIMEOUT_MS, "collision client link up") == 0, "client link down");
g_conn = sc; g_port = cp1;
CHECK(poll_until(c_link_on_port_up, COLLISION_TIMEOUT_MS, "collision server link up") == 0, "server link down");
struct ETCP_LINK* cl = find_link_by_remote_port(cc, sp3);
struct ETCP_LINK* sl = find_link_by_remote_port(sc, cp1);
g_conn = cc; g_sock = c1; g_port = sp3;
CHECK(poll_until(c_link_on_pair_up, COLLISION_TIMEOUT_MS, "collision client link up") == 0, "client link down");
g_conn = sc; g_sock = s3; g_port = cp1;
CHECK(poll_until(c_link_on_pair_up, COLLISION_TIMEOUT_MS, "collision server link up") == 0, "server link down");
struct ETCP_LINK* cl = etcp_link_find_by_addr(c1, &s3_addr, 0);
struct ETCP_LINK* sl = etcp_link_find_by_addr(s3, &c1_addr, 0);
CHECK(cl->is_server != sl->is_server, "collision left both ends in the same role");
CHECK(count_links(sc) == 2 && count_links(cc) == 2, "collision created duplicate links");
CHECK(check_link_traffic(sp3, "collision") == 0, "new link carried no traffic");
CHECK(sl->is_server == (sc->instance->node_id > cc->instance->node_id), "collision elected wrong initiator");
CHECK(sl->remote_link_id == cl->local_link_id && cl->remote_link_id == sl->local_link_id, "collision IDs disagree");
CHECK(check_unique_links(sc) == 0 && check_unique_links(cc) == 0, "duplicate path or broken socket index");
CHECK(check_link_traffic(c1, sp3, "collision") == 0, "new link carried no traffic");
} else {
struct ETCP_LINK* good = (which == COL_B) ? ls3 : lc3;
g_link = good;
@ -480,6 +562,7 @@ static int session_collision(int which) {
CHECK(sc->reinit_count == srv_reinit0 && cc->reinit_count == cli_reinit0, "link add reset the stream");
CHECK(sc->reset_id == srv_epoch && cc->reset_id == cli_epoch, "link add changed session epochs");
CHECK(check_unique_links(sc) == 0 && check_unique_links(cc) == 0, "duplicate path after link add");
printf("[%s] PASSED\n", names[which]);
stop_instances();
return 0;
@ -489,16 +572,66 @@ fail:
return -1;
}
// Один remote IP:port на двух локальных сокетах — два пути, а не коллизия INIT.
static int session_cross_socket(void) {
printf("\n=== Session: same peer endpoint, different local sockets ===\n");
if (start_instances() != 0) { stop_instances(); return -1; }
struct ETCP_CONN* sc = get_conn(server_instance);
struct ETCP_CONN* cc = get_conn(client_instance);
struct ETCP_SOCKET* s0 = server_instance->etcp_sockets;
CHECK(sc && cc && s0, "missing initial connection/socket");
CHECK(sc->instance->node_id < cc->instance->node_id, "fixture requires server to win collision arbitration");
CHECK(run_traffic(TRAFFIC_PACKETS, TRAFFIC_TIMEOUT_MS) == 0, "initial traffic failed");
uint32_t srv_reinit = sc->reinit_count, cli_reinit = cc->reinit_count;
uint64_t srv_epoch = sc->reset_id, cli_epoch = cc->reset_id;
int sp1 = alloc_port(), cp1 = alloc_port();
struct ETCP_SOCKET* s1 = add_udp_socket(server_instance, sp1, "cross_s1");
struct ETCP_SOCKET* c1 = add_udp_socket(client_instance, cp1, "cross_c1");
CHECK(s1 && c1, "failed to add sockets");
struct sockaddr_storage c1_addr, s0_addr, s1_addr;
mk_addr(&c1_addr, cp1); mk_addr(&s0_addr, server_port); mk_addr(&s1_addr, sp1);
CHECK(etcp_link_new(sc, s1, &c1_addr, 0), "outbound S1 -> C1 failed");
CHECK(etcp_link_new(cc, c1, &s0_addr, 0), "outbound C1 -> S0 failed");
struct ETCP_SOCKET* sockets[] = {s0, s1, c1, c1};
struct ETCP_CONN* conns[] = {sc, sc, cc, cc};
uint16_t ports[] = {cp1, cp1, server_port, sp1};
for (int i = 0; i < 4; i++) {
g_sock = sockets[i]; g_conn = conns[i]; g_port = ports[i];
CHECK(poll_until(c_link_on_pair_up, LINK_UP_TIMEOUT_MS, "independent path handshake") == 0,
"path socket=%s -> port=%u did not establish", g_sock->name, g_port);
}
struct ETCP_LINK* in = etcp_link_find_by_addr(s0, &c1_addr, 0);
struct ETCP_LINK* out = etcp_link_find_by_addr(s1, &c1_addr, 0);
struct ETCP_LINK* peer_out = etcp_link_find_by_addr(c1, &s0_addr, 0);
struct ETCP_LINK* peer_in = etcp_link_find_by_addr(c1, &s1_addr, 0);
CHECK(in != out && in->is_server == 1 && out->is_server == 0, "independent paths lost their roles");
CHECK(peer_out->is_server == 0 && peer_in->is_server == 1, "peer paths lost their roles");
CHECK(in->remote_link_id == peer_out->local_link_id && peer_out->remote_link_id == in->local_link_id, "S0/C1 IDs disagree");
CHECK(out->remote_link_id == peer_in->local_link_id && peer_in->remote_link_id == out->local_link_id, "S1/C1 IDs disagree");
CHECK(check_unique_links(sc) == 0 && check_unique_links(cc) == 0, "duplicate independent path");
CHECK(check_link_traffic(c1, server_port, "S0/C1") == 0, "S0 path carried no traffic");
CHECK(check_link_traffic(c1, sp1, "S1/C1") == 0, "S1 path carried no traffic");
CHECK(sc->reinit_count == srv_reinit && cc->reinit_count == cli_reinit, "independent paths reset stream");
CHECK(sc->reset_id == srv_epoch && cc->reset_id == cli_epoch, "independent paths changed epochs");
printf("independent S0/C1 and S1/C1 paths — PASSED\n");
stop_instances();
return 0;
fail:
stop_instances();
return -1;
}
int main(void) {
if (create_temp_configs() != 0) return 1;
debug_config_init();
debug_set_level(DEBUG_LEVEL_ERROR);
if (getenv("UTUN_TEST_DEBUG")) {
debug_set_level(DEBUG_LEVEL_DEBUG);
debug_set_level(DEBUG_LEVEL_WARN);
debug_set_category_level(DEBUG_CATEGORY_ETCP, DEBUG_LEVEL_DEBUG);
debug_set_category_level(DEBUG_CATEGORY_CONNECTION, DEBUG_LEVEL_DEBUG);
debug_set_category_level(DEBUG_CATEGORY_KEEPALIVE, DEBUG_LEVEL_DEBUG);
debug_set_category_level(DEBUG_CATEGORY_BGP, DEBUG_LEVEL_DEBUG);
debug_set_category_level(DEBUG_CATEGORY_DEBUG, DEBUG_LEVEL_DEBUG);
}
@ -509,10 +642,12 @@ int main(void) {
ua = uasync_create();
int rc = 0;
rc |= test_collision_selection() != 0;
rc |= session_a() != 0;
rc |= session_collision(COL_A) != 0;
rc |= session_collision(COL_B) != 0;
rc |= session_collision(COL_C) != 0;
rc |= session_cross_socket() != 0;
if (ua) { uasync_poll(ua, 0); uasync_destroy(ua, 0); ua = NULL; }
cleanup_temp_configs();

126
tests/test_etcp_router.c

@ -1,6 +1,6 @@
// test_etcp_router.c — Integration test for etcp_router with dummynet congestion
// 2 ноды в одном UASYNC, dummynet-прокси между ними.
// Фазы: baseline (без congestion) + congestion cycles (fill→push→drain→verify)
// test_etcp_router.c — интеграционный тест backpressure/recovery роутера через UDP и TCP.
// 2 ноды в одном UASYNC; UDP начинает через dummynet, NCD может достраивать прямые пути.
// Фазы: baseline + ограниченные циклы fill→push→drain→verify, независимо от скорости линков.
// Проверка: bitmap (нет потерь/дубликатов), per-packet payload integrity, порядок.
#include <stdio.h>
#include <stdlib.h>
@ -22,6 +22,7 @@
#include "etcp.h"
#include "etcp_connections.h"
#include "etcp_api.h"
#include "etcp_dump.h"
#include "node_conn_direct.h"
#include "topo_node.h"
#include "etcp_router.h"
@ -43,7 +44,8 @@
#define BASELINE_PACKETS 20
#define PACKETS_PER_CYCLE 20
#define CONGESTION_CYCLES 5
#define BITMAP_SIZE 2000 // с запасом (baseline + cycles * (fill+push))
#define MAX_FILL_PACKETS (ROUTER_MAX_INFLIGHT + ROUTER_MAX_SEND_Q_PACKETS)
#define BITMAP_SIZE (BASELINE_PACKETS + CONGESTION_CYCLES * (MAX_FILL_PACKETS + PACKETS_PER_CYCLE))
#define LOW_BW_KBPS 80 // ~55 pkt/s @ ~180 bytes — медленнее чем можем слать
#define DUMMY_QUEUE_SIZE 300 // чтобы dummynet не дропал
#define MAX_PAYLOAD 200
@ -258,6 +260,31 @@ static int send_one_pkt(uint32_t seq, int data_len) {
}
// ======================== Server handler ========================
static void log_failure_state(uint32_t seq) {
DEBUG_ERROR(DEBUG_CATEGORY_DEBUG,
"router test: tcp=%d seq=%u expected=%u bitmap=%u state=%d cycle=%d sent=%u received=%u cycle_sent=%u",
g_tcp, seq, g_expected_seq, BITMAP_SIZE, g_state, g_cycle + 1, g_total_sent, g_total_rcvd, g_cycle_sent);
struct ETCP_ROUTER_CONN* rc = etcp_router_conn_get(cli, TOPO_GROUP_UTUN, server_node_id, TEST_SVC_ID);
if (rc) {
DEBUG_ERROR(DEBUG_CATEGORY_DEBUG,
"router queues: tx_seq=%u tx_sent=%u tx_acked=%u send_q=%d inflight=%d limit=%u backpressure=%d",
rc->tx_seq, rc->tx_sent, rc->tx_acked, rc->send_q->count, rc->inflight_q->count,
rc->inflight_limit, g_backpressure_seen);
}
if (dn) {
const struct dummynet_stats* sf = dummynet_get_stats(dn, DUMMYNET_FORWARD);
const struct dummynet_stats* sb = dummynet_get_stats(dn, DUMMYNET_BACKWARD);
DEBUG_ERROR(DEBUG_CATEGORY_DEBUG,
"dummynet: port=%d server=%d client=%d fwd_recv=%llu fwd_sent=%llu bwd_recv=%llu bwd_sent=%llu",
g_dn_port, g_srv_port, g_cli_port, (unsigned long long)sf->recv, (unsigned long long)sf->sent,
(unsigned long long)sb->recv, (unsigned long long)sb->sent);
} else DEBUG_ERROR(DEBUG_CATEGORY_DEBUG, "router test: dummynet is absent, set_bw has no effect");
debug_level_t saved = g_debug_config.category_levels[DEBUG_CATEGORY_ETCP_DUMP];
debug_set_category_level(DEBUG_CATEGORY_ETCP_DUMP, DEBUG_LEVEL_DEBUG);
etcp_dump_all(srv); etcp_dump_all(cli);
debug_set_category_level(DEBUG_CATEGORY_ETCP_DUMP, saved);
}
static void srv_handler(struct ETCP_CONN* conn, struct ll_entry* entry) {
// Restart/close notification: [svc_id][src][dst] без payload
if (entry && entry->dgram && entry->len == ROUTER_SVC_HDR_SIZE) {
@ -286,6 +313,7 @@ static void srv_handler(struct ETCP_CONN* conn, struct ll_entry* entry) {
// Check seq in range
if (seq >= BITMAP_SIZE) {
if (!g_fail) log_failure_state(seq);
printf("[FAIL] srv_handler: seq=%u out of bitmap range\n", seq);
queue_dgram_free(entry); queue_entry_free(entry);
g_fail = 3; g_fail_seq = seq; g_done = -1;
@ -385,10 +413,8 @@ static void set_bw(uint32_t kbps) {
}
// ======================== Monitor ========================
static int g_conn_ok = 0;
static int g_conn_delay = 0;
static uint32_t g_baseline_sent = 0;
static double g_cycle_drain_start_ms = 0;
static void monitor(void* arg) {
(void)arg;
@ -399,7 +425,6 @@ static void monitor(void* arg) {
if (conn_established(srv) && conn_established(cli)) {
g_conn_delay++;
if (g_conn_delay < 50) { g_mon_id = uasync_set_timeout(ua, 10, NULL, monitor, "mon"); return; }
g_conn_ok = 1;
t_connect_ms = toc_ms(&t_phase_start);
printf("Phase 1 — Connect: %7.1f ms\n", t_connect_ms);
// Run loopback test
@ -411,6 +436,10 @@ static void monitor(void* arg) {
tic(&t_phase_start);
printf("Phase 2 — Baseline (%d packets, no congestion):\n", BASELINE_PACKETS);
}
if (toc_ms(&t_phase_start) >= CONNECT_TIMEOUT_MS) {
DEBUG_ERROR(DEBUG_CATEGORY_DEBUG, "router test: connection timeout tcp=%d", g_tcp);
log_failure_state(g_total_sent); g_done = -1; return;
}
g_mon_id = uasync_set_timeout(ua, 10, NULL, monitor, "mon");
return;
}
@ -449,51 +478,67 @@ static void monitor(void* arg) {
}
}
if (g_total_rcvd >= BASELINE_PACKETS) {
if (g_total_rcvd == BASELINE_PACKETS && rc->tx_acked == rc->tx_seq && !rc->send_q->count && !rc->inflight_q->count) {
t_baseline_ms = toc_ms(&t_phase_start);
printf(" baseline: %7.1f ms sent=%u rcvd=%u\n", t_baseline_ms, g_baseline_sent, g_total_rcvd);
// Start congestion cycles
g_cycle = 0;
g_state = ST_CYCLE_FILL;
tic(&t_phase_start);
printf("Phase 3 — Congestion cycles (%d cycles, %d pkt/cycle, bw=%u kbps):\n",
CONGESTION_CYCLES, PACKETS_PER_CYCLE, LOW_BW_KBPS);
printf("Phase 3 — Backpressure cycles (%d cycles, bounded fill + %d packets; UDP proxy bw=%u kbps):\n",
CONGESTION_CYCLES, PACKETS_PER_CYCLE, dn ? LOW_BW_KBPS : 0);
}
g_mon_id = uasync_set_timeout(ua, 10, NULL, monitor, "mon");
return;
}
// Phase: cycle fill — send until backpressure
// Заполняем за один callback: ACK не может освободить очередь между отправками.
if (g_state == ST_CYCLE_FILL) {
set_bw(LOW_BW_KBPS);
g_cycle_start_seq = g_total_sent;
g_cycle_sent = 0;
g_backpressure_seen = 0;
g_state = ST_CYCLE_FILL + 100; // sub-state: actively filling
tic(&t_phase_start);
g_mon_id = uasync_set_timeout(ua, 1, NULL, monitor, "mon");
return;
}
if (g_state == ST_CYCLE_FILL + 100) {
// Send one packet per tick until backpressure
int data_len = 32 + (g_total_sent % 64);
int ret = send_one_pkt(g_total_sent, data_len);
if (ret == 0) {
g_total_sent++;
g_cycle_sent++;
} else {
// Backpressure hit — inflight+send_q full
struct ETCP_ROUTER_CONN* rc = etcp_router_conn_get(cli, TOPO_GROUP_UTUN, server_node_id, TEST_SVC_ID);
uint32_t start_tx = rc->tx_seq;
while (g_cycle_sent <= MAX_FILL_PACKETS) {
int data_len = 32 + (g_total_sent % 64);
int ret = send_one_pkt(g_total_sent, data_len);
if (ret == 0) {
g_total_sent++; g_cycle_sent++;
continue;
}
if (rc->send_q->count != ROUTER_MAX_SEND_Q_PACKETS || rc->tx_seq != start_tx + g_cycle_sent) {
DEBUG_ERROR(DEBUG_CATEGORY_DEBUG, "fill failed without backpressure: ret=%d send_q=%d tx_seq=%u expected=%u",
ret, rc->send_q->count, rc->tx_seq, start_tx + g_cycle_sent);
log_failure_state(g_total_sent); g_done = -1; return;
}
// Повторный отказ не должен потреблять seq или менять очереди.
int inflight = rc->inflight_q->count;
for (int retry = 0; retry < 3; retry++) {
if (send_one_pkt(g_total_sent, data_len) == 0 || rc->tx_seq != start_tx + g_cycle_sent
|| rc->send_q->count != ROUTER_MAX_SEND_Q_PACKETS || rc->inflight_q->count != inflight) {
DEBUG_ERROR(DEBUG_CATEGORY_DEBUG, "backpressure retry changed sequence/queues: retry=%d", retry);
log_failure_state(g_total_sent); g_done = -1; return;
}
}
g_backpressure_seen = 1;
t_cycle_fill_ms[g_cycle] = toc_ms(&t_phase_start);
g_cycle_fill_count[g_cycle] = g_cycle_sent;
printf(" cycle %d fill: %7.1f ms sent=%u (backpressure)\n",
g_cycle + 1, t_cycle_fill_ms[g_cycle], g_cycle_sent);
// Transition to push phase
g_state = ST_CYCLE_PUSH;
g_cycle_push_ok = 0;
g_cycle_drop_count[g_cycle] = 0;
tic(&t_phase_start);
break;
}
if (!g_backpressure_seen) {
DEBUG_ERROR(DEBUG_CATEGORY_DEBUG, "fill exceeded capacity without backpressure: accepted=%u bound=%u",
g_cycle_sent, MAX_FILL_PACKETS);
log_failure_state(g_total_sent); g_done = -1; return;
}
t_cycle_fill_ms[g_cycle] = toc_ms(&t_phase_start);
g_cycle_fill_count[g_cycle] = g_cycle_sent;
printf(" cycle %d fill: %7.1f ms sent=%u send_q=%d inflight=%d; four refusals preserved seq=%u\n",
g_cycle + 1, t_cycle_fill_ms[g_cycle], g_cycle_sent, rc->send_q->count, rc->inflight_q->count, rc->tx_seq);
g_state = ST_CYCLE_PUSH;
g_cycle_push_ok = 0;
g_cycle_drop_count[g_cycle] = 4;
g_total_drop += 4;
tic(&t_phase_start);
g_mon_id = uasync_set_timeout(ua, 1, NULL, monitor, "mon");
return;
}
@ -509,7 +554,6 @@ static void monitor(void* arg) {
set_bw(0);
g_state = ST_CYCLE_DRAIN;
tic(&t_phase_start);
g_cycle_drain_start_ms = toc_ms(&t_phase_start);
g_mon_id = uasync_set_timeout(ua, 10, NULL, monitor, "mon");
return;
}
@ -528,10 +572,11 @@ static void monitor(void* arg) {
return;
}
// Phase: cycle drain — wait for all sent packets to arrive
// До следующего цикла ждём доставку, ACK и освобождение обеих очередей.
if (g_state == ST_CYCLE_DRAIN) {
uint32_t cycle_target = g_cycle_start_seq + g_cycle_sent;
if (g_total_rcvd >= cycle_target) {
struct ETCP_ROUTER_CONN* rc = etcp_router_conn_get(cli, TOPO_GROUP_UTUN, server_node_id, TEST_SVC_ID);
if (g_total_rcvd == cycle_target && rc->tx_acked == rc->tx_seq && !rc->send_q->count && !rc->inflight_q->count) {
t_cycle_drain_ms[g_cycle] = toc_ms(&t_phase_start);
printf(" cycle %d drain: %7.1f ms rcvd=%u/%u\n",
g_cycle + 1, t_cycle_drain_ms[g_cycle], g_total_rcvd, cycle_target);
@ -541,6 +586,9 @@ static void monitor(void* arg) {
} else {
g_state = ST_CYCLE_FILL;
}
} else if (toc_ms(&t_phase_start) >= DRAIN_TIMEOUT_MS) {
DEBUG_ERROR(DEBUG_CATEGORY_DEBUG, "router drain timeout: cycle=%d target=%u", g_cycle + 1, cycle_target);
log_failure_state(g_total_sent); g_done = -1; return;
}
g_mon_id = uasync_set_timeout(ua, 10, NULL, monitor, "mon");
return;
@ -606,7 +654,8 @@ static void monitor(void* arg) {
printf("[FAIL] %d gaps in bitmap (losses)\n", gaps);
g_done = -1;
} else {
printf("[PASS] test_etcp_router — %u packets, 0 loss, 0 dup, 0 corrupt\n", g_total_sent);
printf("[PASS] test_etcp_router — %u packets, %d backpressure cycles, 0 loss, 0 dup, 0 corrupt\n",
g_total_sent, CONGESTION_CYCLES);
g_done = 1;
}
g_mon_id = NULL;
@ -622,6 +671,7 @@ static void timeout_cb(void* arg) {
if (!g_done) {
printf("[FAIL] TIMEOUT: state=%d cycle=%d sent=%u rcvd=%u drops=%u\n",
g_state, g_cycle, g_total_sent, g_total_rcvd, g_total_drop);
log_failure_state(g_total_sent);
g_done = -1;
}
if (g_mon_id) { uasync_cancel_timeout(ua, g_mon_id); g_mon_id = NULL; }
@ -658,7 +708,7 @@ int main(void) {
write_configs();
printf("=== test_etcp_router (with dummynet congestion) ===\n");
printf("=== test_etcp_router (%s, bounded backpressure and transport recovery) ===\n", g_tcp ? "TCP" : "UDP");
debug_config_init();
debug_set_level(DEBUG_LEVEL_ERROR);

Loading…
Cancel
Save