diff --git a/doc/proxy_protocol.md b/doc/proxy_protocol.md new file mode 100644 index 00000000..7332da77 --- /dev/null +++ b/doc/proxy_protocol.md @@ -0,0 +1,61 @@ +# Proxy: потоки TCP и датаграммы + +Клиентские входы — TUN (lwIP), SOCKS5 CONNECT и HTTP proxy/CONNECT. Exit открывает +обычный TCP-сокет к назначению. TCP-поток на exit определяется парой `(peer_node_id, +stream_id)`. Клиент принимает ответы и уведомления restart только от настроенного +exit в UTUN-группе. UDP/ICMP-контексты принадлежат конкретному `UTUN_INSTANCE`. + +## TCP + +Общие определения находятся в `src/proxy/proxy_protocol.h`. Передаваемый заголовок: +`svc_id:1, command:1, stream_id:4`; поля маршрута добавляет ETCP-router. + +- CONNECT (1): IPv4:4 и порт:2 в сетевом порядке. +- CONNECTED (2): exit подтвердил успешный TCP connect. Только после этого клиент + сообщает SOCKS success / HTTP 200 и начинает передавать накопленные данные. +- DATA (3): от 1 до 4096 байт. +- CLOSE (4), ERROR (5): прекращение потока. +- FIN (6): конец одного направления; все предшествующие DATA должны быть переданы. +- WINDOW (7): uint32_t, число дополнительно разрешённых байт. + +У каждого направления начальное окно 65536 байт. DATA уменьшает окно отправителя +и получателя. Получатель возвращает кредит после записи в локальный TCP-сокет +либо принятия данных ограниченным send-buffer lwIP. Очередь ETCP сама по себе +не является подтверждением потребления данных конечным TCP. + +FIN ждёт освобождения исходящей очереди, включая данные, остановленные окном +или backpressure роутера. Закрытие одного направления не прекращает другое. +Окончательное освобождение происходит после обоих FIN и передачи оставшихся данных. +Непереданные управляющие сообщения повторяются таймером. Ошибки и переходы +состояния диагностируются категорией `proxy`; состояние сокетов — `socket`. + +Парсеры сохраняют остаток после greeting/request/HTTP headers. HTTP CONNECT ждёт +полного блока заголовков. Во время DNS/connect чтение ограничено очередью tcp_io; +накопленные HTTP headers/body передаются частями размером DATA. + +## UDP и ICMP + +REQUEST и REPLY различаются явно. Ответ допускается только от настроенного exit. +UDP-ответ восстанавливается с нулевой UDP checksum (допустимо для IPv4), IPv4 checksum +пересчитывается. ICMP checksum учитывает нечётный последний байт без чтения за буфером. + +Exit заменяет ICMP id/sequence уникальной среди ожидающих запросов парой и сопоставляет +ответ также с IP назначения. Исходные id/sequence восстанавливаются для клиента. +Это исключает смешивание ping-запросов разных клиентов с одинаковыми исходными ID. + +## Проверки и ограничения + +`test_proxy_packets` проверяет IPv4/UDP/ICMP и payload длиной 0..1501. +`test_tcp_io_flush` проверяет уведомление после последнего синхронного send. +`test_proxy_regressions` использует реальные локальные сокеты и управляемую доставку +ETCP: совместные/раздельные handshake, ранние данные, отложенный CONNECTED, HTTP 502, +большой POST, одинаковые stream ID разных пиров, FIN при закрытом окне, RST и IPv4 options. + +На Linux полный `./check.sh`: 109 passed, 0 failed, 1 skipped; дополнительно прошли +настоящий TUN, full-duplex burst и нагрузка 8 соединений / 128 МиБ. +Windows/FreeBSD в этом прогоне не проверялись. + +SOCKS IPv6 явно отклоняется кодом address type not supported. IPv4-фрагменты +на входе proxy TUN отклоняются с диагностикой; сборка фрагментов не реализована. +Формат TCP дополнен CONNECTED/WINDOW, UDP-команды разделены: клиент и exit необходимо +обновлять вместе; совместимость со старым proxy-протоколом не предусмотрена. diff --git a/lib/tcp_io.c b/lib/tcp_io.c index 7f830648..c2b648d2 100644 --- a/lib/tcp_io.c +++ b/lib/tcp_io.c @@ -282,6 +282,8 @@ static int flush_write_buf(struct tcp_conn* tc) { return 0; } +// Последний send может опустошить очередь при уже выключенном POLLOUT. +// Уведомляем потребителя здесь, не рассчитывая на ещё один write event. static void notify_flushed(struct tcp_conn* tc) { if (tc->destroyed || tc->write_buf || tc->write_queue->head || !tc->on_flushed) return; void (*cb)(struct tcp_conn*, void*) = tc->on_flushed; @@ -316,7 +318,7 @@ static void write_queue_fetch_cb(struct ll_queue* q, void* arg) { if (!e->dgram) { DEBUG_DEBUG(DEBUG_CATEGORY_SOCKET, "tcp_io: CLOSE fd=%d", (int)tc->sock); queue_entry_free(e); queue_resume_callback(q); - uasync_remove_socket_t(tc->ua, tc->sock); + if (tc->socket_id) uasync_remove_socket_t(tc->ua, tc->sock); socket_close_wrapper(tc->sock); tc->socket_id = NULL; tc->sock = SOCKET_INVALID; tc->closed = 1; if (tc->on_closed) tc->on_closed(tc, tc->arg); @@ -394,6 +396,8 @@ static void write_cb(socket_t sock, void* arg) { tc->connected = 1; if (tc->connect_timer) { uasync_cancel_timeout(tc->ua, tc->connect_timer); tc->connect_timer = NULL; } DEBUG_DEBUG(DEBUG_CATEGORY_SOCKET, "tcp_io: connect ok fd=%d", (int)tc->sock); + if (tc->on_connected) tc->on_connected(tc, tc->arg); + if (tc->destroyed) return; } else { DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "tcp_io: connect fail fd=%d err=%d", (int)tc->sock, err); tcp_conn_handle_error(tc, err ? err : -1); @@ -426,6 +430,14 @@ static void error_cb(socket_t sock, void* arg) { // Peer полностью закрыл свой сокет (EPOLLHUP) после того, как FIN уже прочитан (fin_remote=1). // EPOLLHUP — level-triggered и «always reported»: если оставить сокет в epoll, каждая итерация // epoll_wait мгновенно вернёт HUP и error_cb зациклится на 100% CPU. Закрываем сокет. + if (tc->fin_remote && tc->fin_local) { + // Обе TCP-половины закрыты, но владелец ещё может передавать ранее прочитанные данные. + // Убираем level-triggered HUP из poll; CLOSE остаётся решением владельца. + DEBUG_DEBUG(DEBUG_CATEGORY_SOCKET, "tcp_io: both FINs fd=%d — detach poll, await owner drain", (int)tc->sock); + if (tc->socket_id) uasync_remove_socket_t(tc->ua, tc->sock); + tc->socket_id = NULL; tc->write_monitor = 0; + return; + } if (tc->fin_remote) { DEBUG_DEBUG(DEBUG_CATEGORY_SOCKET, "tcp_io: async HUP (graceful) fd=%d — peer fully closed after FIN, closing", (int)tc->sock); diff --git a/lib/tcp_io.h b/lib/tcp_io.h index 39480036..e3dab4c9 100644 --- a/lib/tcp_io.h +++ b/lib/tcp_io.h @@ -70,6 +70,7 @@ struct tcp_conn { struct queue_waiter_handle read_waiter; // Коллбэки + void (*on_connected)(struct tcp_conn* tc, void* arg); // исходящий connect завершён void (*on_fin)(struct tcp_conn* tc, void* arg); // FIN получен от удалённой стороны (fin_remote=1) void (*on_fin_sent)(struct tcp_conn* tc, void* arg); // FIN отправлен удалённой стороне (fin_local=1) void (*on_error)(struct tcp_conn* tc, int err, void* arg); // фатальная ошибка (вызывается после close сокета; обязан вызвать tcp_conn_destroy) diff --git a/src/Makefile.am b/src/Makefile.am index 6ebe70db..089177f2 100644 --- a/src/Makefile.am +++ b/src/Makefile.am @@ -67,6 +67,7 @@ utun_CORE_SOURCES = \ transport_layer/dummynet.c \ ntp_time.c \ ntp_node_time.c \ + proxy/proxy_protocol.h \ proxy/tcp_proxy_client.c \ routing_layer/etcp_router.c \ routing_layer/route_crypto.c \ @@ -171,6 +172,7 @@ libutun_a_SOURCES = \ transport_layer/dummynet.c \ ntp_time.c \ ntp_node_time.c \ + proxy/proxy_protocol.h \ proxy/tcp_proxy_client.c \ routing_layer/etcp_router.c \ routing_layer/route_crypto.c \ diff --git a/src/proxy/proxy_protocol.h b/src/proxy/proxy_protocol.h new file mode 100644 index 00000000..c0ea372f --- /dev/null +++ b/src/proxy/proxy_protocol.h @@ -0,0 +1,157 @@ +// Общий TCP proxy протокол. DATA ограничен окном каждого потока; WINDOW возвращает +// кредит после потребления данных локальным TCP. FIN следует за всеми DATA своего +// направления. CONNECTED подтверждает реальный connect exit → destination. +#ifndef PROXY_PROTOCOL_H +#define PROXY_PROTOCOL_H + +#include +#include +#include "../routing_layer/etcp_router.h" +#include "../routing_layer/topo_node.h" +#include "../../lib/u_async.h" +#include "../../lib/ll_queue.h" +#include "../../lib/mem.h" +#include "../../lib/debug_config.h" + +#define TCP_PROXY_SUBCMD_CONNECT 0x01 +#define TCP_PROXY_SUBCMD_CONNECTED 0x02 +#define TCP_PROXY_SUBCMD_DATA 0x03 +#define TCP_PROXY_SUBCMD_CLOSE 0x04 +#define TCP_PROXY_SUBCMD_ERROR 0x05 +#define TCP_PROXY_SUBCMD_FIN 0x06 +#define TCP_PROXY_SUBCMD_WINDOW 0x07 +#define TCP_PROXY_HDR_SIZE 6 +#define TCP_PROXY_CONNECT_HDR_SIZE 12 +#define TCP_PROXY_RECV_HDR_SIZE (ROUTER_SVC_PAYLOAD_OFF + 5) +#define TCP_PROXY_CHUNK 4096 +#define TCP_PROXY_WINDOW 65536u + +struct proxy_flow { + struct UTUN_INSTANCE* inst; + struct UASYNC* ua; + uint64_t peer; + uint32_t sid; + uint8_t svc; + uint8_t ready; + uint8_t connected_pending; + uint8_t fin_pending; + uint8_t fin_sent; + uint8_t fin_received; + uint32_t tx_credit; + uint32_t rx_credit; + uint32_t consumed; + void* retry; + void (*wake)(void* arg); + void* arg; +}; + +// Сообщение целиком передаётся маршрутизатору, который освобождает entry при любом результате. +static inline int proxy_send(struct UTUN_INSTANCE* inst, uint64_t peer, uint8_t svc, uint8_t cmd, + uint32_t sid, const uint8_t* data, size_t len, int force) { + if (len > TCP_PROXY_CHUNK) { + DEBUG_ERROR(DEBUG_CATEGORY_PROXY, "proxy: oversized message sid=%08x cmd=%u len=%zu", sid, cmd, len); + return -1; + } + struct ll_entry* e = queue_entry_new(0); + if (!e) { DEBUG_ERROR(DEBUG_CATEGORY_PROXY, "proxy: entry allocation failed sid=%08x", sid); return -1; } + e->dgram = u_malloc(TCP_PROXY_HDR_SIZE + len); + if (!e->dgram) { + DEBUG_ERROR(DEBUG_CATEGORY_PROXY, "proxy: payload allocation failed sid=%08x", sid); + queue_entry_free(e); return -1; + } + e->dgram[0] = svc; e->dgram[1] = cmd; + memcpy(e->dgram + 2, &sid, 4); + if (len) memcpy(e->dgram + TCP_PROXY_HDR_SIZE, data, len); + e->len = TCP_PROXY_HDR_SIZE + len; + return etcp_route_send(inst, TOPO_GROUP_UTUN, peer, e, force, 0); +} + +static inline void proxy_flow_flush(struct proxy_flow* f); + +// Повторяем только недоставленные control-сообщения; DATA остаётся у владельца потока. +static inline void proxy_flow_retry(void* arg) { + struct proxy_flow* f = arg; + f->retry = NULL; + proxy_flow_flush(f); + if (f->wake) f->wake(f->arg); +} + +static inline void proxy_flow_init(struct proxy_flow* f, struct UTUN_INSTANCE* inst, struct UASYNC* ua, + uint64_t peer, uint8_t svc, uint32_t sid, void (*wake)(void*), void* arg) { + memset(f, 0, sizeof(*f)); + f->inst = inst; f->ua = ua; f->peer = peer; f->svc = svc; f->sid = sid; + f->tx_credit = TCP_PROXY_WINDOW; f->rx_credit = TCP_PROXY_WINDOW; + f->wake = wake; f->arg = arg; +} + +static inline void proxy_flow_destroy(struct proxy_flow* f) { + if (f->retry) { uasync_cancel_timeout(f->ua, f->retry); f->retry = NULL; } +} + +// Сначала CONNECTED, затем WINDOW и FIN. Force обходит лишь очередь router, не окно потока. +static inline void proxy_flow_flush(struct proxy_flow* f) { + if (f->connected_pending) { + if (proxy_send(f->inst, f->peer, f->svc, TCP_PROXY_SUBCMD_CONNECTED, f->sid, NULL, 0, 1) < 0) goto retry; + f->connected_pending = 0; f->ready = 1; + } + if (f->consumed) { + uint32_t credit = f->consumed; + // Резервируем до синхронной loopback-доставки. + f->consumed = 0; f->rx_credit += credit; + if (proxy_send(f->inst, f->peer, f->svc, TCP_PROXY_SUBCMD_WINDOW, f->sid, (uint8_t*)&credit, 4, 1) < 0) { + f->consumed += credit; f->rx_credit -= credit; goto retry; + } + } + if (f->fin_pending && f->ready) { + if (proxy_send(f->inst, f->peer, f->svc, TCP_PROXY_SUBCMD_FIN, f->sid, NULL, 0, 1) < 0) goto retry; + f->fin_pending = 0; f->fin_sent = 1; + DEBUG_DEBUG(DEBUG_CATEGORY_PROXY, "proxy: FIN queued peer=%016llx sid=%08x", (unsigned long long)f->peer, f->sid); + } + return; +retry: + DEBUG_DEBUG(DEBUG_CATEGORY_PROXY, "proxy: control retry sid=%08x connected=%u credit=%u fin=%u", + f->sid, f->connected_pending, f->consumed, f->fin_pending); + if (!f->retry) f->retry = uasync_set_timeout(f->ua, 5000, f, proxy_flow_retry, "proxy_control"); +} + +// 0 = отправлено, -2 = ждём окно/CONNECTED, -1 = router backpressure/ошибка выделения. +static inline int proxy_flow_send(struct proxy_flow* f, const uint8_t* data, size_t len, int force) { + if (!f->ready || len > f->tx_credit) return -2; + if (!len || len > TCP_PROXY_CHUNK || f->fin_sent) { + DEBUG_ERROR(DEBUG_CATEGORY_PROXY, "proxy: invalid DATA sid=%08x len=%zu fin=%u", f->sid, len, f->fin_sent); + return -1; + } + f->tx_credit -= len; + int ret = proxy_send(f->inst, f->peer, f->svc, TCP_PROXY_SUBCMD_DATA, f->sid, data, len, force); + if (ret < 0) f->tx_credit += len; + return ret; +} + +static inline int proxy_flow_receive(struct proxy_flow* f, size_t len) { + if (!len || len > TCP_PROXY_CHUNK || len > f->rx_credit || f->fin_received) { + DEBUG_ERROR(DEBUG_CATEGORY_PROXY, "proxy: invalid receive sid=%08x len=%zu credit=%u fin=%u", + f->sid, len, f->rx_credit, f->fin_received); + return -1; + } + f->rx_credit -= len; + return 0; +} + +static inline int proxy_flow_window(struct proxy_flow* f, const uint8_t* data, size_t len) { + uint32_t credit = 0; + if (len == 4) memcpy(&credit, data, 4); + if (!credit || credit > TCP_PROXY_WINDOW - f->tx_credit) { + DEBUG_ERROR(DEBUG_CATEGORY_PROXY, "proxy: invalid WINDOW sid=%08x len=%zu credit=%u tx=%u", f->sid, len, credit, f->tx_credit); + return -1; + } + f->tx_credit += credit; + DEBUG_TRACE(DEBUG_CATEGORY_PROXY, "proxy: WINDOW sid=%08x add=%u available=%u", f->sid, credit, f->tx_credit); + return 0; +} + +// Вызывать только после передачи байтов локальному TCP, не после получения ETCP DATA. +static inline void proxy_flow_consume(struct proxy_flow* f, uint32_t len) { + f->consumed += len; + proxy_flow_flush(f); +} +#endif diff --git a/src/proxy/socks_proxy.c b/src/proxy/socks_proxy.c index 9527b74d..4a92c717 100644 --- a/src/proxy/socks_proxy.c +++ b/src/proxy/socks_proxy.c @@ -23,6 +23,8 @@ #include #endif +static void socks_flow_wake(void* arg); +static void socks_flushed_cb(struct tcp_conn* tc, void* arg); static void on_accept_cb(socket_t sock, void* arg); static void on_read_cb(struct ll_queue* q, void* arg); static void on_fin_cb(struct tcp_conn* tc, void* arg); @@ -58,24 +60,10 @@ struct listen_ctx { // ==================================================================== // Отправка сообщений через ETCP // ==================================================================== -static int send_msg(struct UTUN_INSTANCE* inst, uint64_t group_id, 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_PROXY, "socks_proxy: 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_PROXY, "socks_proxy: malloc(%zu) failed", TCP_PROXY_HDR_SIZE + len); queue_entry_free(e); return -1; } - e->dgram[0] = ETCP_RT_ID_TCP_PROXY_SERVER; - e->dgram[1] = subcmd; - memcpy(e->dgram + 2, &sid, 4); - if (len > 0) memcpy(e->dgram + TCP_PROXY_HDR_SIZE, data, len); - if (TCP_PROXY_HDR_SIZE + len > UINT16_MAX) { - DEBUG_ERROR(DEBUG_CATEGORY_PROXY, "socks_proxy: msg too large len=%zu subcmd=%02x", len, subcmd); - queue_dgram_free(e); queue_entry_free(e); return -1; - } - e->len = (uint16_t)(TCP_PROXY_HDR_SIZE + len); - int ret = etcp_route_send(inst, group_id, dst, e, force, 0); - DEBUG_TRACE(DEBUG_CATEGORY_PROXY, "SOCKS send_msg subcmd=%02x sid=%08x len=%zu force=%d → ret=%d", - subcmd, sid, len, force, ret); - return ret; +static int send_msg(struct UTUN_INSTANCE* inst, uint64_t group_id, uint64_t dst, uint8_t subcmd, + uint32_t sid, const uint8_t* data, size_t len, int force) { + (void)group_id; + return proxy_send(inst, dst, ETCP_RT_ID_TCP_PROXY_SERVER, subcmd, sid, data, len, force); } static int send_connect(struct socks_proxy_conn* c) { @@ -88,7 +76,7 @@ static int send_connect(struct socks_proxy_conn* c) { } static int send_data(struct socks_proxy_conn* c, const uint8_t* data, uint16_t len, int force) { - return send_msg(c->inst, TOPO_GROUP_UTUN, c->via_node_id, TCP_PROXY_SUBCMD_DATA, c->stream_id, data, len, force); + return proxy_flow_send(&c->flow, data, len, force); } static void send_close(struct socks_proxy_conn* c) { @@ -97,9 +85,10 @@ static void send_close(struct socks_proxy_conn* c) { else { c->close_pending = 0; c->close_sent = 1; } } -static void send_fin(struct socks_proxy_conn* c) { - DEBUG_TRACE(DEBUG_CATEGORY_PROXY, "SOCKS send FIN sid=%08x", c->stream_id); - send_msg(c->inst, TOPO_GROUP_UTUN, c->via_node_id, TCP_PROXY_SUBCMD_FIN, c->stream_id, NULL, 0, 1); +// Удалить только разобранные байты, сохранив следующий запрос/данные того же recv(). +static void consume_header(struct socks_proxy_conn* c, uint16_t len) { + c->buf_len -= len; + memmove(c->buf, c->buf + len, c->buf_len); } // ==================================================================== @@ -145,7 +134,7 @@ static void process_socks_greeting(struct socks_proxy_conn* c) { DEBUG_DEBUG(DEBUG_CATEGORY_PROXY, "socks_proxy: greeting ver=%d nmethods=%d", ver, nmethods); uint8_t reply[] = { 0x05, 0x00 }; write_to_client(c, reply, 2); - c->buf_len = 0; + consume_header(c, 2 + nmethods); c->state = SOCKS_STATE_REQUEST; } @@ -170,7 +159,7 @@ static void process_socks_request(struct socks_proxy_conn* c) { if (atyp == 1) { memcpy(c->dest_ip, c->buf + 4, 4); memcpy(&c->dest_port, c->buf + 8, 2); - c->buf_len = 0; + consume_header(c, need); socks_finalize(c); return; } else if (atyp == 4) { @@ -183,7 +172,7 @@ static void process_socks_request(struct socks_proxy_conn* c) { uint8_t dlen = c->buf[4]; char domain[256]; memcpy(domain, c->buf + 5, dlen); domain[dlen] = '\0'; memcpy(&c->dest_port, c->buf + 5 + dlen, 2); - c->buf_len = 0; + consume_header(c, need); socks_issue_dns(c, domain, DNSK_SOCKS); return; @@ -204,7 +193,11 @@ static void process_http_request(struct socks_proxy_conn* c) { uint16_t line_len = (uint16_t)((uint8_t*)line_end - c->buf); if (line_len < 8) { DEBUG_ERROR(DEBUG_CATEGORY_PROXY, "socks_proxy: http request too short"); goto error; } - char line[512]; if (line_len > sizeof(line) - 1) line_len = sizeof(line) - 1; + char line[512]; + if (line_len >= sizeof(line)) { DEBUG_WARN(DEBUG_CATEGORY_PROXY, "socks_proxy: request line too long"); goto error; } + char* hdr_end = memmem(c->buf, c->buf_len, "\r\n\r\n", 4); + if (!hdr_end) return; + uint16_t headers_len = (uint16_t)((uint8_t*)hdr_end - c->buf) + 4; memcpy(line, c->buf, line_len); line[line_len] = '\0'; // ===== CONNECT ===== @@ -222,16 +215,12 @@ static void process_http_request(struct socks_proxy_conn* c) { c->dest_port = htons((uint16_t)port); DEBUG_INFO(DEBUG_CATEGORY_PROXY, "socks_proxy: HTTP CONNECT %s:%d sid=%08x", host_port, port, c->stream_id); - c->buf_len = 0; + consume_header(c, headers_len); socks_issue_dns(c, host_port, DNSK_HTTP_CONNECT); return; } // ===== Не-CONNECT: HTTP-прокси (GET, POST, PUT, HEAD, OPTIONS, ...) ===== - char* hdr_end = memmem(c->buf, c->buf_len, "\r\n\r\n", 4); - if (!hdr_end) return; - - uint16_t headers_len = (uint16_t)((uint8_t*)hdr_end - c->buf) + 4; char method[16] = {0}, url[512] = {0}; if (sscanf(line, "%15s %511s", method, url) < 2) { @@ -287,7 +276,7 @@ static void process_http_request(struct socks_proxy_conn* c) { uint16_t rest_off = line_len + 2; uint16_t rest_len = headers_len - rest_off; - uint8_t hdr_buf[2048]; + uint8_t hdr_buf[sizeof(c->buf) + 512]; int hdr_n = snprintf((char*)hdr_buf, sizeof(hdr_buf), "%s %s%s\r\n", method, path_start, version_str); if (hdr_n < 0 || (size_t)hdr_n + rest_len > sizeof(hdr_buf)) { DEBUG_ERROR(DEBUG_CATEGORY_PROXY, "socks_proxy: header reconstruction overflow hdr_n=%d rest=%u", hdr_n, rest_len); @@ -318,43 +307,40 @@ error: { // ==================================================================== // DNS-резолвинг (неблокирующий) + финализация рукопожатия // ==================================================================== +// DNS готов: запрашиваем exit, успех клиенту сообщаем только после CONNECTED. static void socks_finalize(struct socks_proxy_conn* c) { - switch (c->dns_kind) { - case DNSK_SOCKS: { - uint8_t reply[] = { 0x05, 0x00, 0x00, 0x01, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00 }; - write_to_client(c, reply, 10); - break; - } - case DNSK_HTTP_CONNECT: { - uint8_t resp[] = "HTTP/1.1 200 Connection Established\r\n\r\n"; - write_to_client(c, resp, (uint16_t)strlen((char*)resp)); - break; - } - case DNSK_HTTP_PROXY: - default: - break; // reply нет — сразу CONNECT + данные + c->state = SOCKS_STATE_CONNECTING; + if (send_connect(c) < 0) { + DEBUG_ERROR(DEBUG_CATEGORY_PROXY, "socks_proxy: send_connect failed sid=%08x", c->stream_id); + socks_dns_error_and_close(c); } +} - if (send_connect(c) < 0) { - DEBUG_ERROR(DEBUG_CATEGORY_PROXY, "socks_proxy: send_connect failed sid=%08x kind=%d", c->stream_id, c->dns_kind); - tcp_conn_push_close(c->tc); +// Ответ CONNECTED открывает релей и выпускает сохранённый хвост рукопожатия. +static void socks_connected(struct socks_proxy_conn* c) { + if (c->state != SOCKS_STATE_CONNECTING || c->flow.ready) { + DEBUG_WARN(DEBUG_CATEGORY_PROXY, "socks_proxy: unexpected CONNECTED sid=%08x state=%u", c->stream_id, c->state); return; } - - if (c->dns_kind == DNSK_HTTP_PROXY && c->dns_pending) { - uint8_t* pkt = c->dns_pending; - uint16_t len = c->dns_pending_len; + c->flow.ready = 1; + if (c->dns_kind == DNSK_SOCKS) { + const uint8_t reply[] = {5, 0, 0, 1, 0, 0, 0, 0, 0, 0}; + if (write_to_client(c, reply, sizeof(reply)) < 0) { on_error_cb(c->tc, ENOMEM, c); return; } + } else if (c->dns_kind == DNSK_HTTP_CONNECT) { + const uint8_t reply[] = "HTTP/1.1 200 Connection Established\r\n\r\n"; + if (write_to_client(c, reply, sizeof(reply) - 1) < 0) { on_error_cb(c->tc, ENOMEM, c); return; } + } + c->state = SOCKS_STATE_RELAY; + if (c->dns_pending) { + c->tx_buf = c->dns_pending; c->tx_len = c->dns_pending_len; c->dns_pending = NULL; c->dns_pending_len = 0; - int ret = send_data(c, pkt, len, 0); - if (ret == 0) { - u_free(pkt); - } else { - c->tx_buf = pkt; c->tx_len = len; - socks_bp_register(c); - } + } else if (c->buf_len) { + c->tx_buf = u_malloc(c->buf_len); + if (!c->tx_buf) { on_error_cb(c->tc, ENOMEM, c); return; } + memcpy(c->tx_buf, c->buf, c->buf_len); c->tx_len = c->buf_len; c->buf_len = 0; } - - c->state = c->is_http ? HTTP_STATE_RELAY : SOCKS_STATE_RELAY; + DEBUG_INFO(DEBUG_CATEGORY_PROXY, "socks_proxy: established sid=%08x buffered=%u", c->stream_id, c->tx_len); + socks_flow_wake(c); } static void socks_dns_error_and_close(struct socks_proxy_conn* c) { @@ -365,6 +351,7 @@ static void socks_dns_error_and_close(struct socks_proxy_conn* c) { uint8_t err[] = { 0x05, 0x04, 0x00, 0x01, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00 }; write_to_client(c, err, 10); } + c->rem_closed = 1; tcp_conn_push_close(c->tc); } @@ -402,6 +389,7 @@ static void socks_dns_done_cb(const struct adns_result* res, void* arg) { // ==================================================================== static void on_read_cb(struct ll_queue* q, void* arg) { struct socks_proxy_conn* c = (struct socks_proxy_conn*)arg; + if (c->state == SOCKS_STATE_CONNECTING || c->tx_buf || c->rem_closed) return; struct ll_entry* e = queue_data_get(q); if (!e) { queue_resume_callback(q); return; } @@ -421,43 +409,16 @@ static void on_read_cb(struct ll_queue* q, void* arg) { memcpy(c->buf + c->buf_len, e->dgram, e->len); c->buf_len += e->len; memory_pool_free(c->tc->data_pool, e->dgram); queue_entry_free(e); - queue_resume_callback(q); - if (c->is_http) { process_http_request(c); } else { if (c->state == SOCKS_STATE_GREETING) process_socks_greeting(c); - if (c->state == SOCKS_STATE_REQUEST) process_socks_request(c); + if (!c->rem_closed && c->state == SOCKS_STATE_REQUEST) process_socks_request(c); } + if (c->state != SOCKS_STATE_CONNECTING && !c->rem_closed) queue_resume_callback(q); return; } - if (c->state == SOCKS_STATE_CONNECTING || c->state == HTTP_STATE_CONNECTING) { - // DNS в процессе: копим входящий body (HTTP proxy) / дропаем (CONNECT) - if (c->dns_kind == DNSK_HTTP_PROXY) { - size_t nl = (size_t)c->dns_pending_len + e->len; - if (nl <= UINT16_MAX) { - uint8_t* nb = u_realloc(c->dns_pending, nl); - if (nb) { memcpy(nb + c->dns_pending_len, e->dgram, e->len); c->dns_pending = nb; c->dns_pending_len = (uint16_t)nl; } - else DEBUG_ERROR(DEBUG_CATEGORY_PROXY, "socks_proxy: dns_pending realloc failed sid=%08x", c->stream_id); - } else { - DEBUG_ERROR(DEBUG_CATEGORY_PROXY, "socks_proxy: dns_pending overflow sid=%08x", c->stream_id); - } - } else { - DEBUG_WARN(DEBUG_CATEGORY_PROXY, "socks_proxy: data during DNS (kind=%d) sid=%08x — drop len=%u", c->dns_kind, c->stream_id, e->len); - } - memory_pool_free(c->tc->data_pool, e->dgram); queue_entry_free(e); - queue_resume_callback(q); - return; - } - - // RELAY = релей данных в ETCP - if (c->tx_buf) { - DEBUG_ERROR(DEBUG_CATEGORY_PROXY, "socks_proxy: RELAY with pending tx_buf sid=%08x — drop new len=%u", c->stream_id, e->len); - memory_pool_free(c->tc->data_pool, e->dgram); queue_entry_free(e); - queue_resume_callback(q); - return; - } int ret = send_data(c, e->dgram, e->len, 0); DEBUG_TRACE(DEBUG_CATEGORY_PROXY, "SOCKS PROXY SEND sid=%08x len=%u ret=%d is_http=%d", c->stream_id, e->len, ret, c->is_http); if (ret == 0) { @@ -466,90 +427,70 @@ static void on_read_cb(struct ll_queue* q, void* arg) { } else { c->tx_buf = u_malloc(e->len); if (c->tx_buf) { memcpy(c->tx_buf, e->dgram, e->len); c->tx_len = e->len; } - else { DEBUG_ERROR(DEBUG_CATEGORY_PROXY, "socks_proxy: tx_buf malloc=%u failed sid=%08x — drop", e->len, c->stream_id); } + else { DEBUG_ERROR(DEBUG_CATEGORY_PROXY, "socks_proxy: tx_buf allocation failed sid=%08x", c->stream_id); } memory_pool_free(c->tc->data_pool, e->dgram); queue_entry_free(e); + if (!c->tx_buf) { on_error_cb(c->tc, ENOMEM, c); return; } socks_bp_register(c); DEBUG_DEBUG(DEBUG_CATEGORY_PROXY, "SOCKS PROXY BP: sid=%08x tx_buf=%u waiter_reg", c->stream_id, c->tx_len); } } +// Отправка буфера частями сохраняет HTTP headers/body, независимо от размера одного DATA. +static void socks_flow_wake(void* arg) { + struct socks_proxy_conn* c = arg; + if (c->rem_closed || c->freed || !c->flow.ready) return; + while (c->tx_buf) { + uint16_t len = c->tx_len > TCP_PROXY_CHUNK ? TCP_PROXY_CHUNK : c->tx_len; + int ret = send_data(c, c->tx_buf, len, 0); + if (ret < 0) { socks_bp_register(c); return; } + c->tx_len -= len; + if (c->tx_len) memmove(c->tx_buf, c->tx_buf + len, c->tx_len); + else { u_free(c->tx_buf); c->tx_buf = NULL; } + } + queue_resume_callback(c->tc->read_queue); + socks_maybe_relay_fin(c); +} + static void tx_waiter_cb(struct ll_queue* q, void* arg) { (void)q; - struct socks_proxy_conn* c = (struct socks_proxy_conn*)arg; - if (c->rem_closed || c->close_sent) return; - if (!c->tx_buf) { if (c->tc && c->tc->read_queue) queue_resume_callback(c->tc->read_queue); return; } - int ret = send_data(c, c->tx_buf, c->tx_len, 0); - DEBUG_DEBUG(DEBUG_CATEGORY_PROXY, "SOCKS PROXY WAKE: sid=%08x tx_buf=%u ret=%d is_http=%d", c->stream_id, c->tx_len, ret, c->is_http); - if (ret == 0) { - u_free(c->tx_buf); c->tx_buf = NULL; c->tx_len = 0; - if (c->tc && c->tc->read_queue) queue_resume_callback(c->tc->read_queue); - socks_maybe_relay_fin(c); - } else { - socks_bp_register(c); - } + socks_flow_wake(arg); } static void tx_retry_timer_cb(void* arg) { - struct socks_proxy_conn* c = (struct socks_proxy_conn*)arg; + struct socks_proxy_conn* c = arg; c->tx_retry_timer = NULL; - if (c->rem_closed || c->close_sent) return; - if (!c->tx_buf) return; - int ret = send_data(c, c->tx_buf, c->tx_len, 1); - DEBUG_DEBUG(DEBUG_CATEGORY_PROXY, "SOCKS PROXY RETRY: sid=%08x len=%u force=1 ret=%d", c->stream_id, c->tx_len, ret); - if (ret == 0) { - u_free(c->tx_buf); c->tx_buf = NULL; c->tx_len = 0; - if (c->tc && c->tc->read_queue) queue_resume_callback(c->tc->read_queue); - socks_maybe_relay_fin(c); - } else { - c->tx_retry_timer = uasync_set_timeout(c->ua, 5000, c, tx_retry_timer_cb, "socks_retry"); - } + socks_flow_wake(c); } static void socks_bp_register(struct socks_proxy_conn* c) { - etcp_router_on_send_ready(c->inst, TOPO_GROUP_UTUN, c->via_node_id, ETCP_RT_ID_TCP_PROXY_SERVER, &c->tx_waiter, tx_waiter_cb, c); - if (!c->tx_retry_timer) - c->tx_retry_timer = uasync_set_timeout(c->ua, 5000, c, tx_retry_timer_cb, "socks_retry"); + if (c->flow.ready && c->flow.tx_credit >= (c->tx_len > TCP_PROXY_CHUNK ? TCP_PROXY_CHUNK : c->tx_len)) + etcp_router_on_send_ready(c->inst, TOPO_GROUP_UTUN, c->via_node_id, ETCP_RT_ID_TCP_PROXY_SERVER, &c->tx_waiter, tx_waiter_cb, c); + if (!c->tx_retry_timer) c->tx_retry_timer = uasync_set_timeout(c->ua, 5000, c, tx_retry_timer_cb, "socks_retry"); } +// Обе половины закрываются независимо: write_queue не задерживает FIN прочитанного направления. static void socks_maybe_relay_fin(struct socks_proxy_conn* c) { - if (c->fin_deferred && !c->tx_buf) { - c->fin_deferred = 0; - DEBUG_DEBUG(DEBUG_CATEGORY_PROXY, "socks_proxy: relay deferred FIN sid=%08x", c->stream_id); - send_fin(c); - if (c->fin_remote && !c->close_sent && !c->close_pending) send_close(c); + if (c->rem_closed || c->close_queued || !c->flow.ready) return; + if (c->tc->fin_remote && !c->tx_buf && !c->tc->read_queue->head && !c->flow.fin_sent) { + c->flow.fin_pending = 1; + proxy_flow_flush(&c->flow); + } + if (c->flow.fin_sent && c->flow.fin_received && c->tc->fin_local) { + c->close_queued = 1; + tcp_conn_push_close(c->tc); } } -// ==================================================================== -// tcp_io: FIN / Error / Closed -// ==================================================================== static void on_fin_cb(struct tcp_conn* tc, void* arg) { - struct socks_proxy_conn* c = (struct socks_proxy_conn*)arg; - if (c->rem_closed || c->close_sent) return; - if (c->tx_buf) { - // отложенный FIN: дождёмся сброса tx_buf (backpressure) - c->fin_deferred = 1; - DEBUG_DEBUG(DEBUG_CATEGORY_PROXY, "socks_proxy: local FIN deferred (tx_buf pending) sid=%08x", c->stream_id); - return; - } - if (tc->fin_local) { - // оба FIN обменяны: мы уже shutdown-нули свой write (fin_local), peer прислал FIN. - // Закрываем локальный сокет, симметрично exit-стороне (tcp_proxy_server.c on_fin_cb). - DEBUG_DEBUG(DEBUG_CATEGORY_PROXY, "socks_proxy: both FINs done → close local socket sid=%08x", c->stream_id); - send_close(c); - tcp_conn_push_close(tc); - return; - } - if (!tc->write_buf && !tc->write_queue->head) { - DEBUG_DEBUG(DEBUG_CATEGORY_PROXY, "socks_proxy: local FIN → relay FIN sid=%08x", c->stream_id); - send_fin(c); - if (c->fin_remote && !c->close_sent && !c->close_pending) send_close(c); - } else { - // данные ещё в write_queue, отложим FIN - DEBUG_DEBUG(DEBUG_CATEGORY_PROXY, "socks_proxy: local FIN deferred (wq=%d wbuf=%s) sid=%08x", - tc->write_queue->count, tc->write_buf ? "y" : "n", c->stream_id); - tcp_conn_set_flushed(tc, NULL); - } + (void)tc; + socks_maybe_relay_fin(arg); +} + +static void socks_flushed_cb(struct tcp_conn* tc, void* arg) { + struct socks_proxy_conn* c = arg; + if (!tc->write_buf && !tc->write_queue->head) + proxy_flow_consume(&c->flow, TCP_PROXY_WINDOW - c->flow.rx_credit - c->flow.consumed); + socks_maybe_relay_fin(c); } static void on_error_cb(struct tcp_conn* tc, int err, void* arg) { @@ -563,6 +504,7 @@ static void on_error_cb(struct tcp_conn* tc, int err, void* arg) { struct UTUN_INSTANCE* inst = c->inst; uint64_t via = c->via_node_id; uint32_t sid = c->stream_id; + c->rem_closed = 1; socks_proxy_conn_free_soon(c); if (notify && inst) send_msg(inst, TOPO_GROUP_UTUN, via, TCP_PROXY_SUBCMD_ERROR, sid, NULL, 0, 1); } @@ -571,7 +513,7 @@ static void on_closed_cb(struct tcp_conn* tc, void* arg) { struct socks_proxy_conn* c = (struct socks_proxy_conn*)arg; DEBUG_INFO(DEBUG_CATEGORY_PROXY, "socks_proxy: tcp closed sid=%08x state=%d fm_r=%d fm_l=%d rem_cl=%d bytes_client=%u bytes_exit=%u", c->stream_id, c->state, tc->fin_remote, tc->fin_local, c->rem_closed, c->bytes_to_client, c->bytes_from_exit); - if (!c->close_sent && !c->close_pending) send_close(c); + if (!c->flow.fin_sent && !c->rem_closed && !c->close_sent && !c->close_pending) send_close(c); socks_proxy_conn_free_soon(c); } @@ -592,10 +534,12 @@ static void on_accept_cb(socket_t sock, void* arg) { c->ua = ctx->ua; c->inst = ctx->inst; c->via_node_id = ctx->via_node_id; c->state = ctx->is_http ? HTTP_STATE_REQUEST : SOCKS_STATE_GREETING; c->head = ctx->conns; c->count = ctx->conn_count; + proxy_flow_init(&c->flow, ctx->inst, ctx->ua, ctx->via_node_id, ETCP_RT_ID_TCP_PROXY_SERVER, c->stream_id, socks_flow_wake, c); c->tc = tcp_conn_create(ctx->ua, csock, 4096, 4096, 8, 0, 0, on_fin_cb, on_error_cb, c); if (!c->tc) { DEBUG_ERROR(DEBUG_CATEGORY_PROXY, "socks_proxy: tcp_conn_create failed"); u_free(c); socket_close_wrapper(csock); return; } c->tc->on_closed = on_closed_cb; + c->tc->on_fin_sent = on_fin_cb; queue_set_callback(c->tc->read_queue, on_read_cb, c); queue_set_waiter_defer(c->tc->read_queue, 1); @@ -620,9 +564,17 @@ int socks_proxy_handle_etcp(struct socks_proxy_conn** head, int* count, const uint8_t* data, size_t data_len) { struct socks_proxy_conn* c = socks_proxy_find_conn(*head, stream_id); if (!c) return 0; + if (c->rem_closed || c->freed) return 1; + if (subcmd == TCP_PROXY_SUBCMD_CONNECTED) { socks_connected(c); return 1; } + if (subcmd == TCP_PROXY_SUBCMD_WINDOW) { + if (proxy_flow_window(&c->flow, data, data_len) < 0) on_error_cb(c->tc, EPROTO, c); + else socks_flow_wake(c); + return 1; + } if (subcmd == TCP_PROXY_SUBCMD_DATA) { DEBUG_TRACE(DEBUG_CATEGORY_PROXY, "SOCKS DATA <- sid=%08x len=%zu", stream_id, data_len); + if (proxy_flow_receive(&c->flow, data_len) < 0) { on_error_cb(c->tc, EPROTO, c); return 1; } c->bytes_from_exit += (uint32_t)data_len; if (data_len > 0) { if (data_len > c->tc->data_pool->object_size) { @@ -635,11 +587,13 @@ int socks_proxy_handle_etcp(struct socks_proxy_conn** head, int* count, if (e && buf) { memcpy(buf, data, data_len); e->dgram = buf; e->len = (uint16_t)data_len; c->bytes_to_client += (uint32_t)data_len; + tcp_conn_set_flushed(c->tc, socks_flushed_cb); queue_data_put(c->tc->write_queue, e); } else { DEBUG_ERROR(DEBUG_CATEGORY_PROXY, "socks_proxy: handle_data alloc failed sid=%08x wq=%d bytes_in=%u bytes_out=%u", stream_id, c->tc->write_queue->count, c->bytes_from_exit, c->bytes_to_client); if (e) queue_entry_free(e); if (buf) memory_pool_free(c->tc->data_pool, buf); + on_error_cb(c->tc, ENOMEM, c); } } return 1; @@ -654,7 +608,8 @@ int socks_proxy_handle_etcp(struct socks_proxy_conn** head, int* count, } if (subcmd == TCP_PROXY_SUBCMD_ERROR) { - DEBUG_INFO(DEBUG_CATEGORY_PROXY, "socks_proxy: ERROR from exit sid=%08x", stream_id); + DEBUG_WARN(DEBUG_CATEGORY_PROXY, "socks_proxy: ERROR from exit sid=%08x", stream_id); + if (!c->flow.ready) { socks_dns_error_and_close(c); return 1; } c->rem_closed = 1; etcp_router_cancel_send_ready(c->inst, TOPO_GROUP_UTUN, c->via_node_id, ETCP_RT_ID_TCP_PROXY_SERVER, &c->tx_waiter); tcp_conn_push_close(c->tc); @@ -663,9 +618,9 @@ int socks_proxy_handle_etcp(struct socks_proxy_conn** head, int* count, if (subcmd == TCP_PROXY_SUBCMD_FIN) { DEBUG_INFO(DEBUG_CATEGORY_PROXY, "socks_proxy: FIN from exit sid=%08x", stream_id); - c->fin_remote = 1; - tcp_conn_push_fin(c->tc); - if (c->tc->fin_local && !c->close_sent && !c->close_pending) send_close(c); + if (c->flow.fin_received) return 1; + c->fin_remote = 1; c->flow.fin_received = 1; + if (tcp_conn_push_fin(c->tc) < 0) on_error_cb(c->tc, ENOMEM, c); return 1; } @@ -687,6 +642,7 @@ void socks_proxy_conn_free(struct socks_proxy_conn* c) { if (!c) return; if (c->freed) return; c->freed = 1; + proxy_flow_destroy(&c->flow); if (c->free_soon_id) { uasync_call_soon_cancel(c->ua, c->free_soon_id); c->free_soon_id = NULL; } struct socks_proxy_conn** head = c->head; int* count = c->count; DEBUG_INFO(DEBUG_CATEGORY_PROXY, "socks_proxy: FREE sid=%08x state=%d total=%d", c->stream_id, c->state, count ? *count : 0); diff --git a/src/proxy/socks_proxy.h b/src/proxy/socks_proxy.h index c4bf6114..ee59d132 100644 --- a/src/proxy/socks_proxy.h +++ b/src/proxy/socks_proxy.h @@ -8,6 +8,7 @@ extern "C" { #include +#include "proxy_protocol.h" #include "../lib/socket_compat.h" #include "../lib/ll_queue.h" @@ -42,6 +43,7 @@ struct socks_proxy_conn { struct UASYNC* ua; struct UTUN_INSTANCE* inst; uint64_t via_node_id; + struct proxy_flow flow; uint32_t stream_id; uint8_t dest_ip[4]; uint16_t dest_port; @@ -52,8 +54,9 @@ struct socks_proxy_conn { uint8_t is_http; uint8_t state; uint8_t freed; + uint8_t close_queued; void* free_soon_id; // handle отложенного освобождения (uasync_call_soon) - uint8_t buf[1024]; + uint8_t buf[16384]; uint16_t buf_len; uint8_t* tx_buf; // буфер при backpressure (retry в tx_waiter_cb) uint16_t tx_len; diff --git a/src/proxy/tcp_proxy_client.c b/src/proxy/tcp_proxy_client.c index 2c6cb1ec..f05f62f9 100644 --- a/src/proxy/tcp_proxy_client.c +++ b/src/proxy/tcp_proxy_client.c @@ -34,6 +34,8 @@ // ==================================================================== // Предварительные объявления // ==================================================================== +static void tcp_proxy_client_flow_wake(void* arg); +static int tcp_proxy_client_fin_flush(struct tcp_proxy_client_conn* pc); static err_t tcp_proxy_client_recv_cb(void *arg, struct tcp_pcb *pcb, struct pbuf *p, err_t err); static err_t tcp_proxy_client_sent_cb(void *arg, struct tcp_pcb *pcb, uint16_t len); static void tcp_proxy_client_err_cb(void *arg, err_t err); @@ -70,17 +72,9 @@ static struct ll_entry* tcp_proxy_client_entry_from_data(struct memory_pool* poo // Протокол: сборка и отправка сообщений прокси через ETCP // ==================================================================== static int tcp_proxy_client_send_msg(struct UTUN_INSTANCE* inst, uint64_t group_id, 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_PROXY, "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_PROXY, "malloc(%zu) failed subcmd=%02x sid=%08x", TCP_PROXY_HDR_SIZE + len, subcmd, sid); queue_entry_free(e); return -1; } - e->dgram[0] = ETCP_RT_ID_TCP_PROXY_SERVER; - 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, group_id, dst, e, force, 0); + uint32_t sid, const uint8_t* data, size_t len, int force) { + (void)group_id; + return proxy_send(inst, dst, ETCP_RT_ID_TCP_PROXY_SERVER, subcmd, sid, data, len, force); } static int tcp_proxy_client_send_connect(struct tcp_proxy_client_conn* pc) { @@ -94,8 +88,7 @@ static int tcp_proxy_client_send_connect(struct tcp_proxy_client_conn* pc) { } static int tcp_proxy_client_send_data(struct tcp_proxy_client_conn* pc, const uint8_t* data, uint16_t len, int force) { - int ret = tcp_proxy_client_send_msg(pc->proxy->inst, TOPO_GROUP_UTUN, pc->proxy->via_node_id, - TCP_PROXY_SUBCMD_DATA, pc->stream_id, data, len, force); + int ret = proxy_flow_send(&pc->flow, data, len, force); if (ret == 0) pc->bytes_to_exit += len; else { pc->bp_count++; @@ -125,12 +118,6 @@ static int tcp_proxy_client_send_error(struct tcp_proxy_client_conn* pc) { return 0; } -static int tcp_proxy_client_send_fin(struct tcp_proxy_client_conn* pc) { - DEBUG_INFO(DEBUG_CATEGORY_PROXY, "PROXY FIN RELAY sid=%08x", pc->stream_id); - return tcp_proxy_client_send_msg(pc->proxy->inst, TOPO_GROUP_UTUN, pc->proxy->via_node_id, - TCP_PROXY_SUBCMD_FIN, pc->stream_id, NULL, 0, 1); -} - // ==================================================================== // Вывод: lwIP TCP отправляет IP пакеты через этот callback // ==================================================================== @@ -189,53 +176,33 @@ static int tcp_proxy_client_handle_non_tcp(struct tcp_proxy_client* p, uint8_t* // Помощник: передача данных из очереди to_lwip в lwIP TCP // ==================================================================== static void tcp_proxy_client_feed_from_transport(struct tcp_proxy_client_conn *pc) { - if (!pc->to_lwip || !pc->pcb) return; - if (pc->rem_closed) return; - int sent_any = 0; - uint32_t q_pre = queue_entry_count(pc->to_lwip); - while (1) { + if (!pc->to_lwip || !pc->pcb || pc->rem_closed) return; + int sent = 0; + while (pc->to_lwip->head) { uint16_t space = tcp_sndbuf(pc->pcb); - if (space < TCP_MSS / 2) { - if (queue_entry_count(pc->to_lwip) > 0) - DEBUG_DEBUG(DEBUG_CATEGORY_DEBUG, - "PROXY FEED stalled sid=%08x snd_buf=%u to_lwip=%d snd_wnd=%u cwnd=%u", - pc->stream_id, space, queue_entry_count(pc->to_lwip), - pc->pcb->snd_wnd, pc->pcb->cwnd); - break; - } - struct ll_entry *e = queue_data_get(pc->to_lwip); - if (!e) break; - if (e->len <= space) { - err_t ret = tcp_write(pc->pcb, e->dgram, e->len, TCP_WRITE_FLAG_COPY); - if (ret == LERR_OK) { sent_any = 1; queue_dgram_free(e); queue_entry_free(e); } - else { - DEBUG_DEBUG(DEBUG_CATEGORY_DEBUG, - "PROXY FEED write fail sid=%08x len=%u ret=%d snd_buf=%u q=%d", - pc->stream_id, e->len, ret, space, queue_entry_count(pc->to_lwip) + 1); - queue_data_put_first(pc->to_lwip, e); break; - } - } else { - queue_data_put_first(pc->to_lwip, e); - break; + if (!space) break; + struct ll_entry* e = queue_data_get(pc->to_lwip); + uint16_t len = e->len < space ? e->len : space; + err_t ret = tcp_write(pc->pcb, e->dgram, len, TCP_WRITE_FLAG_COPY); + if (ret != LERR_OK) { + DEBUG_DEBUG(DEBUG_CATEGORY_PROXY, "proxy: lwIP write blocked sid=%08x ret=%d space=%u", pc->stream_id, ret, space); + queue_data_put_first(pc->to_lwip, e); break; } - queue_resume_callback(pc->to_lwip); + sent = 1; + e->len -= len; + if (e->len) { memmove(e->dgram, e->dgram + len, e->len); queue_data_put_first(pc->to_lwip, e); } + else { queue_dgram_free(e); queue_entry_free(e); } + // Буфер lwIP тоже ограничен snd_buf; to_lwip ограничен окном proxy. + proxy_flow_consume(&pc->flow, len); } queue_resume_callback(pc->to_lwip); - if (sent_any) { - uint32_t unsent = 0; struct tcp_seg* s; - for (s = pc->pcb->unsent; s; s = s->next) unsent++; - DEBUG_TRACE(DEBUG_CATEGORY_PROXY, "PROXY FEED exit->client sid=%08x fed=%u q=%u->%u snd_wnd=%u cwnd=%u unsent=%u", - pc->stream_id, sent_any, q_pre, queue_entry_count(pc->to_lwip), pc->pcb->snd_wnd, pc->pcb->cwnd, unsent); - tcp_output(pc->pcb); - } - // FIN от exit был отложен до дренажа to_lwip — теперь очередь пуста, шлём FIN локальной стороне - if (pc->fin_deferred && pc->pcb && queue_entry_count(pc->to_lwip) == 0 && - pc->pcb->state != TIME_WAIT && pc->pcb->state != CLOSED) { - pc->fin_deferred = 0; - DEBUG_INFO(DEBUG_CATEGORY_PROXY, "PROXY FIN RELAY (deferred) sid=%08x — to_lwip drained, shutdown write", pc->stream_id); - tcp_shutdown(pc->pcb, 0, 1); - tcp_sent(pc->pcb, NULL); - tcp_poll(pc->pcb, NULL, 0); + if (sent) tcp_output(pc->pcb); + if (pc->fin_deferred && !pc->to_lwip->head) { + err_t ret = tcp_shutdown(pc->pcb, 0, 1); + if (ret == LERR_OK) { + pc->fin_deferred = 0; + DEBUG_DEBUG(DEBUG_CATEGORY_PROXY, "proxy: lwIP FIN queued sid=%08x", pc->stream_id); + } else DEBUG_WARN(DEBUG_CATEGORY_PROXY, "proxy: lwIP FIN retry sid=%08x ret=%d", pc->stream_id, ret); } } @@ -249,6 +216,7 @@ static err_t tcp_proxy_client_accept_cb(void *arg, struct tcp_pcb *newpcb, err_t struct tcp_proxy_client_conn *pc = u_calloc(1, sizeof(struct tcp_proxy_client_conn)); if (!pc) { DEBUG_ERROR(DEBUG_CATEGORY_PROXY, "TCP proxy client: accept alloc failed"); return LERR_MEM; } pc->proxy = p; pc->pcb = newpcb; pc->stream_id = ++p->next_stream_id; + proxy_flow_init(&pc->flow, p->inst, p->ua, p->via_node_id, ETCP_RT_ID_TCP_PROXY_SERVER, pc->stream_id, tcp_proxy_client_flow_wake, pc); int i; for (i = 0; i < p->mapping_count; i++) { @@ -273,14 +241,14 @@ 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"); pc->tx_queue = queue_new(p->ua, 0, 0, 0, "tx_queue"); - if (!pc->tx_queue) { + if (!pc->tx_queue || !pc->to_lwip) { DEBUG_ERROR(DEBUG_CATEGORY_PROXY, "TCP proxy client: tx_queue alloc failed sid=%08x", pc->stream_id); - queue_free(pc->to_lwip); pc->to_lwip = NULL; - tcp_arg(newpcb, NULL); tcp_abort(newpcb); u_free(pc); return LERR_MEM; + queue_free(pc->to_lwip); queue_free(pc->tx_queue); pc->to_lwip = NULL; + tcp_arg(newpcb, NULL); tcp_abort(newpcb); u_free(pc); return LERR_ABRT; } queue_set_callback(pc->tx_queue, tcp_proxy_client_tx_queue_drain_cb, pc); - if (tcp_proxy_client_send_connect(pc) != 0) { DEBUG_ERROR(DEBUG_CATEGORY_PROXY, "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; } + if (tcp_proxy_client_send_connect(pc) != 0) { DEBUG_ERROR(DEBUG_CATEGORY_PROXY, "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); queue_free(pc->tx_queue); u_free(pc); return LERR_ABRT; } pc->next = p->conns; p->conns = pc; p->conn_count++; pc->diag_timer = uasync_set_timeout(p->ua, 5000, pc, tcp_proxy_client_diag_timer_cb, "tpc_diag"); @@ -297,19 +265,25 @@ static err_t tcp_proxy_client_accept_cb(void *arg, struct tcp_pcb *newpcb, err_t // После отправки последних данных (очередь пуста): релеить FIN/CLOSE в exit. // Возвращает 1, если соединение завершено и pc освобождён (дальше pc/q использовать нельзя). static int tcp_proxy_client_fin_flush(struct tcp_proxy_client_conn* pc) { - if (pc->fin_local && !pc->close_sent && !pc->close_pending) { - if (pc->fin_remote || pc->rem_closed) - tcp_proxy_client_send_close(pc); - else - tcp_proxy_client_send_fin(pc); + if (pc->fin_local && !pc->tx_queue->head && !pc->flow.fin_sent) { + pc->flow.fin_pending = 1; + proxy_flow_flush(&pc->flow); } - if (pc->fin_local && pc->fin_remote) { + if (pc->flow.fin_sent && pc->fin_remote && !pc->fin_deferred && !pc->tx_queue->head && !pc->to_lwip->head) { tcp_proxy_client_conn_finish(pc); return 1; } return 0; } +static void tcp_proxy_client_flow_wake(void* arg) { + struct tcp_proxy_client_conn* pc = arg; + if (pc->rem_closed || pc->error) return; + queue_resume_callback(pc->tx_queue); + tcp_proxy_client_feed_from_transport(pc); + tcp_proxy_client_fin_flush(pc); +} + // Дрейн tx_queue → ETCP. tcp_recved вызывается только после успешной отправки: // окно закрывается плавно (штатная backpressure), а при возобновлении открывается // пачкой (за 2 пакета wnd_inflation ≥ TCP_WND_UPDATE_THRESHOLD → немедленный ACK). @@ -330,7 +304,7 @@ static void tcp_proxy_client_tx_queue_drain_cb(struct ll_queue* q, void* arg) { queue_data_put_first(q, e); DEBUG_DEBUG(DEBUG_CATEGORY_PROXY, "PROXY BP stall sid=%08x q=%d rcv_wnd=%u", pc->stream_id, queue_entry_count(q), pc->pcb ? pc->pcb->rcv_wnd : 0); - etcp_router_on_send_ready(pc->proxy->inst, TOPO_GROUP_UTUN, pc->proxy->via_node_id, + if (ret != -2) etcp_router_on_send_ready(pc->proxy->inst, TOPO_GROUP_UTUN, pc->proxy->via_node_id, ETCP_RT_ID_TCP_PROXY_SERVER, &pc->tx_waiter, tcp_proxy_client_pause_resume_cb, pc); if (!pc->tx_retry_timer) @@ -421,19 +395,22 @@ static err_t tcp_proxy_client_recv_cb(void *arg, struct tcp_pcb *pcb, struct pbu tcp_recved(pcb, len); pbuf_free(p); return LERR_OK; } - uint8_t *data = u_malloc(len); - if (!data) { - DEBUG_ERROR(DEBUG_CATEGORY_PROXY, "PROXY recv malloc(%u) failed sid=%08x", len, pc->stream_id); - pbuf_free(p); return LERR_OK; - } - pbuf_copy_partial(p, data, len, 0); - struct ll_entry* e = queue_entry_new_from_pool(pc->proxy->entry_pool); - if (!e) { - DEBUG_ERROR(DEBUG_CATEGORY_PROXY, "PROXY recv entry pool exhausted sid=%08x", pc->stream_id); - u_free(data); pbuf_free(p); return LERR_OK; + struct ll_entry* head = NULL; + struct ll_entry** tail = &head; + for (uint32_t off = 0; off < len;) { + uint16_t chunk = len - off > TCP_PROXY_CHUNK ? TCP_PROXY_CHUNK : len - off; + struct ll_entry* e = queue_entry_new_from_pool(pc->proxy->entry_pool); + if (e) { e->dgram = u_malloc(chunk); e->len = chunk; } + if (!e || !e->dgram) { + DEBUG_ERROR(DEBUG_CATEGORY_PROXY, "proxy: recv allocation failed sid=%08x len=%u", pc->stream_id, len); + if (e) queue_entry_free(e); + while (head) { struct ll_entry* next = head->next; queue_dgram_free(head); queue_entry_free(head); head = next; } + return LERR_MEM; + } + pbuf_copy_partial(p, e->dgram, chunk, off); + e->next = NULL; *tail = e; tail = &e->next; off += chunk; } - e->dgram = data; e->len = len; - queue_data_put(pc->tx_queue, e); + while (head) { struct ll_entry* next = head->next; head->next = NULL; queue_data_put(pc->tx_queue, head); head = next; } pbuf_free(p); return LERR_OK; } @@ -443,6 +420,7 @@ static err_t tcp_proxy_client_sent_cb(void *arg, struct tcp_pcb *pcb, uint16_t l struct tcp_proxy_client_conn *pc = (struct tcp_proxy_client_conn *)arg; if (!pc || !pc->proxy || pc->rem_closed) return LERR_OK; tcp_proxy_client_feed_from_transport(pc); + tcp_proxy_client_fin_flush(pc); return LERR_OK; } @@ -478,6 +456,8 @@ static err_t tcp_proxy_client_poll_cb(void *arg, struct tcp_pcb *pcb) { return LERR_ABRT; } + tcp_proxy_client_feed_from_transport(pc); + if (tcp_proxy_client_fin_flush(pc)) return LERR_OK; if (pc->close_pending) { if (pc->error) tcp_proxy_client_send_error(pc); else tcp_proxy_client_send_close(pc); @@ -545,6 +525,7 @@ static void tcp_proxy_client_tun_input(struct ll_queue* q, void* arg) { static void tcp_proxy_client_conn_free(struct tcp_proxy_client_conn *pc) { if (!pc) return; struct tcp_proxy_client *p = pc->proxy; + proxy_flow_destroy(&pc->flow); uint32_t pending = pc->to_lwip ? queue_entry_count(pc->to_lwip) : 0; DEBUG_INFO(DEBUG_CATEGORY_PROXY, "PROXY FREE sid=%08x total_conns=%d to_lwip_q=%u", pc->stream_id, p->conn_count, pending); @@ -607,8 +588,16 @@ static void tcp_proxy_client_handle_data(struct tcp_proxy_client* p, struct ETCP if (data_len > 0) { pc->bytes_from_exit += (uint32_t)data_len; struct ll_entry* e = tcp_proxy_client_entry_from_data(pc->proxy->entry_pool, entry->dgram + TCP_PROXY_RECV_HDR_SIZE, (uint16_t)data_len); - if (e) queue_data_put(pc->to_lwip, e); - tcp_proxy_client_feed_from_transport(pc); + if (!e || proxy_flow_receive(&pc->flow, data_len) < 0) { + if (e) { queue_dgram_free(e); queue_entry_free(e); } + DEBUG_ERROR(DEBUG_CATEGORY_PROXY, "proxy: receive failed sid=%08x", stream_id); + tcp_proxy_client_send_error(pc); + if (pc->pcb) { tcp_arg(pc->pcb, NULL); tcp_err(pc->pcb, NULL); tcp_abort(pc->pcb); pc->pcb = NULL; } + tcp_proxy_client_conn_free(pc); + } else { + queue_data_put(pc->to_lwip, e); + tcp_proxy_client_feed_from_transport(pc); + } } queue_dgram_free(entry); queue_entry_free(entry); } @@ -640,7 +629,7 @@ static void tcp_proxy_client_handle_error(struct tcp_proxy_client* p, uint32_t s tcp_sent(pc->pcb, NULL); tcp_err(pc->pcb, NULL); tcp_poll(pc->pcb, NULL, 0); - tcp_close(pc->pcb); + tcp_abort(pc->pcb); pc->pcb = NULL; } pc->rem_closed = 1; @@ -652,22 +641,10 @@ static void tcp_proxy_client_handle_fin(struct tcp_proxy_client* p, uint32_t str if (!pc || !pc->pcb) return; DEBUG_INFO(DEBUG_CATEGORY_PROXY, "PROXY FIN FROM exit sid=%08x pcb_state=%u to_lwip=%d — shutdown write (send FIN to local)", stream_id, pc->pcb->state, pc->to_lwip ? queue_entry_count(pc->to_lwip) : -1); - pc->fin_remote = 1; - // Если в to_lwip ещё есть данные от exit — откладываем FIN до их дренажа, - // иначе остаток потеряется (после tcp_shutdown писать в lwIP нельзя). - if (pc->to_lwip && queue_entry_count(pc->to_lwip) > 0 && pc->pcb->state != TIME_WAIT && pc->pcb->state != CLOSED) { - pc->fin_deferred = 1; - tcp_proxy_client_feed_from_transport(pc); - return; - } - if (pc->pcb->state != TIME_WAIT && pc->pcb->state != CLOSED) { - tcp_shutdown(pc->pcb, 0, 1); - tcp_sent(pc->pcb, NULL); - tcp_poll(pc->pcb, NULL, 0); - } - if (pc->fin_local && !pc->close_sent && !pc->close_pending) - tcp_proxy_client_send_close(pc); - if (pc->fin_local) tcp_proxy_client_conn_finish(pc); + if (pc->flow.fin_received) return; + pc->fin_remote = 1; pc->flow.fin_received = 1; pc->fin_deferred = 1; + tcp_proxy_client_feed_from_transport(pc); + tcp_proxy_client_fin_flush(pc); } // ==================================================================== @@ -738,6 +715,18 @@ void tcp_proxy_client_router_recv_cb(struct ETCP_CONN* conn, struct ll_entry* en stream_id, subcmd, entry->dgram + TCP_PROXY_RECV_HDR_SIZE, data_len)) { queue_dgram_free(entry); queue_entry_free(entry); return; } + struct tcp_proxy_client_conn* pc = tcp_proxy_client_find_conn(proxy, stream_id); + if (pc && subcmd == TCP_PROXY_SUBCMD_CONNECTED) { + pc->flow.ready = 1; + tcp_proxy_client_flow_wake(pc); + queue_dgram_free(entry); queue_entry_free(entry); return; + } + if (pc && subcmd == TCP_PROXY_SUBCMD_WINDOW) { + if (proxy_flow_window(&pc->flow, entry->dgram + TCP_PROXY_RECV_HDR_SIZE, data_len) < 0) + tcp_proxy_client_handle_error(proxy, stream_id); + else tcp_proxy_client_flow_wake(pc); + queue_dgram_free(entry); queue_entry_free(entry); return; + } // Существующие lwIP conns if (subcmd == TCP_PROXY_SUBCMD_DATA) { tcp_proxy_client_handle_data(proxy, conn, stream_id, entry); return; } if (subcmd == TCP_PROXY_SUBCMD_CLOSE) { tcp_proxy_client_handle_close(proxy, stream_id); queue_dgram_free(entry); queue_entry_free(entry); return; } @@ -844,6 +833,7 @@ void tcp_proxy_client_destroy(struct tcp_proxy_client* p) { struct tcp_proxy_client_conn* pc = p->conns; while (pc) { struct tcp_proxy_client_conn* next = pc->next; + proxy_flow_destroy(&pc->flow); if (pc->pcb && !memory_pool_is_freed(p->lwip->pcb_pool, pc->pcb) && pc->pcb->state != CLOSED) { tcp_arg(pc->pcb, NULL); tcp_recv(pc->pcb, NULL); diff --git a/src/proxy/tcp_proxy_client.h b/src/proxy/tcp_proxy_client.h index ca385156..ce37f8c1 100644 --- a/src/proxy/tcp_proxy_client.h +++ b/src/proxy/tcp_proxy_client.h @@ -8,6 +8,7 @@ extern "C" { #include +#include "proxy_protocol.h" #include #include "../lib/socket_compat.h" #include "../lib/ll_queue.h" @@ -32,6 +33,7 @@ struct tcp_proxy_client_conn { struct tcp_proxy_client* proxy; struct tcp_pcb* pcb; + struct proxy_flow flow; uint32_t stream_id; struct ll_queue* to_lwip; // DATA от exit → lwIP (управление потоком) diff --git a/src/proxy/tcp_proxy_server.c b/src/proxy/tcp_proxy_server.c index ad77b5c2..e07106b4 100644 --- a/src/proxy/tcp_proxy_server.c +++ b/src/proxy/tcp_proxy_server.c @@ -27,18 +27,18 @@ #endif +static void server_wake(void* arg); +static void on_connected_cb(struct tcp_conn* tc, void* arg); 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 on_closed_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 void retry_timer_cb(void* arg); static int send_msg(struct UTUN_INSTANCE* inst, uint64_t group_id, 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); void tcp_proxy_server_conn_free(struct tcp_proxy_server_conn* rc); static inline int write_pending(struct tcp_conn* tc) { @@ -51,94 +51,50 @@ static inline int write_pending(struct tcp_conn* tc) { static int send_msg(struct UTUN_INSTANCE* inst, uint64_t group_id, 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_PROXY, "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_PROXY, "tcp_proxy_server: malloc(%zu) failed", TCP_PROXY_HDR_SIZE + len); queue_entry_free(e); return -1; } - e->dgram[0] = ETCP_RT_ID_TCP_PROXY_CLIENT; - e->dgram[1] = subcmd; - memcpy(e->dgram + 2, &sid, 4); - if (len > 0) memcpy(e->dgram + TCP_PROXY_HDR_SIZE, data, len); - if (TCP_PROXY_HDR_SIZE + len > UINT16_MAX) { - DEBUG_ERROR(DEBUG_CATEGORY_PROXY, "tcp_proxy_server: msg too large len=%zu subcmd=%02x", len, subcmd); - queue_dgram_free(e); queue_entry_free(e); return -1; - } - e->len = TCP_PROXY_HDR_SIZE + len; - return etcp_route_send(inst, group_id, dst, e, force, 0); + (void)group_id; + return proxy_send(inst, dst, ETCP_RT_ID_TCP_PROXY_CLIENT, subcmd, sid, data, len, 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, TOPO_GROUP_UTUN, 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; +// FIN относится только к прочитанному направлению: сначала передаём весь read_queue. +static void server_maybe_finish(struct tcp_proxy_server_conn* rc) { + struct tcp_conn* tc = rc->tc; + if (!tc || rc->cli_closed || rc->close_queued) return; + if (tc->fin_remote && !tc->read_queue->head && !rc->flow.fin_sent) { + rc->flow.fin_pending = 1; + proxy_flow_flush(&rc->flow); } - 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, TOPO_GROUP_UTUN, 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"); + if (rc->flow.fin_sent && rc->flow.fin_received && tc->fin_local) { + if (tcp_conn_push_close(tc) == 0) rc->close_queued = 1; } } -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, TOPO_GROUP_UTUN, 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 (rc->cli_closed) return; - int pend_w = write_pending(tc); - int pend_r = (rc->tx_buf != NULL) || (tc->read_queue && tc->read_queue->head); - DEBUG_INFO(DEBUG_CATEGORY_PROXY, - "SOCK:DST_FIN fd=%d sid=%08x fin_l=%d pend_w=%d pend_r=%d " - "sent=%u recv=%u relayed=%u drain#=%u data#=%u " - "rq=%d(%zub) wq=%d(%zub) wbuf=%s err=%d tx_buf=%s", - (int)tc->sock, rc->stream_id, tc->fin_local, pend_w, pend_r, - rc->bytes_sent, rc->bytes_recv, rc->bytes_relayed, - rc->drain_count, rc->data_count, - 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", tc->error, rc->tx_buf ? "y" : "n"); - if (tc->fin_local) { - send_close(rc); tcp_conn_push_close(tc); - } else if (pend_w) { - tcp_conn_set_flushed(tc, on_flushed_cb); - } else if (pend_r) { - rc->dst_fin_deferred = 1; - } else { - send_fin(rc); - } + (void)tc; + server_maybe_finish(arg); } +// Кредит возвращается после записи всей принятой порции в OS TCP, включая write_buf. static void on_flushed_cb(struct tcp_conn* tc, void* arg) { - struct tcp_proxy_server_conn* rc = (struct tcp_proxy_server_conn*)arg; + struct tcp_proxy_server_conn* rc = arg; if (rc->cli_closed) return; - if (tc->fin_local) { - DEBUG_DEBUG(DEBUG_CATEGORY_PROXY, "SOCK:FLUSHED fd=%d sid=%08x fin_local=%d — both FINs, closing", - (int)tc->sock, rc->stream_id, tc->fin_local); - send_close(rc); tcp_conn_push_close(tc); return; - } - send_fin(rc); + if (!write_pending(tc)) proxy_flow_consume(&rc->flow, TCP_PROXY_WINDOW - rc->flow.rx_credit - rc->flow.consumed); + server_maybe_finish(rc); +} + +static void on_connected_cb(struct tcp_conn* tc, void* arg) { + struct tcp_proxy_server_conn* rc = arg; + DEBUG_INFO(DEBUG_CATEGORY_PROXY, "proxy exit connected peer=%016llx sid=%08x fd=%d", + (unsigned long long)rc->peer_node_id, rc->stream_id, (int)tc->sock); + rc->flow.connected_pending = 1; + proxy_flow_flush(&rc->flow); + server_wake(rc); +} + +static void server_wake(void* arg) { + struct tcp_proxy_server_conn* rc = arg; + if (!rc->tc || rc->cli_closed) return; + queue_resume_callback(rc->tc->read_queue); + server_maybe_finish(rc); } static void on_closed_cb(struct tcp_conn* tc, void* arg) { @@ -154,11 +110,11 @@ static void on_error_cb(struct tcp_conn* tc, int err, void* arg) { struct tcp_proxy_server_conn* rc = (struct tcp_proxy_server_conn*)arg; DEBUG_ERROR(DEBUG_CATEGORY_PROXY, "SOCK:ERR fd=%d sid=%08x err=%d fin_r=%d fin_l=%d " - "sent=%u recv=%u relayed=%u drain#=%u data#=%u " + "sent=%u relayed=%u drain#=%u data#=%u " "rq=%d wq=%d wbuf=%s conn=%d", tc ? (int)tc->sock : -1, rc->stream_id, err, tc ? tc->fin_remote : 0, tc ? tc->fin_local : 0, - rc->bytes_sent, rc->bytes_recv, rc->bytes_relayed, + rc->bytes_sent, rc->bytes_relayed, rc->drain_count, rc->data_count, tc ? tc->read_queue->count : 0, tc ? tc->write_queue->count : 0, tc && tc->write_buf ? "y" : "n", tc ? tc->connected : 0); @@ -202,99 +158,36 @@ static void diag_timer_cb(void* arg) { // ==================================================================== 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 tcp_conn* tc = rc->tc; - struct ll_entry* e; - - if (rc->tx_buf) { - int ret = send_msg(inst, TOPO_GROUP_UTUN, rc->peer_node_id, TCP_PROXY_SUBCMD_DATA, rc->stream_id, rc->tx_buf, rc->tx_len, 0); - if (!tc->read_queue) { u_free(rc->tx_buf); return; } - DEBUG_DEBUG(DEBUG_CATEGORY_PROXY, "TPS TXBUF RETRY: sid=%08x len=%u ret=%d", rc->stream_id, rc->tx_len, ret); - if (ret == 0) { - u_free(rc->tx_buf); rc->tx_buf = NULL; rc->tx_len = 0; - if (rc->retry_timer) { uasync_cancel_timeout(rc->ua, rc->retry_timer); rc->retry_timer = NULL; } - if (rc->dst_fin_deferred && !q->head) { - rc->dst_fin_deferred = 0; - send_fin(rc); - return; - } - } else { - etcp_router_on_send_ready(inst, TOPO_GROUP_UTUN, rc->peer_node_id, ETCP_RT_ID_TCP_PROXY_CLIENT, - &rc->pause_waiter, pause_resume_cb, rc); - if (!rc->retry_timer) rc->retry_timer = uasync_set_timeout(rc->ua, 5000, rc, retry_timer_cb, "tps_retry"); - DEBUG_DEBUG(DEBUG_CATEGORY_PROXY, "TPS TXBUF BP AGAIN: sid=%08x waiter_reg", rc->stream_id); - return; - } - } - - e = queue_data_get(q); - if (!e) { - if (rc->dst_fin_deferred) { - rc->dst_fin_deferred = 0; - send_fin(rc); - } - queue_resume_callback(q); return; - } - if (rc->cli_closed) { - do { - memory_pool_free(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, TOPO_GROUP_UTUN, rc->peer_node_id, TCP_PROXY_SUBCMD_DATA, rc->stream_id, e->dgram, e->len, 0); - if (!tc->read_queue) { memory_pool_free(tc->data_pool, e->dgram); queue_entry_free(e); return; } - rc->bytes_relayed += e->len; rc->drain_count++; - DEBUG_TRACE(DEBUG_CATEGORY_PROXY, "TPS DRAIN #%u sid=%08x len=%u → ret=%d rq=%d", - rc->drain_count, rc->stream_id, e->len, ret, q->count); + struct tcp_proxy_server_conn* rc = arg; + if (!rc->tc || rc->cli_closed) return; + struct ll_entry* e = queue_data_get(q); + if (!e) { queue_resume_callback(q); server_maybe_finish(rc); return; } + int ret = proxy_flow_send(&rc->flow, e->dgram, e->len, 0); if (ret == 0) { - memory_pool_free(tc->data_pool, e->dgram); queue_entry_free(e); + rc->bytes_relayed += e->len; rc->drain_count++; + memory_pool_free(rc->tc->data_pool, e->dgram); queue_entry_free(e); queue_resume_callback(q); - } else { - rc->tx_buf = u_malloc(e->len); - if (rc->tx_buf) { memcpy(rc->tx_buf, e->dgram, e->len); rc->tx_len = e->len; } - else { DEBUG_ERROR(DEBUG_CATEGORY_PROXY, "TPS tx_buf malloc=%u failed sid=%08x — drop", e->len, rc->stream_id); } - memory_pool_free(tc->data_pool, e->dgram); queue_entry_free(e); - if (!tc->read_queue) return; - DEBUG_DEBUG(DEBUG_CATEGORY_PROXY, "SOCK:BACKPRESSURE fd=%d sid=%08x — buffered %u for retry", - (int)tc->sock, rc->stream_id, rc->tx_len); - etcp_router_on_send_ready(inst, TOPO_GROUP_UTUN, rc->peer_node_id, ETCP_RT_ID_TCP_PROXY_CLIENT, - &rc->pause_waiter, pause_resume_cb, rc); - if (!rc->retry_timer) rc->retry_timer = uasync_set_timeout(rc->ua, 5000, rc, retry_timer_cb, "tps_retry"); + server_maybe_finish(rc); + return; } + queue_data_put_first(q, e); + DEBUG_DEBUG(DEBUG_CATEGORY_PROXY, "proxy exit: paused sid=%08x ret=%d credit=%u queued=%zu", + rc->stream_id, ret, rc->flow.tx_credit, queue_total_bytes(q)); + if (ret != -2) + etcp_router_on_send_ready(rc->ctx->inst, TOPO_GROUP_UTUN, rc->peer_node_id, ETCP_RT_ID_TCP_PROXY_CLIENT, + &rc->pause_waiter, pause_resume_cb, rc); + if (!rc->retry_timer) rc->retry_timer = uasync_set_timeout(rc->ua, 5000, rc, retry_timer_cb, "tps_retry"); } 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 || rc->cli_closed) return; - DEBUG_DEBUG(DEBUG_CATEGORY_PROXY, "TPS RESUME sid=%08x tx_buf=%s rq=%d", - rc->stream_id, rc->tx_buf ? "y" : "n", - rc->tc->read_queue ? rc->tc->read_queue->count : 0); - queue_resume_callback(rc->tc->read_queue); + server_wake(arg); } static void retry_timer_cb(void* arg) { - struct tcp_proxy_server_conn* rc = (struct tcp_proxy_server_conn*)arg; + struct tcp_proxy_server_conn* rc = arg; rc->retry_timer = NULL; - if (!rc->tx_buf || !rc->tc || rc->tc->sock == SOCKET_INVALID || rc->cli_closed) return; - struct UTUN_INSTANCE* inst = rc->ctx ? rc->ctx->inst : NULL; - int ret = send_msg(inst, TOPO_GROUP_UTUN, rc->peer_node_id, TCP_PROXY_SUBCMD_DATA, rc->stream_id, rc->tx_buf, rc->tx_len, 1); - DEBUG_DEBUG(DEBUG_CATEGORY_PROXY, "TPS RETRY TIMER: sid=%08x len=%u force=1 ret=%d", rc->stream_id, rc->tx_len, ret); - if (ret == 0) { - u_free(rc->tx_buf); rc->tx_buf = NULL; rc->tx_len = 0; - if (rc->dst_fin_deferred && rc->tc->read_queue && !rc->tc->read_queue->head) { - rc->dst_fin_deferred = 0; - send_fin(rc); - return; - } - queue_resume_callback(rc->tc->read_queue); - } else { - rc->retry_timer = uasync_set_timeout(rc->ua, 5000, rc, retry_timer_cb, "tps_retry"); - } + server_wake(rc); } // ==================================================================== @@ -319,19 +212,13 @@ void tcp_proxy_server_conn_free(struct tcp_proxy_server_conn* rc) { int total = conn_total(rc); DEBUG_DEBUG(DEBUG_CATEGORY_PROXY, "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_remote : 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, TOPO_GROUP_UTUN, 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->ctx && rc->ctx->inst) etcp_router_cancel_send_ready(rc->ctx->inst, TOPO_GROUP_UTUN, rc->peer_node_id, ETCP_RT_ID_TCP_PROXY_CLIENT, &rc->pause_waiter); - if (rc->tx_buf) { u_free(rc->tx_buf); rc->tx_buf = NULL; rc->tx_len = 0; } + proxy_flow_destroy(&rc->flow); if (rc->retry_timer) { uasync_cancel_timeout(rc->ua, rc->retry_timer); rc->retry_timer = 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->tc) { tcp_conn_destroy(rc->tc); rc->tc = NULL; } DEBUG_DEBUG(DEBUG_CATEGORY_PROXY, "SOCK:FREE u_free rc=%p", (void*)rc); @@ -365,6 +252,7 @@ int tcp_proxy_server_handle_connect(struct UTUN_INSTANCE* inst, struct ll_entry* 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; + proxy_flow_init(&rc->flow, inst, inst->ua, src_node_id, ETCP_RT_ID_TCP_PROXY_CLIENT, stream_id, server_wake, rc); socket_t sock = socket(AF_INET, SOCK_STREAM, 0); if (sock == SOCKET_INVALID) { DEBUG_ERROR(DEBUG_CATEGORY_PROXY, "TCP proxy server: socket() failed"); u_free(rc); queue_dgram_free(entry); queue_entry_free(entry); return -1; } @@ -403,6 +291,8 @@ int tcp_proxy_server_handle_connect(struct UTUN_INSTANCE* inst, struct ll_entry* rc->tc = tcp_conn_create(inst->ua, sock, 1500, 8192, 4, 0, inst->config->global.tcp_recv_buf, on_fin_cb, on_error_cb, rc); if (!rc->tc) { DEBUG_ERROR(DEBUG_CATEGORY_PROXY, "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; } rc->tc->on_closed = on_closed_cb; + rc->tc->on_connected = on_connected_cb; + rc->tc->on_fin_sent = on_fin_cb; queue_set_callback(rc->tc->read_queue, read_queue_drain_cb, rc); queue_set_waiter_defer(rc->tc->read_queue, 1); @@ -418,6 +308,7 @@ int tcp_proxy_server_handle_connect(struct UTUN_INSTANCE* inst, struct ll_entry* } rc->next = ctx->conns; ctx->conns = rc; + if (ret == 0) { rc->tc->connected = 1; on_connected_cb(rc->tc, rc); } tcp_conn_set_connect_timeout(rc->tc, g->tcp_proxy_server_connect_timeout_ms); rc->diag_timer = uasync_set_timeout(rc->ua, 10000, rc, diag_timer_cb, "tps_diag"); @@ -439,6 +330,9 @@ int tcp_proxy_server_handle_data(struct UTUN_INSTANCE* inst, struct ETCP_CONN* c rc->bytes_sent += (uint32_t)data_len; rc->data_count++; DEBUG_TRACE(DEBUG_CATEGORY_PROXY, "TPS DATA #%u sid=%08x len=%zu", rc->data_count, stream_id, data_len); + if (proxy_flow_receive(&rc->flow, data_len) < 0) { + on_error_cb(rc->tc, EPROTO, rc); queue_dgram_free(entry); queue_entry_free(entry); return -1; + } if (data_len > 0) { if (data_len > rc->tc->data_pool->object_size) { DEBUG_ERROR(DEBUG_CATEGORY_PROXY, "tcp_proxy_server: data_len=%zu > pool_sz=%zu, dropping sid=%08x", @@ -450,11 +344,13 @@ int tcp_proxy_server_handle_data(struct UTUN_INSTANCE* inst, struct ETCP_CONN* c if (e && buf) { memcpy(buf, entry->dgram + TCP_PROXY_RECV_HDR_SIZE, data_len); e->dgram = buf; e->len = (uint16_t)data_len; + tcp_conn_set_flushed(rc->tc, on_flushed_cb); queue_data_put(rc->tc->write_queue, e); } else { DEBUG_ERROR(DEBUG_CATEGORY_PROXY, "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); + on_error_cb(rc->tc, ENOMEM, rc); } } queue_dgram_free(entry); queue_entry_free(entry); @@ -480,7 +376,9 @@ void tcp_proxy_server_handle_close(struct UTUN_INSTANCE* inst, uint64_t peer, ui queue_entry_free(e); } } - if (rc->tc->connected) tcp_conn_push_close(rc->tc); else tcp_proxy_server_conn_free(rc); + if (rc->tc->connected) { + if (!rc->close_queued && tcp_conn_push_close(rc->tc) == 0) rc->close_queued = 1; + } else tcp_proxy_server_conn_free(rc); } void tcp_proxy_server_handle_error(struct UTUN_INSTANCE* inst, uint64_t peer, uint32_t stream_id) { @@ -498,10 +396,11 @@ void tcp_proxy_server_handle_fin(struct UTUN_INSTANCE* inst, uint64_t peer, uint if (!inst) return; struct tcp_proxy_server* ctx = &inst->tcp_proxy_server; struct tcp_proxy_server_conn* rc = tcp_proxy_server_find_conn(ctx, peer, stream_id); - if (!rc || !rc->tc || !rc->tc->connected) return; + if (!rc || !rc->tc || rc->flow.fin_received) return; + rc->flow.fin_received = 1; DEBUG_DEBUG(DEBUG_CATEGORY_PROXY, "SOCK:FIN_RECV fd=%d sid=%08x — pushing FIN to wq", (int)rc->tc->sock, stream_id); - tcp_conn_push_fin(rc->tc); + if (tcp_conn_push_fin(rc->tc) < 0) on_error_cb(rc->tc, ENOMEM, rc); } // ==================================================================== @@ -550,6 +449,12 @@ void tcp_proxy_server_recv_cb(struct ETCP_CONN* conn, struct ll_entry* entry) { if (inst && inst->tcp_proxy_server.enabled) { struct tcp_proxy_server_conn* rc = tcp_proxy_server_find_conn(&inst->tcp_proxy_server, peer, stream_id); if (rc) { + if (subcmd == TCP_PROXY_SUBCMD_WINDOW) { + if (proxy_flow_window(&rc->flow, entry->dgram + TCP_PROXY_RECV_HDR_SIZE, entry->len - TCP_PROXY_RECV_HDR_SIZE) < 0) + on_error_cb(rc->tc, EPROTO, rc); + else server_wake(rc); + queue_dgram_free(entry); queue_entry_free(entry); return; + } if (subcmd == TCP_PROXY_SUBCMD_DATA) { tcp_proxy_server_handle_data(inst, conn, entry, stream_id); return; } if (subcmd == TCP_PROXY_SUBCMD_CLOSE) { tcp_proxy_server_handle_close(inst, peer, stream_id); queue_dgram_free(entry); queue_entry_free(entry); return; } if (subcmd == TCP_PROXY_SUBCMD_ERROR) { tcp_proxy_server_handle_error(inst, peer, stream_id); queue_dgram_free(entry); queue_entry_free(entry); return; } diff --git a/src/proxy/tcp_proxy_server.h b/src/proxy/tcp_proxy_server.h index 3d9ec657..a72abed6 100644 --- a/src/proxy/tcp_proxy_server.h +++ b/src/proxy/tcp_proxy_server.h @@ -18,42 +18,27 @@ struct UASYNC; struct ll_entry; struct ETCP_CONN; -#define TCP_PROXY_SUBCMD_CONNECT 0x01 -#define TCP_PROXY_SUBCMD_DATA 0x03 -#define TCP_PROXY_SUBCMD_CLOSE 0x04 -#define TCP_PROXY_SUBCMD_ERROR 0x05 -#define TCP_PROXY_SUBCMD_FIN 0x06 - -#define TCP_PROXY_HDR_SIZE 6 // svc_id(1)+subcmd(1)+stream_id(4) -#define TCP_PROXY_CONNECT_HDR_SIZE 12 // HDR_SIZE + dest_ip(4)+dest_port(2) - -// recv-формат (после router_deliver): [svc_id][src][dst][subcmd][stream_id][data] -#define TCP_PROXY_RECV_HDR_SIZE (ROUTER_SVC_PAYLOAD_OFF + 5) // 22: до data (subcmd+stream_id) +#include "proxy_protocol.h" struct tcp_proxy_server_conn { struct tcp_proxy_server_conn* next; struct tcp_proxy_server* ctx; struct tcp_conn* tc; + struct proxy_flow flow; uint32_t stream_id; uint64_t peer_node_id; uint8_t dest_ip[4]; uint16_t dest_port; + uint8_t close_queued; uint8_t cli_closed; // client sent CLOSE - uint8_t close_pending; // CLOSE/ERROR не доставлен, ретрай - void* close_timer; // таймер повтора CLOSE/ERROR - int close_backoff; // backoff: 50..5000 tb (5ms..500ms) void* diag_timer; // 1-секундный таймер диагностики - uint8_t dst_fin_deferred; // FIN от destination отложен (ждём drain read_queue + tx_buf) uint8_t freed; // 1 = уже освобождён, защита от double-free struct queue_waiter_handle pause_waiter; void* retry_timer; // fallback timer (500ms) for force=1 retry on stall - uint8_t* tx_buf; // буфер при backpressure (retry в pause_resume_cb) - uint16_t tx_len; uint32_t bytes_sent; // байт записано в destination (реальный TCP) - uint32_t bytes_recv; // байт прочитано от destination uint32_t bytes_relayed; // байт отправлено клиенту через ETCP uint32_t drain_count; // счётчик вызовов read_queue_drain_cb uint32_t data_count; // счётчик входящих DATA от клиента diff --git a/tests/Makefile.am b/tests/Makefile.am index 21ba8ee5..73856f23 100644 --- a/tests/Makefile.am +++ b/tests/Makefile.am @@ -50,6 +50,8 @@ check_PROGRAMS = \ test_tcp_proxy_server \ test_udp_proxy \ test_icmp_proxy \ + test_proxy_packets \ + test_proxy_regressions \ test_tcp_proxy_client \ test_socks_http_proxy \ test_socks_client \ @@ -633,10 +635,15 @@ test_services_SOURCES = test_services.c test_services_CFLAGS = -I$(top_srcdir)/src -I$(top_srcdir)/src/routing_layer -I$(top_srcdir)/src/transport_layer -I$(top_srcdir)/lib test_services_LDADD = $(LIBUTUN) $(CRYPTO_LIBS) $(COMMON_LIBS) -check_PROGRAMS += test_proxy_packets +# Граничные длины и checksum UDP/ICMP после восстановления пакета. test_proxy_packets_SOURCES = test_proxy_packets.c test_proxy_packets_LDADD = $(LIBUTUN) $(CRYPTO_LIBS) $(COMMON_LIBS) +# Управляемый транспорт для граничных состояний proxy. +test_proxy_regressions_SOURCES = test_proxy_regressions.c +test_proxy_regressions_LDADD = $(LIBUTUN) $(CRYPTO_LIBS) $(COMMON_LIBS) +test_proxy_regressions_LDFLAGS = -Wl,--wrap=etcp_route_send + check_PROGRAMS += test_tcp_io_flush test_tcp_io_flush_SOURCES = test_tcp_io_flush.c test_tcp_io_flush_LDADD = $(LIBUTUN) $(CRYPTO_LIBS) $(COMMON_LIBS) diff --git a/tests/test_proxy_regressions.c b/tests/test_proxy_regressions.c new file mode 100644 index 00000000..bdc64d49 --- /dev/null +++ b/tests/test_proxy_regressions.c @@ -0,0 +1,292 @@ +// Реальные локальные TCP-сокеты и управляемый транспорт: проверяем границы +// handshake, подтверждение connect, изоляцию потоков и порядок DATA/FIN. +#include +#include +#include +#include "utun_instance.h" +#include "etcp_connections.h" +#include "etcp.h" +#include "proxy/socks_proxy.h" +#include "proxy/udp_proxy.h" +#include "proxy/icmp_proxy.h" +#include "lwip_tcp/lwip_tcp.h" +#include "tun_if.h" +#include "debug_config.h" +#include "mem.h" +#ifndef _WIN32 +#include +#else +#define SHUT_WR SD_SEND +#endif + +#define CHECK(x) do { if (!(x)) { fprintf(stderr, "FAIL line %d: %s\n", __LINE__, #x); exit(1); } } while (0) +struct message { struct UTUN_INSTANCE* inst; uint64_t peer; uint8_t svc, cmd; uint32_t sid; size_t len; uint8_t data[4096]; }; +static struct message messages[512]; +static unsigned message_count; +static struct UASYNC* ua; + +// Перехватывается только транспорт. Парсеры, tcp_io, очереди и lifecycle — настоящие. +int __wrap_etcp_route_send(struct UTUN_INSTANCE* inst, uint64_t group, uint64_t peer, + struct ll_entry* entry, int force, int mode) { + (void)group; (void)force; (void)mode; + CHECK(message_count < 512 && entry->len >= 6 && entry->len <= 4102); + struct message* m = &messages[message_count++]; + m->inst = inst; m->peer = peer; m->svc = entry->dgram[0]; m->cmd = entry->dgram[1]; + memcpy(&m->sid, entry->dgram + 2, 4); m->len = entry->len - 6; + memcpy(m->data, entry->dgram + 6, m->len); + queue_dgram_free(entry); queue_entry_free(entry); + return 0; +} + +static void pump(void) { for (int i = 0; i < 12; i++) uasync_poll(ua, 1); } + +static socket_t listen_local(uint16_t* port) { + socket_t fd = socket(AF_INET, SOCK_STREAM, 0); + CHECK(fd != SOCKET_INVALID); + struct sockaddr_in a = {.sin_family = AF_INET, .sin_addr.s_addr = htonl(INADDR_LOOPBACK)}; + CHECK(bind(fd, (struct sockaddr*)&a, sizeof(a)) == 0 && listen(fd, 8) == 0); + socklen_t len = sizeof(a); CHECK(getsockname(fd, (struct sockaddr*)&a, &len) == 0); + *port = ntohs(a.sin_port); socket_set_nonblocking(fd); + return fd; +} + +static socket_t connect_local(uint16_t port) { + socket_t fd = socket(AF_INET, SOCK_STREAM, 0); + struct sockaddr_in a = {.sin_family = AF_INET, .sin_port = htons(port), .sin_addr.s_addr = htonl(INADDR_LOOPBACK)}; + CHECK(connect(fd, (struct sockaddr*)&a, sizeof(a)) == 0); + socket_set_nonblocking(fd); pump(); return fd; +} + +static void send_bytes(socket_t fd, const void* bytes, size_t len) { + CHECK(send(fd, bytes, len, 0) == (ssize_t)len); pump(); +} + +static unsigned count_cmd(uint8_t cmd) { + unsigned count = 0; + for (unsigned i = 0; i < message_count; i++) if (messages[i].cmd == cmd) count++; + return count; +} + +static void deliver(struct UTUN_INSTANCE* inst, uint64_t peer, uint8_t svc, uint8_t cmd, + uint32_t sid, const void* data, size_t len) { + struct ll_entry* e = queue_entry_new(0); + CHECK(e); + e->len = cmd ? TCP_PROXY_RECV_HDR_SIZE + len : ROUTER_SVC_HDR_SIZE; + e->dgram = u_calloc(1, e->len); CHECK(e->dgram); + e->dgram[0] = svc; + uint64_t group = TOPO_GROUP_UTUN; + memcpy(e->dgram + ROUTER_SVC_SRC_OFF, &peer, 8); + memcpy(e->dgram + ROUTER_SVC_GROUP_OFF, &group, 8); + if (cmd) { + e->dgram[ROUTER_SVC_PAYLOAD_OFF] = cmd; + memcpy(e->dgram + ROUTER_SVC_PAYLOAD_OFF + 1, &sid, 4); + if (len) memcpy(e->dgram + TCP_PROXY_RECV_HDR_SIZE, data, len); + } + struct ETCP_CONN conn = {0}; conn.instance = inst; + if (svc == ETCP_RT_ID_TCP_PROXY_SERVER) tcp_proxy_server_recv_cb(&conn, e); + else tcp_proxy_client_router_recv_cb(&conn, e); + pump(); +} + +static void parser_tests(struct UTUN_INSTANCE* inst) { + uint16_t sp, hp; + socket_t reserved = listen_local(&sp); socket_close_wrapper(reserved); + reserved = listen_local(&hp); socket_close_wrapper(reserved); + char socks[64], http[64]; + snprintf(socks, sizeof(socks), "127.0.0.1:%u", sp); snprintf(http, sizeof(http), "127.0.0.1:%u", hp); + struct tcp_proxy_client* p = tcp_proxy_client_create(inst, ua, NULL, NULL, 1500, 1, NULL, 0, 42, 1, socks, 1, http); + CHECK(p); inst->tcp_proxy_client = p; + socket_t fd = connect_local(sp); + // Greeting, CONNECT и ранние данные намеренно одним write. + const uint8_t request[] = {5,1,0, 5,1,0,1,127,0,0,1,0,80, 'a','b','c'}; + message_count = 0; send_bytes(fd, request, sizeof(request)); + CHECK(count_cmd(TCP_PROXY_SUBCMD_CONNECT) == 1 && count_cmd(TCP_PROXY_SUBCMD_DATA) == 0); + uint8_t reply[256]; CHECK(recv(fd, reply, sizeof(reply), 0) == 2 && reply[1] == 0); + uint32_t sid = p->socks_conns->stream_id; + deliver(inst, 99, ETCP_RT_ID_TCP_PROXY_CLIENT, TCP_PROXY_SUBCMD_CONNECTED, sid, NULL, 0); + deliver(inst, 99, ETCP_RT_ID_TCP_PROXY_CLIENT, 0, 0, NULL, 0); + CHECK(p->socks_conn_count == 1 && !p->socks_conns->flow.ready); + deliver(inst, 42, ETCP_RT_ID_TCP_PROXY_CLIENT, TCP_PROXY_SUBCMD_CONNECTED, sid, NULL, 0); + CHECK(recv(fd, reply, sizeof(reply), 0) == 10 && reply[1] == 0); + CHECK(count_cmd(TCP_PROXY_SUBCMD_DATA) == 1 && messages[message_count-1].len == 3); + CHECK(memcmp(messages[message_count-1].data, "abc", 3) == 0); + // Закрытое окно: FIN обоих направлений не должен уничтожать ожидающую передачу. + p->socks_conns->flow.tx_credit = 0; + message_count = 0; send_bytes(fd, "tail", 4); + CHECK(shutdown(fd, SHUT_WR) == 0); pump(); + CHECK(count_cmd(TCP_PROXY_SUBCMD_FIN) == 0); + deliver(inst, 42, ETCP_RT_ID_TCP_PROXY_CLIENT, TCP_PROXY_SUBCMD_FIN, sid, NULL, 0); + CHECK(p->socks_conn_count == 1); + uint32_t credit = 4; + deliver(inst, 42, ETCP_RT_ID_TCP_PROXY_CLIENT, TCP_PROXY_SUBCMD_WINDOW, sid, &credit, 4); + CHECK(count_cmd(TCP_PROXY_SUBCMD_DATA) == 1 && count_cmd(TCP_PROXY_SUBCMD_FIN) == 1); + CHECK(messages[0].cmd == TCP_PROXY_SUBCMD_DATA && memcmp(messages[0].data, "tail", 4) == 0); + CHECK(p->socks_conn_count == 0); socket_close_wrapper(fd); + + fd = connect_local(hp); message_count = 0; + const char first[] = "CONNECT 127.0.0.1:443 HTTP/1.1\r\n"; + send_bytes(fd, first, sizeof(first)-1); CHECK(message_count == 0); + CHECK(recv(fd, reply, sizeof(reply), 0) < 0); + const char second[] = "Host: 127.0.0.1\r\n\r\nTLS"; + send_bytes(fd, second, sizeof(second)-1); CHECK(count_cmd(TCP_PROXY_SUBCMD_CONNECT) == 1); + sid = p->http_conns->stream_id; + CHECK(recv(fd, reply, sizeof(reply), 0) < 0); + deliver(inst, 42, ETCP_RT_ID_TCP_PROXY_CLIENT, TCP_PROXY_SUBCMD_CONNECTED, sid, NULL, 0); + ssize_t n = recv(fd, reply, sizeof(reply), 0); CHECK(n > 12 && memcmp(reply, "HTTP/1.1 200", 12) == 0); + CHECK(messages[message_count-1].len == 3 && memcmp(messages[message_count-1].data, "TLS", 3) == 0); + deliver(inst, 42, ETCP_RT_ID_TCP_PROXY_CLIENT, TCP_PROXY_SUBCMD_ERROR, sid, NULL, 0); + socket_close_wrapper(fd); + + fd = connect_local(hp); message_count = 0; + send_bytes(fd, first, sizeof(first)-1); send_bytes(fd, second, sizeof(second)-1); + sid = p->http_conns->stream_id; + deliver(inst, 42, ETCP_RT_ID_TCP_PROXY_CLIENT, TCP_PROXY_SUBCMD_ERROR, sid, NULL, 0); + n = recv(fd, reply, sizeof(reply), 0); CHECK(n > 12 && memcmp(reply, "HTTP/1.1 502", 12) == 0); + socket_close_wrapper(fd); + + fd = connect_local(sp); message_count = 0; + const uint8_t auth[] = {5,1,2}; send_bytes(fd, auth, sizeof(auth)); + CHECK(recv(fd, reply, sizeof(reply), 0) == 2 && reply[1] == 255 && message_count == 0); + socket_close_wrapper(fd); + fd = connect_local(sp); message_count = 0; + uint8_t v6[25] = {5,1,0,5,1,0,4}; v6[24] = 80; + send_bytes(fd, v6, sizeof(v6)); + n = recv(fd, reply, sizeof(reply), 0); CHECK(n == 12 && reply[3] == 8 && message_count == 0); + socket_close_wrapper(fd); + // Большие заголовки и body, поступающий до CONNECTED, должны сохраниться и разбиться на DATA. + fd = connect_local(hp); message_count = 0; + char upload[32000], expected[32000]; + int header = snprintf(upload, sizeof(upload), "POST http://127.0.0.1/upload HTTP/1.1\r\nHost: 127.0.0.1\r\nX-Pad: "); + memset(upload + header, 'x', 9000); header += 9000; + header += snprintf(upload + header, sizeof(upload) - header, "\r\nContent-Length: 20000\r\n\r\n"); + memset(upload + header, 'B', 20000); + for (int off = 0; off < header + 20000;) { + int chunk = header + 20000 - off; if (chunk > 1024) chunk = 1024; + send_bytes(fd, upload + off, chunk); off += chunk; + } + CHECK(count_cmd(TCP_PROXY_SUBCMD_CONNECT) == 1 && count_cmd(TCP_PROXY_SUBCMD_DATA) == 0); + sid = p->http_conns->stream_id; + CHECK(p->http_conns->tc->read_queue->count <= 8); + deliver(inst, 42, ETCP_RT_ID_TCP_PROXY_CLIENT, TCP_PROXY_SUBCMD_CONNECTED, sid, NULL, 0); + for (int i = 0; i < 10; i++) pump(); + size_t total = 0; + for (unsigned i = 0; i < message_count; i++) if (messages[i].cmd == TCP_PROXY_SUBCMD_DATA) { + CHECK(messages[i].len <= TCP_PROXY_CHUNK && total + messages[i].len <= sizeof(expected)); + memcpy(expected + total, messages[i].data, messages[i].len); total += messages[i].len; + } + const size_t prefix = strlen("http://127.0.0.1"); + CHECK(total == header + 20000 - prefix); + CHECK(memcmp(expected, "POST ", 5) == 0 && memcmp(expected + 5, upload + 5 + prefix, total - 5) == 0); + deliver(inst, 42, ETCP_RT_ID_TCP_PROXY_CLIENT, TCP_PROXY_SUBCMD_ERROR, sid, NULL, 0); + socket_close_wrapper(fd); + tcp_proxy_client_destroy(p); inst->tcp_proxy_client = NULL; pump(); + puts("[PASS] SOCKS/HTTP split/coalesced headers, CONNECTED/error, source identity and half-close"); +} + +static void server_tests(struct UTUN_INSTANCE* inst) { + uint16_t port; socket_t listener = listen_local(&port); + inst->tcp_proxy_server.enabled = 1; inst->tcp_proxy_server.inst = inst; + uint8_t connect_data[6] = {127,0,0,1}; uint16_t wire_port = htons(port); memcpy(connect_data+4, &wire_port, 2); + deliver(inst, 11, ETCP_RT_ID_TCP_PROXY_SERVER, TCP_PROXY_SUBCMD_CONNECT, 2, connect_data, 6); + socket_t a = accept(listener, NULL, NULL); CHECK(a != SOCKET_INVALID); socket_set_nonblocking(a); + deliver(inst, 22, ETCP_RT_ID_TCP_PROXY_SERVER, TCP_PROXY_SUBCMD_CONNECT, 2, connect_data, 6); + socket_t b = accept(listener, NULL, NULL); CHECK(b != SOCKET_INVALID); socket_set_nonblocking(b); + CHECK(inst->tcp_proxy_server.conn_count == 2); + deliver(inst, 11, ETCP_RT_ID_TCP_PROXY_SERVER, TCP_PROXY_SUBCMD_DATA, 2, "AAA", 3); + deliver(inst, 22, ETCP_RT_ID_TCP_PROXY_SERVER, TCP_PROXY_SUBCMD_DATA, 2, "BBB", 3); + uint8_t buf[16]; CHECK(recv(a, buf, sizeof(buf), 0) == 3 && memcmp(buf, "AAA", 3) == 0); + CHECK(recv(b, buf, sizeof(buf), 0) == 3 && memcmp(buf, "BBB", 3) == 0); + deliver(inst, 33, ETCP_RT_ID_TCP_PROXY_SERVER, TCP_PROXY_SUBCMD_CLOSE, 2, NULL, 0); + CHECK(inst->tcp_proxy_server.conn_count == 2); + deliver(inst, 11, ETCP_RT_ID_TCP_PROXY_SERVER, TCP_PROXY_SUBCMD_CONNECT, 2, connect_data, 6); + CHECK(inst->tcp_proxy_server.conn_count == 2); + deliver(inst, 11, ETCP_RT_ID_TCP_PROXY_SERVER, TCP_PROXY_SUBCMD_CLOSE, 2, NULL, 0); + CHECK(inst->tcp_proxy_server.conn_count == 1); + deliver(inst, 22, ETCP_RT_ID_TCP_PROXY_SERVER, TCP_PROXY_SUBCMD_CLOSE, 2, NULL, 0); + CHECK(inst->tcp_proxy_server.conn_count == 0); + socket_close_wrapper(a); socket_close_wrapper(b); socket_close_wrapper(listener); + inst->tcp_proxy_server.enabled = 0; pump(); + puts("[PASS] exit streams isolated by node + stream ID, duplicate CONNECT rejected"); +} + +static void instance_tests(struct UTUN_INSTANCE* a) { + struct UTUN_INSTANCE* b = u_calloc(1, sizeof(*b)); CHECK(b); b->ua = ua; + CHECK(udp_proxy_init(a, ua) == 0 && udp_proxy_init(b, ua) == 0); + CHECK(icmp_proxy_init(a, ua) == 0 && icmp_proxy_init(b, ua) == 0); + CHECK(a->udp_proxy != b->udp_proxy && a->icmp_proxy != b->icmp_proxy); + struct udp_proxy_ctx* udp = b->udp_proxy; struct icmp_proxy_ctx* icmp = b->icmp_proxy; + udp_proxy_destroy(a); icmp_proxy_destroy(a); + CHECK(b->udp_proxy == udp && b->icmp_proxy == icmp); + CHECK(udp_proxy_init(b, ua) == 0 && b->udp_proxy == udp); + udp_proxy_destroy(b); icmp_proxy_destroy(b); u_free(b); + puts("[PASS] UDP/ICMP contexts and destruction are instance-local"); +} + +// Принимающий callback lwIP запускаем с PCB установленного соединения; дальнейшие +// recv/abort/shutdown проходят через настоящие callbacks proxy и стек lwIP. +static struct tcp_pcb* accept_tun(struct tcp_proxy_client* p) { + struct tcp_pcb* pcb = tcp_new(p->lwip); CHECK(pcb); + pcb->state = ESTABLISHED; pcb->local_port = 80; pcb->remote_port = 12345; + pcb->local_ip = htonl(0xc0000201); pcb->remote_ip = htonl(0x0a000002); + pcb->next = p->lwip->active_pcbs; p->lwip->active_pcbs = pcb; + struct tcp_pcb_listen* listener = (struct tcp_pcb_listen*)p->lwip->listen_pcbs; CHECK(listener && listener->accept); + CHECK(listener->accept(p, pcb, LERR_OK) == LERR_OK); + return pcb; +} + +static void tun_tests(struct UTUN_INSTANCE* inst) { +#ifndef _WIN32 + struct tcp_proxy_client_mapping_config map = {.local_port=80, .remote_port=80, .remote_ip="192.0.2.1"}; + struct tcp_proxy_client* p = tcp_proxy_client_create(inst, ua, "testproxy", "10.0.0.1", 1500, 1, &map, 1, 42, 0, NULL, 0, NULL); + CHECK(p); inst->tcp_proxy_client = p; + message_count = 0; + struct tcp_pcb* pcb = accept_tun(p); CHECK(p->conn_count == 1); + uint32_t sid = p->conns->stream_id; + deliver(inst, 42, ETCP_RT_ID_TCP_PROXY_CLIENT, TCP_PROXY_SUBCMD_CONNECTED, sid, NULL, 0); + p->conns->flow.tx_credit = 0; + struct pbuf* pb = pbuf_alloc(PBUF_RAW, 4); CHECK(pb); pbuf_take(pb, "tail", 4); + CHECK(pcb->recv(pcb->callback_arg, pcb, pb, LERR_OK) == LERR_OK); + pcb->state = CLOSE_WAIT; + CHECK(pcb->recv(pcb->callback_arg, pcb, NULL, LERR_OK) == LERR_OK); + deliver(inst, 42, ETCP_RT_ID_TCP_PROXY_CLIENT, TCP_PROXY_SUBCMD_FIN, sid, NULL, 0); + CHECK(p->conn_count == 1 && p->conns->tx_queue->count == 1); + message_count = 0; uint32_t credit = 4; + deliver(inst, 42, ETCP_RT_ID_TCP_PROXY_CLIENT, TCP_PROXY_SUBCMD_WINDOW, sid, &credit, 4); + CHECK(p->conn_count == 0 && count_cmd(TCP_PROXY_SUBCMD_DATA) == 1 && count_cmd(TCP_PROXY_SUBCMD_FIN) == 1); + CHECK(messages[0].cmd == TCP_PROXY_SUBCMD_DATA && memcmp(messages[0].data, "tail", 4) == 0); + pcb = accept_tun(p); CHECK(p->conn_count == 1); + tcp_abort(pcb); CHECK(p->conn_count == 0); + + // IPv4 options смещают UDP-заголовок; последующий фрагмент не является UDP-запросом. + uint8_t ip[36] = {0x46, 0, 0, 36, 0, 0, 0, 0, 64, 17}; + uint32_t src = htonl(0x0a000002), dst = htonl(0xc0000201); + memcpy(ip+12, &src, 4); memcpy(ip+16, &dst, 4); + uint16_t sport=htons(12345), dport=htons(53), ulen=htons(12); + memcpy(ip+24, &sport, 2); memcpy(ip+26, &dport, 2); memcpy(ip+28, &ulen, 2); memcpy(ip+32, "data", 4); + message_count = 0; + for (int fragment = 0; fragment < 2; fragment++) { + struct ll_entry* e = queue_entry_new(0); CHECK(e); + e->dgram = u_calloc(1, sizeof(ip)+1); CHECK(e->dgram); e->len = sizeof(ip)+1; + ip[7] = fragment; memcpy(e->dgram+1, ip, sizeof(ip)); + queue_data_put(p->tun->output_queue, e); pump(); + } + CHECK(message_count == 1 && messages[0].svc == ETCP_RT_ID_UDP_PROXY && messages[0].len == 12); + CHECK(memcmp(messages[0].data+6, &dport, 2) == 0 && memcmp(messages[0].data+8, "data", 4) == 0); + tcp_proxy_client_destroy(p); inst->tcp_proxy_client = NULL; pump(); + puts("[PASS] TUN pending DATA survives FIN; aborted PCB freed; IPv4 options/fragment validation"); +#endif +} + +int main(void) { + debug_config_init(); debug_set_level(DEBUG_LEVEL_WARN); + if (getenv("UTUN_TEST_DEBUG")) debug_set_category_level(DEBUG_CATEGORY_PROXY, DEBUG_LEVEL_DEBUG); + CHECK(socket_platform_init() == 0); + ua = uasync_create(); CHECK(ua); + struct UTUN_INSTANCE* inst = u_calloc(1, sizeof(*inst)); CHECK(inst); + struct utun_config config = {0}; inst->config = &config; inst->ua = ua; + instance_tests(inst); parser_tests(inst); server_tests(inst); tun_tests(inst); + u_free(inst); pump(); uasync_destroy(ua, 0); + CHECK(u_get_allocated_count() == 0); + socket_platform_cleanup(); + return 0; +}