Browse Source

stcp: add protocol field to NODEINFO addr structs, TCP support in etcp_connect, stcp traffic test

- route_node.h: add NODE_PROTO_UDP/NODE_PROTO_TCP constants, protocol field in NODEINFO_IPV4_ADDR/NODEINFO_IPV6_ADDR
- route_node.c: fill protocol for UDP sockets + TCP servers from config, update dumps
- etcp_connect.c: create stcp_link for TCP-capable addresses, ready/close callbacks, cleanup
- etcp_connections.c: store stcp_server_listen result in instance->stcp_server, fix socket init checks for TCP-only configs
- utun_instance.c: use stcp_link_server_destroy (high-level API) instead of stcp_server_destroy
- tests/test_stcp_traffic.c: 2 nodes via etcp_connect+TCP, 1MB bidirectional traffic (10-1800 byte pkts)
chatgui
Evgeny 3 months ago
parent
commit
90ad64a106
  1. 77
      src/etcp_connect.c
  2. 8
      src/etcp_connections.c
  3. 75
      src/route_node.c
  4. 6
      src/route_node.h
  5. 2
      src/utun_instance.c
  6. 5
      tests/Makefile.am
  7. 233
      tests/test_stcp_traffic.c

77
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;
}

8
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;
}

75
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

6
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 {

2
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;
}

5
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)

233
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 <stdio.h>
#include <string.h>
#include <stdlib.h>
#include <time.h>
#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;
}
Loading…
Cancel
Save