You can not select more than 25 topics Topics must start with a letter or number, can include dashes ('-') and can be up to 35 characters long.
 
 
 
 
 
 

954 lines
47 KiB

// tcp_proxy.c — TCP proxy: lwIP TCP stack ↔ ll_queue ↔ transport → destination
// Raw-fd mode + TUN mode + active connections + EIM NAT + half-close
#include "tcp_proxy.h"
#include "lwip_tcp/lwip_tcp.h"
#include "lwip_tcp/lwip_tcp_priv.h"
#include "lwip_tcp/lwip_tcp_opts.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 "udp_proxy.h"
#include "icmp_proxy.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 <netinet/tcp.h>
#include <arpa/inet.h>
#include <fcntl.h>
#endif
#ifndef MSG_NOSIGNAL
#define MSG_NOSIGNAL 0
#endif
#ifndef INADDR_ANY
#define INADDR_ANY 0
#endif
// ====================================================================
// Forward declarations
// ====================================================================
struct sock_transport;
struct etcp_transport;
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 void sock_transport_connect(struct sock_transport* st);
static void proxy_try_close(struct proxy_conn *pc);
static struct sock_transport* sock_transport_create(struct proxy_conn* pc, struct UASYNC* ua);
static struct etcp_transport* etcp_transport_create(struct proxy_conn* pc, struct UTUN_INSTANCE* inst, uint64_t remote_node_id);
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;
uint16_t send_seq;
uint16_t recv_last_seq;
uint8_t recv_seq_init;
};
// Forward declarations for lwIP callbacks
static err_t proxy_recv_cb(void *arg, struct tcp_pcb *pcb, struct pbuf *p, err_t err);
static err_t proxy_sent_cb(void *arg, struct tcp_pcb *pcb, uint16_t len);
static void proxy_err_cb(void *arg, err_t err);
static err_t proxy_poll_cb(void *arg, struct tcp_pcb *pcb);
static err_t proxy_connected_cb(void *arg, struct tcp_pcb *pcb, err_t err);
static err_t proxy_accept_cb(void *arg, struct tcp_pcb *newpcb, err_t err);
static err_t tcp_output_cb(void *arg, struct pbuf *p, uint32_t src_ip, uint32_t dst_ip);
static void proxy_feed_from_transport(struct proxy_conn *pc);
static struct ll_entry* entry_from_data(struct memory_pool* pool, const uint8_t* data, uint16_t len);
// ====================================================================
// 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: lwIP TCP sends IP packets through this callback
// ====================================================================
static err_t tcp_output_cb(void *arg, struct pbuf *p, uint32_t src_ip, uint32_t dst_ip) {
struct tcp_proxy *proxy = (struct tcp_proxy *)arg;
(void)src_ip; (void)dst_ip;
uint16_t len = p->tot_len;
if (len > 2000) return LERR_BUF;
uint8_t buf[2002];
if (proxy->tun) {
pbuf_copy_partial(p, buf, len, 0);
ssize_t wr = tun_platform_write(proxy->tun, buf, len);
if (wr != (ssize_t)len) DEBUG_ERROR(DEBUG_CATEGORY_TRAFFIC, "PROXY TCP_OUT TUN write failed: ret=%zd len=%u errno=%d", wr, len, errno);
} else if (proxy->ip_fd >= 0) {
buf[0] = (len >> 8) & 0xFF; buf[1] = len & 0xFF;
pbuf_copy_partial(p, buf + 2, len, 0);
ssize_t n = write(proxy->ip_fd, buf, 2 + len); (void)n;
}
DEBUG_INFO(DEBUG_CATEGORY_TRAFFIC, "PROXY TCP_OUT %u.%u.%u.%u:%u len=%u",
(uint8_t)(src_ip), (uint8_t)(src_ip>>8), (uint8_t)(src_ip>>16), (uint8_t)(src_ip>>24),
(uint8_t)(dst_ip), (uint8_t)(dst_ip>>8), (uint8_t)(dst_ip>>16), (uint8_t)(dst_ip>>24), len);
return LERR_OK;
}
// ====================================================================
// Non-TCP handler (UDP/ICMP proxy)
// ====================================================================
static int tcp_proxy_handle_non_tcp(struct tcp_proxy* p, uint8_t* buf, size_t len) {
if (len < 20) return 0;
uint8_t ip_ver = (buf[0] >> 4) & 0xF;
if (ip_ver != 4) return 0;
uint8_t proto = buf[9];
if (proto == IPPROTO_UDP) {
if (!p->has_remote_mappings) return 0;
if (len < 28) return 0;
uint32_t src_ip, dst_ip; uint16_t src_port, dst_port;
memcpy(&src_ip, buf + 12, 4); memcpy(&dst_ip, buf + 16, 4);
memcpy(&src_port, buf + 20, 2); memcpy(&dst_port, buf + 22, 2);
udp_proxy_send_to_exit(p->inst, p->via_node_id, src_ip, src_port, dst_ip, dst_port, buf + 28, len - 28);
return 1;
}
if (proto == IPPROTO_ICMP) {
if (!p->has_remote_mappings) return 0;
if (len < 28) return 0;
uint8_t icmp_type = buf[20];
if (icmp_type != 8) return 0;
uint32_t src_ip, dst_ip;
memcpy(&src_ip, buf + 12, 4); memcpy(&dst_ip, buf + 16, 4);
uint16_t icmp_id, icmp_seq;
memcpy(&icmp_id, buf + 24, 2); memcpy(&icmp_seq, buf + 26, 2);
icmp_proxy_send_to_exit(p->inst, p->via_node_id, dst_ip, src_ip, icmp_id, icmp_seq, buf + 28, len - 28);
return 1;
}
return 0;
}
// ====================================================================
// Queue callback: автоматически кормит данные из transport_to_uip в lwIP
// ====================================================================
static void transport_to_uip_cb(struct ll_queue* q, void* arg) {
proxy_feed_from_transport((struct proxy_conn*)arg);
queue_resume_callback(q);
}
// ====================================================================
// Helper: feed data from transport_to_uip queue to lwIP TCP
// ====================================================================
static void proxy_feed_from_transport(struct proxy_conn *pc) {
if (!pc->transport_to_uip || !pc->pcb) return;
int sent_any = 0;
int entries_fed = 0;
while (1) {
uint16_t space = tcp_sndbuf(pc->pcb);
if (space < TCP_MSS / 2) { DEBUG_INFO(DEBUG_CATEGORY_TRAFFIC, "PROXY FEED → stream=%016llx space=%u break", (unsigned long long)pc->remote_stream_id, space); break; }
struct ll_entry *e = queue_data_get(pc->transport_to_uip);
if (!e) break;
uint16_t len = e->len;
if (len > space) len = space;
err_t ret = tcp_write(pc->pcb, e->dgram, len, TCP_WRITE_FLAG_COPY);
DEBUG_INFO(DEBUG_CATEGORY_TRAFFIC, "PROXY FEED → stream=%016llx len=%u ret=%d sndbuf=%u", (unsigned long long)pc->remote_stream_id, len, ret, pc->pcb->snd_buf);
if (ret == LERR_OK) {
sent_any = 1;
if (len >= e->len) { queue_dgram_free(e); queue_entry_free(e); }
else { memmove(e->dgram, e->dgram + len, e->len - len); e->len -= len; queue_data_put_first(pc->transport_to_uip, e); break; }
} else {
queue_data_put_first(pc->transport_to_uip, e);
break;
}
queue_resume_callback(pc->transport_to_uip);
}
queue_resume_callback(pc->transport_to_uip);
if (sent_any) tcp_output(pc->pcb);
}
// ====================================================================
// lwIP TCP callbacks
// ====================================================================
static err_t proxy_accept_cb(void *arg, struct tcp_pcb *newpcb, err_t err) {
struct tcp_proxy *p = (struct tcp_proxy *)arg;
if (err != LERR_OK || !newpcb) return LERR_ABRT;
struct proxy_conn *pc = u_calloc(1, sizeof(struct proxy_conn));
if (!pc) return LERR_MEM;
pc->proxy = p; pc->pcb = newpcb; pc->closing_tun = 0; pc->closing_rem = 0;
struct tcp_proxy_mapping *m = find_mapping_by_port(p, newpcb->local_port);
if (m) { memcpy(pc->dest_ip, m->remote_ip, 4); pc->dest_port = m->remote_port; }
else { memcpy(pc->dest_ip, &newpcb->local_ip, 4); pc->dest_port = htons(newpcb->local_port); }
pc->tun_ip = newpcb->local_ip;
tcp_arg(newpcb, pc);
tcp_recv(newpcb, proxy_recv_cb);
tcp_sent(newpcb, proxy_sent_cb);
tcp_err(newpcb, proxy_err_cb);
tcp_poll(newpcb, proxy_poll_cb, 2);
tcp_nagle_disable(newpcb);
pc->uip_to_transport = queue_new(p->ua, 0, 0, 0, "uip_to_transport");
pc->transport_to_uip = queue_new(p->ua, 0, 0, 0, "transport_to_uip");
if (p->via_node_id != 0 && p->via_node_id != p->inst->node_id) {
struct etcp_transport *et = etcp_transport_create(pc, p->inst, p->via_node_id);
if (!et) { u_free(pc); return LERR_MEM; }
pc->transport = &et->base;
} else {
struct sock_transport *st = sock_transport_create(pc, p->ua);
if (!st) { u_free(pc); return LERR_MEM; }
pc->transport = &st->base; st->conn = pc;
sock_transport_connect(st);
}
pc->next = p->conns; p->conns = pc; p->conn_count++;
DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "TCP proxy: new passive conn 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));
return LERR_OK;
}
static err_t proxy_recv_cb(void *arg, struct tcp_pcb *pcb, struct pbuf *p, err_t err) {
struct proxy_conn *pc = (struct proxy_conn *)arg;
if (!pc) { if (p) pbuf_free(p); return LERR_OK; }
if (p == NULL || err != LERR_OK) {
DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "PROXY FIN stream=%016llx tun=1 active=%d", (unsigned long long)(pc->remote_stream_id ? pc->remote_stream_id : 0), pc->active);
{ uint16_t wnd_gap = TCP_WND_MAX(pcb) - pcb->rcv_wnd;
if (wnd_gap > 0) tcp_recved(pcb, wnd_gap); }
pc->closing_tun = 1;
return LERR_OK;
}
if (pc->active) {
DEBUG_ERROR(DEBUG_CATEGORY_TRAFFIC, "RECV active pc=%p len=%u", (void*)pc, p->tot_len);
struct ll_entry *e = entry_from_data(pc->proxy->entry_pool, p->payload, p->tot_len);
if (e) queue_data_put(pc->uip_to_transport, e);
tcp_recved(pcb, p->tot_len);
pbuf_free(p);
} else if (pc->transport && !pc->closing_rem) {
uint16_t len = p->tot_len;
uint8_t *data = u_malloc(len);
if (data) {
pbuf_copy_partial(p, data, len, 0);
int ret = pc->transport->ops->send(pc->transport, data, len);
if (ret < 0) {
struct ll_entry *e = entry_from_data(pc->proxy->entry_pool, data, len);
if (e) queue_data_put(pc->uip_to_transport, e);
}
u_free(data);
}
tcp_recved(pcb, len);
pbuf_free(p);
}
return LERR_OK;
}
static err_t proxy_sent_cb(void *arg, struct tcp_pcb *pcb, uint16_t len) {
(void)len;
struct proxy_conn *pc = (struct proxy_conn *)arg;
if (!pc) return LERR_OK;
proxy_feed_from_transport(pc);
return LERR_OK;
}
static void proxy_err_cb(void *arg, err_t err) {
struct proxy_conn *pc = (struct proxy_conn *)arg;
if (!pc) return;
DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "TCP proxy: tcp error %d", err);
pc->closing_tun = 1;
if (pc->transport && !pc->closing_sent) pc->transport->ops->close(pc->transport);
}
static void proxy_try_close(struct proxy_conn *pc) {
if (!pc->closing_tun || pc->closing_sent || pc->closing_rem) return;
if (!pc->transport || !pc->pcb) return;
uint32_t pending = pc->transport_to_uip ? queue_entry_count(pc->transport_to_uip) : 0;
uint32_t pending_out = pc->uip_to_transport ? queue_entry_count(pc->uip_to_transport) : 0;
int done = (pending == 0 && pending_out == 0 && pc->pcb->unsent == NULL && pc->pcb->unacked == NULL);
uint32_t sndbuf = tcp_sndbuf(pc->pcb);
if (done || sndbuf == 0) {
DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "PROXY CLOSE stream=%016llx tun=1 pending=%u pending_out=%u done=%d sndbuf=%u cwnd=%u unsent=%p unacked=%p",
(unsigned long long)pc->remote_stream_id, pending, pending_out, done, sndbuf,
pc->pcb->cwnd, (void*)pc->pcb->unsent, (void*)pc->pcb->unacked);
pc->transport->ops->close(pc->transport);
}
}
static err_t proxy_poll_cb(void *arg, struct tcp_pcb *pcb) {
struct proxy_conn *pc = (struct proxy_conn *)arg;
if (!pc) return LERR_OK;
proxy_try_close(pc);
if (pc->closing_tun && (pc->closing_rem || pc->closing_sent) && pc->pcb) {
DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "PROXY CLEANUP stream=%016llx tun=1 rem=%d sent=%d cwnd=%u sndbuf=%u unsent=%p unacked=%p",
(unsigned long long)pc->remote_stream_id, pc->closing_rem, pc->closing_sent,
pc->pcb->cwnd, pc->pcb->snd_buf, (void*)pc->pcb->unsent, (void*)pc->pcb->unacked);
struct tcp_pcb *save = pc->pcb;
tcp_arg(save, NULL); pc->pcb = NULL;
{ uint16_t wnd_gap = TCP_WND_MAX(save) - save->rcv_wnd;
if (wnd_gap > 0) tcp_recved(save, wnd_gap); }
while (save->unsent) { struct tcp_seg *seg = save->unsent; save->unsent = seg->next; u_free(seg); }
tcp_close(save);
struct tcp_proxy *proxy = pc->proxy;
struct proxy_conn **prev = &proxy->conns;
while (*prev) { if (*prev == pc) { *prev = pc->next; proxy->conn_count--; break; } prev = &(*prev)->next; }
if (pc->transport) { pc->transport->ops->destroy(pc->transport); pc->transport = NULL; }
if (pc->eim_mapping) { pc->eim_mapping->delete_at_tb = get_time_tb() + proxy->eim_timeout_tb; pc->eim_mapping = NULL; }
if (pc->uip_to_transport) { struct ll_entry *e2; while ((e2 = queue_data_get(pc->uip_to_transport))) { queue_dgram_free(e2); queue_entry_free(e2); } queue_free(pc->uip_to_transport); }
if (pc->transport_to_uip) { struct ll_entry *e2; while ((e2 = queue_data_get(pc->transport_to_uip))) { queue_dgram_free(e2); queue_entry_free(e2); } queue_free(pc->transport_to_uip); }
u_free(pc);
}
return LERR_OK;
}
static err_t proxy_connected_cb(void *arg, struct tcp_pcb *pcb, err_t err) {
struct proxy_conn *pc = (struct proxy_conn *)arg;
if (!pc) return LERR_ABRT;
if (err != LERR_OK) {
DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "TCP proxy: active connect failed: %d", err);
pc->closing_tun = 1;
return LERR_ABRT;
}
DEBUG_ERROR(DEBUG_CATEGORY_TRAFFIC, "TCP proxy: active connect established to %d.%d.%d.%d:%d cwnd=%u snd_wnd=%u snd_buf=%u",
pc->dest_ip[0], pc->dest_ip[1], pc->dest_ip[2], pc->dest_ip[3], ntohs(pc->dest_port),
pcb->cwnd, pcb->snd_wnd, pcb->snd_buf);
proxy_feed_from_transport(pc);
return LERR_OK;
}
// ====================================================================
// Ensure dynamic listen pcb for transparent outbound TCP proxying
// ====================================================================
static void tcp_proxy_ensure_outbound_listen(struct tcp_proxy* p, uint16_t dport_net) {
if (!p->has_remote_mappings) return;
uint16_t dport_host = ntohs(dport_net);
struct tcp_pcb* lp = p->lwip->listen_pcbs;
while (lp) { if (lp->local_port == dport_host) return; lp = lp->next; }
struct tcp_pcb* lpcb = tcp_new(p->lwip);
if (!lpcb) return;
tcp_bind(lpcb, INADDR_ANY, dport_net);
struct tcp_pcb* listen_pcb = tcp_listen(lpcb);
if (listen_pcb) {
tcp_arg(listen_pcb, p); tcp_accept(listen_pcb, proxy_accept_cb);
DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "TCP proxy: dynamic listen on port %u", dport_host);
}
}
// ====================================================================
// Input: IP packet → lwip_tcp
// ====================================================================
static void tcp_proxy_raw_read(int fd, void* arg) {
(void)fd; struct tcp_proxy* p = (struct tcp_proxy*)arg;
ssize_t n;
if (p->raw_need == 0) {
/* Ждём 2-байтный заголовок длины */
n = read(p->ip_fd, p->raw_buf + p->raw_pos, 2 - p->raw_pos);
if (n <= 0) { if (n < 0 && errno == EAGAIN) return; p->raw_pos = 0; return; }
p->raw_pos += n;
if (p->raw_pos < 2) return;
p->raw_need = ((uint16_t)p->raw_buf[0] << 8) | p->raw_buf[1];
if (p->raw_need == 0 || p->raw_need > 2000) { p->raw_pos = 0; return; }
p->raw_pos = 0;
}
/* Ждём p->raw_need байт данных */
n = read(p->ip_fd, p->raw_buf + 2 + p->raw_pos, p->raw_need - p->raw_pos);
if (n <= 0) { if (n < 0 && errno == EAGAIN) return; p->raw_pos = p->raw_need = 0; return; }
p->raw_pos += n;
if (p->raw_pos < p->raw_need) return;
/* Пакет получен полностью */
uint8_t* pkt = p->raw_buf + 2;
size_t remaining = p->raw_need;
p->raw_pos = 0; p->raw_need = 0;
if (!tcp_proxy_handle_non_tcp(p, pkt, remaining)) {
uint8_t proto = pkt[9];
if (proto == IPPROTO_TCP) {
uint16_t ip_hdr_len = (pkt[0] & 0x0F) * 4;
uint16_t ip_total = ((uint16_t)pkt[2] << 8) | pkt[3];
if (ip_hdr_len >= 20 && ip_total >= ip_hdr_len && (size_t)ip_total <= remaining) {
uint16_t dport_net; memcpy(&dport_net, pkt + ip_hdr_len + 2, 2);
tcp_proxy_ensure_outbound_listen(p, dport_net);
uint32_t src_ip, dst_ip;
memcpy(&src_ip, pkt + 12, 4); memcpy(&dst_ip, pkt + 16, 4);
uint16_t tcp_len = ip_total - ip_hdr_len;
DEBUG_INFO(DEBUG_CATEGORY_TRAFFIC, "PROXY TCP_IN %u.%u.%u.%u:%u len=%u",
(uint8_t)(src_ip), (uint8_t)(src_ip>>8), (uint8_t)(src_ip>>16), (uint8_t)(src_ip>>24), ntohs(dport_net), tcp_len);
struct pbuf *pb = pbuf_alloc(PBUF_RAW, tcp_len);
if (pb) { pbuf_take(pb, pkt + ip_hdr_len, tcp_len); lwip_tcp_input(p->lwip, pb, src_ip, dst_ip); }
}
}
}
}
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 > 20 && len <= 2000) {
if (!tcp_proxy_handle_non_tcp(p, ip, len)) {
uint8_t proto = ip[9];
if (proto == IPPROTO_TCP) {
uint16_t ip_hdr_len = (ip[0] & 0x0F) * 4;
uint16_t ip_total = ((uint16_t)ip[2] << 8) | ip[3];
if (ip_hdr_len >= 20 && ip_total >= ip_hdr_len && len >= ip_total) {
uint16_t dport_net; memcpy(&dport_net, ip + ip_hdr_len + 2, 2);
tcp_proxy_ensure_outbound_listen(p, dport_net);
uint32_t src_ip, dst_ip;
memcpy(&src_ip, ip + 12, 4); memcpy(&dst_ip, ip + 16, 4);
uint16_t tcp_len = ip_total - ip_hdr_len;
struct pbuf *pb = pbuf_alloc(PBUF_RAW, tcp_len);
if (pb) { pbuf_take(pb, ip + ip_hdr_len, tcp_len); lwip_tcp_input(p->lwip, pb, src_ip, dst_ip); }
}
}
}
}
}
queue_dgram_free(entry); queue_entry_free(entry); queue_resume_callback(q);
}
// ====================================================================
// 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);
{ int one = 1; setsockopt(st->sock, IPPROTO_TCP, TCP_NODELAY, &one, sizeof(one)); }
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_ERROR(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_rem) pc->closing_rem = 1;
if(pc) pc->transport = NULL;
}
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; proxy_feed_from_transport(pc); } 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;
DEBUG_ERROR(DEBUG_CATEGORY_TRAFFIC, "WRITE_CB st=%p connected=%d connect_called=%d", (void*)st, st->connected, st->connect_called);
if(!st->connected) {
if(!st->connect_called) return;
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;
// Create a listen PCB for this dynamic EIM port
uint16_t dyn_port = local.sin_port;
struct tcp_pcb *lpcb = tcp_new(p->lwip);
if (lpcb) {
tcp_bind(lpcb, INADDR_ANY, dyn_port);
struct tcp_pcb *listen_pcb = tcp_listen(lpcb);
if (listen_pcb) { tcp_arg(listen_pcb, p); tcp_accept(listen_pcb, proxy_accept_cb); }
}
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_TRAFFIC, "ERROR_CB st=%p sock=%d connect_called=%d connected=%d", (void*)st, st->sock, st->connect_called, st->connected);
// For non-blocking connect, EPOLLERR can fire before EPOLLOUT.
// Check SO_ERROR: if connect succeeded (err==0), trigger write callback.
if (st->connect_called && !st->connected) {
int err = 0; socklen_t len = sizeof(err);
if (getsockopt(st->sock, SOL_SOCKET, SO_ERROR, &err, &len) == 0) {
if (err == 0) {
sock_transport_write_callback(st->sock, st);
return;
}
DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "TCP proxy: connect error fd=%d err=%d (%s)", st->sock, err, strerror(err));
sock_transport_close(&st->base);
return;
}
}
DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "TCP proxy: socket error fd=%d", st->sock);
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);
memcpy(e->dgram + 10, &et->send_seq, 2);
if (len > 0) memcpy(e->dgram + TCP_PROXY_HDR_SIZE, data, len);
e->len = TCP_PROXY_HDR_SIZE + len;
et->send_seq++;
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); memset(e->dgram + 10, 0, 2); e->len = TCP_PROXY_HDR_SIZE;
etcp_route_send(et->inst, et->remote_node_id, e); } else queue_entry_free(e);
}
et->connected = 0;
if (et->conn) { et->conn->transport = NULL; et->conn->closing_sent = 1;
DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "PROXY SENT=1 stream=%016llx (transport close)", (unsigned long long)et->conn->remote_stream_id);
}
}
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 = ++pc->proxy->next_stream_id;
pc->remote_stream_id = et->stream_id;
et->send_seq = 0; et->recv_seq_init = 0;
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);
memset(e->dgram + 10, 0, 2);
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(struct tcp_proxy* p, uint64_t stream_id) {
struct proxy_conn* pc;
for (pc = p->conns; pc; pc = pc->next) if (pc->remote_stream_id == stream_id) return pc;
return NULL;
}
static void handle_connected(struct tcp_proxy* p, uint64_t stream_id, struct ll_entry* entry) {
struct proxy_conn* pc = find_pc_by_stream(p, stream_id);
if (!pc || !pc->transport) {
DEBUG_WARN(DEBUG_CATEGORY_SOCKET, "PROXY CONNECTED stream=%016llx pc=%p transport=%p — drop", (unsigned long long)stream_id, (void*)pc, pc ? pc->transport : NULL);
queue_dgram_free(entry); queue_entry_free(entry);
return;
}
struct etcp_transport* et = (struct etcp_transport*)pc->transport;
if (et->base.ops != &etcp_transport_ops) { queue_dgram_free(entry); queue_entry_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, "PROXY CONNECTED stream=%016llx status=OK", (unsigned long long)stream_id);
struct ll_entry* e;
while ((e = queue_data_get(pc->uip_to_transport)) != NULL) {
if (e->len > 0) {
etcp_transport_send(&et->base, e->dgram, e->len);
}
queue_dgram_free(e); queue_entry_free(e);
}
}
else { DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "TCP proxy: remote refused stream=%016llx", (unsigned long long)stream_id);
DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "PROXY REM=1 stream=%016llx (refused)", (unsigned long long)stream_id);
pc->closing_rem = 1; }
}
queue_dgram_free(entry); queue_entry_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_dgram_free(entry); queue_entry_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 : NULL;
struct tcp_proxy* proxy = inst ? inst->tcp_proxy : NULL;
if (subcmd == TCP_PROXY_SUBCMD_CONNECT) {
uint64_t src_node_id = conn ? conn->peer_node_id : (inst ? inst->node_id : 0);
DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "RP CONNECT recv stream=%016llx", (unsigned long long)stream_id);
remote_proxy_handle_connect(inst, entry, stream_id, src_node_id);
return;
}
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_dgram_free(entry); queue_entry_free(entry); return; }
}
}
if (proxy) {
if (subcmd == TCP_PROXY_SUBCMD_CONNECTED) { handle_connected(proxy, stream_id, entry); return; }
if (subcmd == TCP_PROXY_SUBCMD_DATA) {
struct proxy_conn* pc = find_pc_by_stream(proxy, stream_id);
if (pc) {
if (pc->transport && pc->transport->ops == &etcp_transport_ops) {
struct etcp_transport* et = (struct etcp_transport*)pc->transport;
uint16_t seq; memcpy(&seq, entry->dgram + 10, 2);
if (!et->recv_seq_init) {
et->recv_last_seq = seq; et->recv_seq_init = 1;
} else {
int16_t delta = (int16_t)(seq - et->recv_last_seq);
if (delta <= 0) {
DEBUG_WARN(DEBUG_CATEGORY_SOCKET, "TCP proxy: duplicate DATA seq=%u stream=%016llx, discard", seq, (unsigned long long)stream_id);
queue_dgram_free(entry); queue_entry_free(entry); return;
}
if (delta > 1) {
DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "TCP proxy: seq gap seq=%u last=%u stream=%016llx, tearing down", seq, et->recv_last_seq, (unsigned long long)stream_id);
pc->closing_rem = 1;
if (pc->transport) { pc->transport->ops->destroy(pc->transport); pc->transport = NULL; }
{
int n_active = 0, n_tw = 0;
struct tcp_pcb* pcb2;
for (pcb2 = proxy->lwip->active_pcbs; pcb2; pcb2 = pcb2->next) n_active++;
for (pcb2 = proxy->lwip->tw_pcbs; pcb2; pcb2 = pcb2->next) n_tw++;
DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "PROXY CLEANUP done stream=%016llx active_pcbs=%d tw_pcbs=%d",
(unsigned long long)pc->remote_stream_id, n_active, n_tw);
}
queue_dgram_free(entry); queue_entry_free(entry); return;
}
et->recv_last_seq = seq;
}
size_t data_len = entry->len - TCP_PROXY_HDR_SIZE;
DEBUG_INFO(DEBUG_CATEGORY_TRAFFIC, "PROXY DATA ← stream=%016llx seq=%u len=%zu", (unsigned long long)stream_id, seq, data_len);
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);
proxy_feed_from_transport(pc);
}
}
}
}
if (subcmd == TCP_PROXY_SUBCMD_CLOSE) {
struct proxy_conn* pc = find_pc_by_stream(proxy, stream_id);
if (pc) {
DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "PROXY REM=1 stream=%016llx (CLOSE from exit)", (unsigned long long)stream_id);
pc->closing_rem = 1;
if (pc->transport) { pc->transport->ops->destroy(pc->transport); pc->transport = NULL; }
}
}
}
queue_dgram_free(entry); queue_entry_free(entry);
}
// ====================================================================
// 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, uint64_t via_node_id)
{
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->via_node_id = via_node_id;
p->has_remote_mappings = (via_node_id != 0 && inst && via_node_id != inst->node_id) ? 1 : 0;
p->ip_fd_mode = !use_tun;
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; }
}
p->lwip = lwip_tcp_init(ua, tcp_output_cb, p);
if (!p->lwip) {
if (p->ip_fd_id) uasync_remove_socket(ua, p->ip_fd_id);
if (p->tun) tun_close(p->tun);
memory_pool_destroy(p->entry_pool); u_free(p); return NULL;
}
if (mapping_count > 0) {
int j; for (j = 0; j < mapping_count; j++) {
uint16_t port_net = htons(mappings[j].local_port);
struct tcp_pcb *lpcb = tcp_new(p->lwip);
if (!lpcb) continue;
tcp_bind(lpcb, INADDR_ANY, port_net);
struct tcp_pcb *listen_pcb = tcp_listen(lpcb);
if (listen_pcb) {
tcp_arg(listen_pcb, p);
tcp_accept(listen_pcb, proxy_accept_cb);
DEBUG_ERROR(DEBUG_CATEGORY_TRAFFIC, "TCP proxy: listen on %d ok lpcb=%p ctx_listen=%p",
mappings[j].local_port, (void*)listen_pcb, (void*)p->lwip->listen_pcbs);
}
struct tcp_proxy_mapping* m = u_calloc(1, sizeof(struct tcp_proxy_mapping));
if (m) { m->local_port = mappings[j].local_port; 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->next = p->mappings; p->mappings = m; }
}
}
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);
udp_proxy_init(inst, ua);
icmp_proxy_init(inst, ua);
}
}
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);
udp_proxy_destroy(p->inst);
icmp_proxy_destroy(p->inst);
if (p->lwip) { lwip_tcp_destroy(p->lwip); p->lwip = 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))) { 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))) { queue_dgram_free(e); queue_entry_free(e); } queue_free(pc->transport_to_uip); }
u_free(pc); pc = next; }
p->conns = NULL;
struct tcp_proxy_mapping* m = p->mappings;
while (m) { struct tcp_proxy_mapping* next = m->next; 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);
u_free(p);
}
// ====================================================================
// Active connections API
// ====================================================================
struct proxy_conn* tcp_proxy_active_open(struct tcp_proxy* p, const char* dest_ip, uint16_t dest_port) {
if (!p || !dest_ip) return NULL;
unsigned int a0, a1, a2, a3;
if (sscanf(dest_ip, "%u.%u.%u.%u", &a0, &a1, &a2, &a3) != 4) return NULL;
struct proxy_conn* pc = u_calloc(1, sizeof(struct proxy_conn));
if (!pc) return NULL;
pc->proxy = p; pc->active = 1; pc->closing_tun = 0; pc->closing_rem = 0;
pc->dest_port = htons(dest_port);
pc->dest_ip[0] = (uint8_t)a0; pc->dest_ip[1] = (uint8_t)a1; pc->dest_ip[2] = (uint8_t)a2; pc->dest_ip[3] = (uint8_t)a3;
pc->tun_ip = htonl((a0 << 24) | (a1 << 16) | (a2 << 8) | a3);
pc->uip_to_transport = queue_new(p->ua, 0, 0, 0, "active_uip_to");
pc->transport_to_uip = queue_new(p->ua, 0, 0, 0, "active_to_uip");
if (!pc->uip_to_transport || !pc->transport_to_uip) { u_free(pc); return NULL; }
struct tcp_pcb *pcb = tcp_new(p->lwip);
if (!pcb) { queue_free(pc->uip_to_transport); queue_free(pc->transport_to_uip); u_free(pc); return NULL; }
pc->pcb = pcb;
tcp_arg(pcb, pc);
tcp_recv(pcb, proxy_recv_cb);
tcp_sent(pcb, proxy_sent_cb);
tcp_err(pcb, proxy_err_cb);
tcp_poll(pcb, proxy_poll_cb, 2);
tcp_nagle_disable(pcb);
tcp_bind(pcb, INADDR_ANY, 0);
pcb->local_ip = htonl(0x7F000001); // use 127.0.0.1 as source for raw-fd mode
uint32_t ip = htonl((a0 << 24) | (a1 << 16) | (a2 << 8) | a3);
tcp_connect(pcb, ip, htons(dest_port), proxy_connected_cb);
pc->next = p->conns; p->conns = pc; p->conn_count++;
DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "TCP proxy: active_open to %s:%d", dest_ip, dest_port);
return pc;
}
int tcp_proxy_active_send(struct tcp_proxy* p, struct proxy_conn* pc, const uint8_t* data, size_t len) {
(void)p;
if (!pc || !pc->active || !data || len == 0) 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);
if (pc->pcb && pc->pcb->state == ESTABLISHED) proxy_feed_from_transport(pc);
return 0;
}
ssize_t tcp_proxy_active_recv(struct tcp_proxy* p, struct proxy_conn* pc, uint8_t* buf, size_t len) {
(void)p;
if (!pc || !pc->active || !buf || len == 0) 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, struct proxy_conn* pc) {
(void)p;
if (!pc || !pc->active) return -1;
pc->closing_tun = 1;
tcp_close(pc->pcb);
return 0;
}
int tcp_proxy_active_send_done(struct tcp_proxy* p, struct proxy_conn* pc) {
(void)p;
if (!pc || !pc->active || !pc->pcb) return 1;
if (pc->closing_tun) return 1;
int q = queue_entry_count(pc->transport_to_uip);
if (q > 0) return 0;
if (pc->pcb->unsent || pc->pcb->unacked) return 0;
return 1;
}