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.
969 lines
49 KiB
969 lines
49 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) { |
|
uint16_t ip_total = ((uint16_t)buf[2] << 8) | buf[3]; |
|
uint32_t src_ip, dst_ip; memcpy(&src_ip, buf + 12, 4); memcpy(&dst_ip, buf + 16, 4); |
|
DEBUG_ERROR(DEBUG_CATEGORY_TRAFFIC, "TUN_WRITE_ERR ret=%zd len=%u ip_total=%u vhl=%02x proto=%u %u.%u.%u.%u→%u.%u.%u.%u errno=%d", |
|
wr, len, ip_total, buf[0], buf[9], |
|
(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), errno); |
|
{ char hx[128]; int p = 0; size_t n = len < 32 ? len : 32; |
|
for (size_t i = 0; i < n && p < 124; i++) p += snprintf(hx + p, sizeof(hx) - p, "%02x", buf[i]); |
|
DEBUG_ERROR(DEBUG_CATEGORY_TRAFFIC, "TUN_WRITE_ERR hex[0..%zu]: %s", n, hx); } |
|
} |
|
} 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 { |
|
DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "PROXY FEED tcp_write failed stream=%016llx len=%u ret=%d", (unsigned long long)pc->remote_stream_id, len, ret); |
|
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); DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "PROXY proxy_recv_cb: pc=NULL"); 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); } |
|
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); } |
|
} |
|
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) { |
|
DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "PROXY send failed stream=%016llx ret=%d len=%u", (unsigned long long)pc->remote_stream_id, ret, len); |
|
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) { DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "PROXY proxy_err_cb: pc=NULL err=%d", err); 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) { DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "PROXY proxy_poll_cb: pc=NULL"); 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) { DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "PROXY send not connected stream=%016llx len=%zu", (unsigned long long)et->stream_id, len); return -1; } |
|
struct ll_entry* e = queue_entry_new(0); |
|
if (!e) { DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "PROXY send alloc entry fail stream=%016llx", (unsigned long long)et->stream_id); return -1; } |
|
e->dgram = u_malloc(TCP_PROXY_HDR_SIZE + len); |
|
if (!e->dgram) { DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "PROXY send alloc dgram fail stream=%016llx len=%zu", (unsigned long long)et->stream_id, len); 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) { DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "PROXY CONNECTED wrong ops stream=%016llx", (unsigned long long)stream_id); 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; |
|
}
|
|
|