Browse Source
etcp_router: routable service packets with dst/src node_id - etcp_router_bind/unbind: per-instance service handler registry - etcp_route_send: send to node by id via BGP routing - multi-hop forwarding through intermediate nodes - loopback optimization (dst==self dispatches directly) remote_proxy: exit node creating OS sockets on demand - handles CONNECT/DATA/CLOSE via etcp_router - non-blocking socket connect with uasync callbacks - echo back response data to requesting node tcp_proxy: etcp_transport for remote proxy - via_node_id in mapping selects between sock/etcp transport - unified handler for both tcp_proxy and remote_proxy roles - config: forward = port -> ip:port via NODE_HEX [remote_proxy] section tests: test_etcp_router (20 pkts bidirectional), test_remote_proxy (echo)congestion
15 changed files with 1997 additions and 6 deletions
@ -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 <string.h> |
||||
|
||||
// Обработчик 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); |
||||
} |
||||
@ -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 <stdint.h> |
||||
#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
|
||||
@ -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 <stdlib.h> |
||||
#include <string.h> |
||||
#include <errno.h> |
||||
#ifndef _WIN32 |
||||
#include <unistd.h> |
||||
#include <sys/socket.h> |
||||
#include <netinet/in.h> |
||||
#include <arpa/inet.h> |
||||
#include <fcntl.h> |
||||
#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"); |
||||
} |
||||
@ -0,0 +1,58 @@
|
||||
// remote_proxy.h — Удаленный TCP прокси (exit node)
|
||||
// Создаёт OS сокеты к адресатам по запросам от других нод через etcp_router
|
||||
#ifndef REMOTE_PROXY_H |
||||
#define REMOTE_PROXY_H |
||||
|
||||
#include <stdint.h> |
||||
#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 |
||||
@ -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 <stdlib.h> |
||||
#include <string.h> |
||||
#include <errno.h> |
||||
#ifndef _WIN32 |
||||
#include <unistd.h> |
||||
#include <sys/socket.h> |
||||
#include <netinet/in.h> |
||||
#include <arpa/inet.h> |
||||
#include <fcntl.h> |
||||
#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; |
||||
} |
||||
@ -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 <stdint.h> |
||||
#include <stddef.h> |
||||
#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
|
||||
@ -0,0 +1,340 @@
|
||||
// test_etcp_router.c — Unit test for etcp_router: bind/unbind, send, loopback, forward, reply
|
||||
#include <stdio.h> |
||||
#include <stdlib.h> |
||||
#include <string.h> |
||||
#include "../lib/platform_compat.h" |
||||
#include "test_utils.h" |
||||
#ifdef _WIN32 |
||||
#include <windows.h> |
||||
#include <direct.h> |
||||
#else |
||||
#include <unistd.h> |
||||
#endif |
||||
#include <time.h> |
||||
#include <sys/stat.h> |
||||
|
||||
#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; |
||||
} |
||||
@ -0,0 +1,175 @@
|
||||
// test_remote_proxy.c — Test remote_proxy: CONNECT → socket → echo → verify
|
||||
#include <stdio.h> |
||||
#include <stdlib.h> |
||||
#include <string.h> |
||||
#include <errno.h> |
||||
#include <time.h> |
||||
#include <unistd.h> |
||||
#include <signal.h> |
||||
#include <sys/socket.h> |
||||
#include <netinet/in.h> |
||||
#include <arpa/inet.h> |
||||
#include <sys/wait.h> |
||||
#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; |
||||
} |
||||
Loading…
Reference in new issue