diff --git a/src/proxy/tcp_proxy_client.c b/src/proxy/tcp_proxy_client.c index 9ffb2965..20378994 100644 --- a/src/proxy/tcp_proxy_client.c +++ b/src/proxy/tcp_proxy_client.c @@ -229,7 +229,7 @@ static err_t tcp_proxy_client_accept_cb(void *arg, struct tcp_pcb *newpcb, err_t pc->to_lwip = queue_new(p->ua, 0, 0, 0, "to_lwip"); - if (tcp_proxy_client_send_connect(pc) != 0) { DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "TCP proxy client: send_connect failed sid=%08x to %d.%d.%d.%d:%d", pc->stream_id, pc->dest_ip[0], pc->dest_ip[1], pc->dest_ip[2], pc->dest_ip[3], ntohs(pc->dest_port)); u_free(pc); return LERR_MEM; } + if (tcp_proxy_client_send_connect(pc) != 0) { DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "TCP proxy client: send_connect failed sid=%08x to %d.%d.%d.%d:%d", pc->stream_id, pc->dest_ip[0], pc->dest_ip[1], pc->dest_ip[2], pc->dest_ip[3], ntohs(pc->dest_port)); tcp_arg(newpcb, NULL); tcp_abort(newpcb); queue_free(pc->to_lwip); u_free(pc); return LERR_MEM; } pc->next = p->conns; p->conns = pc; p->conn_count++; DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "TCP proxy client: new conn sid=%08x local_port=%u -> %d.%d.%d.%d:%d total_conns=%d snd_wnd=%u mss=%u", @@ -579,6 +579,8 @@ void tcp_proxy_client_destroy(struct tcp_proxy_client* p) { while ((e = queue_data_get(pc->to_lwip))) { queue_dgram_free(e); queue_entry_free(e); } queue_free(pc->to_lwip); pc->to_lwip = NULL; } + if (pc->tx_buf) { u_free(pc->tx_buf); pc->tx_buf = NULL; pc->tx_len = 0; } + if (p->inst) etcp_router_waiter_cancel(p->inst, p->via_node_id, &pc->tx_waiter); u_free(pc); pc = next; } p->conns = NULL; diff --git a/src/proxy/tcp_proxy_server.c b/src/proxy/tcp_proxy_server.c index 251dcace..0d792866 100644 --- a/src/proxy/tcp_proxy_server.c +++ b/src/proxy/tcp_proxy_server.c @@ -37,6 +37,8 @@ static void tcp_proxy_server_close_retry_cb(void* arg); static void tcp_proxy_server_sock_send(struct tcp_proxy_server_conn* rc, const uint8_t* data, size_t len); void tcp_proxy_server_conn_free(struct tcp_proxy_server_conn* rc); +// Отправка сообщения через ETCP-маршрут: выделяет queue_entry, заполняет заголовок TCP_PROXY_HDR_SIZE, вызывает etcp_route_send. +// Вызывается из всех send-функций (close/error/data). static int tcp_proxy_server_send_msg(struct UTUN_INSTANCE* inst, uint64_t dst, uint8_t subcmd, uint32_t sid, const uint8_t* data, size_t len) { struct ll_entry* e = queue_entry_new(0); @@ -51,6 +53,7 @@ static int tcp_proxy_server_send_msg(struct UTUN_INSTANCE* inst, uint64_t dst, u return etcp_route_send(inst, dst, e); } +// Таймер-коллбэк повтора CLOSE/ERROR с экспоненциальным backoff (50..5000 tb). Запускается из send_close/send_error при неудаче etcp_route_send. static void tcp_proxy_server_close_retry_cb(void* arg) { struct tcp_proxy_server_conn* rc = (struct tcp_proxy_server_conn*)arg; if (!rc) return; @@ -67,6 +70,7 @@ static void tcp_proxy_server_close_retry_cb(void* arg) { rc->close_pending = 0; } +// Отправка TCP_PROXY_SUBCMD_CLOSE пиру; при неудаче — close_pending=1 + retry-таймер. Вызывается из sock_read_cb (EOF), sock_retry_cb (завершение send). static void tcp_proxy_server_send_close(struct tcp_proxy_server_conn* rc) { struct UTUN_INSTANCE* inst = rc->ctx ? rc->ctx->inst : NULL; if (!inst) return; @@ -78,6 +82,7 @@ static void tcp_proxy_server_send_close(struct tcp_proxy_server_conn* rc) { } } +// Отправка TCP_PROXY_SUBCMD_ERROR пиру; при неудаче — close_pending=1 + retry-таймер. Вызывается из sock_write_cb (ошибка connect), sock_error_cb. static void tcp_proxy_server_send_error(struct tcp_proxy_server_conn* rc) { struct UTUN_INSTANCE* inst = rc->ctx ? rc->ctx->inst : NULL; if (!inst) return; @@ -89,6 +94,7 @@ static void tcp_proxy_server_send_error(struct tcp_proxy_server_conn* rc) { } } +// Подсчёт всех коннектов в ctx->conns (для диагностики). static int tcp_proxy_server_conn_total(struct tcp_proxy_server_conn* rc) { if (!rc || !rc->ctx) return 0; int n = 0; struct tcp_proxy_server_conn* c; @@ -96,6 +102,8 @@ static int tcp_proxy_server_conn_total(struct tcp_proxy_server_conn* rc) { return n; } +// Неблокирующий send в TCP-сокет с обработкой EAGAIN/EWOULDBLOCK/EINTR. Возвращает: 0=done, 1=need_retry, -1=error. +// Вызывается из sock_send и sock_retry_cb. // ==================================================================== // Неблокирующий send с EAGAIN и таймером повтора // ==================================================================== @@ -118,6 +126,7 @@ static int tcp_proxy_server_sock_try_send(struct tcp_proxy_server_conn* rc) { return 0; } +// Таймер-коллбэк повтора sock_try_send с backoff (10..5000 tb). При успешном завершении — shutdown(SHUT_WR) если cli_closed, send_close если sock_closed, conn_free если обе стороны закрыты. static void tcp_proxy_server_sock_retry_cb(void* arg) { struct tcp_proxy_server_conn* rc = (struct tcp_proxy_server_conn*)arg; if (!rc || !rc->out_buf) { rc->out_timer = NULL; return; } @@ -137,6 +146,8 @@ static void tcp_proxy_server_sock_retry_cb(void* arg) { } } +// Буферизованная неблокирующая отправка данных в TCP-сокет. Если уже есть out_buf — дописывает (realloc). Иначе sock_try_send, при EAGAIN ставит retry-таймер. +// Вызывается из sock_write_cb (сброс после connect) и handle_data (приём DATA от клиента). static void tcp_proxy_server_sock_send(struct tcp_proxy_server_conn* rc, const uint8_t* data, size_t len) { if (!rc || !data || len == 0) return; if (rc->out_buf) { @@ -160,6 +171,9 @@ static void tcp_proxy_server_sock_send(struct tcp_proxy_server_conn* rc, const u } } +// [queue] Коллбэк освобождения очереди ETCP (backpressure relief). Срабатывает когда etcp_router_waiter сигнализирует что очередь освободилась. +// Отправляет сохранённый в pause_buf буфер через ETCP, затем восстанавливает uasync-чтение из сокета. +// Цепочка: queue waiter (threshold) → etcp_router → эта функция. static void tcp_proxy_server_pause_waiter_cb(struct ll_queue* q, void* arg) { (void)q; struct tcp_proxy_server_conn* rc = (struct tcp_proxy_server_conn*)arg; @@ -179,8 +193,12 @@ static void tcp_proxy_server_pause_waiter_cb(struct ll_queue* q, void* arg) { } // ==================================================================== -// Socket callbacks +// Socket callbacks — вызываются из uasync при событиях на TCP-сокете // ==================================================================== + +// Сокет-коллбэк чтения. Читает данные из TCP-сокета, отправляет через ETCP пиру (DATA). +// При ошибке ETCP-send (backpressure) — буферизует в pause_buf, снимает сокет с uasync, регистрирует waiter в очереди. +// При EOF (n=0) — sock_closed=1, send_close. static void tcp_proxy_server_sock_read_cb(socket_t sock, void* arg) { (void)sock; struct tcp_proxy_server_conn* rc = (struct tcp_proxy_server_conn*)arg; if (!rc || rc->sock == SOCKET_INVALID) return; @@ -220,6 +238,7 @@ static void tcp_proxy_server_sock_read_cb(socket_t sock, void* arg) { } } +// Сокет-коллбэк записи. Срабатывает при завершении неблокирующего connect(). Проверяет SO_ERROR, при успехе — connected=1, сбрасывает pending_buf в сокет. static void tcp_proxy_server_sock_write_cb(socket_t sock, void* arg) { (void)sock; struct tcp_proxy_server_conn* rc = (struct tcp_proxy_server_conn*)arg; if (!rc || rc->sock == SOCKET_INVALID || rc->connected) return; @@ -246,6 +265,7 @@ static void tcp_proxy_server_sock_write_cb(socket_t sock, void* arg) { } } +// Сокет-коллбэк ошибки. Асинхронная ошибка сокета — отправляет ERROR пиру и освобождает коннект через conn_free. static void tcp_proxy_server_sock_error_cb(socket_t sock, void* arg) { (void)sock; struct tcp_proxy_server_conn* rc = (struct tcp_proxy_server_conn*)arg; if (!rc || rc->sock == SOCKET_INVALID) return; @@ -255,6 +275,8 @@ static void tcp_proxy_server_sock_error_cb(socket_t sock, void* arg) { tcp_proxy_server_conn_free(rc); } +// Полное освобождение коннекта: удаление из ctx->conns, закрытие сокета, отмена таймеров (out/close), освобождение буферов (pending/out/pause), отмена queue waiter. +// Вызывается из всех точек завершения/очистки коннекта. void tcp_proxy_server_conn_free(struct tcp_proxy_server_conn* rc) { if (!rc) return; int total = tcp_proxy_server_conn_total(rc); @@ -278,14 +300,18 @@ void tcp_proxy_server_conn_free(struct tcp_proxy_server_conn* rc) { } // ==================================================================== -// Public API +// Public API — входные точки из ETCP-диспетчера (tcp_proxy_client_etcp_recv_cb) // ==================================================================== + +// Поиск коннекта по stream_id в цепочке ctx->conns. 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; } +// Входная точка: обработка CONNECT от клиента. Создаёт коннект, инициирует неблокирующий connect() к dest_ip:dest_port, регистрирует сокет в uasync. +// Вызывается из ETCP-диспетчера при приёме пакета с subcmd=CONNECT. 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; } @@ -307,7 +333,7 @@ int tcp_proxy_server_handle_connect(struct UTUN_INSTANCE* inst, struct ll_entry* (int)rc->sock, stream_id, dest_ip[0],dest_ip[1],dest_ip[2],dest_ip[3],ntohs(dest_port), ctx->conn_count); socket_set_nonblocking(rc->sock); rc->read_id = uasync_add_socket_t(rc->ua, rc->sock, tcp_proxy_server_sock_read_cb, tcp_proxy_server_sock_write_cb, tcp_proxy_server_sock_error_cb, rc); - if (!rc->read_id) { DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "TCP proxy server: uasync_add_socket_t failed"); u_free(rc); queue_dgram_free(entry); queue_entry_free(entry); return -1; } + if (!rc->read_id) { DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "TCP proxy server: uasync_add_socket_t failed"); ctx->conn_count--; socket_close_wrapper(rc->sock); u_free(rc); queue_dgram_free(entry); queue_entry_free(entry); return -1; } struct sockaddr_in addr; memset(&addr, 0, sizeof(addr)); addr.sin_family = AF_INET; memcpy(&addr.sin_addr.s_addr, dest_ip, 4); addr.sin_port = dest_port; @@ -316,7 +342,7 @@ int tcp_proxy_server_handle_connect(struct UTUN_INSTANCE* inst, struct ll_entry* 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)); tcp_proxy_server_send_msg(inst, src_node_id, TCP_PROXY_SUBCMD_ERROR, stream_id, NULL, 0); - uasync_remove_socket_t(rc->ua, rc->sock); socket_close_wrapper(rc->sock); u_free(rc); + uasync_remove_socket_t(rc->ua, rc->sock); socket_close_wrapper(rc->sock); ctx->conn_count--; u_free(rc); queue_dgram_free(entry); queue_entry_free(entry); return -1; } @@ -325,6 +351,8 @@ int tcp_proxy_server_handle_connect(struct UTUN_INSTANCE* inst, struct ll_entry* return 0; } +// Входная точка: обработка DATA от клиента. Пересылает данные в целевой TCP-сокет через sock_send; если connect не завершён — буферизует в pending_buf. +// Вызывается из ETCP-диспетчера при приёме DATA. int tcp_proxy_server_handle_data(struct UTUN_INSTANCE* inst, struct ETCP_CONN* conn, struct ll_entry* entry, uint32_t stream_id) { if (!inst) { queue_dgram_free(entry); queue_entry_free(entry); return -1; } struct tcp_proxy_server* ctx = &inst->tcp_proxy_server; @@ -352,6 +380,8 @@ int tcp_proxy_server_handle_data(struct UTUN_INSTANCE* inst, struct ETCP_CONN* c return 0; } +// Входная точка: обработка CLOSE от клиента. cli_closed=1, shutdown(SHUT_WR) сокета, при sock_closed && !out_buf — conn_free. +// Вызывается из ETCP-диспетчера при приёме CLOSE. 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; @@ -372,6 +402,8 @@ void tcp_proxy_server_handle_close(struct UTUN_INSTANCE* inst, uint32_t stream_i DEBUG_WARN(DEBUG_CATEGORY_SOCKET, "TCP proxy server: CLOSE sid=%08x — no conn, dropping", stream_id); } +// Инициализация TCP proxy server: читает enabled из конфига, биндит tcp_proxy_client_etcp_recv_cb на ETCP_ID_TCP_PROXY. +// Если клиент не включён — инициализирует также UDP/ICMP proxy. Вызывается из utun_instance_init. int tcp_proxy_server_init(struct UTUN_INSTANCE* inst) { if (!inst) return -1; struct tcp_proxy_server* ctx = &inst->tcp_proxy_server; @@ -390,6 +422,7 @@ int tcp_proxy_server_init(struct UTUN_INSTANCE* inst) { return 0; } +// Деинициализация: освобождает все коннекты через conn_free, сбрасывает ctx, вызывает udp/icmp proxy destroy. Вызывается из utun_instance_destroy. void tcp_proxy_server_destroy(struct UTUN_INSTANCE* inst) { if (!inst) return; struct tcp_proxy_server* ctx = &inst->tcp_proxy_server;