diff --git a/src/Makefile.am b/src/Makefile.am index b0af8400..3bcea869 100644 --- a/src/Makefile.am +++ b/src/Makefile.am @@ -31,7 +31,11 @@ utun_CORE_SOURCES = \ firewall.c \ eim_nat.c \ nat_transport.c \ - dummynet.c + dummynet.c \ + tcp_proxy.c \ + etcp_router.c \ + remote_proxy.c \ + uip/uip.c # Platform-specific TUN libs (Windows only) utun_TUN_LIBS = @TUN_LIBS@ @@ -62,6 +66,7 @@ endif # Include paths utun_CORE_CFLAGS = \ -I$(top_srcdir)/lib \ + -I$(top_srcdir)/src/uip \ -I$(top_srcdir)/tinycrypt/lib/include \ -I$(top_srcdir)/tinycrypt/lib/source \ -g \ diff --git a/src/config_parser.c b/src/config_parser.c index 54aa00a4..22282d87 100644 --- a/src/config_parser.c +++ b/src/config_parser.c @@ -44,7 +44,9 @@ typedef enum { SECTION_FIREWALL, SECTION_CONTROL, SECTION_ALLOWED_KEYS, - SECTION_NAT + SECTION_NAT, + SECTION_TCP_PROXY, + SECTION_REMOTE_PROXY } section_type_t; static char* trim(char *str) { @@ -436,6 +438,71 @@ static int parse_control(const char *key, const char *value, struct global_confi return 0; } +static int parse_tcp_proxy(const char *key, const char *value, struct global_config *global) { + if (strcmp(key, "enabled") == 0) { + global->tcp_proxy_enabled = strcasecmp(value, "yes") == 0 || strcasecmp(value, "1") == 0 || strcasecmp(value, "true") == 0; + return 0; + } + if (strcmp(key, "tun_name") == 0) { + strncpy(global->tcp_proxy_tun_name, value, sizeof(global->tcp_proxy_tun_name) - 1); + return 0; + } + if (strcmp(key, "tun_ip") == 0) { + strncpy(global->tcp_proxy_tun_ip, value, sizeof(global->tcp_proxy_tun_ip) - 1); + return 0; + } + if (strcmp(key, "mtu") == 0) { + global->tcp_proxy_mtu = atoi(value); + return 0; + } + if (strcmp(key, "eim_timeout") == 0) { + global->tcp_proxy_eim_timeout = atoi(value); + return 0; + } + if (strcmp(key, "forward") == 0) { + if (global->tcp_proxy_mapping_count >= MAX_TCP_PROXY_MAPPINGS) { + DEBUG_ERROR(DEBUG_CATEGORY_CONFIG, "Too many tcp_proxy forward rules (max %d)", MAX_TCP_PROXY_MAPPINGS); + return -1; + } + char buf[256]; + strncpy(buf, value, sizeof(buf) - 1); + buf[sizeof(buf) - 1] = '\0'; + trim(buf); + + char* local_port_str = strtok(buf, " "); + char* arrow = strtok(NULL, " "); + char* remote_str = strtok(NULL, ""); + if (!local_port_str || !arrow || !remote_str || strcmp(arrow, "->") != 0) { + DEBUG_ERROR(DEBUG_CATEGORY_CONFIG, "Invalid tcp_proxy forward format: %s (expected: port -> ip:port [via NODE_HEX])", value); + return -1; + } + int local_port = atoi(local_port_str); + if (local_port <= 0 || local_port > 65535) { DEBUG_ERROR(DEBUG_CATEGORY_CONFIG, "Invalid tcp_proxy forward local port: %s", local_port_str); return -1; } + + char* remote_ip_str = strtok(remote_str, ":"); + char* remote_port_str = strtok(NULL, ""); + if (!remote_ip_str || !remote_port_str) { + DEBUG_ERROR(DEBUG_CATEGORY_CONFIG, "Invalid tcp_proxy forward remote: %s (expected ip:port)", remote_str); + return -1; + } + int remote_port = atoi(remote_port_str); + if (remote_port <= 0 || remote_port > 65535) { DEBUG_ERROR(DEBUG_CATEGORY_CONFIG, "Invalid tcp_proxy forward remote port: %s", remote_port_str); return -1; } + + uint64_t via_node_id = 0; + char* via_ptr = strstr(remote_port_str, "via "); + if (via_ptr) via_node_id = strtoull(via_ptr + 4, NULL, 16); + + int idx = global->tcp_proxy_mapping_count; + global->tcp_proxy_mappings[idx].local_port = (uint16_t)local_port; + strncpy(global->tcp_proxy_mappings[idx].remote_ip, remote_ip_str, sizeof(global->tcp_proxy_mappings[idx].remote_ip) - 1); + global->tcp_proxy_mappings[idx].remote_port = (uint16_t)remote_port; + global->tcp_proxy_mappings[idx].via_node_id = via_node_id; + global->tcp_proxy_mapping_count++; + return 0; + } + return 0; +} + static int parse_nat(const char *key, const char *value, struct global_config *global) { if (strcmp(key, "tun_ifname") == 0) { strncpy(global->nat_tun_ifname, value, sizeof(global->nat_tun_ifname) - 1); @@ -627,6 +694,8 @@ static section_type_t parse_section_header(const char *line, char *name, size_t if (strcasecmp(section, "control") == 0) return SECTION_CONTROL; if (strcasecmp(section, "allowed_keys") == 0) return SECTION_ALLOWED_KEYS; if (strcasecmp(section, "nat") == 0) return SECTION_NAT; + if (strcasecmp(section, "tcp_proxy") == 0) return SECTION_TCP_PROXY; + if (strcasecmp(section, "remote_proxy") == 0) return SECTION_REMOTE_PROXY; char *colon = strchr(section, ':'); if (!colon) return SECTION_UNKNOWN; @@ -791,6 +860,15 @@ static struct utun_config* parse_config_internal(FILE *fp, const char *filename) DEBUG_WARN(DEBUG_CATEGORY_CONFIG, "%s:%d: Invalid NAT key '%s'", filename, line_num, key); } break; + case SECTION_TCP_PROXY: + cfg->global.tcp_proxy_enabled = 1; + if (parse_tcp_proxy(key, value, &cfg->global) < 0) { + DEBUG_WARN(DEBUG_CATEGORY_CONFIG, "%s:%d: Invalid tcp_proxy key '%s'", filename, line_num, key); + } + break; + case SECTION_REMOTE_PROXY: + cfg->global.remote_proxy_enabled = 1; + break; default: DEBUG_WARN(DEBUG_CATEGORY_CONFIG, "%s:%d: Key outside section: %s", filename, line_num, key); break; diff --git a/src/config_parser.h b/src/config_parser.h index 1c1f25e2..508b53f0 100644 --- a/src/config_parser.h +++ b/src/config_parser.h @@ -82,6 +82,14 @@ struct CFG_ALLOWED_KEY { struct CFG_ALLOWED_KEY *next; }; +#define MAX_TCP_PROXY_MAPPINGS 32 +struct tcp_proxy_mapping_config { + uint16_t local_port; + char remote_ip[64]; + uint16_t remote_port; + uint64_t via_node_id; // 0=локальный прокси, !=0=удаленный прокси через эту ноду +}; + struct global_config { char name[16]; // Instance name char my_private_key_hex[MAX_KEY_LEN]; @@ -143,6 +151,18 @@ struct global_config { uint16_t external_port; } nat_forwards[MAX_NAT_FORWARDS]; int nat_forward_count; + + // TCP proxy configuration ([tcp_proxy] section) + int tcp_proxy_enabled; + char tcp_proxy_tun_name[16]; + char tcp_proxy_tun_ip[64]; + int tcp_proxy_mtu; + int tcp_proxy_eim_timeout; + struct tcp_proxy_mapping_config tcp_proxy_mappings[MAX_TCP_PROXY_MAPPINGS]; + int tcp_proxy_mapping_count; + + // Remote proxy (exit node) configuration ([remote_proxy] section) + int remote_proxy_enabled; }; struct utun_config { diff --git a/src/etcp_api.h b/src/etcp_api.h index 8c6c43c3..9bf8726f 100644 --- a/src/etcp_api.h +++ b/src/etcp_api.h @@ -24,6 +24,8 @@ #define ETCP_ID_DATA 0x00 // Пакет для передачи адресату #define ETCP_ID_ROUTE_ENTRY 0x01 // Элемент роутинг-таблицы #define ETCP_ID_NAT 0x02 // NAT трафик между узлами +#define ETCP_ID_SVC_ROUTE 0x03 // Маршрутизируемые сервисные пакеты (etcp_router) +#define ETCP_ID_TCP_PROXY 0x04 // TCP proxy через удаленный узел (remote_proxy) // Forward declarations struct ETCP_CONN; diff --git a/src/etcp_router.c b/src/etcp_router.c new file mode 100644 index 00000000..c23ea80e --- /dev/null +++ b/src/etcp_router.c @@ -0,0 +1,150 @@ +// etcp_router.c — Сервисный слой маршрутизации поверх ETCP +// Маршрутизирует сервисные пакеты до целевой ноды, на промежуточных нодах ретранслирует +#include "etcp_router.h" +#include "etcp.h" +#include "utun_instance.h" +#include "route_bgp.h" +#include "../lib/debug_config.h" +#include "../lib/mem.h" +#include "../lib/ll_queue.h" +#include + +// Обработчик ETCP_ID_SVC_ROUTE — вызывается etcp_int_recv на каждом узле +static void etcp_router_recv_cb(struct ETCP_CONN* conn, struct ll_entry* entry) { + if (!conn || !entry) { + if (entry) { queue_entry_free(entry); queue_dgram_free(entry); } + return; + } + struct UTUN_INSTANCE* inst = conn->instance; + if (!inst || entry->len < SVC_ROUTE_HDR_SIZE) { + DEBUG_WARN(DEBUG_CATEGORY_ETCP, "etcp_router: invalid packet inst=%p len=%zu min=%d", + (void*)inst, entry ? entry->len : 0, SVC_ROUTE_HDR_SIZE); + queue_entry_free(entry); queue_dgram_free(entry); + return; + } + uint8_t svc_id = entry->dgram[1]; + uint64_t dst_node_id = 0; memcpy(&dst_node_id, entry->dgram + 2, 8); + if (dst_node_id == inst->node_id) { + if (svc_id >= SVC_ROUTE_MAX_BINDINGS || !inst->router_bindings.callbacks[svc_id]) { + DEBUG_WARN(DEBUG_CATEGORY_ETCP, "etcp_router: no handler for svc_id=%u dst=%016llx self=%016llx", + svc_id, (unsigned long long)dst_node_id, (unsigned long long)inst->node_id); + queue_entry_free(entry); queue_dgram_free(entry); + return; + } + size_t payload_len = entry->len - SVC_ROUTE_HDR_SIZE; + struct ll_entry* svc_entry = queue_entry_new(0); + if (!svc_entry) { queue_entry_free(entry); queue_dgram_free(entry); return; } + svc_entry->len = 1 + payload_len; + svc_entry->dgram = u_malloc(svc_entry->len); + if (!svc_entry->dgram) { queue_entry_free(svc_entry); queue_entry_free(entry); queue_dgram_free(entry); return; } + svc_entry->dgram[0] = svc_id; + if (payload_len > 0) memcpy(svc_entry->dgram + 1, entry->dgram + SVC_ROUTE_HDR_SIZE, payload_len); + DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "etcp_router: delivering svc_id=%u len=%zu to handler %p src=%016llx", + svc_id, payload_len, (void*)inst->router_bindings.callbacks[svc_id], + (unsigned long long)(*(uint64_t*)(entry->dgram + 10))); + queue_entry_free(entry); queue_dgram_free(entry); + inst->router_bindings.callbacks[svc_id](conn, svc_entry); + } else { + struct ETCP_CONN* next = route_bgp_find_conn_for_node(inst->bgp, dst_node_id); + if (!next) { + DEBUG_WARN(DEBUG_CATEGORY_ETCP, "etcp_router: no route to %016llx svc_id=%u, dropping", + (unsigned long long)dst_node_id, svc_id); + queue_entry_free(entry); queue_dgram_free(entry); + return; + } + DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "etcp_router: forwarding svc_id=%u → %016llx via %s", + svc_id, (unsigned long long)dst_node_id, next->log_name); + etcp_send(next, entry); + } +} + +// ==================================================================== +// Public API +// ==================================================================== + +int etcp_router_init(struct UTUN_INSTANCE* inst) { + if (!inst) { DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "etcp_router_init: NULL instance"); return -1; } + memset(&inst->router_bindings, 0, sizeof(inst->router_bindings)); + int ret = etcp_bind(inst, ETCP_ID_SVC_ROUTE, etcp_router_recv_cb); + if (ret != 0) { DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "etcp_router_init: etcp_bind failed, ret=%d", ret); return -1; } + DEBUG_INFO(DEBUG_CATEGORY_ETCP, "etcp_router initialized for node %016llx", (unsigned long long)inst->node_id); + return 0; +} + +void etcp_router_destroy(struct UTUN_INSTANCE* inst) { + if (!inst) return; + etcp_unbind(inst, ETCP_ID_SVC_ROUTE); + memset(&inst->router_bindings, 0, sizeof(inst->router_bindings)); + DEBUG_INFO(DEBUG_CATEGORY_ETCP, "etcp_router destroyed for node %016llx", (unsigned long long)inst->node_id); +} + +int etcp_router_bind(struct UTUN_INSTANCE* inst, uint8_t svc_id, etcp_recv_fn callback) { + if (!inst) { DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "etcp_router_bind: NULL instance"); return -1; } + if (svc_id >= SVC_ROUTE_MAX_BINDINGS) { DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "etcp_router_bind: invalid svc_id %u", svc_id); return -1; } + if (!callback) { DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "etcp_router_bind: NULL callback for svc_id=%u", svc_id); return -1; } + if (inst->router_bindings.callbacks[svc_id]) DEBUG_WARN(DEBUG_CATEGORY_ETCP, "etcp_router_bind: overwriting svc_id=%u", svc_id); + inst->router_bindings.callbacks[svc_id] = callback; + DEBUG_INFO(DEBUG_CATEGORY_ETCP, "etcp_router_bind: svc_id=%u → cb=%p", svc_id, (void*)callback); + return 0; +} + +int etcp_router_unbind(struct UTUN_INSTANCE* inst, uint8_t svc_id) { + if (!inst) { DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "etcp_router_unbind: NULL instance"); return -1; } + if (svc_id >= SVC_ROUTE_MAX_BINDINGS) { DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "etcp_router_unbind: invalid svc_id %u", svc_id); return -1; } + if (!inst->router_bindings.callbacks[svc_id]) { DEBUG_WARN(DEBUG_CATEGORY_ETCP, "etcp_router_unbind: svc_id=%u not bound", svc_id); return -1; } + inst->router_bindings.callbacks[svc_id] = NULL; + DEBUG_INFO(DEBUG_CATEGORY_ETCP, "etcp_router_unbind: svc_id=%u", svc_id); + return 0; +} + +int etcp_route_send(struct UTUN_INSTANCE* inst, uint64_t dst_node_id, struct ll_entry* entry) { + if (!inst || !entry) { + DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "etcp_route_send: NULL inst=%p entry=%p", (void*)inst, (void*)entry); + if (entry) { queue_entry_free(entry); queue_dgram_free(entry); } + return -1; + } + if (!entry->dgram || entry->len < 1) { + DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "etcp_route_send: empty entry"); + queue_entry_free(entry); queue_dgram_free(entry); + return -1; + } + uint8_t svc_id = entry->dgram[0]; + size_t payload_len = entry->len - 1; + + // Loopback — dispatch прямо локально + if (dst_node_id == inst->node_id) { + DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "etcp_route_send: loopback svc_id=%u len=%zu", svc_id, payload_len); + if (svc_id < SVC_ROUTE_MAX_BINDINGS && inst->router_bindings.callbacks[svc_id]) + inst->router_bindings.callbacks[svc_id](NULL, entry); + else { queue_entry_free(entry); queue_dgram_free(entry); } + return 0; + } + + // Упаковываем в routing header + size_t total_len = SVC_ROUTE_HDR_SIZE + payload_len; + uint8_t* new_dgram = u_malloc(total_len); + if (!new_dgram) { queue_entry_free(entry); queue_dgram_free(entry); return -1; } + new_dgram[0] = ETCP_ID_SVC_ROUTE; + new_dgram[1] = svc_id; + memcpy(new_dgram + 2, &dst_node_id, 8); + memcpy(new_dgram + 10, &inst->node_id, 8); + if (payload_len > 0) memcpy(new_dgram + SVC_ROUTE_HDR_SIZE, entry->dgram + 1, payload_len); + + struct ll_entry* new_entry = queue_entry_new(0); + if (!new_entry) { u_free(new_dgram); queue_entry_free(entry); queue_dgram_free(entry); return -1; } + new_entry->dgram = new_dgram; + new_entry->len = total_len; + + queue_entry_free(entry); queue_dgram_free(entry); + + struct ETCP_CONN* conn = route_bgp_find_conn_for_node(inst->bgp, dst_node_id); + if (!conn) { + DEBUG_WARN(DEBUG_CATEGORY_ETCP, "etcp_route_send: no route to %016llx svc_id=%u, dropping", + (unsigned long long)dst_node_id, svc_id); + queue_entry_free(new_entry); queue_dgram_free(new_entry); + return -1; + } + DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "etcp_route_send: svc_id=%u → %016llx via %s len=%zu", + svc_id, (unsigned long long)dst_node_id, conn->log_name, total_len); + return etcp_send(conn, new_entry); +} diff --git a/src/etcp_router.h b/src/etcp_router.h new file mode 100644 index 00000000..84c708d3 --- /dev/null +++ b/src/etcp_router.h @@ -0,0 +1,38 @@ +// etcp_router.h — Сервисный слой маршрутизации поверх ETCP +// Позволяет отправлять сервисные пакеты конкретной ноде по node_id с многошаговой маршрутизацией +// Формат transit-пакета: [ETCP_ID_SVC_ROUTE:1] [svc_id:1] [dst_node_id:8] [src_node_id:8] [payload...] +#ifndef ETCP_ROUTER_H +#define ETCP_ROUTER_H + +#include +#include "etcp_api.h" + +#define SVC_ROUTE_HDR_SIZE 18 // ETCP_ID_SVC_ROUTE(1) + svc_id(1) + dst_node_id(8) + src_node_id(8) +#define SVC_ROUTE_MAX_BINDINGS 256 + +// Bindings для сервисов внутри etcp_router (аналогично ETCP_BINDINGS) +struct ETCP_ROUTER_BINDINGS { + etcp_recv_fn callbacks[SVC_ROUTE_MAX_BINDINGS]; +}; + +// Инициализация: etcp_bind(ETCP_ID_SVC_ROUTE) + очистка bindings +int etcp_router_init(struct UTUN_INSTANCE* inst); + +// Деинициализация: etcp_unbind(ETCP_ID_SVC_ROUTE) +void etcp_router_destroy(struct UTUN_INSTANCE* inst); + +// Зарегистрировать обработчик сервиса. +// Когда пакет достигает целевой ноды, вызывается callback(conn, entry) +// entry содержит: [svc_id:1] [payload...] — как отправлено через etcp_route_send +int etcp_router_bind(struct UTUN_INSTANCE* inst, uint8_t svc_id, etcp_recv_fn callback); + +// Отписаться от сервиса +int etcp_router_unbind(struct UTUN_INSTANCE* inst, uint8_t svc_id); + +// Отправить сервисный пакет узлу по node_id через overlay. +// entry формат: [svc_id:1] [payload...] +// Функция заворачивает в routing header, находит next hop через BGP, отправляет. +// Принимает ownership entry. +int etcp_route_send(struct UTUN_INSTANCE* inst, uint64_t dst_node_id, struct ll_entry* entry); + +#endif // ETCP_ROUTER_H diff --git a/src/remote_proxy.c b/src/remote_proxy.c new file mode 100644 index 00000000..488b7920 --- /dev/null +++ b/src/remote_proxy.c @@ -0,0 +1,229 @@ +// remote_proxy.c — Удаленный TCP прокси (exit node) +#include "remote_proxy.h" +#include "tcp_proxy.h" +#include "etcp.h" +#include "etcp_api.h" +#include "etcp_router.h" +#include "utun_instance.h" +#include "../lib/u_async.h" +#include "../lib/debug_config.h" +#include "../lib/ll_queue.h" +#include "../lib/mem.h" +#include +#include +#include +#ifndef _WIN32 +#include +#include +#include +#include +#include +#endif +#ifndef MSG_NOSIGNAL +#define MSG_NOSIGNAL 0 +#endif + +static struct remote_proxy_ctx* g_rp_ctx = NULL; + +static void rp_sock_read_cb (socket_t sock, void* arg); +static void rp_sock_write_cb(socket_t sock, void* arg); +static void rp_sock_error_cb(socket_t sock, void* arg); + +static int rp_send_msg(struct UTUN_INSTANCE* inst, uint64_t dst, uint8_t subcmd, + uint64_t sid, const uint8_t* data, size_t len) { + struct ll_entry* e = queue_entry_new(0); + if (!e) return -1; + e->dgram = u_malloc(TCP_PROXY_HDR_SIZE + len); + if (!e->dgram) { queue_entry_free(e); return -1; } + e->dgram[0] = ETCP_ID_TCP_PROXY; + e->dgram[1] = subcmd; + memcpy(e->dgram + 2, &sid, 8); + if (len > 0) memcpy(e->dgram + TCP_PROXY_HDR_SIZE, data, len); + e->len = TCP_PROXY_HDR_SIZE + len; + return etcp_route_send(inst, dst, e); +} + +static int rp_send_connected(struct UTUN_INSTANCE* inst, uint64_t dst, uint64_t sid, + uint16_t local_port, uint8_t status) { + uint8_t buf[3]; + memcpy(buf, &local_port, 2); buf[2] = status; + return rp_send_msg(inst, dst, TCP_PROXY_SUBCMD_CONNECTED, sid, buf, 3); +} + +// ==================================================================== +// Socket callbacks +// ==================================================================== +static void rp_sock_read_cb(socket_t sock, void* arg) { + (void)sock; struct remote_proxy_conn* rc = (struct remote_proxy_conn*)arg; + if (!rc || rc->sock == SOCKET_INVALID) return; + uint8_t buf[8192]; ssize_t n = recv(rc->sock, buf, sizeof(buf), 0); + if (n > 0) { + struct UTUN_INSTANCE* inst = rc->ctx ? rc->ctx->inst : NULL; + if (inst) rp_send_msg(inst, rc->peer_node_id, TCP_PROXY_SUBCMD_DATA, rc->stream_id, buf, (size_t)n); + } else if (n == 0) { + DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "remote_proxy: EOF stream=%016llx", (unsigned long long)rc->stream_id); + struct UTUN_INSTANCE* inst = rc->ctx ? rc->ctx->inst : NULL; + if (inst) rp_send_msg(inst, rc->peer_node_id, TCP_PROXY_SUBCMD_CLOSE, rc->stream_id, NULL, 0); + rc->connected = -1; + } else if (errno != EAGAIN && errno != EWOULDBLOCK && errno != EINTR) { + DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "remote_proxy: recv error %s", strerror(errno)); + } +} + +static void rp_sock_write_cb(socket_t sock, void* arg) { + (void)sock; struct remote_proxy_conn* rc = (struct remote_proxy_conn*)arg; + if (!rc || rc->sock == SOCKET_INVALID) return; + if (!rc->connected) { + if (!rc->connect_called) return; + int err = 0; socklen_t len = sizeof(err); + if (getsockopt(rc->sock, SOL_SOCKET, SO_ERROR, &err, &len) == 0 && err == 0) { + rc->connected = 1; + struct sockaddr_in local; socklen_t llen = sizeof(local); + uint16_t local_port = 0; + if (getsockname(rc->sock, (struct sockaddr*)&local, &llen) == 0) local_port = local.sin_port; + DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "remote_proxy: connected stream=%016llx to %d.%d.%d.%d:%d local_port=%d", + (unsigned long long)rc->stream_id, rc->dest_ip[0], rc->dest_ip[1], rc->dest_ip[2], rc->dest_ip[3], + ntohs(rc->dest_port), ntohs(local_port)); + struct UTUN_INSTANCE* inst = rc->ctx ? rc->ctx->inst : NULL; + if (inst) rp_send_connected(inst, rc->peer_node_id, rc->stream_id, local_port, TCP_PROXY_CONNECTED_OK); + } else { + DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "remote_proxy: connect failed %s", err ? strerror(err) : "unknown"); + struct UTUN_INSTANCE* inst = rc->ctx ? rc->ctx->inst : NULL; + if (inst) rp_send_connected(inst, rc->peer_node_id, rc->stream_id, 0, TCP_PROXY_CONNECTED_REFUSED); + rc->connected = -1; + } + } +} + +static void rp_sock_error_cb(socket_t sock, void* arg) { + (void)sock; struct remote_proxy_conn* rc = (struct remote_proxy_conn*)arg; + DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "remote_proxy: socket error stream=%016llx", (unsigned long long)rc->stream_id); + if (rc) rc->connected = -1; +} + +static void rp_conn_free(struct remote_proxy_conn* rc) { + if (!rc) return; + if (rc->sock != SOCKET_INVALID) { + if (rc->read_id) { uasync_remove_socket_t(rc->ua, rc->sock); rc->read_id = NULL; } + socket_close_wrapper(rc->sock); + } + u_free(rc); +} + +// ==================================================================== +// Standalone handler (registered when tcp_proxy doesn't have remote mappings) +// ==================================================================== +static void rp_standalone_recv_cb(struct ETCP_CONN* conn, struct ll_entry* entry) { + struct UTUN_INSTANCE* i = conn ? conn->instance : (g_rp_ctx ? g_rp_ctx->inst : NULL); + if (!i || !entry || entry->dgram == NULL || entry->len < TCP_PROXY_HDR_SIZE) { + if (entry) { queue_entry_free(entry); queue_dgram_free(entry); } + return; + } + uint8_t subcmd = entry->dgram[1]; + uint64_t sid = 0; memcpy(&sid, entry->dgram + 2, 8); + uint64_t src = conn ? conn->peer_node_id : i->node_id; + if (subcmd == TCP_PROXY_SUBCMD_CONNECT) { remote_proxy_handle_connect(i, entry, sid, src); return; } + if (subcmd == TCP_PROXY_SUBCMD_DATA) { remote_proxy_handle_data(i, entry, sid); return; } + if (subcmd == TCP_PROXY_SUBCMD_CLOSE) { remote_proxy_handle_close(i, sid); queue_entry_free(entry); queue_dgram_free(entry); return; } + queue_entry_free(entry); queue_dgram_free(entry); +} + +// ==================================================================== +// Public API +// ==================================================================== +struct remote_proxy_conn* remote_proxy_find_conn(struct remote_proxy_ctx* ctx, uint64_t stream_id) { + struct remote_proxy_conn* c; + for (c = ctx->conns; c; c = c->next) if (c->stream_id == stream_id) return c; + return NULL; +} + +int remote_proxy_handle_connect(struct UTUN_INSTANCE* inst, struct ll_entry* entry, + uint64_t stream_id, uint64_t src_node_id) { + if (!inst || !inst->remote_proxy.enabled) { queue_entry_free(entry); queue_dgram_free(entry); return -1; } + struct remote_proxy_ctx* ctx = &inst->remote_proxy; + if (entry->len < TCP_PROXY_CONNECT_HDR_SIZE) { queue_entry_free(entry); queue_dgram_free(entry); return -1; } + uint8_t* dest_ip = entry->dgram + TCP_PROXY_HDR_SIZE; + uint16_t dest_port = 0; memcpy(&dest_port, dest_ip + 4, 2); + + struct remote_proxy_conn* rc = u_calloc(1, sizeof(struct remote_proxy_conn)); + if (!rc) { queue_entry_free(entry); queue_dgram_free(entry); return -1; } + rc->ctx = ctx; rc->stream_id = stream_id; rc->peer_node_id = src_node_id; + memcpy(rc->dest_ip, dest_ip, 4); rc->dest_port = dest_port; + rc->ua = inst->ua; rc->sock = SOCKET_INVALID; + + rc->sock = socket(AF_INET, SOCK_STREAM, 0); + if (rc->sock == SOCKET_INVALID) { DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "remote_proxy: socket() failed"); rp_conn_free(rc); queue_entry_free(entry); queue_dgram_free(entry); return -1; } + socket_set_nonblocking(rc->sock); + rc->read_id = uasync_add_socket_t(rc->ua, rc->sock, rp_sock_read_cb, rp_sock_write_cb, rp_sock_error_cb, rc); + if (!rc->read_id) { DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "remote_proxy: uasync_add_socket_t failed"); rp_conn_free(rc); queue_entry_free(entry); queue_dgram_free(entry); return -1; } + + struct sockaddr_in addr; memset(&addr, 0, sizeof(addr)); + addr.sin_family = AF_INET; memcpy(&addr.sin_addr.s_addr, dest_ip, 4); addr.sin_port = dest_port; + DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "remote_proxy: connecting stream=%016llx to %d.%d.%d.%d:%d", + (unsigned long long)stream_id, dest_ip[0], dest_ip[1], dest_ip[2], dest_ip[3], ntohs(dest_port)); + rc->connect_called = 1; + int ret = connect(rc->sock, (struct sockaddr*)&addr, sizeof(addr)); + if (ret < 0 && errno != EINPROGRESS) { + DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "remote_proxy: connect() failed: %s", strerror(errno)); + rp_send_connected(inst, src_node_id, stream_id, 0, TCP_PROXY_CONNECTED_REFUSED); + rp_conn_free(rc); queue_entry_free(entry); queue_dgram_free(entry); return -1; + } + + rc->next = ctx->conns; ctx->conns = rc; + queue_entry_free(entry); queue_dgram_free(entry); + return 0; +} + +int remote_proxy_handle_data(struct UTUN_INSTANCE* inst, struct ll_entry* entry, uint64_t stream_id) { + if (!inst) { queue_entry_free(entry); queue_dgram_free(entry); return -1; } + struct remote_proxy_ctx* ctx = &inst->remote_proxy; + struct remote_proxy_conn* rc = remote_proxy_find_conn(ctx, stream_id); + if (!rc || rc->sock == SOCKET_INVALID || rc->connected != 1) { queue_entry_free(entry); queue_dgram_free(entry); return -1; } + size_t data_len = entry->len - TCP_PROXY_HDR_SIZE; + uint8_t* data = entry->dgram + TCP_PROXY_HDR_SIZE; + ssize_t n = send(rc->sock, data, data_len, MSG_NOSIGNAL); + if (n < 0 && errno != EAGAIN && errno != EWOULDBLOCK) { + DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "remote_proxy: send error %s", strerror(errno)); + } + queue_entry_free(entry); queue_dgram_free(entry); + return 0; +} + +void remote_proxy_handle_close(struct UTUN_INSTANCE* inst, uint64_t stream_id) { + if (!inst) return; + struct remote_proxy_ctx* ctx = &inst->remote_proxy; + struct remote_proxy_conn** prev = &ctx->conns; + while (*prev) { + struct remote_proxy_conn* rc = *prev; + if (rc->stream_id == stream_id) { + DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "remote_proxy: close stream=%016llx", (unsigned long long)stream_id); + *prev = rc->next; rp_conn_free(rc); return; + } + prev = &rc->next; + } +} + +int remote_proxy_init(struct UTUN_INSTANCE* inst) { + if (!inst) return -1; + struct remote_proxy_ctx* ctx = &inst->remote_proxy; + memset(ctx, 0, sizeof(*ctx)); + ctx->enabled = inst->config && inst->config->global.remote_proxy_enabled; + ctx->inst = inst; + if (!ctx->enabled) return 0; + g_rp_ctx = ctx; + // Register standalone handler for when tcp_proxy is not using remote mappings + // (tcp_proxy_create will overwrite this handler if it has remote mappings) + etcp_router_bind(inst, ETCP_ID_TCP_PROXY, rp_standalone_recv_cb); + DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "remote_proxy initialized on node %016llx", (unsigned long long)inst->node_id); + return 0; +} + +void remote_proxy_destroy(struct UTUN_INSTANCE* inst) { + if (!inst) return; + struct remote_proxy_ctx* ctx = &inst->remote_proxy; + struct remote_proxy_conn* rc = ctx->conns; + while (rc) { struct remote_proxy_conn* next = rc->next; rp_conn_free(rc); rc = next; } + ctx->conns = NULL; ctx->enabled = 0; + if (g_rp_ctx == ctx) g_rp_ctx = NULL; + DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "remote_proxy destroyed"); +} diff --git a/src/remote_proxy.h b/src/remote_proxy.h new file mode 100644 index 00000000..96ab4009 --- /dev/null +++ b/src/remote_proxy.h @@ -0,0 +1,58 @@ +// remote_proxy.h — Удаленный TCP прокси (exit node) +// Создаёт OS сокеты к адресатам по запросам от других нод через etcp_router +#ifndef REMOTE_PROXY_H +#define REMOTE_PROXY_H + +#include +#include "../lib/socket_compat.h" + +struct UTUN_INSTANCE; +struct UASYNC; +struct ll_entry; +struct ETCP_CONN; +struct remote_proxy_ctx; + +// Sub-commands for ETCP_ID_TCP_PROXY protocol +#define TCP_PROXY_SUBCMD_CONNECT 0x01 +#define TCP_PROXY_SUBCMD_CONNECTED 0x02 +#define TCP_PROXY_SUBCMD_DATA 0x03 +#define TCP_PROXY_SUBCMD_CLOSE 0x04 + +#define TCP_PROXY_CONNECTED_OK 0 +#define TCP_PROXY_CONNECTED_REFUSED 1 + +#define TCP_PROXY_HDR_SIZE 10 +#define TCP_PROXY_CONNECT_HDR_SIZE 16 +#define TCP_PROXY_CONNECTED_HDR_SIZE 13 + +struct remote_proxy_conn { + struct remote_proxy_conn* next; + struct remote_proxy_ctx* ctx; + uint64_t stream_id; + uint64_t peer_node_id; + socket_t sock; + void* read_id; + struct UASYNC* ua; + int connected; + int connect_called; + uint8_t dest_ip[4]; + uint16_t dest_port; +}; + +struct remote_proxy_ctx { + int enabled; + struct remote_proxy_conn* conns; + struct UTUN_INSTANCE* inst; +}; + +int remote_proxy_init(struct UTUN_INSTANCE* inst); +void remote_proxy_destroy(struct UTUN_INSTANCE* inst); + +struct remote_proxy_conn* remote_proxy_find_conn(struct remote_proxy_ctx* ctx, uint64_t stream_id); + +int remote_proxy_handle_connect(struct UTUN_INSTANCE* inst, struct ll_entry* entry, + uint64_t stream_id, uint64_t src_node_id); +int remote_proxy_handle_data(struct UTUN_INSTANCE* inst, struct ll_entry* entry, uint64_t stream_id); +void remote_proxy_handle_close(struct UTUN_INSTANCE* inst, uint64_t stream_id); + +#endif diff --git a/src/tcp_proxy.c b/src/tcp_proxy.c new file mode 100644 index 00000000..16a0c857 --- /dev/null +++ b/src/tcp_proxy.c @@ -0,0 +1,720 @@ +// tcp_proxy.c — TCP proxy: uIP TCP ↔ ll_queue ↔ OS socket → destination +// Raw-fd mode + TUN mode + active connections + EIM NAT + half-close +#include "tcp_proxy.h" +#include "config_parser.h" +#include "tun_if.h" +#include "utun_instance.h" +#include "etcp.h" +#include "etcp_router.h" +#include "remote_proxy.h" +#include "uip/uip.h" +#include "../lib/u_async.h" +#include "../lib/debug_config.h" +#include "../lib/ll_queue.h" +#include "../lib/memory_pool.h" +#include "../lib/mem.h" +#include +#include +#include +#ifndef _WIN32 +#include +#include +#include +#include +#include +#endif + +#ifndef MSG_NOSIGNAL +#define MSG_NOSIGNAL 0 +#endif + +#define BUF ((struct uip_tcpip_hdr *)&uip_buf[0]) +#define PROXY_TIMER_TB 100 // default 10ms (was 1000=100ms, faster for tests) + +// Globals (uIP is a singleton) +static uip_ipaddr_t g_pkt_dest_ip; +static uint16_t g_pkt_dest_port; +static struct tcp_proxy* g_tcp_proxy = NULL; +static int g_timer_period_tb = PROXY_TIMER_TB; // can be shortened for tests + +// ==================================================================== +// Forward declarations +// ==================================================================== +static int sock_transport_send(struct tcp_proxy_transport* t, const uint8_t* data, size_t len); +static void sock_transport_close(struct tcp_proxy_transport* t); +static void sock_transport_destroy(struct tcp_proxy_transport* t); +static void sock_transport_read_callback(socket_t sock, void* arg); +static void sock_transport_write_callback(socket_t sock, void* arg); +static void sock_transport_error_callback(socket_t sock, void* arg); +static int etcp_transport_send(struct tcp_proxy_transport* t, const uint8_t* data, size_t len); +static void etcp_transport_close(struct tcp_proxy_transport* t); +static void etcp_transport_destroy(struct tcp_proxy_transport* t); + +static struct tcp_proxy_transport_ops sock_transport_ops = { + .send = sock_transport_send, .close = sock_transport_close, .destroy = sock_transport_destroy, +}; +static struct tcp_proxy_transport_ops etcp_transport_ops = { + .send = etcp_transport_send, .close = etcp_transport_close, .destroy = etcp_transport_destroy, +}; + +struct sock_transport { + struct tcp_proxy_transport base; + socket_t sock; + void* read_id; + void* write_id; + struct proxy_conn* conn; + struct UASYNC* ua; + int connected; + int connect_called; +}; + +struct etcp_transport { + struct tcp_proxy_transport base; + struct proxy_conn* conn; + struct UTUN_INSTANCE* inst; + uint64_t remote_node_id; + uint64_t stream_id; + int connected; +}; + +// ==================================================================== +// Mapping +// ==================================================================== +static struct tcp_proxy_mapping* find_mapping_by_port(struct tcp_proxy* p, uint16_t port_net) { + struct tcp_proxy_mapping* m; + for(m = p->mappings; m; m = m->next) if(m->local_port == port_net) return m; + return NULL; +} + +// ==================================================================== +// ll_entry helper +// ==================================================================== +static struct ll_entry* entry_from_data(struct memory_pool* pool, const uint8_t* data, uint16_t len) { + struct ll_entry* e = queue_entry_new_from_pool(pool); + if(!e) { DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "entry_from_data: pool exhausted"); return NULL; } + e->len = 0; e->dgram = NULL; + if(len > 0) { + uint8_t* buf = u_malloc(len); + if(!buf) { queue_entry_free(e); return NULL; } + memcpy(buf, data, len); e->dgram = buf; e->len = len; + } + return e; +} + +// ==================================================================== +// Output: write IP packet to TUN or raw fd +// ==================================================================== +static inline void tcp_proxy_write_output(struct tcp_proxy* p) { + if(uip_len == 0) return; + if(p->tun) tun_platform_write(p->tun, uip_buf, uip_len); + else write(p->ip_fd, uip_buf, uip_len); +} + +// ==================================================================== +// OS socket transport (passive connections) +// ==================================================================== +static struct sock_transport* sock_transport_create(struct proxy_conn* pc, struct UASYNC* ua) { + struct sock_transport* st = u_calloc(1, sizeof(struct sock_transport)); + if(!st) return NULL; + st->base.ops = &sock_transport_ops; st->ua = ua; st->conn = pc; + st->sock = SOCKET_INVALID; + st->sock = socket(AF_INET, SOCK_STREAM, 0); + if(st->sock == SOCKET_INVALID) { DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "socket() failed: %s", strerror(errno)); u_free(st); return NULL; } + socket_set_nonblocking(st->sock); + // Register all callbacks at once (uasync_add_socket_t can only be called once per fd) + st->read_id = uasync_add_socket_t(ua, st->sock, sock_transport_read_callback, sock_transport_write_callback, sock_transport_error_callback, st); + if(!st->read_id) { DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "uasync_add_socket_t failed"); socket_close_wrapper(st->sock); u_free(st); return NULL; } + return st; +} + +static void sock_transport_connect(struct sock_transport* st) { + struct proxy_conn* pc = st->conn; + struct sockaddr_in addr; memset(&addr, 0, sizeof(addr)); + addr.sin_family = AF_INET; memcpy(&addr.sin_addr.s_addr, pc->dest_ip, 4); addr.sin_port = pc->dest_port; + DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "TCP proxy: connecting to %d.%d.%d.%d:%d", pc->dest_ip[0], pc->dest_ip[1], pc->dest_ip[2], pc->dest_ip[3], ntohs(pc->dest_port)); + st->connect_called = 1; + int ret = connect(st->sock, (struct sockaddr*)&addr, sizeof(addr)); + if(ret < 0 && errno != EINPROGRESS) { DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "connect() failed: %s", strerror(errno)); sock_transport_close(&st->base); return; } +} + +static int sock_transport_send(struct tcp_proxy_transport* t, const uint8_t* data, size_t len) { + struct sock_transport* st = (struct sock_transport*)t; + if(!st->connected) return -1; + ssize_t n = send(st->sock, data, len, MSG_NOSIGNAL); + if(n < 0) { if(errno == EAGAIN || errno == EWOULDBLOCK) return -1; DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "send() failed: %s", strerror(errno)); sock_transport_close(t); return -1; } + if((size_t)n < len) return -1; + return 0; +} + +static void sock_transport_flush(struct sock_transport* st) { + struct proxy_conn* pc = st->conn; + while(st->connected) { + struct ll_entry* e = queue_data_get(pc->uip_to_transport); + if(!e) { queue_resume_callback(pc->uip_to_transport); break; } + ssize_t n = send(st->sock, e->dgram, e->len, MSG_NOSIGNAL); + if(n < 0) { if(errno == EAGAIN || errno == EWOULDBLOCK) { queue_data_put_first(pc->uip_to_transport, e); queue_resume_callback(pc->uip_to_transport); return; } + DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "send() in flush failed: %s", strerror(errno)); sock_transport_close(&st->base); queue_dgram_free(e); queue_entry_free(e); queue_resume_callback(pc->uip_to_transport); return; } + if((size_t)n < e->len) { memmove(e->dgram, e->dgram + n, e->len - n); e->len -= n; queue_data_put_first(pc->uip_to_transport, e); queue_resume_callback(pc->uip_to_transport); return; } + queue_dgram_free(e); queue_entry_free(e); queue_resume_callback(pc->uip_to_transport); + } +} + +static void sock_transport_close(struct tcp_proxy_transport* t) { + struct sock_transport* st = (struct sock_transport*)t; + struct proxy_conn* pc = st->conn; + if(st->sock == SOCKET_INVALID) return; + DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "TCP proxy: closing transport to %d.%d.%d.%d:%d", pc->dest_ip[0], pc->dest_ip[1], pc->dest_ip[2], pc->dest_ip[3], ntohs(pc->dest_port)); + if(st->read_id) { uasync_remove_socket_t(st->ua, st->sock); st->read_id = NULL; } + if(st->write_id) { uasync_remove_socket_t(st->ua, st->sock); st->write_id = NULL; } + socket_close_wrapper(st->sock); st->sock = SOCKET_INVALID; st->connected = 0; + if(pc && !pc->closing) pc->closing = 1; +} + +static void sock_transport_destroy(struct tcp_proxy_transport* t) { + struct sock_transport* st = (struct sock_transport*)t; + if(st->sock != SOCKET_INVALID) sock_transport_close(t); + u_free(st); +} + +// ==================================================================== +// Socket callbacks +// ==================================================================== +static void sock_transport_read_callback(socket_t sock, void* arg) { + (void)sock; struct sock_transport* st = (struct sock_transport*)arg; struct proxy_conn* pc = st->conn; + uint8_t buf[8192]; ssize_t n = recv(st->sock, buf, sizeof(buf), 0); + if(n > 0) { size_t offset = 0; + while(offset < (size_t)n) { size_t chunk = n - offset; if(chunk > 1460) chunk = 1460; + struct ll_entry* e = entry_from_data(pc->proxy->entry_pool, buf + offset, chunk); + if(e) { queue_data_put(pc->transport_to_uip, e); offset += chunk; } else break; } + } else if(n == 0) { DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "TCP proxy: transport EOF"); sock_transport_close(&st->base); } + else if(errno != EAGAIN && errno != EWOULDBLOCK && errno != EINTR) { DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "recv() error: %s", strerror(errno)); sock_transport_close(&st->base); } +} + +static void sock_transport_write_callback(socket_t sock, void* arg) { + (void)sock; struct sock_transport* st = (struct sock_transport*)arg; + if(!st->connected) { + if(!st->connect_called) return; // connect not yet initiated + int err = 0; socklen_t len = sizeof(err); + if(getsockopt(st->sock, SOL_SOCKET, SO_ERROR, &err, &len) == 0 && err == 0) { st->connected = 1; struct proxy_conn* pc = st->conn; + DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "TCP proxy: transport connected"); + struct sockaddr_in local; socklen_t llen = sizeof(local); + if(getsockname(st->sock, (struct sockaddr*)&local, &llen) == 0 && local.sin_port != 0) { + struct tcp_proxy* p = pc->proxy; struct tcp_proxy_mapping* dm = u_calloc(1, sizeof(struct tcp_proxy_mapping)); + if(dm) { dm->local_port = local.sin_port; memcpy(dm->remote_ip, pc->dest_ip, 4); dm->remote_port = pc->dest_port; + dm->dynamic = 1; dm->created_tb = get_time_tb(); dm->next = p->mappings; p->mappings = dm; + pc->eim_mapping = dm; uip_listen(local.sin_port); + DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "EIM NAT: mapped port %d -> %d.%d.%d.%d:%d", + ntohs(local.sin_port), pc->dest_ip[0], pc->dest_ip[1], pc->dest_ip[2], pc->dest_ip[3], ntohs(pc->dest_port)); } + } + sock_transport_flush(st); + } else { DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "TCP proxy: connect failed: %s", err ? strerror(err) : "unknown"); sock_transport_close(&st->base); } + return; } + sock_transport_flush(st); +} + +static void sock_transport_error_callback(socket_t sock, void* arg) { + (void)sock; struct sock_transport* st = (struct sock_transport*)arg; + DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "TCP proxy: socket error"); sock_transport_close(&st->base); +} + +// ==================================================================== +// ETCP transport (remote proxy via etcp_router) +// ==================================================================== +static int etcp_transport_send(struct tcp_proxy_transport* t, const uint8_t* data, size_t len) { + struct etcp_transport* et = (struct etcp_transport*)t; + if (!et->connected) return -1; + struct ll_entry* e = queue_entry_new(0); + if (!e) return -1; + e->dgram = u_malloc(TCP_PROXY_HDR_SIZE + len); + if (!e->dgram) { queue_entry_free(e); return -1; } + e->dgram[0] = ETCP_ID_TCP_PROXY; e->dgram[1] = TCP_PROXY_SUBCMD_DATA; + memcpy(e->dgram + 2, &et->stream_id, 8); + if (len > 0) memcpy(e->dgram + TCP_PROXY_HDR_SIZE, data, len); + e->len = TCP_PROXY_HDR_SIZE + len; + int ret = etcp_route_send(et->inst, et->remote_node_id, e); + return ret; +} + +static void etcp_transport_close(struct tcp_proxy_transport* t) { + struct etcp_transport* et = (struct etcp_transport*)t; + if (!et->connected && et->stream_id == 0) return; + struct ll_entry* e = queue_entry_new(0); + if (e) { + e->dgram = u_malloc(TCP_PROXY_HDR_SIZE); + if (e->dgram) { e->dgram[0] = ETCP_ID_TCP_PROXY; e->dgram[1] = TCP_PROXY_SUBCMD_CLOSE; + memcpy(e->dgram + 2, &et->stream_id, 8); e->len = TCP_PROXY_HDR_SIZE; + etcp_route_send(et->inst, et->remote_node_id, e); } else queue_entry_free(e); + } + et->connected = 0; +} + +static void etcp_transport_destroy(struct tcp_proxy_transport* t) { + struct etcp_transport* et = (struct etcp_transport*)t; + if (et->connected) etcp_transport_close(t); + u_free(et); +} + +static struct etcp_transport* etcp_transport_create(struct proxy_conn* pc, struct UTUN_INSTANCE* inst, uint64_t remote_node_id) { + struct etcp_transport* et = u_calloc(1, sizeof(struct etcp_transport)); + if (!et) return NULL; + et->base.ops = &etcp_transport_ops; et->conn = pc; et->inst = inst; + et->remote_node_id = remote_node_id; + et->stream_id = ++g_tcp_proxy->next_stream_id; + pc->remote_stream_id = et->stream_id; + uint8_t conn_buf[6]; + memcpy(conn_buf, pc->dest_ip, 4); memcpy(conn_buf + 4, &pc->dest_port, 2); + struct ll_entry* e = queue_entry_new(0); + if (!e) { u_free(et); return NULL; } + e->dgram = u_malloc(TCP_PROXY_HDR_SIZE + 6); + if (!e->dgram) { queue_entry_free(e); u_free(et); return NULL; } + e->dgram[0] = ETCP_ID_TCP_PROXY; e->dgram[1] = TCP_PROXY_SUBCMD_CONNECT; + memcpy(e->dgram + 2, &et->stream_id, 8); + memcpy(e->dgram + TCP_PROXY_HDR_SIZE, conn_buf, 6); + e->len = TCP_PROXY_HDR_SIZE + 6; + DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "TCP proxy: remote connect stream=%016llx to %d.%d.%d.%d:%d via node %016llx", + (unsigned long long)et->stream_id, pc->dest_ip[0], pc->dest_ip[1], pc->dest_ip[2], pc->dest_ip[3], + ntohs(pc->dest_port), (unsigned long long)remote_node_id); + if (etcp_route_send(inst, remote_node_id, e) != 0) { u_free(et); return NULL; } + return et; +} + +// ==================================================================== +// Stream lookup helpers +// ==================================================================== +static struct proxy_conn* find_pc_by_stream(uint64_t stream_id) { + struct proxy_conn* pc; + for (pc = g_tcp_proxy->conns; pc; pc = pc->next) if (pc->remote_stream_id == stream_id) return pc; + return NULL; +} + +static void handle_connected(uint64_t stream_id, struct ll_entry* entry) { + struct proxy_conn* pc = find_pc_by_stream(stream_id); + if (!pc || !pc->transport) { + queue_entry_free(entry); queue_dgram_free(entry); + return; + } + struct etcp_transport* et = (struct etcp_transport*)pc->transport; + if (et->base.ops != &etcp_transport_ops) { queue_entry_free(entry); queue_dgram_free(entry); return; } + if (entry->len >= TCP_PROXY_CONNECTED_HDR_SIZE) { + uint8_t status = entry->dgram[TCP_PROXY_HDR_SIZE + 2]; + if (status == TCP_PROXY_CONNECTED_OK) { et->connected = 1; + DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "TCP proxy: remote connected stream=%016llx", (unsigned long long)stream_id); } + else { DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "TCP proxy: remote refused stream=%016llx", (unsigned long long)stream_id); + pc->closing = 1; } + } + queue_entry_free(entry); queue_dgram_free(entry); +} + +// ==================================================================== +// Unified etcp_router handler (dispatches to tcp_proxy or remote_proxy) +// ==================================================================== +void tcp_proxy_etcp_recv_cb(struct ETCP_CONN* conn, struct ll_entry* entry) { + if (!entry || !entry->dgram || entry->len < TCP_PROXY_HDR_SIZE) { + if (entry) { queue_entry_free(entry); queue_dgram_free(entry); } + return; + } + uint8_t subcmd = entry->dgram[1]; + uint64_t stream_id; memcpy(&stream_id, entry->dgram + 2, 8); + struct UTUN_INSTANCE* inst = conn ? conn->instance : (g_tcp_proxy ? g_tcp_proxy->inst : NULL); + + // CONNECT → remote_proxy + if (subcmd == TCP_PROXY_SUBCMD_CONNECT) { + uint64_t src_node_id = conn ? conn->peer_node_id : (inst ? inst->node_id : 0); + remote_proxy_handle_connect(inst, entry, stream_id, src_node_id); + return; + } + // DATA / CLOSE for remote_proxy (check if stream belongs to remote_proxy) + if (inst && inst->remote_proxy.enabled) { + struct remote_proxy_conn* rc = remote_proxy_find_conn(&inst->remote_proxy, stream_id); + if (rc) { + if (subcmd == TCP_PROXY_SUBCMD_DATA) { remote_proxy_handle_data(inst, entry, stream_id); return; } + if (subcmd == TCP_PROXY_SUBCMD_CLOSE) { remote_proxy_handle_close(inst, stream_id); queue_entry_free(entry); queue_dgram_free(entry); return; } + } + } + // CONNECTED / DATA / CLOSE for tcp_proxy + if (g_tcp_proxy) { + if (subcmd == TCP_PROXY_SUBCMD_CONNECTED) { handle_connected(stream_id, entry); return; } + if (subcmd == TCP_PROXY_SUBCMD_DATA) { + struct proxy_conn* pc = find_pc_by_stream(stream_id); + if (pc) { + size_t data_len = entry->len - TCP_PROXY_HDR_SIZE; + if (data_len > 0) { + struct ll_entry* e = entry_from_data(pc->proxy->entry_pool, entry->dgram + TCP_PROXY_HDR_SIZE, (uint16_t)data_len); + if (e) queue_data_put(pc->transport_to_uip, e); + } + } + } + if (subcmd == TCP_PROXY_SUBCMD_CLOSE) { + struct proxy_conn* pc = find_pc_by_stream(stream_id); + if (pc) pc->closing = 1; + } + } + queue_entry_free(entry); queue_dgram_free(entry); +} + +// ==================================================================== +// uIP application callback +// ==================================================================== +void tcp_proxy_appcall(void) +{ + struct proxy_conn* pc = (struct proxy_conn*)uip_conn->appstate; + + if(uip_connected()) { + if(!pc) { + pc = u_calloc(1, sizeof(struct proxy_conn)); + if(!pc) { uip_abort(); return; } + pc->proxy = g_tcp_proxy; uip_conn->appstate = pc; pc->uip_conn = uip_conn; pc->closing = 0; pc->half_closed = 0; + + // Check if this was an active open (uip_connect): appstate was pre-set + if(uip_conn->appstate && ((struct proxy_conn*)uip_conn->appstate)->active) { + pc = (struct proxy_conn*)uip_conn->appstate; + pc->uip_conn = uip_conn; uip_conn->appstate = pc; + } else { + // Passive: look up mapping or use wildcard saved data + struct tcp_proxy_mapping* m = find_mapping_by_port(pc->proxy, uip_conn->lport); + if(m) { memcpy(pc->dest_ip, m->remote_ip, 4); pc->dest_port = m->remote_port; } + else { memcpy(pc->dest_ip, g_pkt_dest_ip, 4); pc->dest_port = g_pkt_dest_port; } + // tun_ip = IP from incoming packet dest (what client connected to) + memcpy(pc->tun_ip, g_pkt_dest_ip, 4); + } + DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "TCP proxy: new %s conn to %d.%d.%d.%d:%d (client port %d)", + pc->active ? "active" : "passive", + pc->dest_ip[0], pc->dest_ip[1], pc->dest_ip[2], pc->dest_ip[3], ntohs(pc->dest_port), ntohs(uip_conn->rport)); + } + } + + if(!pc) { DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "tcp_proxy_appcall: no proxy_conn"); uip_abort(); return; } + + // --- Passive connection setup --- + if(uip_connected() && !pc->active) { + if(!pc->uip_to_transport) { + pc->uip_to_transport = queue_new(pc->proxy->ua, 0, "uip_to_transport"); + pc->transport_to_uip = queue_new(pc->proxy->ua, 0, "transport_to_uip"); + if(!pc->uip_to_transport || !pc->transport_to_uip) { DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "Failed to create queues"); uip_abort(); return; } + struct tcp_proxy_mapping* m = find_mapping_by_port(pc->proxy, uip_conn->lport); + uint64_t via_node = m ? m->via_node_id : 0; + if(via_node != 0 && g_tcp_proxy && g_tcp_proxy->inst) { + struct etcp_transport* et = etcp_transport_create(pc, g_tcp_proxy->inst, via_node); + if(!et) { uip_abort(); return; } + pc->transport = &et->base; + } else { + struct sock_transport* st = sock_transport_create(pc, pc->proxy->ua); + if(!st) { uip_abort(); return; } + pc->transport = &st->base; st->conn = pc; + sock_transport_connect(st); + } + pc->next = pc->proxy->conns; pc->proxy->conns = pc; pc->proxy->conn_count++; + } + return; + } + + // --- Active connection: established --- + if(uip_connected() && pc->active) { + pc->next = pc->proxy->conns; pc->proxy->conns = pc; pc->proxy->conn_count++; + return; + } + + // --- Close / abort --- + if(uip_aborted() || uip_timedout()) { + DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "TCP proxy: uIP aborted/timedout"); + if(!pc->closing) { pc->closing = 1; if(pc->transport) pc->transport->ops->close(pc->transport); } + return; + } + + if(uip_closed()) { + if(pc->active) { pc->closing = 1; return; } + // Passive: client sent FIN — flush queued data, then signal EOF to transport + DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "TCP proxy: uIP half-closed (client FIN)"); + pc->half_closed = 1; + if(pc->transport) { + if(pc->transport->ops == &sock_transport_ops) { + struct sock_transport* st = (struct sock_transport*)pc->transport; + if(st->sock != SOCKET_INVALID) { + sock_transport_flush(st); + shutdown(st->sock, SHUT_WR); + } + } else if(pc->transport->ops == &etcp_transport_ops) { + pc->transport->ops->close(pc->transport); + } + } + return; + } + + // --- New data from remote --- + if(uip_newdata()) { + uint16_t dlen = uip_datalen(); if(dlen == 0) return; + if(pc->active) { + // Active: data from remote → queue for test to read + struct ll_entry* e = entry_from_data(pc->proxy->entry_pool, uip_appdata, dlen); + if(e) queue_data_put(pc->uip_to_transport, e); + } else if(pc->transport && !pc->half_closed) { + // Passive: data from client → forward to transport + int ret = pc->transport->ops->send(pc->transport, uip_appdata, dlen); + if(ret < 0) { + struct ll_entry* e = entry_from_data(pc->proxy->entry_pool, uip_appdata, dlen); + if(e) { queue_data_put(pc->uip_to_transport, e); } + } + } + return; + } + + // --- Retransmit --- + if(uip_rexmit()) { if(pc->rexmit_len > 0) uip_send(pc->rexmit_buf, pc->rexmit_len); return; } + + // --- ACK --- + if(uip_acked()) { pc->rexmit_len = 0; } + + // --- Poll: send data from queue to uIP --- + if(uip_poll() || uip_acked()) { + if(!uip_outstanding(uip_conn) && pc->transport_to_uip) { + struct ll_entry* e = queue_data_get(pc->transport_to_uip); + if(e && e->len > 0) { + uint16_t mss = uip_mss(); uint16_t send_len = (e->len < mss) ? e->len : mss; + memcpy(pc->rexmit_buf, e->dgram, send_len); pc->rexmit_len = send_len; + uip_send(pc->rexmit_buf, send_len); + if(e->len <= send_len) { queue_dgram_free(e); queue_entry_free(e); } + else { memmove(e->dgram, e->dgram + send_len, e->len - send_len); e->len -= send_len; queue_data_put_first(pc->transport_to_uip, e); } + queue_resume_callback(pc->transport_to_uip); + } else { + queue_resume_callback(pc->transport_to_uip); + if(pc->closing && !uip_outstanding(uip_conn) && !pc->close_sent) { uip_close(); pc->close_sent = 1; } + } + } else if(pc->closing && !uip_outstanding(uip_conn) && !pc->close_sent) { uip_close(); pc->close_sent = 1; } + return; + } +} + +// ==================================================================== +// Raw-fd IP packet input (no TUN) +// ==================================================================== +static void tcp_proxy_raw_read(int fd, void* arg) { + (void)fd; struct tcp_proxy* p = (struct tcp_proxy*)arg; + uint8_t buf[UIP_BUFSIZE]; ssize_t n = read(p->ip_fd, buf, sizeof(buf)); + if(n <= 0) return; + memcpy(uip_buf, buf, n); uip_len = n; + uip_ipaddr_copy(g_pkt_dest_ip, BUF->destipaddr); g_pkt_dest_port = BUF->destport; + uip_ipaddr_copy(uip_hostaddr, BUF->destipaddr); + uip_input(); + tcp_proxy_write_output(p); +} + +// ==================================================================== +// TUN input (when use_tun=1) +// ==================================================================== +static void tcp_proxy_tun_input(struct ll_queue* q, void* arg) { + struct tcp_proxy* p = (struct tcp_proxy*)arg; + struct ll_entry* entry = queue_data_get(q); + if(!entry) return; + if(entry->dgram && entry->len > 1) { + uint8_t* ip = entry->dgram + 1; size_t len = entry->len - 1; + if(len > 0 && len <= UIP_BUFSIZE) { + memcpy(uip_buf, ip, len); uip_len = len; + uip_ipaddr_copy(g_pkt_dest_ip, BUF->destipaddr); g_pkt_dest_port = BUF->destport; + uip_ipaddr_copy(uip_hostaddr, BUF->destipaddr); + uip_input(); + tcp_proxy_write_output(p); + } + } + queue_dgram_free(entry); queue_entry_free(entry); queue_resume_callback(q); +} + +// ==================================================================== +// Periodic timer: drive uIP, cleanup, expire EIM +// ==================================================================== +static void tcp_proxy_periodic(void* arg) { + struct tcp_proxy* p = (struct tcp_proxy*)arg; + int i; + for(i = 0; i < UIP_CONNS; i++) { + struct proxy_conn* pc = (struct proxy_conn*)uip_conns[i].appstate; + if(pc) uip_ipaddr_copy(uip_hostaddr, pc->tun_ip); + uip_periodic(i); + tcp_proxy_write_output(p); + } + + // Clean up fully closed connections + struct proxy_conn** prev = &p->conns; struct proxy_conn* pc = p->conns; + while(pc) { struct uip_conn* uc = (struct uip_conn*)pc->uip_conn; + if(pc->closing && uc && uc->tcpstateflags == UIP_CLOSED) { + *prev = pc->next; p->conn_count--; + if(uc) uc->appstate = NULL; + DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "TCP proxy: freeing %s conn to %d.%d.%d.%d:%d", + pc->active ? "active" : "passive", pc->dest_ip[0], pc->dest_ip[1], pc->dest_ip[2], pc->dest_ip[3], ntohs(pc->dest_port)); + if(pc->eim_mapping) { pc->eim_mapping->delete_at_tb = get_time_tb() + p->eim_timeout_tb; pc->eim_mapping = NULL; } + if(pc->transport) pc->transport->ops->destroy(pc->transport); + if(pc->uip_to_transport) { struct ll_entry* e; while((e = queue_data_get(pc->uip_to_transport)) != NULL) { queue_dgram_free(e); queue_entry_free(e); } queue_free(pc->uip_to_transport); } + if(pc->transport_to_uip) { struct ll_entry* e; while((e = queue_data_get(pc->transport_to_uip)) != NULL) { queue_dgram_free(e); queue_entry_free(e); } queue_free(pc->transport_to_uip); } + u_free(pc); pc = *prev; + } else { prev = &pc->next; pc = pc->next; } + } + + // Expire EIM dynamic mappings + uint64_t now = get_time_tb(); struct tcp_proxy_mapping** mp = &p->mappings; + while(*mp) { struct tcp_proxy_mapping* m = *mp; + if(m->dynamic && m->delete_at_tb > 0 && now >= m->delete_at_tb) { + *mp = m->next; uip_unlisten(m->local_port); + DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "EIM NAT: expired mapping port %d", ntohs(m->local_port)); u_free(m); + } else mp = &m->next; + } + + p->uip_timer_id = uasync_set_timeout(p->ua, g_timer_period_tb, p, tcp_proxy_periodic, "uip_periodic"); +} + +// ==================================================================== +// Public API +// ==================================================================== +struct tcp_proxy* tcp_proxy_create(struct UTUN_INSTANCE* inst, struct UASYNC* ua, + const char* tun_name, const char* tun_ip, int mtu, int test_mode, + struct tcp_proxy_mapping_config* mappings, int mapping_count, + int eim_timeout_sec, int use_tun, int ip_fd) +{ + if(!ua) return NULL; + + struct tcp_proxy* p = u_calloc(1, sizeof(struct tcp_proxy)); + if(!p) return NULL; + p->inst = inst; p->ua = ua; p->ip_fd = -1; p->next_stream_id = 1; + p->eim_timeout_tb = eim_timeout_sec > 0 ? eim_timeout_sec * 10000 : 300000; + + p->entry_pool = memory_pool_init(sizeof(struct ll_entry)); + if(!p->entry_pool) { u_free(p); return NULL; } + + if(use_tun) { + if(!tun_name || !tun_ip) { memory_pool_destroy(p->entry_pool); u_free(p); return NULL; } + p->tun = tun_init_nat(ua, tun_name, tun_ip, mtu > 0 ? mtu : 1500, test_mode); + if(!p->tun) { DEBUG_ERROR(DEBUG_CATEGORY_TUN, "tcp_proxy: failed to create TUN %s", tun_name); memory_pool_destroy(p->entry_pool); u_free(p); return NULL; } + queue_set_callback(p->tun->output_queue, tcp_proxy_tun_input, p); + } else { + p->ip_fd = ip_fd; + p->ip_fd_id = uasync_add_socket(ua, ip_fd, tcp_proxy_raw_read, NULL, NULL, p); + if(!p->ip_fd_id) { DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "tcp_proxy: uasync_add_socket for ip_fd failed"); memory_pool_destroy(p->entry_pool); u_free(p); return NULL; } + } + + uip_init(); + uip_ipaddr_t addr; uip_ipaddr(addr, 127,0,0,1); uip_sethostaddr(addr); + + if(mapping_count > 0) { + uip_set_promiscuous(0); uip_listen_all(0); + int j; for(j = 0; j < mapping_count; j++) { + uint16_t port_net = htons(mappings[j].local_port); uip_listen(port_net); + struct tcp_proxy_mapping* m = u_calloc(1, sizeof(struct tcp_proxy_mapping)); + if(m) { m->local_port = port_net; struct in_addr ra; ra.s_addr = inet_addr(mappings[j].remote_ip); + memcpy(m->remote_ip, &ra.s_addr, 4); m->remote_port = htons(mappings[j].remote_port); m->dynamic = 0; + m->via_node_id = mappings[j].via_node_id; + if(m->via_node_id != 0) p->has_remote_mappings = 1; + m->next = p->mappings; p->mappings = m; + DEBUG_INFO(DEBUG_CATEGORY_TUN, "TCP proxy mapping: %d -> %s:%d %s", mappings[j].local_port, mappings[j].remote_ip, mappings[j].remote_port, + m->via_node_id ? "(remote)" : ""); } + } + } else { uip_set_promiscuous(1); uip_listen_all(1); } + + g_tcp_proxy = p; + p->uip_timer_id = uasync_set_timeout(ua, g_timer_period_tb, p, tcp_proxy_periodic, "uip_periodic"); + + // Register etcp_router handler if we have remote proxy mappings + if(p->has_remote_mappings && inst) { + if(etcp_router_bind(inst, ETCP_ID_TCP_PROXY, tcp_proxy_etcp_recv_cb) != 0) { + DEBUG_WARN(DEBUG_CATEGORY_SOCKET, "TCP proxy: etcp_router_bind for remote proxy failed"); + p->has_remote_mappings = 0; + } else { + DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "TCP proxy: etcp_router bind registered for ID=0x%02x", ETCP_ID_TCP_PROXY); + } + } + + DEBUG_INFO(DEBUG_CATEGORY_TUN, "TCP proxy created: use_tun=%d mappings=%d remote=%d", use_tun, mapping_count, p->has_remote_mappings); + return p; +} + +void tcp_proxy_destroy(struct tcp_proxy* p) { + if(!p) return; + if(p->has_remote_mappings && p->inst) etcp_router_unbind(p->inst, ETCP_ID_TCP_PROXY); + if(p->uip_timer_id) { uasync_cancel_timeout(p->ua, p->uip_timer_id); p->uip_timer_id = NULL; } + struct proxy_conn* pc = p->conns; + while(pc) { struct proxy_conn* next = pc->next; + if(pc->transport) pc->transport->ops->destroy(pc->transport); + if(pc->uip_to_transport) { struct ll_entry* e; while((e = queue_data_get(pc->uip_to_transport)) != NULL) { queue_dgram_free(e); queue_entry_free(e); } queue_free(pc->uip_to_transport); } + if(pc->transport_to_uip) { struct ll_entry* e; while((e = queue_data_get(pc->transport_to_uip)) != NULL) { queue_dgram_free(e); queue_entry_free(e); } queue_free(pc->transport_to_uip); } + if(pc->uip_conn) ((struct uip_conn*)pc->uip_conn)->appstate = NULL; + u_free(pc); pc = next; } + struct tcp_proxy_mapping* m = p->mappings; + while(m) { struct tcp_proxy_mapping* next = m->next; if(m->dynamic) uip_unlisten(m->local_port); u_free(m); m = next; } + if(p->ip_fd_id) { uasync_remove_socket(p->ua, p->ip_fd_id); p->ip_fd_id = NULL; } + if(p->tun) tun_close(p->tun); + if(p->entry_pool) memory_pool_destroy(p->entry_pool); + if(g_tcp_proxy == p) g_tcp_proxy = NULL; + u_free(p); +} + +// ==================================================================== +// Active connections API +// ==================================================================== +int tcp_proxy_active_open(struct tcp_proxy* p, const char* dest_ip, uint16_t dest_port) { + if(!p || !dest_ip) return -1; + uip_ipaddr_t addr; + unsigned int a0, a1, a2, a3; + if(sscanf(dest_ip, "%u.%u.%u.%u", &a0, &a1, &a2, &a3) != 4) return -1; + uip_ipaddr(&addr, a0, a1, a2, a3); + + struct proxy_conn* pc = u_calloc(1, sizeof(struct proxy_conn)); + if(!pc) return -1; + pc->proxy = p; pc->active = 1; pc->closing = 0; pc->half_closed = 0; + pc->dest_port = htons(dest_port); + pc->dest_ip[0] = a0; pc->dest_ip[1] = a1; pc->dest_ip[2] = a2; pc->dest_ip[3] = a3; + pc->tun_ip[0] = a0; pc->tun_ip[1] = a1; pc->tun_ip[2] = a2; pc->tun_ip[3] = a3; + + struct uip_conn* conn = uip_connect(&addr, htons(dest_port)); + if(!conn) { u_free(pc); return -1; } + pc->uip_conn = conn; conn->appstate = pc; + + pc->uip_to_transport = queue_new(p->ua, 0, "active_uip_to"); + pc->transport_to_uip = queue_new(p->ua, 0, "active_to_uip"); + if(!pc->uip_to_transport || !pc->transport_to_uip) { u_free(pc); return -1; } + + DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "TCP proxy: active_open to %s:%d (conn_idx=%ld)", dest_ip, dest_port, conn - uip_conns); + return (int)(conn - uip_conns); +} + +int tcp_proxy_active_send(struct tcp_proxy* p, int conn_idx, const uint8_t* data, size_t len) { + if(!p || conn_idx < 0 || conn_idx >= UIP_CONNS || !data || len == 0) return -1; + struct uip_conn* uc = &uip_conns[conn_idx]; + struct proxy_conn* pc = (struct proxy_conn*)uc->appstate; + if(!pc || !pc->active) return -1; + struct ll_entry* e = entry_from_data(p->entry_pool, data, (uint16_t)(len > 65535 ? 65535 : len)); + if(!e) return -1; + queue_data_put(pc->transport_to_uip, e); + return 0; +} + +ssize_t tcp_proxy_active_recv(struct tcp_proxy* p, int conn_idx, uint8_t* buf, size_t len) { + if(!p || conn_idx < 0 || conn_idx >= UIP_CONNS || !buf || len == 0) return -1; + struct uip_conn* uc = &uip_conns[conn_idx]; + struct proxy_conn* pc = (struct proxy_conn*)uc->appstate; + if(!pc || !pc->active) return -1; + struct ll_entry* e = queue_data_get(pc->uip_to_transport); + if(!e) { queue_resume_callback(pc->uip_to_transport); return 0; } + size_t copylen = e->len < len ? e->len : len; + memcpy(buf, e->dgram, copylen); + queue_dgram_free(e); queue_entry_free(e); queue_resume_callback(pc->uip_to_transport); + return copylen; +} + +int tcp_proxy_active_close(struct tcp_proxy* p, int conn_idx) { + if(!p || conn_idx < 0 || conn_idx >= UIP_CONNS) return -1; + struct uip_conn* uc = &uip_conns[conn_idx]; + struct proxy_conn* pc = (struct proxy_conn*)uc->appstate; + if(!pc || !pc->active) return -1; + pc->closing = 1; + pc->close_sent = 1; + uip_close(); + return 0; +} + +int tcp_proxy_active_send_done(struct tcp_proxy* p, int conn_idx) { + if(!p || conn_idx < 0 || conn_idx >= UIP_CONNS) return 1; + struct uip_conn* uc = &uip_conns[conn_idx]; + struct proxy_conn* pc = (struct proxy_conn*)uc->appstate; + if(!pc || !pc->active) return 1; + if(pc->closing) return 1; + // Check if there's data pending in the send queue or outstanding in uIP + if(uc->len > 0) return 0; // outstanding data + if(queue_entry_count(pc->transport_to_uip) > 0) return 0; + return 1; +} diff --git a/src/tcp_proxy.h b/src/tcp_proxy.h new file mode 100644 index 00000000..7b67e90f --- /dev/null +++ b/src/tcp_proxy.h @@ -0,0 +1,104 @@ +// tcp_proxy.h — TCP proxy: uIP TCP ↔ ll_queue ↔ OS socket → destination +// Modes: TUN (use_tun=1) or raw-fd (use_tun=0, socketpair etc.) +// Supports port forwarding + EIM NAT + half-close + active connections + remote proxy +#ifndef TCP_PROXY_H +#define TCP_PROXY_H + +#include +#include +#include "../lib/socket_compat.h" + +struct UASYNC; +struct UTUN_INSTANCE; +struct ETCP_CONN; +struct tun_if; +struct ll_queue; +struct ll_entry; +struct memory_pool; +struct tcp_proxy_mapping_config; +struct remote_proxy_ctx; + +// Transport abstraction (local OS socket or remote via ETCP) +struct tcp_proxy_transport; +struct tcp_proxy_transport_ops { + int (*send)(struct tcp_proxy_transport* t, const uint8_t* data, size_t len); + void (*close)(struct tcp_proxy_transport* t); + void (*destroy)(struct tcp_proxy_transport* t); +}; +struct tcp_proxy_transport { + struct tcp_proxy_transport_ops* ops; +}; + +// Port mapping (static from config or dynamic from EIM NAT) +struct tcp_proxy_mapping { + struct tcp_proxy_mapping* next; + uint16_t local_port; + uint8_t remote_ip[4]; + uint16_t remote_port; + int dynamic; + uint64_t created_tb; + uint64_t delete_at_tb; + uint64_t via_node_id; // 0=local proxy, !=0=remote via this node +}; + +// One proxy connection: uIP TCP ↔ ll_queues ↔ transport ↔ destination +struct proxy_conn { + struct proxy_conn* next; + struct tcp_proxy* proxy; + void* uip_conn; + struct tcp_proxy_transport* transport; + struct ll_queue* uip_to_transport; // data: uIP → transport (passive) or remote→test (active) + struct ll_queue* transport_to_uip; // data: transport → uIP (passive) or test→remote (active) + uint8_t dest_ip[4]; // proxy target IP (where to forward OS socket) + uint8_t tun_ip[4]; // TUN IP (what client connected to, used for uip_hostaddr) + uint16_t dest_port; + uint16_t rexmit_len; + uint8_t rexmit_buf[1500]; + int closing; + int half_closed; + int active; // 1 = outgoing connection (uip_connect) + int close_sent; // 1 = uip_close() already called + struct tcp_proxy_mapping* eim_mapping; + uint64_t remote_stream_id; // stream ID for remote proxy (0=local) +}; + +// Main TCP proxy module +struct tcp_proxy { + struct UTUN_INSTANCE* inst; + struct UASYNC* ua; + struct tun_if* tun; // TUN interface (NULL in raw-fd mode) + int ip_fd; // raw IP fd (only when tun==NULL) + void* ip_fd_id; // uasync socket id for ip_fd + struct proxy_conn* conns; + int conn_count; + void* uip_timer_id; + struct memory_pool* entry_pool; + struct tcp_proxy_mapping* mappings; + int eim_timeout_tb; + uint64_t next_stream_id; // counter for remote proxy stream IDs + int has_remote_mappings; // 1=at least one mapping uses via_node_id +}; + +// ========== API ========== + +// Create TCP proxy +// use_tun=1: TUN mode (needs root), uses tun_name/tun_ip/mtu/test_mode, ip_fd ignored +// use_tun=0: raw-fd mode, uses ip_fd for IP packet I/O, TUN params ignored +struct tcp_proxy* tcp_proxy_create(struct UTUN_INSTANCE* inst, struct UASYNC* ua, + const char* tun_name, const char* tun_ip, int mtu, int test_mode, + struct tcp_proxy_mapping_config* mappings, int mapping_count, + int eim_timeout_sec, int use_tun, int ip_fd); + +void tcp_proxy_destroy(struct tcp_proxy* p); + +// Active (outgoing) connections — for test client or EIM NAT +int tcp_proxy_active_open (struct tcp_proxy* p, const char* dest_ip, uint16_t dest_port); +int tcp_proxy_active_send (struct tcp_proxy* p, int conn_idx, const uint8_t* data, size_t len); +ssize_t tcp_proxy_active_recv (struct tcp_proxy* p, int conn_idx, uint8_t* buf, size_t len); +int tcp_proxy_active_close(struct tcp_proxy* p, int conn_idx); +int tcp_proxy_active_send_done(struct tcp_proxy* p, int conn_idx); // 1=all queued data sent + +// etcp_router handler for incoming proxy messages (called from etcp_router dispatch) +void tcp_proxy_etcp_recv_cb(struct ETCP_CONN* conn, struct ll_entry* entry); + +#endif // TCP_PROXY_H diff --git a/src/utun_instance.c b/src/utun_instance.c index 159e4c8a..eaf9c9cf 100644 --- a/src/utun_instance.c +++ b/src/utun_instance.c @@ -1,5 +1,6 @@ // utun_instance.c - Root instance implementation #include "utun_instance.h" +#include "etcp_router.h" #include "config_parser.h" #include "config_updater.h" #include "tun_if.h" @@ -128,7 +129,38 @@ static int instance_init_common(struct UTUN_INSTANCE* instance, struct UASYNC* u DEBUG_INFO(DEBUG_CATEGORY_CONFIG, "Firewall initialized: %d rules, bypass_all=%d", instance->fw.count, instance->fw.bypass_all); } - + + // etcp_router — сервисная маршрутизация (после BGP, до TCP proxy/NAT) + if (etcp_router_init(instance) != 0) { + DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "Failed to initialize etcp_router"); + return -1; + } + + // Remote proxy (exit node, optional) — must be before tcp_proxy so + // tcp_proxy_create can overwrite the handler if it has remote mappings + if (remote_proxy_init(instance) != 0) { + DEBUG_WARN(DEBUG_CATEGORY_SOCKET, "Failed to initialize remote_proxy (non-fatal)"); + } + + // TCP proxy (from [tcp_proxy] config section) + if (config->global.tcp_proxy_enabled) { + const char* tun_name = config->global.tcp_proxy_tun_name[0] ? config->global.tcp_proxy_tun_name : "tun_tcp"; + const char* tun_ip = config->global.tcp_proxy_tun_ip[0] ? config->global.tcp_proxy_tun_ip : "10.99.0.1"; + int mtu = config->global.tcp_proxy_mtu > 0 ? config->global.tcp_proxy_mtu : 1500; + instance->tcp_proxy = tcp_proxy_create(instance, ua, tun_name, tun_ip, mtu, 0, + config->global.tcp_proxy_mappings, config->global.tcp_proxy_mapping_count, + config->global.tcp_proxy_eim_timeout, 1, -1); + if (instance->tcp_proxy) { + DEBUG_INFO(DEBUG_CATEGORY_TUN, "TCP proxy enabled: TUN=%s IP=%s MTU=%d mappings=%d", + tun_name, tun_ip, mtu, config->global.tcp_proxy_mapping_count); + } else { + DEBUG_ERROR(DEBUG_CATEGORY_TUN, "Failed to create TCP proxy"); + return -1; + } + } else { + instance->tcp_proxy = NULL; + } + return 0; } @@ -283,7 +315,20 @@ void utun_instance_destroy(struct UTUN_INSTANCE *instance) { tun_close(instance->tun); instance->tun = NULL; } - + + // Cleanup TCP proxy module + if (instance->tcp_proxy) { + DEBUG_INFO(DEBUG_CATEGORY_TUN, "Destroying TCP proxy module"); + tcp_proxy_destroy(instance->tcp_proxy); + instance->tcp_proxy = NULL; + } + + // Cleanup remote proxy + remote_proxy_destroy(instance); + + // Cleanup etcp_router + etcp_router_destroy(instance); + // Cleanup routing module routing_destroy(instance); diff --git a/src/utun_instance.h b/src/utun_instance.h index 2c317540..d9840846 100644 --- a/src/utun_instance.h +++ b/src/utun_instance.h @@ -11,6 +11,9 @@ #include "firewall.h" #include "eim_nat.h" #include "nat_transport.h" +#include "tcp_proxy.h" +#include "etcp_router.h" +#include "remote_proxy.h" // Forward declarations struct utun_config; @@ -21,6 +24,7 @@ typedef void (*etcp_new_conn_fn)(struct ETCP_CONN* conn, void* arg); struct ETCP_SOCKET; struct tun_if; struct ETCP_BINDINGS; +struct ETCP_ROUTER_BINDINGS; struct ROUTE_BGP; struct control_server; struct PING_CONTEXT; @@ -92,6 +96,15 @@ struct UTUN_INSTANCE { // Socket initialization status: 0=OK, 1=partial (some sockets failed), -1=error (none created) int socket_init_status; + + // TCP proxy (optional, NULL if not enabled) + struct tcp_proxy* tcp_proxy; + + // etcp_router bindings (per-instance service routing) + struct ETCP_ROUTER_BINDINGS router_bindings; + + // Remote proxy (exit node) + struct remote_proxy_ctx remote_proxy; }; // Functions diff --git a/tests/Makefile.am b/tests/Makefile.am index a0194452..74c314e2 100644 --- a/tests/Makefile.am +++ b/tests/Makefile.am @@ -31,6 +31,9 @@ check_PROGRAMS = \ test_nat_transport \ test_nat_stress \ test_etcp_reinit_inflight \ + test_tcp_proxy \ + test_etcp_router \ + test_remote_proxy \ bench_timeout_heap \ bench_uasync_timeouts @@ -107,6 +110,10 @@ ETCP_FULL_OBJS = \ $(top_builddir)/src/utun-control_server.o \ $(TUN_PLATFORM_OBJ) \ $(top_builddir)/src/utun-utun_instance.o \ + $(top_builddir)/src/utun-tcp_proxy.o \ + $(top_builddir)/src/utun-etcp_router.o \ + $(top_builddir)/src/utun-remote_proxy.o \ + $(top_builddir)/src/uip/utun-uip.o \ $(ETCP_CORE_OBJS) # Windows-specific libraries (advapi32 for CryptGenRandom, ws2_32 for sockets) @@ -177,8 +184,15 @@ test_etcp_reinit_inflight_SOURCES = test_etcp_reinit_inflight.c test_etcp_reinit_inflight_CFLAGS = -I$(top_srcdir)/src -I$(top_srcdir)/lib -I$(top_srcdir)/tinycrypt/lib/include -I$(top_srcdir)/tinycrypt/lib/source test_etcp_reinit_inflight_LDADD = $(top_builddir)/src/utun-dummynet.o $(ETCP_FULL_OBJS) $(SECURE_CHANNEL_OBJS) $(CRYPTO_LIBS) $(COMMON_LIBS) -test_etcp_minimal_SOURCES = test_etcp_minimal.c -test_etcp_minimal_CFLAGS = -I$(top_srcdir)/src -I$(top_srcdir)/lib -I$(top_srcdir)/tinycrypt/lib/include -I$(top_srcdir)/tinycrypt/lib/source +test_tcp_proxy_SOURCES = test_tcp_proxy.c +test_tcp_proxy_LDADD = $(ETCP_FULL_OBJS) $(SECURE_CHANNEL_OBJS) $(CRYPTO_LIBS) $(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_remote_proxy_SOURCES = test_remote_proxy.c +test_remote_proxy_LDADD = $(ETCP_FULL_OBJS) $(SECURE_CHANNEL_OBJS) $(CRYPTO_LIBS) $(COMMON_LIBS) + test_etcp_minimal_LDADD = $(ETCP_FULL_OBJS) $(SECURE_CHANNEL_OBJS) $(CRYPTO_LIBS) $(COMMON_LIBS) test_etcp_100_packets_SOURCES = test_etcp_100_packets.c diff --git a/tests/test_etcp_router.c b/tests/test_etcp_router.c new file mode 100644 index 00000000..ba64e9e3 --- /dev/null +++ b/tests/test_etcp_router.c @@ -0,0 +1,340 @@ +// test_etcp_router.c — Unit test for etcp_router: bind/unbind, send, loopback, forward, reply +#include +#include +#include +#include "../lib/platform_compat.h" +#include "test_utils.h" +#ifdef _WIN32 +#include +#include +#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/config_parser.h" +#include "../src/utun_instance.h" +#include "../src/routing.h" +#include "../src/tun_if.h" +#include "../src/secure_channel.h" +#include "../src/pkt_normalizer.h" +#include "../src/route_bgp.h" +#include "../lib/u_async.h" +#include "../lib/ll_queue.h" +#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 + +static char temp_dir[] = "/tmp/utun_test_XXXXXX"; +static char server_conf[256], client_conf[256]; +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 void* g_mon_id = NULL; + +// Server config — node 0x1111... +static const char* srv_cfg = + "[global]\n" + "my_node_id=0x1111111111111111\n" + "my_private_key=67b705a92b41bcaae105af2d6a17743faa7b26ccebba8b3b9b0af05e9cd1d5fb\n" + "my_public_key=1c55e4ccae7c4470707759086738b10681bf88b81f198cc2ab54a647d1556e17c65e6b1833e0c771e5a39382c03067c388915a4c732191bc130480f20f8e00b9\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=4813d31d28b7e9829247f488c6be7672f2bdf61b2508333128e386d1759afed2\n" + "my_public_key=c594f33c91f3a2222795c2c110c527bf214ad1009197ce14556cb13df3c461b3c373bed8f205a8dd1fc0c364f90bf471d7c6f5db49564c33e4235d268569ac71\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=1c55e4ccae7c4470707759086738b10681bf88b81f198cc2ab54a647d1556e17c65e6b1833e0c771e5a39382c03067c388915a4c732191bc130480f20f8e00b9\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; + +// ======================== Helpers ======================== +static int conn_established(struct UTUN_INSTANCE* inst) { + if (!inst) return 0; + for (struct ETCP_CONN* c = inst->connections; c; c = c->next) { + 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; } + if (ok) return 1; + } + } + } + return 0; +} + +static struct ETCP_CONN* first_conn(struct UTUN_INSTANCE* inst) { + return inst ? inst->connections : NULL; +} + +// ======================== 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) { queue_entry_free(entry); queue_dgram_free(entry); } + 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 == 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++; + } + queue_entry_free(entry); queue_dgram_free(entry); + + // 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); reply_sent++; } + } + } else { + queue_entry_free(entry); queue_dgram_free(entry); + } +} + +// ======================== 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_entry_free(entry); queue_dgram_free(entry); } + 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++; + } + } + queue_entry_free(entry); queue_dgram_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) { + (void)conn; + if (entry && entry->dgram && entry->len >= 6) { + loop_rcvd++; + if (entry->dgram[1] == 0xAA && entry->dgram[2] == 0xBB && entry->dgram[3] == 0xCC) loop_ok = 1; + } + if (entry) { queue_entry_free(entry); queue_dgram_free(entry); } +} + +static int test_loopback(void) { + loop_rcvd = loop_ok = 0; + etcp_router_bind(srv, 0xF0, loop_handler); + 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); // loopback + etcp_router_unbind(srv, 0xF0); + if (loop_rcvd != 1 || !loop_ok) { + printf("[FAIL] loopback: rcvd=%d ok=%d\n", loop_rcvd, loop_ok); + return 1; + } + printf(" loopback: OK\n"); + 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"); + return 0; +} + +// ======================== Monitor & sender ======================== +static uint64_t dedup_bgp_log = 0; +static void monitor(void* arg) { + (void)arg; + if (g_test_done) { g_mon_id = NULL; return; } + + static int conn_ok = 0, conn_delay = 0; + if (!conn_ok) { + 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 + 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); + } + } + + // 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) fwd_sent++; + else { queue_entry_free(e); queue_dgram_free(e); } + } + } + + 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); + } + + if (g_phase == 2 && reply_rcvd >= TOTAL_PACKETS) { + printf(" reply phase done: rcvd=%d\n", reply_rcvd); + g_test_done = 1; + return; + } + + g_mon_id = uasync_set_timeout(ua, 10, NULL, monitor, "mon"); +} + +static void timeout(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_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); + + printf("=== test_etcp_router ===\n"); + + debug_config_init(); + debug_set_level(DEBUG_LEVEL_ERROR); + debug_set_categories(DEBUG_CATEGORY_ALL); + utun_instance_set_tun_init_enabled(0); + srand((unsigned)time(NULL)); + + ua = uasync_create(); + if (!ua) { printf("[FAIL] uasync_create\n"); goto done; } + + srv = utun_instance_create(ua, server_conf); + if (!srv || utun_instance_init(srv) < 0) { printf("[FAIL] server create\n"); goto done; } + + cli = utun_instance_create(ua, client_conf); + if (!cli || utun_instance_init(cli) < 0) { printf("[FAIL] client create\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"); + + while (!g_test_done) uasync_poll(ua, 100); + + 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 (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; +} diff --git a/tests/test_remote_proxy.c b/tests/test_remote_proxy.c new file mode 100644 index 00000000..98dfe359 --- /dev/null +++ b/tests/test_remote_proxy.c @@ -0,0 +1,175 @@ +// test_remote_proxy.c — Test remote_proxy: CONNECT → socket → echo → verify +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include "../lib/platform_compat.h" +#include "test_utils.h" +#include "../src/etcp.h" +#include "../src/etcp_api.h" +#include "../src/etcp_router.h" +#include "../src/remote_proxy.h" +#include "../src/config_parser.h" +#include "../src/utun_instance.h" +#include "../src/routing.h" +#include "../src/tun_if.h" +#include "../lib/u_async.h" +#include "../lib/ll_queue.h" +#include "../lib/debug_config.h" +#include "../lib/mem.h" + +#define ECHO_PORT 19991 +#define TEST_TIMEOUT_MS 8000 +#define PAYLOAD_SIZE 64 + +static char temp_dir[] = "/tmp/utun_test_XXXXXX"; +static char cfg_path[256]; +static struct UTUN_INSTANCE* inst = NULL; +static struct UASYNC* ua = NULL; +static int g_ok = 0, g_done = 0; +static void* g_mon_id = NULL; +static pid_t echo_pid = 0; +static uint8_t send_buf[PAYLOAD_SIZE], recv_buf[PAYLOAD_SIZE]; +static int connected_ok = 0; +static uint64_t stream_id = 1; + +static const char* cfg = + "[global]\n" + "my_node_id=0x1111111111111111\n" + "my_private_key=67b705a92b41bcaae105af2d6a17743faa7b26ccebba8b3b9b0af05e9cd1d5fb\n" + "my_public_key=1c55e4ccae7c4470707759086738b10681bf88b81f198cc2ab54a647d1556e17c65e6b1833e0c771e5a39382c03067c388915a4c732191bc130480f20f8e00b9\n" + "tun_ip=10.99.0.1/24\n" + "tun_ifname=tun99\n" + "[remote_proxy]\n" + "enabled=yes\n"; + +static void echo_server(void) { + int srv = socket(AF_INET, SOCK_STREAM, 0); + if (srv < 0) _exit(1); + int opt = 1; setsockopt(srv, SOL_SOCKET, SO_REUSEADDR, &opt, sizeof(opt)); + struct sockaddr_in addr = {.sin_family = AF_INET, .sin_port = htons(ECHO_PORT)}; + addr.sin_addr.s_addr = inet_addr("127.0.0.1"); + if (bind(srv, (struct sockaddr*)&addr, sizeof(addr)) < 0 || listen(srv, 1) < 0) { close(srv); _exit(1); } + int cli = accept(srv, NULL, NULL); + if (cli < 0) { close(srv); _exit(1); } + uint8_t buf[8192]; ssize_t n; + while ((n = recv(cli, buf, sizeof(buf), 0)) > 0) { + ssize_t sent = 0; + while (sent < n) { ssize_t s = send(cli, buf + sent, n - sent, 0); if (s < 0) goto done; sent += s; } + } +done: close(cli); close(srv); +} + +static void test_handler(struct ETCP_CONN* conn, struct ll_entry* entry) { + struct UTUN_INSTANCE* i = conn ? conn->instance : inst; + if (!i || !entry || entry->dgram == NULL || entry->len < TCP_PROXY_HDR_SIZE) { + if (entry) { queue_entry_free(entry); queue_dgram_free(entry); } return; + } + uint8_t subcmd = entry->dgram[1]; + uint64_t sid = 0; memcpy(&sid, entry->dgram + 2, 8); + uint64_t src = conn ? conn->peer_node_id : i->node_id; + if (subcmd == TCP_PROXY_SUBCMD_CONNECT && sid == stream_id) { remote_proxy_handle_connect(i, entry, sid, src); return; } + if (sid != stream_id) { queue_entry_free(entry); queue_dgram_free(entry); return; } + if (subcmd == TCP_PROXY_SUBCMD_CONNECTED) { + uint8_t status = (entry->len >= TCP_PROXY_CONNECTED_HDR_SIZE) ? entry->dgram[TCP_PROXY_HDR_SIZE + 2] : 1; + connected_ok = (status == TCP_PROXY_CONNECTED_OK) ? 1 : -1; + } + queue_entry_free(entry); queue_dgram_free(entry); +} + +static void monitor(void* arg) { + (void)arg; + if (g_done) { g_mon_id = NULL; return; } + static int phase = 0; + if (phase == 0) { + phase = 1; + etcp_router_bind(inst, ETCP_ID_TCP_PROXY, test_handler); + uint32_t ip = inet_addr("127.0.0.1"); + uint16_t port = htons(ECHO_PORT); + uint8_t payload[6]; memcpy(payload, &ip, 4); memcpy(payload + 4, &port, 2); + struct ll_entry* e = queue_entry_new(0); + if (e) { + e->dgram = u_malloc(TCP_PROXY_HDR_SIZE + 6); + e->dgram[0] = ETCP_ID_TCP_PROXY; e->dgram[1] = TCP_PROXY_SUBCMD_CONNECT; + memcpy(e->dgram + 2, &stream_id, 8); + memcpy(e->dgram + TCP_PROXY_HDR_SIZE, payload, 6); + e->len = TCP_PROXY_HDR_SIZE + 6; + etcp_route_send(inst, inst->node_id, e); + } + } + if (phase == 1 && (connected_ok == 1 || connected_ok == -1)) { + if (connected_ok != 1) { printf("[FAIL] connect refused\n"); g_done = -1; return; } + // Get proxy conn, detach from uasync so we can use socket directly + struct remote_proxy_conn* rc = remote_proxy_find_conn(&inst->remote_proxy, stream_id); + if (!rc || rc->sock == SOCKET_INVALID) { printf("[FAIL] no proxy conn\n"); g_done = -1; return; } + uasync_remove_socket_t(rc->ua, rc->sock); rc->read_id = NULL; + // Send data through OS socket directly + for (int i = 0; i < PAYLOAD_SIZE; i++) send_buf[i] = (uint8_t)(rand() & 0xFF); + ssize_t n = send(rc->sock, send_buf, PAYLOAD_SIZE, MSG_NOSIGNAL); + if (n != PAYLOAD_SIZE) { printf("[FAIL] send %zd\n", n); g_done = -1; return; } + phase = 2; + } + if (phase == 2) { + struct remote_proxy_conn* rc = remote_proxy_find_conn(&inst->remote_proxy, stream_id); + if (rc && rc->sock != SOCKET_INVALID) { + ssize_t n = recv(rc->sock, recv_buf, PAYLOAD_SIZE, 0); + if (n > 0) { + if ((size_t)n == PAYLOAD_SIZE && memcmp(send_buf, recv_buf, PAYLOAD_SIZE) == 0) { + printf("[PASS] test_remote_proxy — %zd bytes echoed\n", n); g_ok = 1; + } else { + printf("[FAIL] echo mismatch: got %zd expected %d\n", n, PAYLOAD_SIZE); + } + g_done = 1; return; + } + } + } + g_mon_id = uasync_set_timeout(ua, 10, NULL, monitor, "mon"); +} + +static void timeout_cb(void* arg) { + (void)arg; + if (!g_done) { printf("[FAIL] timeout: conn=%d\n", connected_ok); g_done = -1; } + if (g_mon_id) { uasync_cancel_timeout(ua, g_mon_id); g_mon_id = NULL; } +} + +int main(void) { + echo_pid = fork(); + if (echo_pid == 0) { echo_server(); _exit(0); } + if (echo_pid < 0) { perror("fork"); return 1; } + usleep(100000); + + if (test_mkdtemp(temp_dir) != 0) { kill(echo_pid, SIGTERM); waitpid(echo_pid,NULL,0); return 1; } + snprintf(cfg_path, sizeof(cfg_path), "%s/test.conf", temp_dir); + FILE* f = fopen(cfg_path, "w"); + if (!f) { kill(echo_pid, SIGTERM); waitpid(echo_pid,NULL,0); test_rmdir(temp_dir); return 1; } + fprintf(f, "%s", cfg); fclose(f); + + printf("=== test_remote_proxy ===\n"); + debug_config_init(); debug_set_level(DEBUG_LEVEL_ERROR); debug_set_categories(DEBUG_CATEGORY_ALL); + utun_instance_set_tun_init_enabled(0); + srand((unsigned)time(NULL)); + + ua = uasync_create(); + inst = utun_instance_create(ua, cfg_path); + if (!inst) { printf("[FAIL] instance create\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_cb, "to"); + while (!g_done) uasync_poll(ua, 100); + if (to_id) uasync_cancel_timeout(ua, to_id); + +done: + if (g_mon_id) uasync_cancel_timeout(ua, g_mon_id); + if (inst) { inst->running = 0; utun_instance_destroy(inst); } + if (ua) uasync_destroy(ua, 0); + test_unlink(cfg_path); test_rmdir(temp_dir); + if (echo_pid) { kill(echo_pid, SIGTERM); waitpid(echo_pid, NULL, 0); } + return g_ok ? 0 : 1; +}