diff --git a/src/conn_mgr.c b/src/conn_mgr.c index 857fb6fa..d7541324 100644 --- a/src/conn_mgr.c +++ b/src/conn_mgr.c @@ -36,6 +36,7 @@ static uint16_t cm_get_node_min_rtt(struct NODEINFO_Q* nq); static int cm_has_direct_ip(struct NODEINFO_Q* nq); static int cm_has_local_addr(struct NODEINFO_Q* nq); static void cm_entry_destroy(struct CONN_MGR_ENTRY* entry); +static void cm_exchange_probe_retry_cb(void* arg); static uint64_t cm_get_node_max_probe_time(struct NODEINFO_Q* nq); static int cm_is_rtt_fresh(struct NODEINFO_Q* nq, uint64_t now_tb); @@ -92,6 +93,14 @@ void conn_mgr_destroy(struct CONN_MGR* mgr) { if (mgr->candidate_ping_timer) { uasync_cancel_timeout(mgr->instance->ua, mgr->candidate_ping_timer); mgr->candidate_ping_timer = NULL; } etcp_router_unbind(mgr->instance, ETCP_ID_CONN_MGR); for (size_t i = 0; i < mgr->entry_count; i++) cm_entry_destroy(&mgr->entries[i]); + struct cm_exchange_pending* ep = mgr->exchange_pending; + while (ep) { + struct cm_exchange_pending* next = ep->next; + if (ep->timeout_timer) uasync_cancel_timeout(mgr->instance->ua, ep->timeout_timer); + u_free(ep); + ep = next; + } + mgr->exchange_pending = NULL; u_free(mgr->entries); mgr->initialized = 0; DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "conn_mgr: destroyed, entries=%zu", (size_t)mgr->entry_count); @@ -102,6 +111,16 @@ static void cm_entry_destroy(struct CONN_MGR_ENTRY* entry) { if (!entry || entry->state == CONN_MGR_STATE_DISCONNECTED) return; if (entry->idle_timer) { uasync_cancel_timeout(entry->mgr->instance->ua, entry->idle_timer); entry->idle_timer = NULL; } if (entry->main.timer) { uasync_cancel_timeout(entry->mgr->instance->ua, entry->main.timer); entry->main.timer = NULL; } + struct cm_exchange_pending** pp = &entry->mgr->exchange_pending; + while (*pp) { + if ((*pp)->entry == entry) { + if ((*pp)->timeout_timer && (*pp)->timeout_timer != entry->main.timer) { + uasync_cancel_timeout(entry->mgr->instance->ua, (*pp)->timeout_timer); + (*pp)->timeout_timer = NULL; + } + struct cm_exchange_pending* ep = *pp; *pp = ep->next; u_free(ep); + } else pp = &(*pp)->next; + } entry->state = CONN_MGR_STATE_DISCONNECTED; entry->conn_type = CONN_TYPE_NONE; } @@ -758,6 +777,26 @@ static void cm_compute_intermediaries(struct CONN_MGR_ENTRY* entry, struct CONN_ cm_deliver_result(entry, CONN_MGR_OK); } +static void cm_exchange_probe_retry_cb(void* arg) { + struct cm_exchange_pending* ep = (struct cm_exchange_pending*)arg; + struct CONN_MGR* mgr = ep->entry->mgr; + ep->timeout_timer = NULL; + uint8_t need_probe = 0; + uint64_t now_tb = get_time_tb(); + for (uint8_t i = 0; i < ep->cached_resp.my_count && i < 4; i++) { + struct NODEINFO_Q* nq = nodeinfo_find_by_id(mgr->instance->bgp, ep->cached_resp.my_candidates[i].node_id); + if (!nq || !cm_is_rtt_fresh(nq, now_tb)) { need_probe++; break; } + } + if (need_probe == 0) { + cm_compute_intermediaries(ep->entry, &ep->cached_resp); + struct cm_exchange_pending** pp = &mgr->exchange_pending; + while (*pp) { if (*pp == ep) { *pp = ep->next; break; } pp = &(*pp)->next; } + u_free(ep); + } else { + ep->timeout_timer = uasync_set_timeout(mgr->instance->ua, 2000, ep, cm_exchange_probe_retry_cb, "cm_exch_probe"); + } +} + static void cm_handle_interm_exchange_resp(struct CONN_MGR* mgr, const uint8_t* data, size_t len) { const struct CONN_MGR_INTERM_EXCHANGE_RESP* resp = (const struct CONN_MGR_INTERM_EXCHANGE_RESP*)data; size_t min_size = offsetof(struct CONN_MGR_INTERM_EXCHANGE_RESP, your_candidates); @@ -782,7 +821,7 @@ static void cm_handle_interm_exchange_resp(struct CONN_MGR* mgr, const uint8_t* struct cm_exchange_pending** pp = &mgr->exchange_pending; while (*pp) { if (*pp == ep) { *pp = ep->next; break; } pp = &(*pp)->next; } u_free(ep); return; } - uasync_set_timeout(mgr->instance->ua, 2000, ep, NULL, "cm_exch_probe"); + ep->timeout_timer = uasync_set_timeout(mgr->instance->ua, 2000, ep, cm_exchange_probe_retry_cb, "cm_exch_probe"); } static void cm_handle_interm_selected(struct CONN_MGR* mgr, const uint8_t* data, size_t len) { diff --git a/src/utun_instance.c b/src/utun_instance.c index b1dc7685..ea6afe84 100644 --- a/src/utun_instance.c +++ b/src/utun_instance.c @@ -12,6 +12,7 @@ #include "etcp_connections.h" #include "etcp.h" #include "conn_mgr.h" +#include "stcp_server.h" #include "control_server.h" #include "msg_transport.h" #include "../lib/u_async.h" @@ -400,6 +401,12 @@ void utun_instance_destroy(struct UTUN_INSTANCE *instance) { // Cleanup firewall fw_free(&instance->fw); + // Cleanup TCP server + if (instance->stcp_server) { + stcp_server_destroy(instance->stcp_server); + instance->stcp_server = NULL; + } + // Cleanup config if (instance->config) { DEBUG_INFO(DEBUG_CATEGORY_MEMORY, "[INSTANCE_DESTROY] Freeing configuration"); diff --git a/tests/Makefile.am b/tests/Makefile.am index e6626fdd..807855f5 100644 --- a/tests/Makefile.am +++ b/tests/Makefile.am @@ -193,7 +193,7 @@ test_lwip_tcp_LDADD = \ $(COMMON_LIBS) test_etcp_router_SOURCES = test_etcp_router.c -test_etcp_router_LDADD = $(ETCP_FULL_OBJS) $(SECURE_CHANNEL_OBJS) $(CRYPTO_LIBS) $(COMMON_LIBS) +test_etcp_router_LDADD = $(top_builddir)/src/utun-dummynet.o $(ETCP_FULL_OBJS) $(SECURE_CHANNEL_OBJS) $(CRYPTO_LIBS) $(COMMON_LIBS) test_etcp_router_unit_SOURCES = test_etcp_router_unit.c test_etcp_router_unit_LDADD = $(ETCP_FULL_OBJS) $(SECURE_CHANNEL_OBJS) $(CRYPTO_LIBS) $(COMMON_LIBS) diff --git a/tests/test_etcp_router.c b/tests/test_etcp_router.c index b3d38e32..902976dc 100644 --- a/tests/test_etcp_router.c +++ b/tests/test_etcp_router.c @@ -1,8 +1,16 @@ -// test_etcp_router.c — Unit test for etcp_router: bind/unbind, send, loopback, forward, reply +// test_etcp_router.c — Integration test for etcp_router with dummynet congestion +// 2 ноды в одном UASYNC, dummynet-прокси между ними. +// Фазы: baseline (без congestion) + congestion cycles (fill→push→drain→verify) +// Проверка: bitmap (нет потерь/дубликатов), per-packet payload integrity, порядок. #include #include #include +#include +#include +#include + #include "../lib/platform_compat.h" +#include "../lib/socket_compat.h" #include "test_utils.h" #ifdef _WIN32 #include @@ -10,13 +18,12 @@ #else #include #endif -#include -#include #include "../src/etcp.h" #include "../src/etcp_connections.h" #include "../src/etcp_api.h" #include "../src/etcp_router.h" +#include "../src/dummynet.h" #include "../src/config_parser.h" #include "../src/utun_instance.h" #include "../src/routing.h" @@ -29,58 +36,168 @@ #include "../lib/debug_config.h" #include "../lib/mem.h" -#define TEST_SVC_ID 0x01 -#define TEST_TIMEOUT_MS 5000 -#define TOTAL_PACKETS 20 -#define MAX_PAYLOAD 200 - +// ======================== Configuration ======================== +#define TEST_SVC_ID 0x01 +#define BASELINE_PACKETS 20 +#define PACKETS_PER_CYCLE 20 +#define CONGESTION_CYCLES 5 +#define BITMAP_SIZE 2000 // с запасом (baseline + cycles * (fill+push)) +#define LOW_BW_KBPS 80 // ~55 pkt/s @ ~180 bytes — медленнее чем можем слать +#define DUMMY_QUEUE_SIZE 300 // чтобы dummynet не дропал +#define MAX_PAYLOAD 200 +#define CONNECT_TIMEOUT_MS 10000 +#define TEST_TIMEOUT_MS 120000 +#define DRAIN_TIMEOUT_MS 30000 // per-cycle drain timeout + +// ======================== Globals ======================== static char temp_dir[] = "/tmp/utun_test_XXXXXX"; -static char server_conf[256], client_conf[256]; +static char server_conf[512], client_conf[512]; static struct UTUN_INSTANCE* srv = NULL; static struct UTUN_INSTANCE* cli = NULL; static struct UASYNC* ua = NULL; - -static int g_ok = 0; -static int g_test_done = 0; -static int g_phase = 0; // 0=wait_conn, 1=fwd, 2=reply +static struct dummynet* dn = NULL; + +static int g_dn_port = 0, g_srv_port = 0, g_cli_port = 0; +static const uint64_t server_node_id = 0x1111111111111111ULL; +static const uint64_t client_node_id = 0x2222222222222222ULL; + +// Bitmap: 0=not received, 1=received +static uint8_t rcvd_bitmap[BITMAP_SIZE]; +// Corruption/dup tracking +static int g_fail = 0; // 0=ok, 1=corrupt, 2=dup, 3=out_of_range, 4=seq_mismatch +static uint32_t g_fail_seq = 0; +static uint32_t g_fail_detail = 0; + +// Counters +static uint32_t g_total_sent = 0; // total successfully sent (etcp_route_send==0) +static uint32_t g_total_rcvd = 0; // total received by server handler +static uint32_t g_total_drop = 0; // total etcp_route_send returned -1 + +// Expected next seq for in-order delivery check +static uint32_t g_expected_seq = 0; + +// State machine +enum { + ST_WAIT_CONN, + ST_BASELINE_SEND, + ST_BASELINE_WAIT, + ST_CYCLE_FILL, + ST_CYCLE_PUSH, + ST_CYCLE_DRAIN, + ST_FINAL +}; +static int g_state = ST_WAIT_CONN; +static int g_cycle = 0; +static uint32_t g_cycle_start_seq = 0; +static uint32_t g_cycle_sent = 0; // packets sent in current cycle (fill+push) +static uint32_t g_cycle_push_ok = 0; // successful pushes in PUSH phase +static int g_backpressure_seen = 0; // saw etcp_route_send=-1 in FILL + +// Timing +static struct timespec t_phase_start; +static double t_connect_ms = 0, t_baseline_ms = 0; +static double t_cycle_fill_ms[CONGESTION_CYCLES]; +static double t_cycle_push_ms[CONGESTION_CYCLES]; +static double t_cycle_drain_ms[CONGESTION_CYCLES]; +static uint32_t g_cycle_fill_count[CONGESTION_CYCLES]; +static uint32_t g_cycle_drop_count[CONGESTION_CYCLES]; + +// Monitor static void* g_mon_id = NULL; +static int g_done = 0; // 0=running, 1=pass, -1=fail + +// ======================== Timing helpers ======================== +static void tic(struct timespec* t) { + clock_gettime(CLOCK_MONOTONIC, t); +} + +static double toc_ms(const struct timespec* start) { + struct timespec end; + clock_gettime(CLOCK_MONOTONIC, &end); + return (end.tv_sec - start->tv_sec) * 1000.0 + (end.tv_nsec - start->tv_nsec) / 1000000.0; +} + +// ======================== Port allocation ======================== +static int alloc_consecutive_ports(int* base_port) { + for (int attempt = 0; attempt < 200; attempt++) { + socket_t s = socket_create_udp(AF_INET); + if (s == SOCKET_INVALID) return -1; + struct sockaddr_in a; + memset(&a, 0, sizeof(a)); + a.sin_family = AF_INET; + a.sin_addr.s_addr = inet_addr("127.0.0.1"); + a.sin_port = 0; // auto-assign + if (bind(s, (struct sockaddr*)&a, sizeof(a)) != 0) { socket_close_wrapper(s); continue; } + socklen_t alen = sizeof(a); + getsockname(s, (struct sockaddr*)&a, &alen); + int port = ntohs(a.sin_port); + socket_close_wrapper(s); + // Need port-1, port, port+1 all free + int ok = 1; + for (int p = port - 1; p <= port + 1 && ok; p++) { + if (p < 1024) { ok = 0; break; } + socket_t ts = socket_create_udp(AF_INET); + if (ts == SOCKET_INVALID) { ok = 0; break; } + struct sockaddr_in ta; + memset(&ta, 0, sizeof(ta)); + ta.sin_family = AF_INET; + ta.sin_addr.s_addr = inet_addr("127.0.0.1"); + ta.sin_port = htons((uint16_t)p); + if (bind(ts, (struct sockaddr*)&ta, sizeof(ta)) != 0) ok = 0; + socket_close_wrapper(ts); + } + if (ok) { *base_port = port; return 0; } + } + return -1; +} + +// ======================== Config generation ======================== +static const char* srv_priv = "38240cb82199e504686507f11f6eaa4f740fde6f0c425c495e49a523019a5d68"; +static const char* srv_pub = "ce8871f07fa056c636d297115f231b08c29cdf94e0d440fce83a07c34416d36a"; +static const char* cli_priv = "704f2e012c8fa8768130cb0f988a997dccb628372bc5ceccacc78dcbfec5916f"; +static const char* cli_pub = "b3193173def895bd0fcea6f86af077c7d77216f10395275f627ac18242ec0f01"; -// Server config — node 0x1111... -static const char* srv_cfg = - "[global]\n" - "my_node_id=0x1111111111111111\n" - "my_private_key=38240cb82199e504686507f11f6eaa4f740fde6f0c425c495e49a523019a5d68\n" - "my_public_key=ce8871f07fa056c636d297115f231b08c29cdf94e0d440fce83a07c34416d36a\n" - "tun_ip=10.99.0.1/24\n" - "tun_ifname=tun99\n" - "[server: s1]\n" - "addr=127.0.0.1:9041\n" - "type=public\n" - "[allowed_keys]\n" - "allow_all=1\n"; - -// Client config — node 0x2222... -static const char* cli_cfg = - "[global]\n" - "my_node_id=0x2222222222222222\n" - "my_private_key=704f2e012c8fa8768130cb0f988a997dccb628372bc5ceccacc78dcbfec5916f\n" - "my_public_key=b3193173def895bd0fcea6f86af077c7d77216f10395275f627ac18242ec0f01\n" - "tun_ip=10.99.0.2/24\n" - "tun_ifname=tun98\n" - "[server: s1]\n" - "addr=127.0.0.1:9042\n" - "type=public\n" - "[client: c1]\n" - "keepalive=1\n" - "peer_public_key=ce8871f07fa056c636d297115f231b08c29cdf94e0d440fce83a07c34416d36a\n" - "link=s1:127.0.0.1:9041\n"; - -// ======================== Test state ======================== -static int fwd_sent = 0, fwd_rcvd = 0; -static int reply_sent = 0, reply_rcvd = 0; -static uint8_t expected_data[MAX_PAYLOAD]; -static uint64_t server_node_id = 0x1111111111111111ULL; -static uint64_t client_node_id = 0x2222222222222222ULL; +static void write_configs(void) { + // Server: listens on g_srv_port + snprintf(server_conf, sizeof(server_conf), "%s/server.conf", temp_dir); + FILE* f = fopen(server_conf, "w"); + if (!f) { fprintf(stderr, "fopen server fail\n"); exit(1); } + fprintf(f, + "[global]\n" + "my_node_id=0x1111111111111111\n" + "my_private_key=%s\n" + "my_public_key=%s\n" + "tun_ip=10.99.0.1/24\n" + "tun_ifname=tun99\n" + "[server: s1]\n" + "addr=127.0.0.1:%d\n" + "type=public\n" + "[allowed_keys]\n" + "allow_all=1\n", + srv_priv, srv_pub, g_srv_port); + fclose(f); + + // Client: listens on g_cli_port, connects to dummynet on g_dn_port + snprintf(client_conf, sizeof(client_conf), "%s/client.conf", temp_dir); + f = fopen(client_conf, "w"); + if (!f) { fprintf(stderr, "fopen client fail\n"); exit(1); } + fprintf(f, + "[global]\n" + "my_node_id=0x2222222222222222\n" + "my_private_key=%s\n" + "my_public_key=%s\n" + "tun_ip=10.99.0.2/24\n" + "tun_ifname=tun98\n" + "[server: s1]\n" + "addr=127.0.0.1:%d\n" + "type=public\n" + "[client: c1]\n" + "keepalive=1\n" + "peer_public_key=%s\n" + "link=s1:127.0.0.1:%d\n", + cli_priv, cli_pub, g_cli_port, srv_pub, g_dn_port); + fclose(f); +} // ======================== Helpers ======================== static int conn_established(struct UTUN_INSTANCE* inst) { @@ -89,7 +206,8 @@ static int conn_established(struct UTUN_INSTANCE* inst) { for (struct ETCP_LINK* l = c->links; l; l = l->next) { if (l->initialized && c->crypto_ctx.initialized) { int ok = 0; - for (int i = 0; i < SC_SESSION_KEY_SIZE; i++) if (c->crypto_ctx.session_key[i] != 0) { ok = 1; break; } + for (int i = 0; i < SC_SESSION_KEY_SIZE; i++) + if (c->crypto_ctx.session_key[i] != 0) { ok = 1; break; } if (ok) return 1; } } @@ -97,15 +215,38 @@ static int conn_established(struct UTUN_INSTANCE* inst) { return 0; } -static struct ETCP_CONN* first_conn(struct UTUN_INSTANCE* inst) { - return inst ? inst->connections : NULL; +static void gen_payload(uint32_t seq, uint8_t* buf, int len) { + for (int i = 0; i < len; i++) + buf[i] = (uint8_t)((i ^ seq ^ 0xAA) & 0xFF); +} + +// Build and send one router packet. Returns 0 on success, -1 on failure. +static int send_one_pkt(uint32_t seq, int data_len) { + if (data_len > MAX_PAYLOAD) data_len = MAX_PAYLOAD; + uint8_t buf[10 + MAX_PAYLOAD]; + buf[0] = TEST_SVC_ID; + buf[1] = 0x01; // DATA subcmd + memcpy(buf + 2, &seq, 4); + memcpy(buf + 6, &data_len, 4); + gen_payload(seq, buf + 10, data_len); + + struct ll_entry* e = queue_entry_new(0); + if (!e) return -1; + e->dgram = u_malloc(10 + data_len); + if (!e->dgram) { queue_entry_free(e); return -1; } + memcpy(e->dgram, buf, 10 + data_len); + e->len = 10 + data_len; + + return etcp_route_send(cli, server_node_id, e, 0); } // ======================== Server handler ======================== static void srv_handler(struct ETCP_CONN* conn, struct ll_entry* entry) { (void)conn; - if (!entry || !entry->dgram || entry->len < 10) { // svc_id(1) + subcmd(1) + seq(4) + data_len(4) + if (!entry || !entry->dgram || entry->len < 10) { + printf("[FAIL] srv_handler: bad entry len=%zu\n", entry ? entry->len : 0); if (entry) { queue_dgram_free(entry); queue_entry_free(entry); } + g_fail = 4; g_done = -1; return; } uint8_t subcmd = entry->dgram[1]; @@ -113,66 +254,64 @@ static void srv_handler(struct ETCP_CONN* conn, struct ll_entry* entry) { uint32_t data_len = 0; memcpy(&data_len, entry->dgram + 6, 4); uint8_t* payload = entry->dgram + 10; - if (subcmd == 0x01) { // DATA - if (seq != (uint32_t)fwd_rcvd) { - printf("[FAIL] server: seq mismatch expected=%u got=%u\n", (uint32_t)fwd_rcvd, seq); - g_test_done = -1; - } else if (data_len > 0 && memcmp(payload, expected_data, data_len) != 0) { - printf("[FAIL] server: data mismatch at seq=%u\n", seq); - g_test_done = -1; - } else { - fwd_rcvd++; - } + if (subcmd != 0x01) { + printf("[FAIL] srv_handler: unexpected subcmd=%u seq=%u\n", subcmd, seq); queue_dgram_free(entry); queue_entry_free(entry); + g_fail = 4; g_done = -1; + return; + } - // Send reply back - if (!g_test_done && reply_sent < TOTAL_PACKETS) { - uint8_t buf[10 + MAX_PAYLOAD]; - buf[0] = TEST_SVC_ID; - buf[1] = 0x02; // REPLY subcmd - uint32_t rseq = reply_sent; - memcpy(buf + 2, &rseq, 4); - memcpy(buf + 6, &data_len, 4); - memcpy(buf + 10, expected_data, data_len); - struct ll_entry* re = queue_entry_new(0); - if (re) { re->dgram = u_malloc(10 + data_len); memcpy(re->dgram, buf, 10 + data_len); re->len = 10 + data_len; - etcp_route_send(srv, client_node_id, re, 0); reply_sent++; } - } - } else { + // Check seq in range + if (seq >= BITMAP_SIZE) { + 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; + return; } -} -// ======================== Client handler ======================== -static void cli_handler(struct ETCP_CONN* conn, struct ll_entry* entry) { - (void)conn; - if (!entry || !entry->dgram || entry->len < 10) { - if (entry) { queue_dgram_free(entry); queue_entry_free(entry); } + // Check in-order delivery (router should guarantee this) + if (seq != g_expected_seq) { + printf("[FAIL] srv_handler: seq mismatch expected=%u got=%u (loss or reorder)\n", + g_expected_seq, seq); + queue_dgram_free(entry); queue_entry_free(entry); + g_fail = 4; g_fail_seq = seq; g_fail_detail = g_expected_seq; + g_done = -1; return; } - uint8_t subcmd = entry->dgram[1]; - uint32_t seq = 0; memcpy(&seq, entry->dgram + 2, 4); - uint32_t data_len = 0; memcpy(&data_len, entry->dgram + 6, 4); - uint8_t* payload = entry->dgram + 10; - if (subcmd == 0x02) { // REPLY - if (seq != (uint32_t)reply_rcvd) { - printf("[FAIL] client: reply seq mismatch expected=%u got=%u\n", (uint32_t)reply_rcvd, seq); - g_test_done = -1; - } else if (data_len > 0 && memcmp(payload, expected_data, data_len) != 0) { - printf("[FAIL] client: reply data mismatch at seq=%u\n", seq); - g_test_done = -1; - } else { - reply_rcvd++; + // Check payload integrity + if (data_len > 0 && data_len <= MAX_PAYLOAD) { + uint8_t expected[MAX_PAYLOAD]; + gen_payload(seq, expected, data_len); + if (memcmp(payload, expected, data_len) != 0) { + printf("[FAIL] srv_handler: data corrupt at seq=%u\n", seq); + queue_dgram_free(entry); queue_entry_free(entry); + g_fail = 1; g_fail_seq = seq; + g_done = -1; + return; } } - queue_dgram_free(entry); queue_entry_free(entry); + + // Check duplicate + if (rcvd_bitmap[seq] != 0) { + printf("[FAIL] srv_handler: duplicate seq=%u\n", seq); + queue_dgram_free(entry); queue_entry_free(entry); + g_fail = 2; g_fail_seq = seq; + g_done = -1; + return; + } + + rcvd_bitmap[seq] = 1; + g_total_rcvd++; + g_expected_seq++; + + queue_dgram_free(entry); + queue_entry_free(entry); } // ======================== Loopback test (no ETCP) ======================== -static int loop_rcvd = 0; -static int loop_ok = 0; -static void loop_handler(struct ETCP_CONN* conn, struct ll_entry* entry) { +static int loop_rcvd = 0, loop_ok = 0; +static void loop_handler(struct ETCP_CONN* conn, struct ll_entry* entry) { (void)conn; if (entry && entry->dgram && entry->len >= 6) { loop_rcvd++; @@ -187,7 +326,7 @@ static int test_loopback(void) { uint8_t data[6] = { 0xF0, 0xAA, 0xBB, 0xCC, 0x00, 0x00 }; struct ll_entry* e = queue_entry_new(0); e->dgram = u_malloc(6); memcpy(e->dgram, data, 6); e->len = 6; - etcp_route_send(srv, srv->node_id, e, 0); // loopback + etcp_route_send(srv, srv->node_id, e, 0); etcp_router_unbind(srv, 0xF0); if (loop_rcvd != 1 || !loop_ok) { printf("[FAIL] loopback: rcvd=%d ok=%d\n", loop_rcvd, loop_ok); @@ -197,99 +336,267 @@ static int test_loopback(void) { return 0; } -// ======================== Test API bind/unbind ======================== -static int test_api(void) { - if (etcp_router_bind(srv, 0xEE, loop_handler) != 0) { printf("[FAIL] bind\n"); return 1; } - if (etcp_router_bind(srv, 0xEE, loop_handler) != 0) { printf("[FAIL] rebind (overwrite)\n"); etcp_router_unbind(srv, 0xEE); return 1; } - if (etcp_router_unbind(srv, 0xEE) != 0) { printf("[FAIL] unbind\n"); return 1; } - if (etcp_router_unbind(srv, 0xEE) == 0) { printf("[FAIL] double unbind should fail\n"); return 1; } - // Invalid args: - if (etcp_router_bind(NULL, 0, NULL) == 0) { printf("[FAIL] bind null\n"); return 1; } - if (etcp_router_bind(srv, 0xEE, NULL) == 0) { printf("[FAIL] bind null cb\n"); return 1; } - printf(" api bind/unbind: OK\n"); +// ======================== Dummynet setup ======================== +static int setup_dummynet(void) { + dn = dummynet_create(ua, "127.0.0.1", (uint16_t)g_dn_port); + if (!dn) { printf("[FAIL] dummynet_create on port %d\n", g_dn_port); return -1; } + // Forward: dummynet → server + dummynet_set_direction(dn, DUMMYNET_FORWARD, 0, 0, 0, DUMMY_QUEUE_SIZE, 0, + "127.0.0.1", (uint16_t)g_srv_port); + // Backward: dummynet → client + dummynet_set_direction(dn, DUMMYNET_BACKWARD, 0, 0, 0, DUMMY_QUEUE_SIZE, 0, + "127.0.0.1", (uint16_t)g_cli_port); return 0; } -// ======================== Monitor & sender ======================== -static uint64_t dedup_bgp_log = 0; +static void set_bw(uint32_t kbps) { + dummynet_set_direction(dn, DUMMYNET_FORWARD, 0, 0, kbps, DUMMY_QUEUE_SIZE, 0, + "127.0.0.1", (uint16_t)g_srv_port); + dummynet_set_direction(dn, DUMMYNET_BACKWARD, 0, 0, kbps, DUMMY_QUEUE_SIZE, 0, + "127.0.0.1", (uint16_t)g_cli_port); +} + +// ======================== 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; - if (g_test_done) { g_mon_id = NULL; return; } + if (g_done) { g_mon_id = NULL; return; } - static int conn_ok = 0, conn_delay = 0; - if (!conn_ok) { + // Phase: wait for connection + if (g_state == ST_WAIT_CONN) { if (conn_established(srv) && conn_established(cli)) { - conn_delay++; - if (conn_delay < 40) { g_mon_id = uasync_set_timeout(ua, 10, NULL, monitor, "mon"); return; } - conn_ok = 1; - printf(" connections established\n"); - // Run API tests + loopback - if (test_api() != 0) { g_test_done = -1; return; } - if (test_loopback() != 0) { g_test_done = -1; return; } - // Bind handlers + 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 + if (test_loopback() != 0) { g_done = -1; return; } + // Bind data handler etcp_router_bind(srv, TEST_SVC_ID, srv_handler); - etcp_router_bind(cli, TEST_SVC_ID, cli_handler); - // Generate shared random data - for (int i = 0; i < MAX_PAYLOAD; i++) expected_data[i] = (uint8_t)(rand() & 0xFF); - g_phase = 1; - printf(" sending %d packets client→server...\n", TOTAL_PACKETS); - fflush(stdout); + // Start baseline + g_state = ST_BASELINE_SEND; + tic(&t_phase_start); + printf("Phase 2 — Baseline (%d packets, no congestion):\n", BASELINE_PACKETS); + } + g_mon_id = uasync_set_timeout(ua, 10, NULL, monitor, "mon"); + return; + } + + // Phase: baseline send (all at once, no congestion) + if (g_state == ST_BASELINE_SEND) { + while (g_baseline_sent < BASELINE_PACKETS) { + int data_len = 16 + (g_baseline_sent % 32); + if (send_one_pkt(g_total_sent, data_len) == 0) { + g_total_sent++; + g_baseline_sent++; + } else { + printf("[FAIL] baseline: send failed at seq=%u\n", g_total_sent); + g_done = -1; + return; + } } + g_state = ST_BASELINE_WAIT; + g_mon_id = uasync_set_timeout(ua, 10, NULL, monitor, "mon"); + return; } - // Phase 1: send forward - if (g_phase == 1 && fwd_sent < TOTAL_PACKETS) { - int data_len = 16 + (fwd_sent % 32); - uint8_t buf[10 + MAX_PAYLOAD]; - buf[0] = TEST_SVC_ID; - buf[1] = 0x01; // DATA subcmd - uint32_t seq = fwd_sent; - memcpy(buf + 2, &seq, 4); - memcpy(buf + 6, &data_len, 4); - memcpy(buf + 10, expected_data, data_len); - struct ll_entry* e = queue_entry_new(0); - if (e) { e->dgram = u_malloc(10 + data_len); memcpy(e->dgram, buf, 10 + data_len); e->len = 10 + data_len; - if (etcp_route_send(cli, server_node_id, e, 0) == 0) fwd_sent++; - else { queue_dgram_free(e); queue_entry_free(e); } + // Phase: baseline wait + if (g_state == ST_BASELINE_WAIT) { + if (g_total_rcvd >= BASELINE_PACKETS) { + 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); } + g_mon_id = uasync_set_timeout(ua, 10, NULL, monitor, "mon"); + return; + } + + // Phase: cycle fill — send until backpressure + 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 + 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); + } + g_mon_id = uasync_set_timeout(ua, 1, NULL, monitor, "mon"); + return; + } + + // Phase: cycle push — push 20 more packets through congestion + if (g_state == ST_CYCLE_PUSH) { + if (g_cycle_push_ok >= PACKETS_PER_CYCLE) { + t_cycle_push_ms[g_cycle] = toc_ms(&t_phase_start); + printf(" cycle %d push: %7.1f ms pushed=%u drops=%u\n", + g_cycle + 1, t_cycle_push_ms[g_cycle], g_cycle_push_ok, + g_cycle_drop_count[g_cycle]); + // Release shaper + 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; + } + // Try to send one packet + 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++; + g_cycle_push_ok++; + } else { + g_total_drop++; + g_cycle_drop_count[g_cycle]++; + } + g_mon_id = uasync_set_timeout(ua, 500, NULL, monitor, "mon"); // 50ms = ROUTER_SEND_RESUME_TB + return; } - if (g_phase == 1 && fwd_sent >= TOTAL_PACKETS && fwd_rcvd >= TOTAL_PACKETS && reply_sent >= TOTAL_PACKETS) { - g_phase = 2; - printf(" forward phase done: sent=%d rcvd=%d reply_sent=%d\n", fwd_sent, fwd_rcvd, reply_sent); + // Phase: cycle drain — wait for all sent packets to arrive + if (g_state == ST_CYCLE_DRAIN) { + uint32_t cycle_target = g_cycle_start_seq + g_cycle_sent; + if (g_total_rcvd >= cycle_target) { + 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); + g_cycle++; + if (g_cycle >= CONGESTION_CYCLES) { + g_state = ST_FINAL; + } else { + g_state = ST_CYCLE_FILL; + } + } + g_mon_id = uasync_set_timeout(ua, 10, NULL, monitor, "mon"); + return; } - if (g_phase == 2 && reply_rcvd >= TOTAL_PACKETS) { - printf(" reply phase done: rcvd=%d\n", reply_rcvd); - g_test_done = 1; + // Phase: final verification + if (g_state == ST_FINAL) { + printf("Phase 4 — Final verify:\n"); + // Scan bitmap for gaps + int gaps = 0; + uint32_t first_gap = 0, last_gap = 0; + for (uint32_t i = 0; i < g_total_sent; i++) { + if (rcvd_bitmap[i] == 0) { + if (gaps == 0) first_gap = i; + last_gap = i; + gaps++; + } + } + // Dummynet stats + const struct dummynet_stats* sf = dummynet_get_stats(dn, DUMMYNET_FORWARD); + const struct dummynet_stats* sb = dummynet_get_stats(dn, DUMMYNET_BACKWARD); + printf(" total: sent=%u rcvd=%u drops=%u\n", g_total_sent, g_total_rcvd, g_total_drop); + printf(" bitmap: gaps=%d (first=%u last=%u)\n", gaps, first_gap, last_gap); + printf(" dummynet fwd: recv=%" PRIu64 " sent=%" PRIu64 " dropped=%" PRIu64 " lost=%" PRIu64 " qmax=%u\n", + sf->recv, sf->sent, sf->dropped, sf->lost, sf->queue_max); + printf(" dummynet bwd: recv=%" PRIu64 " sent=%" PRIu64 " dropped=%" PRIu64 " lost=%" PRIu64 " qmax=%u\n", + sb->recv, sb->sent, sb->dropped, sb->lost, sb->queue_max); + + // Router stats + struct ETCP_ROUTER_CONN* rconn = etcp_router_conn_get(cli, server_node_id, TEST_SVC_ID); + if (rconn) { + printf(" router: tx_seq=%u rx_seq=%u tx_acked=%u c_pkts_sent=%u c_pkts_rcvd=%u c_dup_dropped=%u c_oob_dropped=%u\n", + rconn->tx_seq, rconn->rx_seq, rconn->tx_acked, + rconn->c_pkts_sent, rconn->c_pkts_rcvd, + rconn->c_dup_dropped, rconn->c_oob_dropped); + } + + // Per-cycle summary + printf(" per-cycle:\n"); + double total_fill = 0, total_push = 0, total_drain = 0; + for (int i = 0; i < CONGESTION_CYCLES; i++) { + printf(" cycle %d: fill=%6.1fms push=%6.1fms drain=%7.1fms fill_pkt=%u push_drops=%u\n", + i + 1, t_cycle_fill_ms[i], t_cycle_push_ms[i], t_cycle_drain_ms[i], + g_cycle_fill_count[i], g_cycle_drop_count[i]); + total_fill += t_cycle_fill_ms[i]; + total_push += t_cycle_push_ms[i]; + total_drain += t_cycle_drain_ms[i]; + } + printf(" totals: fill=%6.1fms push=%6.1fms drain=%7.1fms\n", total_fill, total_push, total_drain); + + // Verdict + if (g_total_sent != g_total_rcvd) { + printf("[FAIL] sent(%u) != rcvd(%u): %u packets lost\n", g_total_sent, g_total_rcvd, g_total_sent - g_total_rcvd); + g_done = -1; + } else if (gaps > 0) { + 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); + g_done = 1; + } + g_mon_id = NULL; return; } g_mon_id = uasync_set_timeout(ua, 10, NULL, monitor, "mon"); } -static void timeout(void* arg) { +// ======================== Timeout ======================== +static void timeout_cb(void* arg) { (void)arg; - if (!g_test_done) { printf("[FAIL] timeout: sent=%d rcvd=%d reply_sent=%d reply_rcvd=%d\n", - fwd_sent, fwd_rcvd, reply_sent, reply_rcvd); - g_test_done = -1; } + 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); + g_done = -1; + } if (g_mon_id) { uasync_cancel_timeout(ua, g_mon_id); g_mon_id = NULL; } } // ======================== Main ======================== int main(void) { if (test_mkdtemp(temp_dir) != 0) { fprintf(stderr, "mkdtemp fail\n"); return 1; } - snprintf(server_conf, sizeof(server_conf), "%s/server.conf", temp_dir); - snprintf(client_conf, sizeof(client_conf), "%s/client.conf", temp_dir); - FILE* f = fopen(server_conf, "w"); - if (!f) { fprintf(stderr, "fopen fail\n"); test_rmdir(temp_dir); return 1; } - fprintf(f, "%s", srv_cfg); fclose(f); - f = fopen(client_conf, "w"); - if (!f) { fprintf(stderr, "fopen fail\n"); test_unlink(server_conf); test_rmdir(temp_dir); return 1; } - fprintf(f, "%s", cli_cfg); fclose(f); + // Allocate consecutive ports + int base_port; + if (alloc_consecutive_ports(&base_port) != 0) { + printf("[FAIL] cannot allocate consecutive ports\n"); + test_rmdir(temp_dir); + return 1; + } + g_dn_port = base_port; + g_srv_port = base_port + 1; + g_cli_port = base_port - 1; + printf("Ports: dummynet=%d server=%d client=%d\n", g_dn_port, g_srv_port, g_cli_port); + + write_configs(); - printf("=== test_etcp_router ===\n"); + printf("=== test_etcp_router (with dummynet congestion) ===\n"); debug_config_init(); debug_set_level(DEBUG_LEVEL_ERROR); @@ -297,6 +604,7 @@ int main(void) { utun_instance_set_tun_init_enabled(0); srand((unsigned)time(NULL)); + tic(&t_phase_start); ua = uasync_create(); if (!ua) { printf("[FAIL] uasync_create\n"); goto done; } @@ -306,35 +614,24 @@ int main(void) { cli = utun_instance_create(ua, client_conf); if (!cli || utun_instance_init(cli) < 0) { printf("[FAIL] client create\n"); goto done; } + // Create dummynet between client and server + if (setup_dummynet() != 0) { printf("[FAIL] dummynet setup\n"); goto done; } + g_mon_id = uasync_set_timeout(ua, 100, NULL, monitor, "mon"); - void* to_id = uasync_set_timeout(ua, TEST_TIMEOUT_MS * 10, NULL, timeout, "to"); + void* to_id = uasync_set_timeout(ua, TEST_TIMEOUT_MS * 10, NULL, timeout_cb, "to"); - while (!g_test_done) uasync_poll(ua, 100); + while (!g_done) uasync_poll(ua, 10); if (to_id) uasync_cancel_timeout(ua, to_id); if (g_mon_id) uasync_cancel_timeout(ua, g_mon_id); - if (g_test_done == 1) { - // Check all statistics - if (fwd_sent == TOTAL_PACKETS && fwd_rcvd == TOTAL_PACKETS && - reply_sent == TOTAL_PACKETS && reply_rcvd == TOTAL_PACKETS) { - printf("[PASS] test_etcp_router — %d packets forwarded, %d replied\n", fwd_rcvd, reply_rcvd); - g_ok = 1; - } else { - printf("[FAIL] incomplete: fwd_sent=%d fwd_rcvd=%d reply_sent=%d reply_rcvd=%d\n", - fwd_sent, fwd_rcvd, reply_sent, reply_rcvd); - } - } else { - printf("[FAIL] test did not complete\n"); - } - etcp_router_unbind(srv, TEST_SVC_ID); - etcp_router_unbind(cli, TEST_SVC_ID); done: + if (dn) dummynet_destroy(dn); if (srv) { srv->running = 0; utun_instance_destroy(srv); } if (cli) { cli->running = 0; utun_instance_destroy(cli); } if (ua) uasync_destroy(ua, 0); test_unlink(server_conf); test_unlink(client_conf); test_rmdir(temp_dir); - return g_ok ? 0 : 1; + return g_done == 1 ? 0 : 1; }