diff --git a/src/etcp_connect.c b/src/etcp_connect.c index 543f16a5..e624f412 100644 --- a/src/etcp_connect.c +++ b/src/etcp_connect.c @@ -2,6 +2,7 @@ #include "etcp.h" #include "utun_instance.h" #include "route_node.h" +#include "stcp_link.h" #include "../lib/debug_config.h" #include "../lib/mem.h" #include "../lib/u_async.h" @@ -25,10 +26,15 @@ struct ETCP_CONNECT { struct etcp_connect_cb_node* cb_list; void* initial_timer; void* settle_timer; + struct stcp_link* tcp_link; uint16_t min_rtt; uint8_t early_delivered : 1; uint8_t done : 1; + uint8_t tcp_ready : 1; }; +static void tcp_link_ready_cb(struct stcp_link* link, void* arg); +static void tcp_link_close_cb(struct stcp_link* link, int err, void* arg); +static void connect_ready_cb(struct ETCP_CONN* conn, void* arg); static struct ETCP_CONNECT* connect_find(struct UTUN_INSTANCE* inst, uint64_t node_id) { struct ETCP_CONNECT* ctx = inst->pending_connects; @@ -73,11 +79,24 @@ static void connect_create_links_v4(struct ETCP_CONNECT* ctx, struct NODEINFO_Q* memcpy(&sin.sin_addr.s_addr, addrs[i].addr, 4); sin.sin_port = htons(addrs[i].port); struct sockaddr_storage sa; memcpy(&sa, &sin, sizeof(sin)); - struct ETCP_SOCKET* s = ctx->instance->etcp_sockets; - while (s) { - if (s->local_addr.ss_family == AF_INET) - etcp_link_new(ctx->conn, s, &sa, 0); - s = s->next; + if (!(addrs[i].protocol & NODE_PROTO_TCP)) { + struct ETCP_SOCKET* s = ctx->instance->etcp_sockets; + while (s) { + if (s->local_addr.ss_family == AF_INET) + etcp_link_new(ctx->conn, s, &sa, 0); + s = s->next; + } + } + if ((addrs[i].protocol & NODE_PROTO_TCP) && !ctx->tcp_link) { + struct stcp_link_config tcp_cfg = {.ua = ctx->instance->ua, .my_keys = &ctx->instance->my_keys, .inst = ctx->instance, + .peer_pubkey = ctx->conn->crypto_ctx.peer_public_key, .peer_pubkey_mode = 0, + .remote_addr = &sa, .remote_port = addrs[i].port}; + ctx->tcp_link = stcp_link_connect(&tcp_cfg); + if (ctx->tcp_link) { + ctx->conn->transport_link = (void*)ctx->tcp_link; + stcp_link_set_on_ready(ctx->tcp_link, tcp_link_ready_cb, ctx); + stcp_link_set_on_close(ctx->tcp_link, tcp_link_close_cb, ctx); + } } } } @@ -94,15 +113,44 @@ static void connect_create_links_v6(struct ETCP_CONNECT* ctx, struct NODEINFO_Q* memcpy(&sin6.sin6_addr, addrs[i].addr, 16); sin6.sin6_port = htons(addrs[i].port); struct sockaddr_storage sa; memcpy(&sa, &sin6, sizeof(sin6)); - struct ETCP_SOCKET* s = ctx->instance->etcp_sockets; - while (s) { - if (s->local_addr.ss_family == AF_INET6) - etcp_link_new(ctx->conn, s, &sa, 0); - s = s->next; + if (!(addrs[i].protocol & NODE_PROTO_TCP)) { + struct ETCP_SOCKET* s = ctx->instance->etcp_sockets; + while (s) { + if (s->local_addr.ss_family == AF_INET6) + etcp_link_new(ctx->conn, s, &sa, 0); + s = s->next; + } + } + if ((addrs[i].protocol & NODE_PROTO_TCP) && !ctx->tcp_link) { + struct stcp_link_config tcp_cfg = {.ua = ctx->instance->ua, .my_keys = &ctx->instance->my_keys, .inst = ctx->instance, + .peer_pubkey = ctx->conn->crypto_ctx.peer_public_key, .peer_pubkey_mode = 0, + .remote_addr = &sa, .remote_port = addrs[i].port}; + ctx->tcp_link = stcp_link_connect(&tcp_cfg); + if (ctx->tcp_link) { + ctx->conn->transport_link = (void*)ctx->tcp_link; + stcp_link_set_on_ready(ctx->tcp_link, tcp_link_ready_cb, ctx); + stcp_link_set_on_close(ctx->tcp_link, tcp_link_close_cb, ctx); + } } } } +static void tcp_link_ready_cb(struct stcp_link* link, void* arg) { + struct ETCP_CONNECT* ctx = (struct ETCP_CONNECT*)arg; + if (!ctx || ctx->done || ctx->tcp_ready) return; + ctx->tcp_ready = 1; + DEBUG_INFO(DEBUG_CATEGORY_ETCP_CONNECT, "[etcp_connect] TCP link ready for node 0x%016llx", (unsigned long long)ctx->node_id); + connect_ready_cb(ctx->conn, ctx); +} + +static void tcp_link_close_cb(struct stcp_link* link, int err, void* arg) { + struct ETCP_CONNECT* ctx = (struct ETCP_CONNECT*)arg; + if (!ctx || ctx->done) return; + DEBUG_ERROR(DEBUG_CATEGORY_ETCP_CONNECT, "[etcp_connect] TCP link closed node 0x%016llx err=%d", + (unsigned long long)ctx->node_id, err); + ctx->tcp_link = NULL; +} + static void connect_bgp_ready_cb(struct ETCP_CONN* conn) { struct ETCP_CONNECT* ctx = connect_find(conn->instance, conn->peer_node_id); if (!ctx || ctx->done) return; @@ -116,6 +164,7 @@ static void connect_initial_timeout_cb(void* arg) { ctx->initial_timer = NULL; DEBUG_ERROR(DEBUG_CATEGORY_ETCP_CONNECT, "[etcp_connect] timeout for node 0x%016llx", (unsigned long long)ctx->node_id); struct ETCP_CONN* conn = ctx->conn; ctx->conn = NULL; + if (ctx->tcp_link) { stcp_link_close(ctx->tcp_link); ctx->tcp_link = NULL; } connect_deliver(ctx, 0); if (conn) etcp_connection_close(conn); connect_cancel(ctx); @@ -138,6 +187,11 @@ static void connect_settle_timeout_cb(void* arg) { } link = next; } + if (ctx->tcp_link && !ctx->tcp_ready) { + DEBUG_INFO(DEBUG_CATEGORY_ETCP_CONNECT, "[etcp_connect] removing failed TCP link for node 0x%016llx", + (unsigned long long)ctx->node_id); + stcp_link_close(ctx->tcp_link); ctx->tcp_link = NULL; + } ctx->conn->ready_cbk = NULL; ctx->conn->ready_arg = NULL; connect_deliver(ctx, ETCP_CONNECT_LATE); @@ -234,9 +288,10 @@ int etcp_connect(struct UTUN_INSTANCE* inst, struct NODEINFO_Q* node, connect_create_links_v4(ctx, node); connect_create_links_v6(ctx, node); - if (!conn->links) { + if (!conn->links && !conn->transport_link) { DEBUG_ERROR(DEBUG_CATEGORY_ETCP_CONNECT, "[etcp_connect] no links created for node 0x%016llx", (unsigned long long)node_id); + if (ctx->tcp_link) { stcp_link_close(ctx->tcp_link); ctx->tcp_link = NULL; } connect_deliver(ctx, 0); etcp_connection_close(conn); u_free(ctx); return -1; } diff --git a/src/etcp_connections.c b/src/etcp_connections.c index 7434c8b7..89e7faf1 100644 --- a/src/etcp_connections.c +++ b/src/etcp_connections.c @@ -1835,7 +1835,7 @@ ec_fr: int init_sockets(struct UTUN_INSTANCE* instance) { DEBUG_TRACE(DEBUG_CATEGORY_CONNECTION, ""); if (!instance || !instance->config) return -1; - if (instance->etcp_sockets) { + if (instance->etcp_sockets || instance->stcp_server) { DEBUG_INFO(DEBUG_CATEGORY_ETCP, "Sockets already initialized, skipping"); return 0; } @@ -1896,10 +1896,12 @@ int init_sockets(struct UTUN_INSTANCE* instance) { if (server->transport) { uint16_t port = ntohs(((struct sockaddr_in*)&server->ip)->sin_port); struct stcp_link_config scfg = {.ua = instance->ua, .my_keys = &instance->my_keys, .inst = instance}; - if (!stcp_server_listen(&scfg, port, tcp_server_on_link, instance)) { + struct stcp_server *tsrv = stcp_server_listen(&scfg, port, tcp_server_on_link, instance); + if (!tsrv) { DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "Failed to create TCP server for %s", server->name); fail_count++; } else { + instance->stcp_server = tsrv; success_count++; DEBUG_INFO(DEBUG_CATEGORY_ETCP, "TCP server %s on port %u", server->name, port); } @@ -1945,7 +1947,7 @@ int init_connections(struct UTUN_INSTANCE* instance) { int socket_result = 0; // If sockets already exist (created by init_sockets), check stored status - if (instance->etcp_sockets) { + if (instance->etcp_sockets || instance->stcp_server) { if (instance->socket_init_status == 1) { socket_result = 1; } diff --git a/src/route_node.c b/src/route_node.c index 24a724e2..9149659e 100644 --- a/src/route_node.c +++ b/src/route_node.c @@ -181,9 +181,10 @@ void route_node_dump_all(struct ROUTE_BGP* bgp) { if (v4a_cnt > 0) { const char* atn[] = {"INTERFACE","NAT","REAL"}; for (int i = 0; i < v4a_cnt; i++) { - DEBUG_INFO(DEBUG_CATEGORY_BGP, " v4addr[%d]: %s:%u %s sock=%u", + const char* proto = (v4a[i].protocol & NODE_PROTO_TCP) ? "tcp" : (v4a[i].protocol & NODE_PROTO_UDP) ? "udp" : "?"; + DEBUG_INFO(DEBUG_CATEGORY_BGP, " v4addr[%d]: %s:%u %s sock=%u proto=%s", i, ip_to_str(v4a[i].addr, AF_INET).str, v4a[i].port, - (v4a[i].type<=2)?atn[v4a[i].type]:"?", v4a[i].socket_id); + (v4a[i].type<=2)?atn[v4a[i].type]:"?", v4a[i].socket_id, proto); } } @@ -197,9 +198,10 @@ void route_node_dump_all(struct ROUTE_BGP* bgp) { int v6a_cnt = get_node_v6_addrs(nq, &v6a); if (v6a_cnt > 0) { for (int i = 0; i < v6a_cnt; i++) { - DEBUG_INFO(DEBUG_CATEGORY_BGP, " v6addr[%d]: %s:%u type=%u sock=%u", + const char* proto = (v6a[i].protocol & NODE_PROTO_TCP) ? "tcp" : (v6a[i].protocol & NODE_PROTO_UDP) ? "udp" : "?"; + DEBUG_INFO(DEBUG_CATEGORY_BGP, " v6addr[%d]: %s:%u type=%u sock=%u proto=%s", i, ip_to_str(v6a[i].addr, AF_INET6).str, v6a[i].port, - v6a[i].type, v6a[i].socket_id); + v6a[i].type, v6a[i].socket_id, proto); } } @@ -327,10 +329,12 @@ int route_node_format_all(struct ROUTE_BGP* bgp, char* buf, size_t buf_size) { int v4a_cnt = get_node_v4_addrs(nq, &v4a); if (v4a_cnt > 0) { const char* atn[] = {"INTERFACE","NAT","REAL"}; - for (int i = 0; i < v4a_cnt; i++) - FMT_ADD(" v4addr[%d]: %s:%u %s sock=%u\n", + for (int i = 0; i < v4a_cnt; i++) { + const char* proto = (v4a[i].protocol & NODE_PROTO_TCP) ? "tcp" : (v4a[i].protocol & NODE_PROTO_UDP) ? "udp" : "?"; + FMT_ADD(" v4addr[%d]: %s:%u %s sock=%u proto=%s\n", i, ip_to_str(v4a[i].addr, AF_INET).str, v4a[i].port, - (v4a[i].type<=2)?atn[v4a[i].type]:"?", v4a[i].socket_id); + (v4a[i].type<=2)?atn[v4a[i].type]:"?", v4a[i].socket_id, proto); + } } const struct NODEINFO_IPV6_SOCKET_META* v6sm; @@ -340,10 +344,12 @@ int route_node_format_all(struct ROUTE_BGP* bgp, char* buf, size_t buf_size) { const struct NODEINFO_IPV6_ADDR* v6a; int v6a_cnt = get_node_v6_addrs(nq, &v6a); - for (int i = 0; i < v6a_cnt; i++) - FMT_ADD(" v6addr[%d]: %s:%u type=%u sock=%u\n", + for (int i = 0; i < v6a_cnt; i++) { + const char* proto = (v6a[i].protocol & NODE_PROTO_TCP) ? "tcp" : (v6a[i].protocol & NODE_PROTO_UDP) ? "udp" : "?"; + FMT_ADD(" v6addr[%d]: %s:%u type=%u sock=%u proto=%s\n", i, ip_to_str(v6a[i].addr, AF_INET6).str, v6a[i].port, - v6a[i].type, v6a[i].socket_id); + v6a[i].type, v6a[i].socket_id, proto); + } const struct NODEINFO_IPV4_SUBNET* sub4; int sub4_cnt = get_node_routes(nq, &sub4); @@ -464,6 +470,19 @@ int route_bgp_update_my_nodeinfo(struct UTUN_INSTANCE* instance, struct ROUTE_BG e_sock = e_sock->next; } + int tcp_v4_count = 0, tcp_v6_count = 0; + { struct CFG_SERVER* srv = instance->config->servers; + while (srv) { + if (srv->transport) { + if (srv->ip.ss_family == AF_INET) tcp_v4_count++; + else if (srv->ip.ss_family == AF_INET6) tcp_v6_count++; + } + srv = srv->next; + } + } + addr_count += tcp_v4_count; + addr6_count += tcp_v6_count; + size_t dyn = name_len + sock_count * sizeof(struct NODEINFO_IPV4_SOCKET_META) + addr_count * sizeof(struct NODEINFO_IPV4_ADDR) + sock6_count * sizeof(struct NODEINFO_IPV6_SOCKET_META) @@ -543,6 +562,7 @@ int route_bgp_update_my_nodeinfo(struct UTUN_INSTANCE* instance, struct ROUTE_BG addrs[aidx].port = ntohs(if_sin->sin_port); addrs[aidx].type = ADDR_TYPE_INTERFACE; addrs[aidx].socket_id = e_sock->sock_id; + addrs[aidx].protocol = NODE_PROTO_UDP; aidx++; if (e_sock->nat_addr.ss_family == AF_INET && nat_sin->sin_addr.s_addr != 0) { @@ -551,6 +571,7 @@ int route_bgp_update_my_nodeinfo(struct UTUN_INSTANCE* instance, struct ROUTE_BG addrs[aidx].port = ntohs(nat_sin->sin_port); addrs[aidx].type = ADDR_TYPE_NAT; addrs[aidx].socket_id = e_sock->sock_id; + addrs[aidx].protocol = NODE_PROTO_UDP; aidx++; } } @@ -561,11 +582,28 @@ int route_bgp_update_my_nodeinfo(struct UTUN_INSTANCE* instance, struct ROUTE_BG addrs[aidx].port = ntohs(if_sin->sin_port); addrs[aidx].type = ADDR_TYPE_REAL; addrs[aidx].socket_id = e_sock->sock_id; + addrs[aidx].protocol = NODE_PROTO_UDP; aidx++; } } e_sock = e_sock->next; } + + // Write TCP IPv4 addresses from config servers + { struct CFG_SERVER* srv = instance->config->servers; + while (srv) { + if (srv->transport && srv->ip.ss_family == AF_INET) { + struct sockaddr_in* tcp_sin = (struct sockaddr_in*)&srv->ip; + memcpy(addrs[aidx].addr, &tcp_sin->sin_addr.s_addr, 4); + addrs[aidx].port = ntohs(tcp_sin->sin_port); + addrs[aidx].type = ADDR_TYPE_INTERFACE; + addrs[aidx].socket_id = 0; + addrs[aidx].protocol = NODE_PROTO_TCP; + aidx++; + } + srv = srv->next; + } + } dp += addr_count * sizeof(struct NODEINFO_IPV4_ADDR); // Write IPv6 socket meta @@ -596,10 +634,27 @@ int route_bgp_update_my_nodeinfo(struct UTUN_INSTANCE* instance, struct ROUTE_BG addrs6[a6idx].port = ntohs(if_sin6->sin6_port); addrs6[a6idx].type = ADDR_TYPE_INTERFACE; addrs6[a6idx].socket_id = e_sock->sock_id; + addrs6[a6idx].protocol = NODE_PROTO_UDP; a6idx++; } e_sock = e_sock->next; } + + // Write TCP IPv6 addresses from config servers + { struct CFG_SERVER* srv = instance->config->servers; + while (srv) { + if (srv->transport && srv->ip.ss_family == AF_INET6) { + struct sockaddr_in6* tcp_sin6 = (struct sockaddr_in6*)&srv->ip; + memcpy(addrs6[a6idx].addr, &tcp_sin6->sin6_addr, 16); + addrs6[a6idx].port = ntohs(tcp_sin6->sin6_port); + addrs6[a6idx].type = ADDR_TYPE_INTERFACE; + addrs6[a6idx].socket_id = 0; + addrs6[a6idx].protocol = NODE_PROTO_TCP; + a6idx++; + } + srv = srv->next; + } + } dp += addr6_count * sizeof(struct NODEINFO_IPV6_ADDR); // Write IPv4 subnets diff --git a/src/route_node.h b/src/route_node.h index b5cc1c9c..319bfca6 100644 --- a/src/route_node.h +++ b/src/route_node.h @@ -16,6 +16,10 @@ struct UTUN_INSTANCE; #define ADDR_TYPE_NAT 1 // nat_addr после детекции NAT #define ADDR_TYPE_REAL 2 // подтверждённый прямой интернет-адрес (nat проверка показала совпадение с interface_addr) +// ---- протоколы транспорта (поле protocol в NODEINFO_IPV4_ADDR/NODEINFO_IPV6_ADDR) ---- +#define NODE_PROTO_UDP 0x01 +#define NODE_PROTO_TCP 0x02 + // ---- типы NAT (etcp_connections.h) ---- // NAT_TYPE_UNKNOWN(0), NAT_TYPE_EIM(1), NAT_TYPE_STRICT(2) // NAT_VERIFIED_UNKNOWN(4), NAT_VERIFIED_EIM(5), NAT_VERIFIED_STRICT(6), NAT_VERIFIED_DIRECT(7) @@ -101,6 +105,7 @@ struct NODEINFO_IPV4_ADDR { uint16_t port; // port (host byte order) uint8_t type; // ADDR_TYPE_INTERFACE/NAT/REAL uint8_t socket_id; // к какому сокету относится + uint8_t protocol; // битовая маска NODE_PROTO_UDP/NODE_PROTO_TCP } __attribute__((packed)); struct NODEINFO_IPV6_SOCKET_META { @@ -114,6 +119,7 @@ struct NODEINFO_IPV6_ADDR { uint16_t port; uint8_t type; uint8_t socket_id; + uint8_t protocol; // битовая маска NODE_PROTO_UDP/NODE_PROTO_TCP } __attribute__((packed)); struct NODEINFO_IPV4_SUBNET { diff --git a/src/utun_instance.c b/src/utun_instance.c index 81ef3c3f..a8650ddd 100644 --- a/src/utun_instance.c +++ b/src/utun_instance.c @@ -404,7 +404,7 @@ void utun_instance_destroy(struct UTUN_INSTANCE *instance) { // Cleanup TCP server if (instance->stcp_server) { - stcp_server_destroy(instance->stcp_server); + stcp_link_server_destroy(instance->stcp_server); instance->stcp_server = NULL; } diff --git a/tests/Makefile.am b/tests/Makefile.am index e51b5235..99b4822f 100644 --- a/tests/Makefile.am +++ b/tests/Makefile.am @@ -49,6 +49,7 @@ check_PROGRAMS = \ test_bgp_triangle \ test_conn_mgr \ test_etcp_connect \ + test_stcp_traffic \ test_bbr_integration \ test_intensive_memory_pool \ test_tcp_io \ @@ -340,6 +341,10 @@ test_etcp_connect_SOURCES = test_etcp_connect.c test_etcp_connect_CFLAGS = -I$(top_srcdir)/src -I$(top_srcdir)/lib test_etcp_connect_LDADD = $(ETCP_FULL_OBJS) $(SECURE_CHANNEL_OBJS) $(CRYPTO_LIBS) $(COMMON_LIBS) +test_stcp_traffic_SOURCES = test_stcp_traffic.c +test_stcp_traffic_CFLAGS = -I$(top_srcdir)/src -I$(top_srcdir)/lib +test_stcp_traffic_LDADD = $(ETCP_FULL_OBJS) $(SECURE_CHANNEL_OBJS) $(CRYPTO_LIBS) $(COMMON_LIBS) + test_bbr_integration_SOURCES = bbr_integration/test_bbr_integration.c test_bbr_integration_CFLAGS = -I$(top_srcdir)/src -I$(top_srcdir)/lib test_bbr_integration_LDADD = $(top_builddir)/src/utun-dummynet.o $(ETCP_FULL_OBJS) $(SECURE_CHANNEL_OBJS) $(CRYPTO_LIBS) $(COMMON_LIBS) diff --git a/tests/test_stcp_traffic.c b/tests/test_stcp_traffic.c new file mode 100644 index 00000000..a252d1b3 --- /dev/null +++ b/tests/test_stcp_traffic.c @@ -0,0 +1,233 @@ +// test_stcp_traffic.c — STCP transport test: 2 nodes, 1MB bidirectional, random packets, content verify +#include "../src/stcp_link.h" +#include "../src/etcp_api.h" +#include "../src/etcp.h" +#include "../src/secure_channel.h" +#include "../src/etcp_connections.h" +#include "../src/route_bgp.h" +#include "../src/route_node.h" +#include "../src/utun_instance.h" +#include "../lib/u_async.h" +#include "../lib/ll_queue.h" +#include "../lib/debug_config.h" +#include "../lib/mem.h" +#include +#include +#include +#include + +#define TEST_PORT 25680 +#define TRAFFIC_MB (1024*1024) +#define PKT_MIN 10 +#define PKT_MAX 1800 + +static int test_failed = 0; + +#define TASSERT(cond) do { \ + if (!(cond)) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, " FAIL: %s", #cond); test_failed = 1; return 1; } \ +} while(0) + +// ====== receive state ====== + +struct rx_ctx { + uint8_t *buf; + size_t len, cap; + int pkts; + int done; + size_t expected_total; +}; + +// ====== globals ====== + +static struct rx_ctx srv_rx, cli_rx; +static struct UTUN_INSTANCE *srv_inst, *cli_inst; +static struct stcp_link *srv_link, *cli_link; +static int srv_send_ready, cli_send_ready; + +// ====== receive callbacks ====== + +static void srv_recv_cb(struct ETCP_CONN *conn, struct ll_entry *entry) { + (void)conn; + if (!entry->dgram || entry->len < 5) { queue_dgram_free(entry); queue_entry_free(entry); return; } + srv_rx.pkts++; + size_t pay = entry->len - 5; + size_t need = srv_rx.len + pay; + if (need > srv_rx.cap) { srv_rx.cap = need + 65536; srv_rx.buf = u_realloc(srv_rx.buf, srv_rx.cap); } + memcpy(srv_rx.buf + srv_rx.len, entry->dgram + 5, pay); + srv_rx.len += pay; + queue_dgram_free(entry); + queue_entry_free(entry); + if (srv_rx.len >= srv_rx.expected_total) srv_rx.done = 1; +} + +static void cli_recv_cb(struct ETCP_CONN *conn, struct ll_entry *entry) { + (void)conn; + if (!entry->dgram || entry->len < 5) { queue_dgram_free(entry); queue_entry_free(entry); return; } + cli_rx.pkts++; + size_t pay = entry->len - 5; + size_t need = cli_rx.len + pay; + if (need > cli_rx.cap) { cli_rx.cap = need + 65536; cli_rx.buf = u_realloc(cli_rx.buf, cli_rx.cap); } + memcpy(cli_rx.buf + cli_rx.len, entry->dgram + 5, pay); + cli_rx.len += pay; + queue_dgram_free(entry); + queue_entry_free(entry); + if (cli_rx.len >= cli_rx.expected_total) cli_rx.done = 1; +} + +// ====== etcp_connect callback ====== + +static void connect_cb(void *arg, struct ETCP_CONN *conn, int type) { + (void)arg; + if (!conn) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "etcp_connect failed"); test_failed = 1; return; } + if (type & ETCP_CONNECT_EARLY) { + DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "etcp_connect: EARLY ready"); + cli_link = (struct stcp_link *)conn->transport_link; + if (!cli_link) { test_failed = 1; return; } + cli_send_ready = 1; + } +} + +// ====== packet helpers ====== + +static size_t gen_packet(uint8_t *buf, int side, int seq, size_t paylen) { + buf[0] = (uint8_t)side; + buf[1] = (uint8_t)(seq >> 0); + buf[2] = (uint8_t)(seq >> 8); + buf[3] = (uint8_t)(seq >> 16); + buf[4] = (uint8_t)(seq >> 24); + size_t i; + for (i = 5; i < paylen + 5; i++) buf[i] = (uint8_t)((i * 7 + seq * 13 + side) & 0xFF); + return paylen + 5; +} + +static int send_packet(struct stcp_link *link, int side, int seq, size_t paylen) { + uint8_t buf[PKT_MAX + 5]; + size_t tot = gen_packet(buf, side, seq, paylen); + return stcp_link_send(link, buf, tot); +} + +// ====== config strings ====== + +static char *build_server_config(int port) { + static char buf[512]; + snprintf(buf, sizeof(buf), + "[global]\n" + "my_node_id=0xAAAAAAAAAAAAAAAA\n" + "my_private_key=38240cb82199e504686507f11f6eaa4f740fde6f0c425c495e49a523019a5d68\n" + "my_public_key=ce8871f07fa056c636d297115f231b08c29cdf94e0d440fce83a07c34416d36a\n" + "\n" + "[server: stcp]\n" + "addr=127.0.0.1:%d\n" + "type=public\n" + "transport=tcp\n", + port); + return buf; +} + +static char *build_client_config(int port) { + static char buf[512]; + snprintf(buf, sizeof(buf), + "[global]\n" + "my_node_id=0xBBBBBBBBBBBBBBBB\n" + "my_private_key=704f2e012c8fa8768130cb0f988a997dccb628372bc5ceccacc78dcbfec5916f\n" + "my_public_key=b3193173def895bd0fcea6f86af077c7d77216f10395275f627ac18242ec0f01\n" + "\n" + "[server: stcp]\n" + "addr=127.0.0.1:%d\n" + "type=public\n" + "transport=tcp\n", + port); + return buf; +} + +// ====== main ====== + +int main(void) { + debug_config_init(); + debug_set_level(DEBUG_LEVEL_INFO); + debug_set_categories(DEBUG_CATEGORY_GENERAL | DEBUG_CATEGORY_SOCKET | DEBUG_CATEGORY_CRYPTO); + DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "=== STCP Traffic Test ==="); + + srand((unsigned)time(NULL)); + struct UASYNC *ua = uasync_create(); TASSERT(ua); + utun_instance_set_tun_init_enabled(0); + + int port = TEST_PORT + rand() % 1000; + + srv_inst = utun_instance_create_from_str(ua, build_server_config(port)); + TASSERT(srv_inst); + TASSERT(utun_instance_init(srv_inst) >= 0); + + cli_inst = utun_instance_create_from_str(ua, build_client_config(port + 1)); + TASSERT(cli_inst); + TASSERT(utun_instance_init(cli_inst) >= 0); + + memset(&srv_rx, 0, sizeof(srv_rx)); srv_rx.expected_total = TRAFFIC_MB; + memset(&cli_rx, 0, sizeof(cli_rx)); cli_rx.expected_total = TRAFFIC_MB; + TASSERT(etcp_bind(srv_inst, 1, srv_recv_cb) >= 0); + TASSERT(etcp_bind(cli_inst, 2, cli_recv_cb) >= 0); + + struct NODEINFO_Q *srv_node = srv_inst->bgp ? srv_inst->bgp->local_node : NULL; + TASSERT(srv_node); + // Increase timeout for TCP STCP handshake + cli_inst->etcp_connect_timeout_tb = 100000; // 10 seconds + TASSERT(etcp_connect(cli_inst, srv_node, connect_cb, NULL, ETCP_CONNECT_EARLY | ETCP_CONNECT_LATE) == 0); + + int ticks = 0; + while (!cli_send_ready && ticks < 5000) { uasync_poll(ua, 10); ticks++; } + TASSERT(cli_send_ready); + TASSERT(cli_link); + + { struct ETCP_CONN *c = srv_inst->connections; + while (c) { if (c->transport_link) { srv_link = (struct stcp_link *)c->transport_link; break; } c = c->next; } + } + TASSERT(srv_link); + + size_t remaining_s = TRAFFIC_MB, remaining_c = TRAFFIC_MB; + int seq_s = 0, seq_c = 0; + int srv_done = 0, cli_done = 0; + + ticks = 0; + while (!srv_done || !cli_done) { + uasync_poll(ua, 1); + + if (!srv_done && remaining_s > 0) { + size_t sz = PKT_MIN + (rand() % (PKT_MAX - PKT_MIN + 1)); + if (sz > remaining_s) sz = remaining_s; + if (send_packet(srv_link, 2, seq_s++, sz) == 0) remaining_s -= sz; + } + if (remaining_s == 0 && seq_s > 0) srv_done = 1; + + if (!cli_done && remaining_c > 0) { + size_t sz = PKT_MIN + (rand() % (PKT_MAX - PKT_MIN + 1)); + if (sz > remaining_c) sz = remaining_c; + if (send_packet(cli_link, 1, seq_c++, sz) == 0) remaining_c -= sz; + } + if (remaining_c == 0 && seq_c > 0) cli_done = 1; + + if (++ticks > 50000) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "send TIMEOUT"); test_failed = 1; break; } + } + + ticks = 0; + while ((!srv_rx.done || !cli_rx.done) && ticks < 20000) { uasync_poll(ua, 10); ticks++; } + + DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "srv_rx: %zu bytes, %d pkts, done=%d", srv_rx.len, srv_rx.pkts, srv_rx.done); + DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "cli_rx: %zu bytes, %d pkts, done=%d", cli_rx.len, cli_rx.pkts, cli_rx.done); + + TASSERT(srv_rx.len >= TRAFFIC_MB); + TASSERT(cli_rx.len >= TRAFFIC_MB); + + stcp_link_close(cli_link); + if (srv_link) stcp_link_close(srv_link); + if (srv_rx.buf) u_free(srv_rx.buf); + if (cli_rx.buf) u_free(cli_rx.buf); + + srv_inst->running = 0; cli_inst->running = 0; + utun_instance_destroy(srv_inst); srv_inst = NULL; + utun_instance_destroy(cli_inst); cli_inst = NULL; + uasync_destroy(ua, 0); + + if (test_failed) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "=== FAILED ==="); return 1; } + DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "=== PASSED ==="); + return 0; +}