// tcp_proxy_server.c — TCP прокси-сервер (exit node) #include "tcp_proxy_server.h" #include "tcp_proxy_client.h" #include "udp_proxy.h" #include "icmp_proxy.h" #include "etcp.h" #include "etcp_api.h" #include "etcp_router.h" #include "utun_instance.h" #include "config_parser.h" #include "../lib/u_async.h" #include "../lib/debug_config.h" #include "../lib/ll_queue.h" #include "../lib/mem.h" #include "../lib/tcp_io.h" #include #include #include #ifndef _WIN32 #include #include #include #include #include #endif static struct tcp_proxy_server* g_tcp_proxy_server_ctx = NULL; static void on_fin_cb(struct tcp_conn* tc, void* arg); static void on_error_cb(struct tcp_conn* tc, int err, void* arg); static void on_flushed_cb(struct tcp_conn* tc, void* arg); static void read_queue_drain_cb(struct ll_queue* q, void* arg); static void pause_resume_cb(struct ll_queue* q, void* arg); static void close_retry_cb(void* arg); static void diag_timer_cb(void* arg); static int send_msg(struct UTUN_INSTANCE* inst, uint64_t dst, uint8_t subcmd, uint32_t sid, const uint8_t* data, size_t len, int force); static void send_close(struct tcp_proxy_server_conn* rc); static void send_error(struct tcp_proxy_server_conn* rc); void tcp_proxy_server_conn_free(struct tcp_proxy_server_conn* rc); static inline int write_pending(struct tcp_conn* tc) { return tc && (tc->write_buf || tc->write_queue->head); } // ==================================================================== // Отправка сообщений через ETCP // ==================================================================== static int send_msg(struct UTUN_INSTANCE* inst, uint64_t dst, uint8_t subcmd, uint32_t sid, const uint8_t* data, size_t len, int force) { struct ll_entry* e = queue_entry_new(0); if (!e) { DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "tcp_proxy_server: queue_entry_new failed subcmd=%02x sid=%08x", subcmd, sid); return -1; } e->dgram = u_malloc(TCP_PROXY_HDR_SIZE + len); if (!e->dgram) { DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "tcp_proxy_server: malloc(%zu) failed", TCP_PROXY_HDR_SIZE + len); queue_entry_free(e); return -1; } e->dgram[0] = ETCP_ID_TCP_PROXY; e->dgram[1] = subcmd; memcpy(e->dgram + 2, &sid, 4); 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, force); } static void close_retry_cb(void* arg) { struct tcp_proxy_server_conn* rc = (struct tcp_proxy_server_conn*)arg; if (!rc) return; rc->close_timer = NULL; if (!rc->close_pending) return; struct UTUN_INSTANCE* inst = rc->ctx ? rc->ctx->inst : NULL; if (!inst) return; uint8_t subcmd = (rc->tc && rc->tc->error) ? TCP_PROXY_SUBCMD_ERROR : TCP_PROXY_SUBCMD_CLOSE; if (send_msg(inst, rc->peer_node_id, subcmd, rc->stream_id, NULL, 0, 1) < 0) { rc->close_backoff = rc->close_backoff < 5000 ? rc->close_backoff * 2 : 5000; rc->close_timer = uasync_set_timeout(rc->ua, rc->close_backoff, rc, close_retry_cb, "tps_close_retry"); return; } rc->close_pending = 0; } static void send_close(struct tcp_proxy_server_conn* rc) { struct UTUN_INSTANCE* inst = rc->ctx ? rc->ctx->inst : NULL; if (!inst) return; if (send_msg(inst, rc->peer_node_id, TCP_PROXY_SUBCMD_CLOSE, rc->stream_id, NULL, 0, 1) < 0) { rc->close_pending = 1; rc->close_backoff = 50; if (!rc->close_timer) rc->close_timer = uasync_set_timeout(rc->ua, rc->close_backoff, rc, close_retry_cb, "tps_close_retry"); } } static void send_error(struct tcp_proxy_server_conn* rc) { struct UTUN_INSTANCE* inst = rc->ctx ? rc->ctx->inst : NULL; if (!inst) return; if (send_msg(inst, rc->peer_node_id, TCP_PROXY_SUBCMD_ERROR, rc->stream_id, NULL, 0, 1) < 0) { rc->close_pending = 1; rc->close_backoff = 50; if (!rc->close_timer) rc->close_timer = uasync_set_timeout(rc->ua, rc->close_backoff, rc, close_retry_cb, "tps_close_retry"); } } static void send_fin(struct tcp_proxy_server_conn* rc) { struct UTUN_INSTANCE* inst = rc->ctx ? rc->ctx->inst : NULL; if (!inst) return; send_msg(inst, rc->peer_node_id, TCP_PROXY_SUBCMD_FIN, rc->stream_id, NULL, 0, 1); } // ==================================================================== // Коллбэки tcp_io // ==================================================================== static void on_fin_cb(struct tcp_conn* tc, void* arg) { struct tcp_proxy_server_conn* rc = (struct tcp_proxy_server_conn*)arg; if (write_pending(tc)) tcp_conn_set_flushed(tc, on_flushed_cb); else send_fin(rc); } static void on_flushed_cb(struct tcp_conn* tc, void* arg) { struct tcp_proxy_server_conn* rc = (struct tcp_proxy_server_conn*)arg; if (rc->cli_closed) { shutdown(tc->sock, SHUT_WR); tcp_proxy_server_conn_free(rc); return; } send_fin(rc); } static void on_error_cb(struct tcp_conn* tc, int err, void* arg) { (void)tc; (void)err; struct tcp_proxy_server_conn* rc = (struct tcp_proxy_server_conn*)arg; send_error(rc); tcp_proxy_server_conn_free(rc); } static void diag_timer_cb(void* arg) { struct tcp_proxy_server_conn* rc = (struct tcp_proxy_server_conn*)arg; struct tcp_conn* tc = rc->tc; if (!tc || tc->sock == SOCKET_INVALID) return; int ent_f = tc->entry_pool->free_count; int dat_f = tc->data_pool->free_count; int rcv_buf = 0, snd_buf = 0; socklen_t optlen = sizeof(int); getsockopt(tc->sock, SOL_SOCKET, SO_RCVBUF, &rcv_buf, &optlen); getsockopt(tc->sock, SOL_SOCKET, SO_SNDBUF, &snd_buf, &optlen); DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "SOCK:DIAG fd=%d sid=%08x conn=%d fin=%d " "rq=%d(%zub) wq=%d(%zub) wbuf=%s " "ent_f=%d dat_f=%d tcp_rb=%zu tcp_sb=%zu", (int)tc->sock, rc->stream_id, tc->connected, tc->fin, tc->read_queue->count, queue_total_bytes(tc->read_queue), tc->write_queue->count, queue_total_bytes(tc->write_queue), tc->write_buf ? "y" : "n", ent_f, dat_f, (size_t)rcv_buf, (size_t)snd_buf); rc->diag_timer = uasync_set_timeout(rc->ua, 10000, rc, diag_timer_cb, "tps_diag"); } // ==================================================================== // Дрейн read_queue → ETCP (автозабор + backpressure) // ==================================================================== static void read_queue_drain_cb(struct ll_queue* q, void* arg) { struct tcp_proxy_server_conn* rc = (struct tcp_proxy_server_conn*)arg; struct UTUN_INSTANCE* inst = rc->ctx ? rc->ctx->inst : NULL; struct ll_entry* e = queue_data_get(q); if (!e) { queue_resume_callback(q); return; } if (rc->cli_closed) { do { memory_pool_free(rc->tc->data_pool, e->dgram); queue_entry_free(e); e = queue_data_get(q); } while (e); queue_resume_callback(q); return; } int ret = send_msg(inst, rc->peer_node_id, TCP_PROXY_SUBCMD_DATA, rc->stream_id, e->dgram, e->len, 0); DEBUG_TRACE(DEBUG_CATEGORY_SOCKET, "SOCK:RECV fd=%d sid=%08x len=%u total=%d", (int)rc->tc->sock, rc->stream_id, e->len, rc->ctx ? rc->ctx->conn_count : 0); memory_pool_free(rc->tc->data_pool, e->dgram); queue_entry_free(e); if (ret == 0) { queue_resume_callback(q); } else { DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "SOCK:BACKPRESSURE fd=%d sid=%08x — pausing drain", (int)rc->tc->sock, rc->stream_id); etcp_router_on_send_ready(inst, rc->peer_node_id, ETCP_ID_TCP_PROXY, &rc->pause_waiter, pause_resume_cb, rc); } } static void pause_resume_cb(struct ll_queue* q, void* arg) { (void)q; struct tcp_proxy_server_conn* rc = (struct tcp_proxy_server_conn*)arg; if (!rc->tc || rc->tc->sock == SOCKET_INVALID) return; queue_resume_callback(rc->tc->read_queue); } // ==================================================================== // Управление жизненным циклом коннекта // ==================================================================== static int conn_total(struct tcp_proxy_server_conn* rc) { if (!rc || !rc->ctx) return 0; int n = 0; struct tcp_proxy_server_conn* c; for (c = rc->ctx->conns; c; c = c->next) n++; return n; } void tcp_proxy_server_conn_free(struct tcp_proxy_server_conn* rc) { if (!rc) return; DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "SOCK:FREE enter rc=%p freed=%d sid=%08x", (void*)rc, rc->freed, rc->stream_id); if (rc->freed) { DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "SOCK:FREE double rc=%p — IGNORED", (void*)rc); return; } rc->freed = 1; int total = conn_total(rc); DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "SOCK:FREE fd=%d sid=%08x total=%d cli_closed=%d fin=%d error=%d", rc->tc ? (int)rc->tc->sock : -1, rc->stream_id, total, rc->cli_closed, rc->tc ? rc->tc->fin : 0, rc->tc ? rc->tc->error : 0); if (rc->close_pending && rc->ctx && rc->ctx->inst) { rc->close_pending = 0; uint8_t subcmd = (rc->tc && rc->tc->error) ? TCP_PROXY_SUBCMD_ERROR : TCP_PROXY_SUBCMD_CLOSE; send_msg(rc->ctx->inst, rc->peer_node_id, subcmd, rc->stream_id, NULL, 0, 1); } if (rc->ctx) { struct tcp_proxy_server_conn** prev = &rc->ctx->conns; while (*prev) { if (*prev == rc) { *prev = rc->next; rc->ctx->conn_count--; break; } prev = &(*prev)->next; } } if (rc->tc) { tcp_conn_destroy(rc->tc); rc->tc = NULL; } if (rc->close_timer) { uasync_cancel_timeout(rc->ua, rc->close_timer); rc->close_timer = NULL; } if (rc->diag_timer) { uasync_cancel_timeout(rc->ua, rc->diag_timer); rc->diag_timer = NULL; } if (rc->ctx && rc->ctx->inst) etcp_router_cancel_send_ready(rc->ctx->inst, rc->peer_node_id, ETCP_ID_TCP_PROXY, &rc->pause_waiter); DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "SOCK:FREE u_free rc=%p", (void*)rc); u_free(rc); } struct tcp_proxy_server_conn* tcp_proxy_server_find_conn(struct tcp_proxy_server* ctx, uint32_t stream_id) { struct tcp_proxy_server_conn* c; for (c = ctx->conns; c; c = c->next) if (c->stream_id == stream_id) return c; return NULL; } // ==================================================================== // Входные точки из ETCP-диспетчера // ==================================================================== int tcp_proxy_server_handle_connect(struct UTUN_INSTANCE* inst, struct ll_entry* entry, uint32_t stream_id, uint64_t src_node_id) { if (!inst || !inst->tcp_proxy_server.enabled) { if (entry) { queue_dgram_free(entry); queue_entry_free(entry); } return -1; } struct tcp_proxy_server* ctx = &inst->tcp_proxy_server; if (entry->len < TCP_PROXY_CONNECT_HDR_SIZE) { DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "TCP proxy server: CONNECT too short len=%u", entry->len); queue_dgram_free(entry); queue_entry_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 tcp_proxy_server_conn* rc = u_calloc(1, sizeof(struct tcp_proxy_server_conn)); if (!rc) { DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "TCP proxy server: u_calloc failed sid=%08x", stream_id); queue_dgram_free(entry); queue_entry_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; socket_t sock = socket(AF_INET, SOCK_STREAM, 0); if (sock == SOCKET_INVALID) { DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "TCP proxy server: socket() failed"); u_free(rc); queue_dgram_free(entry); queue_entry_free(entry); return -1; } ctx->conn_count++; DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "SOCK:NEW fd=%d sid=%08x dest=%d.%d.%d.%d:%d total=%d", (int)sock, stream_id, dest_ip[0],dest_ip[1],dest_ip[2],dest_ip[3],ntohs(dest_port), ctx->conn_count); socket_set_nonblocking(sock); rc->tc = tcp_conn_create(inst->ua, sock, 1500, 8192, 32, 8, on_fin_cb, on_error_cb, rc); if (!rc->tc) { DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "TCP proxy server: tcp_conn_create failed"); socket_close_wrapper(sock); ctx->conn_count--; u_free(rc); queue_dgram_free(entry); queue_entry_free(entry); return -1; } queue_set_callback(rc->tc->read_queue, read_queue_drain_cb, rc); queue_set_waiter_defer(rc->tc->read_queue, 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; int ret = connect(sock, (struct sockaddr*)&addr, sizeof(addr)); if (ret < 0 && errno != EINPROGRESS) { DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "TCP proxy server: connect() to %d.%d.%d.%d:%d failed: %s", dest_ip[0], dest_ip[1], dest_ip[2], dest_ip[3], ntohs(dest_port), strerror(errno)); send_msg(inst, src_node_id, TCP_PROXY_SUBCMD_ERROR, stream_id, NULL, 0, 1); ctx->conn_count--; tcp_conn_destroy(rc->tc); u_free(rc); queue_dgram_free(entry); queue_entry_free(entry); return -1; } rc->next = ctx->conns; ctx->conns = rc; rc->diag_timer = uasync_set_timeout(rc->ua, 10000, rc, diag_timer_cb, "tps_diag"); queue_dgram_free(entry); queue_entry_free(entry); return 0; } int tcp_proxy_server_handle_data(struct UTUN_INSTANCE* inst, struct ETCP_CONN* conn, struct ll_entry* entry, uint32_t stream_id) { (void)conn; if (!inst) { queue_dgram_free(entry); queue_entry_free(entry); return -1; } struct tcp_proxy_server* ctx = &inst->tcp_proxy_server; struct tcp_proxy_server_conn* rc = tcp_proxy_server_find_conn(ctx, stream_id); if (!rc || !rc->tc || rc->tc->sock == SOCKET_INVALID || rc->tc->error || rc->cli_closed) { DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "TPS handle_data: no/closed conn sid=%08x, dropping", stream_id); queue_dgram_free(entry); queue_entry_free(entry); return -1; } size_t data_len = entry->len - TCP_PROXY_HDR_SIZE; if (data_len > 0) { struct ll_entry* e = queue_entry_new_from_pool(rc->tc->entry_pool); uint8_t* buf = memory_pool_alloc(rc->tc->data_pool); if (e && buf) { memcpy(buf, entry->dgram + TCP_PROXY_HDR_SIZE, data_len); e->dgram = buf; e->len = (uint16_t)data_len; queue_data_put(rc->tc->write_queue, e); } else { DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "tcp_proxy_server: handle_data alloc failed sid=%08x", stream_id); if (e) queue_entry_free(e); if (buf) memory_pool_free(rc->tc->data_pool, buf); } } queue_dgram_free(entry); queue_entry_free(entry); return 0; } void tcp_proxy_server_handle_close(struct UTUN_INSTANCE* inst, uint32_t stream_id) { if (!inst) return; struct tcp_proxy_server* ctx = &inst->tcp_proxy_server; struct tcp_proxy_server_conn* rc = tcp_proxy_server_find_conn(ctx, stream_id); if (!rc) { DEBUG_WARN(DEBUG_CATEGORY_SOCKET, "TCP proxy server: CLOSE sid=%08x — no conn", stream_id); return; } DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "SOCK:CLOSE_RECV fd=%d sid=%08x total=%d fin=%d write_pend=%d", rc->tc ? (int)rc->tc->sock : -1, stream_id, conn_total(rc), rc->tc ? rc->tc->fin : 0, rc->tc ? write_pending(rc->tc) : 0); rc->cli_closed = 1; if (!rc->tc) { tcp_proxy_server_conn_free(rc); return; } tcp_conn_pause_read(rc->tc); { struct ll_entry* e; while ((e = queue_data_get(rc->tc->read_queue))) { memory_pool_free(rc->tc->data_pool, e->dgram); queue_entry_free(e); } queue_resume_callback(rc->tc->read_queue); } if (write_pending(rc->tc)) { tcp_conn_set_flushed(rc->tc, on_flushed_cb); } else { if (rc->tc->connected) shutdown(rc->tc->sock, SHUT_WR); tcp_proxy_server_conn_free(rc); } } void tcp_proxy_server_handle_error(struct UTUN_INSTANCE* inst, uint32_t stream_id) { if (!inst) return; struct tcp_proxy_server* ctx = &inst->tcp_proxy_server; struct tcp_proxy_server_conn* rc = tcp_proxy_server_find_conn(ctx, stream_id); if (!rc) { DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "TCP proxy server: ERROR sid=%08x — no conn", stream_id); return; } DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "SOCK:ERROR_RECV fd=%d sid=%08x total=%d fin=%d write_pend=%d", rc->tc ? (int)rc->tc->sock : -1, stream_id, conn_total(rc), rc->tc ? rc->tc->fin : 0, rc->tc ? write_pending(rc->tc) : 0); tcp_proxy_server_handle_close(inst, stream_id); } void tcp_proxy_server_handle_fin(struct UTUN_INSTANCE* inst, uint32_t stream_id) { if (!inst) return; struct tcp_proxy_server* ctx = &inst->tcp_proxy_server; struct tcp_proxy_server_conn* rc = tcp_proxy_server_find_conn(ctx, stream_id); if (!rc || !rc->tc || !rc->tc->connected) return; shutdown(rc->tc->sock, SHUT_WR); } // ==================================================================== // Инициализация / деинициализация // ==================================================================== int tcp_proxy_server_init(struct UTUN_INSTANCE* inst) { if (!inst) return -1; struct tcp_proxy_server* ctx = &inst->tcp_proxy_server; memset(ctx, 0, sizeof(*ctx)); if (!inst->config) return 0; ctx->enabled = inst->config->global.tcp_proxy_server_enabled; ctx->inst = inst; if (!ctx->enabled) return 0; g_tcp_proxy_server_ctx = ctx; etcp_router_bind(inst, ETCP_ID_TCP_PROXY, tcp_proxy_client_etcp_recv_cb); if (!inst->config->global.tcp_proxy_client_enabled) { udp_proxy_init(inst, inst->ua); icmp_proxy_init(inst, inst->ua); } DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "TCP proxy server initialized node=%016llx", (unsigned long long)inst->node_id); return 0; } void tcp_proxy_server_destroy(struct UTUN_INSTANCE* inst) { if (!inst) return; struct tcp_proxy_server* ctx = &inst->tcp_proxy_server; DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "TCP proxy server destroying: conn_count=%d", ctx->conn_count); struct tcp_proxy_server_conn* rc = ctx->conns; while (rc) { struct tcp_proxy_server_conn* next = rc->next; tcp_proxy_server_conn_free(rc); rc = next; } ctx->conns = NULL; ctx->enabled = 0; if (g_tcp_proxy_server_ctx == ctx) g_tcp_proxy_server_ctx = NULL; udp_proxy_destroy(inst); icmp_proxy_destroy(inst); DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "TCP proxy server destroyed"); }