Browse Source
tcp_server_on_link sets conn->peer_node_id after etcp_connection_create which adds the entry with peer_node_id=0. instance_find_conn couldnt find TCP connections until the queue entry was re-indexed.topo_upd
5 changed files with 1 additions and 2685 deletions
@ -1,286 +0,0 @@
|
||||
// stcp_link.c — STCP link management implementation
|
||||
#include "stcp_link.h" |
||||
#include "stcp.h" |
||||
#include "stcp_server.h" |
||||
#include "stcp_client.h" |
||||
#include "secure_channel.h" |
||||
#include "etcp.h" |
||||
#include "etcp_connections.h" |
||||
#include "etcp_router.h" |
||||
#include "utun_instance.h" |
||||
#include "../lib/debug_config.h" |
||||
#include "../lib/ll_queue.h" |
||||
#include "../lib/mem.h" |
||||
#include "../lib/memory_pool.h" |
||||
#include <stdlib.h> |
||||
#include <string.h> |
||||
|
||||
struct stcp_server { |
||||
struct stcp_server *next; // linked list in UTUN_INSTANCE
|
||||
struct stcp_server *srv; // stcp_server from stcp_server.h
|
||||
struct stcp_link_config cfg; |
||||
stcp_server_on_link_cb on_link; |
||||
void *on_link_arg; |
||||
}; |
||||
|
||||
struct stcp_link { |
||||
struct stcp_link_config cfg; |
||||
struct stcp_client *cli; |
||||
struct stcp_conn *conn; |
||||
|
||||
struct ETCP_CONN *etcp_conn; // parent ETCP_CONN (for rx dispatch / etcp_send compat)
|
||||
struct ETCP_LINK *etcp_link; // owning ETCP_LINK (for etcp_conn_input pkt->link)
|
||||
|
||||
uint8_t ready; |
||||
uint8_t closing; |
||||
|
||||
stcp_link_cb on_ready_cb; |
||||
void *ready_arg; |
||||
void (*on_close_cb)(struct stcp_link *link, int err, void *arg); |
||||
void *close_arg; |
||||
|
||||
struct ll_queue *tx_queue; // owned by this link
|
||||
}; |
||||
|
||||
// ====== rx dispatch ======
|
||||
|
||||
static void link_rx_cb(struct ll_queue *q, void *arg) { |
||||
struct stcp_link *link = (struct stcp_link *)arg; |
||||
struct ll_entry *e = queue_data_get(q); |
||||
if (!e) { queue_resume_callback(q); return; } |
||||
|
||||
if (link->etcp_conn && link->etcp_link && link->cfg.inst) { |
||||
if (e->len >= 1 && e->dgram[0] == ETCP_KEEPALIVE) { |
||||
struct ETCP_LINK *l = link->etcp_link; |
||||
if (e->len >= 3) { uint16_t pp = e->dgram[1] | ((uint16_t)e->dgram[2] << 8); l->keepalive_timeout = (uint32_t)pp * KA_TIMEOUT_MULT; } |
||||
l->recv_keepalive = 1; l->remote_keepalive = 1; l->link_status = 1; l->keepalive_recv_count++; |
||||
l->last_recv_local_time = get_time_tb(); |
||||
queue_dgram_free(e); queue_entry_free(e); |
||||
} else { |
||||
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "stcp_link rx → etcp_conn_input len=%zu link=%p", e->len, (void*)link->etcp_link); |
||||
struct ETCP_DGRAM *pkt = memory_pool_alloc(link->cfg.inst->pkt_pool); |
||||
if (pkt) { |
||||
pkt->link = link->etcp_link; pkt->data_len = (uint16_t)e->len; pkt->noencrypt_len = 0; |
||||
if (e->len > 0) memcpy(pkt->data, e->dgram, e->len); |
||||
etcp_conn_input(pkt); |
||||
} |
||||
queue_dgram_free(e); queue_entry_free(e); |
||||
} |
||||
} else { |
||||
struct UTUN_INSTANCE *inst = link->cfg.inst; |
||||
uint8_t id = (e->dgram && e->len > 0) ? e->dgram[0] : 0; |
||||
if (inst && inst->api_bindings.callbacks[id]) |
||||
inst->api_bindings.callbacks[id](link->etcp_conn ? link->etcp_conn : NULL, e); |
||||
else if (inst && inst->api_bindings.callbacks[0]) |
||||
inst->api_bindings.callbacks[0](link->etcp_conn ? link->etcp_conn : NULL, e); |
||||
else { queue_dgram_free(e); queue_entry_free(e); } |
||||
} |
||||
queue_resume_callback(q); |
||||
} |
||||
|
||||
// ====== server accept → link ======
|
||||
|
||||
static void server_accept_cb(struct stcp_conn *conn, void *arg) { |
||||
struct stcp_server *ss = (struct stcp_server *)arg; |
||||
struct stcp_link *link = u_calloc(1, sizeof(struct stcp_link)); |
||||
if (!link) { stcp_conn_free(conn); return; } |
||||
link->cfg = ss->cfg; |
||||
link->ready = 1; |
||||
link->conn = conn; |
||||
|
||||
struct ll_queue *rx = queue_new(conn->ua, 0, 0, 0, "srx"); |
||||
queue_set_callback(rx, link_rx_cb, link); |
||||
stcp_conn_set_rx_queue(conn, rx); |
||||
link->tx_queue = queue_new(conn->ua, 0, 0, 0, "stx"); |
||||
queue_set_threshold(link->tx_queue, 0, 0); |
||||
queue_set_waiter_defer(link->tx_queue, 1); |
||||
stcp_conn_set_tx_queue(conn, link->tx_queue); |
||||
|
||||
DEBUG_INFO(DEBUG_CATEGORY_ETCP, "stcp_link: server accepted connection"); |
||||
if (ss->on_link) ss->on_link(link, ss->on_link_arg); |
||||
} |
||||
|
||||
// ====== client connect → link ======
|
||||
|
||||
static void client_ready_cb(struct stcp_conn *conn, void *arg) { |
||||
struct stcp_link *link = (struct stcp_link *)arg; |
||||
if (!conn) return; |
||||
link->ready = 1; |
||||
link->conn = conn; |
||||
|
||||
struct ll_queue *rx = queue_new(conn->ua, 0, 0, 0, "crx"); |
||||
queue_set_callback(rx, link_rx_cb, link); |
||||
stcp_conn_set_rx_queue(conn, rx); |
||||
link->tx_queue = queue_new(conn->ua, 0, 0, 0, "ctx"); |
||||
queue_set_threshold(link->tx_queue, 0, 0); |
||||
queue_set_waiter_defer(link->tx_queue, 1); |
||||
stcp_conn_set_tx_queue(conn, link->tx_queue); |
||||
|
||||
DEBUG_INFO(DEBUG_CATEGORY_ETCP, "stcp_link: client handshake OK, link=%p etcp_link=%p", |
||||
(void*)link, (void*)link->etcp_link); |
||||
if (link->etcp_link) etcp_link_enter_ready_tcp(link->etcp_link); |
||||
if (link->on_ready_cb) link->on_ready_cb(link, link->ready_arg); |
||||
} |
||||
|
||||
// ====== API ======
|
||||
|
||||
struct stcp_server *stcp_server_listen(struct stcp_link_config *cfg, uint16_t port, |
||||
stcp_server_on_link_cb on_link, void *arg) { |
||||
if (!cfg || !cfg->ua || !cfg->my_keys) return NULL; |
||||
struct stcp_server *ss = u_calloc(1, sizeof(struct stcp_server)); |
||||
if (!ss) return NULL; |
||||
ss->cfg = *cfg; |
||||
ss->on_link = on_link; |
||||
ss->on_link_arg = arg; |
||||
ss->srv = stcp_server_create(cfg->ua, port, cfg->my_keys, server_accept_cb, ss, NULL, NULL, cfg->listen_family); |
||||
if (!ss->srv) { u_free(ss); return NULL; } |
||||
DEBUG_INFO(DEBUG_CATEGORY_ETCP, "port=%u", port); |
||||
return ss; |
||||
} |
||||
|
||||
void stcp_link_server_destroy(struct stcp_server *ss) { |
||||
if (!ss) return; |
||||
if (ss->srv) stcp_server_destroy(ss->srv); |
||||
u_free(ss); |
||||
} |
||||
|
||||
struct stcp_link *stcp_link_connect(struct stcp_link_config *cfg) { |
||||
if (!cfg || !cfg->ua || !cfg->my_keys || !cfg->peer_pubkey || !cfg->remote_addr) |
||||
return NULL; |
||||
|
||||
int family = cfg->remote_addr->ss_family; |
||||
char addr_str[64]; |
||||
uint16_t port; |
||||
if (family == AF_INET) { |
||||
struct sockaddr_in *sa = (struct sockaddr_in *)cfg->remote_addr; |
||||
inet_ntop(AF_INET, &sa->sin_addr, addr_str, sizeof(addr_str)); |
||||
port = cfg->remote_port ? cfg->remote_port : ntohs(sa->sin_port); |
||||
} else if (family == AF_INET6) { |
||||
struct sockaddr_in6 *sa6 = (struct sockaddr_in6 *)cfg->remote_addr; |
||||
getnameinfo((struct sockaddr*)sa6, sizeof(*sa6), addr_str, sizeof(addr_str), NULL, 0, NI_NUMERICHOST); |
||||
port = cfg->remote_port ? cfg->remote_port : ntohs(sa6->sin6_port); |
||||
} else { |
||||
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "stcp_link: unsupported address family %d", family); |
||||
return NULL; |
||||
} |
||||
|
||||
struct stcp_link *link = u_calloc(1, sizeof(struct stcp_link)); |
||||
if (!link) return NULL; |
||||
link->cfg = *cfg; |
||||
|
||||
uint8_t pubkey[SC_PUBKEY_SIZE]; |
||||
if (cfg->peer_pubkey_mode) { |
||||
struct secure_channel sc_tmp; |
||||
sc_init_ctx(&sc_tmp, cfg->my_keys); |
||||
if (sc_set_peer_public_key(&sc_tmp, cfg->peer_pubkey, 1) != SC_OK) { |
||||
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "invalid peer pubkey hex"); |
||||
u_free(link); return NULL; |
||||
} |
||||
memcpy(pubkey, sc_tmp.peer_public_key, SC_PUBKEY_SIZE); |
||||
} else { |
||||
memcpy(pubkey, cfg->peer_pubkey, SC_PUBKEY_SIZE); |
||||
} |
||||
|
||||
link->cli = stcp_client_connect(cfg->ua, addr_str, port, cfg->my_keys, pubkey, |
||||
client_ready_cb, link, NULL, NULL); |
||||
if (!link->cli) { u_free(link); return NULL; } |
||||
|
||||
DEBUG_INFO(DEBUG_CATEGORY_ETCP, "stcp_link: connecting to %s:%u pubkey=%016llx", |
||||
addr_str, port, (unsigned long long)*(const uint64_t*)pubkey); |
||||
return link; |
||||
} |
||||
|
||||
static void stcp_link_close_impl(void *arg) { |
||||
struct stcp_link *link = (struct stcp_link *)arg; |
||||
if (link->on_close_cb) link->on_close_cb(link, 0, link->close_arg); |
||||
if (link->conn) { |
||||
if (link->conn->rx_queue) queue_free(link->conn->rx_queue); |
||||
if (link->tx_queue) queue_free(link->tx_queue); |
||||
stcp_conn_free(link->conn); |
||||
} |
||||
if (link->cli) stcp_client_destroy(link->cli); |
||||
u_free(link); |
||||
} |
||||
|
||||
void stcp_link_close(struct stcp_link *link) { |
||||
if (!link) return; |
||||
if (link->closing) return; |
||||
link->closing = 1; |
||||
uasync_call_soon(link->cfg.ua, link, stcp_link_close_impl); |
||||
} |
||||
|
||||
int stcp_link_send(struct stcp_link *link, const uint8_t *data, size_t len) { |
||||
if (!link || !link->ready) return -1; |
||||
struct ll_entry *e = queue_entry_new(0); |
||||
if (!e) return -1; |
||||
e->dgram = u_malloc(len ? len : 1); |
||||
if (!e->dgram) { queue_entry_free(e); return -1; } |
||||
if (len) memcpy(e->dgram, data, len); |
||||
e->len = (uint16_t)len; |
||||
queue_data_put(link->conn->tx_queue, e); |
||||
return 0; |
||||
} |
||||
|
||||
int stcp_link_is_ready(struct stcp_link *link) { |
||||
return link ? link->ready : 0; |
||||
} |
||||
|
||||
struct ETCP_CONN *stcp_link_get_etcp_conn(struct stcp_link *link) { |
||||
return link ? link->etcp_conn : NULL; |
||||
} |
||||
|
||||
void stcp_link_set_etcp_conn(struct stcp_link *link, struct ETCP_CONN *conn) { |
||||
if (!link) return; |
||||
link->etcp_conn = conn; |
||||
} |
||||
|
||||
void stcp_link_set_etcp_link(struct stcp_link *link, struct ETCP_LINK *elink) { |
||||
if (!link) return; |
||||
link->etcp_link = elink; |
||||
} |
||||
|
||||
const struct sockaddr_storage *stcp_link_get_remote_addr(struct stcp_link *link) { |
||||
return link && link->cfg.remote_addr ? link->cfg.remote_addr : NULL; |
||||
} |
||||
|
||||
void stcp_link_set_on_ready(struct stcp_link *link, stcp_link_cb cb, void *arg) { |
||||
if (!link) return; |
||||
link->on_ready_cb = cb; |
||||
link->ready_arg = arg; |
||||
} |
||||
|
||||
void stcp_link_set_on_close(struct stcp_link *link, void (*cb)(struct stcp_link *link, int err, void *arg), void *arg) { |
||||
if (!link) return; |
||||
link->on_close_cb = cb; |
||||
link->close_arg = arg; |
||||
} |
||||
|
||||
const uint8_t *stcp_link_get_peer_pubkey(struct stcp_link *link) { |
||||
return link && link->conn ? link->conn->peer_pubkey : NULL; |
||||
} |
||||
|
||||
void stcp_server_list_add(struct UTUN_INSTANCE *inst, struct stcp_server *srv) { |
||||
if (!inst || !srv) return; |
||||
srv->next = inst->stcp_servers; |
||||
inst->stcp_servers = srv; |
||||
DEBUG_INFO(DEBUG_CATEGORY_ETCP, "stcp_server_list_add: %p total=%d", (void*)srv, stcp_server_list_count(inst)); |
||||
} |
||||
|
||||
void stcp_server_list_destroy_all(struct UTUN_INSTANCE *inst) { |
||||
if (!inst) return; |
||||
int count = 0; |
||||
while (inst->stcp_servers) { |
||||
struct stcp_server *next = inst->stcp_servers->next; |
||||
stcp_link_server_destroy(inst->stcp_servers); |
||||
inst->stcp_servers = next; |
||||
count++; |
||||
} |
||||
DEBUG_INFO(DEBUG_CATEGORY_ETCP, "stcp_server_list_destroy_all: destroyed %d servers", count); |
||||
} |
||||
|
||||
int stcp_server_list_count(struct UTUN_INSTANCE *inst) { |
||||
if (!inst) return 0; |
||||
int n = 0; |
||||
for (struct stcp_server *s = inst->stcp_servers; s; s = s->next) n++; |
||||
return n; |
||||
} |
||||
@ -1,153 +0,0 @@
|
||||
// test_etcp_stcp.c — integration: etcp_send + etcp_bind over STCP link
|
||||
#include "stcp_link.h" |
||||
#include "etcp_api.h" |
||||
#include "etcp.h" |
||||
#include "etcp_connections.h" |
||||
#include "secure_channel.h" |
||||
#include "topo_group.h" |
||||
#include "topo_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 25678 |
||||
|
||||
static int test_failed = 0; |
||||
static struct SC_MYKEYS s_keys, c_keys; |
||||
static struct UTUN_INSTANCE *srv_inst, *cli_inst; |
||||
static struct ETCP_CONN *srv_conn, *cli_conn; |
||||
static int g_srv_recv_count = 0; |
||||
static uint8_t g_srv_recv_data[256]; |
||||
|
||||
#define TASSERT(cond) do { \ |
||||
if (!(cond)) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, " FAIL: %s", #cond); test_failed = 1; return test_failed; } \
|
||||
} while(0) |
||||
|
||||
static void etcp_recv_cb(struct ETCP_CONN *conn, struct ll_entry *entry) { |
||||
(void)conn; |
||||
g_srv_recv_count++; |
||||
if (entry->dgram && entry->len) { |
||||
size_t n = entry->len < 256 ? entry->len : 255; |
||||
memcpy(g_srv_recv_data, entry->dgram, n); |
||||
} |
||||
queue_dgram_free(entry); |
||||
queue_entry_free(entry); |
||||
} |
||||
|
||||
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_conn = conn; |
||||
} |
||||
} |
||||
|
||||
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; |
||||
} |
||||
|
||||
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, "=== etcp_send/bind over STCP link ==="); |
||||
|
||||
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); |
||||
|
||||
TASSERT(etcp_bind(srv_inst, 0, etcp_recv_cb) >= 0); |
||||
|
||||
struct TOPO_GROUP_NODE *srv_node = topo_groups_get_default(srv_inst->topo_groups) ? topo_groups_get_default(srv_inst->topo_groups)->local_node : NULL; |
||||
TASSERT(srv_node); |
||||
{ struct TOPO_NODE* srv_ni = topo_node_registry_find(srv_inst->topo_groups, srv_node->node_id); TASSERT(srv_ni); |
||||
struct TOPO_NODE* cli_ni = u_calloc(1, sizeof(struct TOPO_NODE)); |
||||
TASSERT(cli_ni); memcpy(cli_ni, srv_ni, sizeof(struct TOPO_NODE)); cli_ni->group_ref_count = 0; |
||||
cli_ni->v4_sock_meta = NULL; cli_ni->v4_addrs = NULL; cli_ni->v6_sock_meta = NULL; cli_ni->v6_addrs = NULL; |
||||
struct TOPO_ADDR4* a = memory_pool_alloc(cli_inst->topo_groups->v4_addr_pool); |
||||
TASSERT(a); a->addr[0]=127; a->addr[1]=0; a->addr[2]=0; a->addr[3]=1; a->port = port; |
||||
a->type = TOPO_ADDR_INTERFACE; a->socket_id = 0; a->protocol = TOPO_PROTO_TCP; |
||||
a->next = cli_ni->v4_addrs; cli_ni->v4_addrs = a; |
||||
topo_node_registry_store(cli_inst->topo_groups, cli_ni); } |
||||
cli_inst->etcp_connect_timeout_tb = 100000; |
||||
TASSERT(etcp_connect(cli_inst, srv_node, connect_cb, NULL, ETCP_CONNECT_EARLY | ETCP_CONNECT_LATE) == 0); |
||||
|
||||
int ticks = 0; |
||||
while (!cli_conn && ticks < 5000) { uasync_poll(ua, 10); ticks++; } |
||||
TASSERT(cli_conn); |
||||
|
||||
{ struct ll_entry* e = srv_inst->connections->head; |
||||
while (e) { struct conn_queue_entry* ce = (struct conn_queue_entry*)e->data; |
||||
struct ETCP_CONN *c = ce->conn; |
||||
struct ETCP_LINK *l = c->links; while (l) { if (l->is_tcp) { srv_conn = c; break; } l = l->next; } |
||||
if (srv_conn) break; |
||||
e = e->next; } |
||||
} |
||||
TASSERT(srv_conn); |
||||
|
||||
const char *msg = "hello via etcp_send!"; |
||||
struct ll_entry *e = queue_entry_new(0); |
||||
e->dgram = u_malloc(strlen(msg) + 1); |
||||
strcpy((char *)e->dgram, msg); |
||||
e->len = (uint16_t)strlen(msg); |
||||
TASSERT(etcp_send(cli_conn, e) == 0); |
||||
|
||||
ticks = 0; |
||||
while (g_srv_recv_count < 1 && ticks < 2000) { uasync_poll(ua, 10); ticks++; } |
||||
TASSERT(g_srv_recv_count == 1); |
||||
TASSERT(memcmp(g_srv_recv_data, msg, strlen(msg)) == 0); |
||||
|
||||
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; |
||||
} |
||||
@ -1,234 +0,0 @@
|
||||
// test_stcp_traffic.c — STCP transport test: 2 nodes, 1MB bidirectional, random packets, content verify
|
||||
#include "stcp_link.h" |
||||
#include "etcp_api.h" |
||||
#include "etcp.h" |
||||
#include "etcp_connections.h" |
||||
#include "secure_channel.h" |
||||
#include "topo_group.h" |
||||
#include "topo_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) |
||||
|
||||
struct rx_ctx { |
||||
uint8_t *buf; |
||||
size_t len, cap; |
||||
int pkts; |
||||
int done; |
||||
size_t expected_total; |
||||
}; |
||||
|
||||
static struct rx_ctx srv_rx, cli_rx; |
||||
static struct UTUN_INSTANCE *srv_inst, *cli_inst; |
||||
static struct ETCP_CONN *srv_conn, *cli_conn; |
||||
static int srv_send_ready, cli_send_ready; |
||||
|
||||
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; |
||||
} |
||||
|
||||
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_conn = conn; |
||||
cli_send_ready = 1; |
||||
} |
||||
} |
||||
|
||||
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_etcp_packet(struct ETCP_CONN *conn, int side, int seq, size_t paylen) { |
||||
uint8_t buf[PKT_MAX + 5]; |
||||
size_t tot = gen_packet(buf, side, seq, paylen); |
||||
struct ll_entry *e = queue_entry_new(0); |
||||
if (!e) return -1; |
||||
e->dgram = u_malloc(tot); |
||||
if (!e->dgram) { queue_entry_free(e); return -1; } |
||||
memcpy(e->dgram, buf, tot); |
||||
e->len = (uint16_t)tot; |
||||
return etcp_send(conn, e); |
||||
} |
||||
|
||||
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; |
||||
} |
||||
|
||||
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_CATEGORY_ETCP); |
||||
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 TOPO_GROUP_NODE *srv_node = topo_groups_get_default(srv_inst->topo_groups) ? topo_groups_get_default(srv_inst->topo_groups)->local_node : NULL; |
||||
TASSERT(srv_node); |
||||
{ struct TOPO_NODE* srv_ni = topo_node_registry_find(srv_inst->topo_groups, srv_node->node_id); TASSERT(srv_ni); |
||||
struct TOPO_NODE* cli_ni = u_calloc(1, sizeof(struct TOPO_NODE)); |
||||
TASSERT(cli_ni); memcpy(cli_ni, srv_ni, sizeof(struct TOPO_NODE)); cli_ni->group_ref_count = 0; |
||||
cli_ni->v4_sock_meta = NULL; cli_ni->v4_addrs = NULL; cli_ni->v6_sock_meta = NULL; cli_ni->v6_addrs = NULL; |
||||
struct TOPO_ADDR4* a = memory_pool_alloc(cli_inst->topo_groups->v4_addr_pool); |
||||
TASSERT(a); a->addr[0]=127; a->addr[1]=0; a->addr[2]=0; a->addr[3]=1; a->port = port; |
||||
a->type = TOPO_ADDR_INTERFACE; a->socket_id = 0; a->protocol = TOPO_PROTO_TCP; |
||||
a->next = cli_ni->v4_addrs; cli_ni->v4_addrs = a; |
||||
topo_node_registry_store(cli_inst->topo_groups, cli_ni); } |
||||
cli_inst->etcp_connect_timeout_tb = 100000; |
||||
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_conn); |
||||
|
||||
{ struct ll_entry* e = srv_inst->connections->head; |
||||
while (e) { struct conn_queue_entry* ce = (struct conn_queue_entry*)e->data; |
||||
struct ETCP_CONN *c = ce->conn; |
||||
struct ETCP_LINK *l = c->links; while (l) { if (l->is_tcp) { srv_conn = c; break; } l = l->next; } |
||||
if (srv_conn) break; |
||||
e = e->next; } |
||||
} |
||||
TASSERT(srv_conn); |
||||
|
||||
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_etcp_packet(srv_conn, 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_etcp_packet(cli_conn, 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); |
||||
|
||||
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…
Reference in new issue