From ba7146f90af503d5fedb6334f34f1fd0a7fde1d9 Mon Sep 17 00:00:00 2001 From: evgeny Date: Mon, 5 Oct 2026 09:27:23 +0200 Subject: [PATCH] Fix HTTP proxy framing, keep-alive and stream lifecycle --- cross-build-win.sh | 2 +- src/Makefile.am | 4 + src/proxy/http_proxy.c | 380 +++++++++++++++++++++++ src/proxy/http_proxy.h | 66 ++++ src/proxy/socks_proxy.c | 549 +++++++++++++++++++++++---------- src/proxy/socks_proxy.h | 12 +- src/proxy/socks_proxy_doc.md | 38 ++- src/proxy/tcp_proxy_server.c | 18 +- tests/Makefile.am | 6 +- tests/test_http_proxy.c | 160 ++++++++++ tests/test_proxy_regressions.c | 239 +++++++++++++- 11 files changed, 1290 insertions(+), 184 deletions(-) create mode 100644 src/proxy/http_proxy.c create mode 100644 src/proxy/http_proxy.h create mode 100644 tests/test_http_proxy.c diff --git a/cross-build-win.sh b/cross-build-win.sh index acd934cb..3ee74cca 100755 --- a/cross-build-win.sh +++ b/cross-build-win.sh @@ -204,7 +204,7 @@ SRC_SOURCES=( control_server.c firewall.c eim_nat.c nat_transport.c transport_layer/dummynet.c etcp_router.c route_crypto.c - proxy/udp_proxy.c proxy/socks_proxy.c + proxy/udp_proxy.c proxy/socks_proxy.c proxy/http_proxy.c lwip_tcp/lwip_pbuf.c ) diff --git a/src/Makefile.am b/src/Makefile.am index 245be663..bd3b6b71 100644 --- a/src/Makefile.am +++ b/src/Makefile.am @@ -84,6 +84,8 @@ utun_CORE_SOURCES = \ proxy/tcp_proxy_server.c \ proxy/udp_proxy.c \ proxy/socks_proxy.c \ + proxy/http_proxy.c \ + proxy/http_proxy.h \ proxy/icmp_proxy.c \ lwip_tcp/lwip_pbuf.c \ lwip_tcp/lwip_tcp.c \ @@ -203,6 +205,8 @@ libutun_a_SOURCES = \ proxy/tcp_proxy_server.c \ proxy/udp_proxy.c \ proxy/socks_proxy.c \ + proxy/http_proxy.c \ + proxy/http_proxy.h \ proxy/icmp_proxy.c \ lwip_tcp/lwip_pbuf.c \ lwip_tcp/lwip_tcp.c \ diff --git a/src/proxy/http_proxy.c b/src/proxy/http_proxy.c new file mode 100644 index 00000000..aa37c5f5 --- /dev/null +++ b/src/proxy/http_proxy.c @@ -0,0 +1,380 @@ +// HTTP headers/body codec: строгие границы сообщений и обработка hop-by-hop полей. +#include "http_proxy.h" +#include "../../lib/debug_config.h" +#include "../../lib/strbuf.h" +#include "../../lib/platform_compat.h" +#include +#include +#include + +// Все протокольные ошибки имеют причину; значения credentials в лог не попадают. +static int http_error(int status, const char* reason) { + DEBUG_WARN(DEBUG_CATEGORY_PROXY, "HTTP codec: status=%d reason=%s", status, reason); + return -status; +} + +// ASCII comparison не зависит от locale. +static int ascii_equal(const char* a, size_t len, const char* b) { + if (len != strlen(b)) return 0; + for (size_t i = 0; i < len; i++) { + unsigned char c = (unsigned char)a[i]; if (c >= 'A' && c <= 'Z') c += 'a' - 'A'; + unsigned char d = (unsigned char)b[i]; if (d >= 'A' && d <= 'Z') d += 'a' - 'A'; + if (c != d) return 0; + } + return 1; +} + +// RFC token, включая имена методов и полей. +static int token_char(unsigned char c) { + return (c >= 'a' && c <= 'z') || (c >= 'A' && c <= 'Z') || (c >= '0' && c <= '9') || + (c && strchr("!#$%&'*+-.^_`|~", c)); +} + +// Проверяет token list и ищет элемент без учёта регистра. +static int list_contains(const char* value, const char* name) { + int found = 0; + while (*value) { + while (*value == ' ' || *value == '\t' || *value == ',') value++; + if (!*value) break; + const char* start = value; + while (token_char((unsigned char)*value)) value++; + if (start == value) return -1; + if (ascii_equal(start, value - start, name)) found = 1; + while (*value == ' ' || *value == '\t') value++; + if (*value && *value != ',') return -1; + } + return found; +} + +// Upgrade содержит protocol-name[/protocol-version], с регистрозависимым сравнением. +static int upgrade_contains(const char* value, const char* selected) { + int found = 0, count = 0; + while (*value) { + while (*value == ' ' || *value == '\t' || *value == ',') value++; + if (!*value) break; + const char* start = value; + while (token_char((unsigned char)*value)) value++; + if (value == start) return -1; + if (*value == '/') { + const char* version = ++value; + while (token_char((unsigned char)*value)) value++; + if (version == value) return -1; + } + if (selected && strlen(selected) == (size_t)(value - start) && !memcmp(start, selected, value - start)) found = 1; + count++; + while (*value == ' ' || *value == '\t') value++; + if (*value && *value != ',') return -1; + } + return count ? found : -1; +} + +// Без переполнения и принятия знака/суффикса. +static int parse_number(const char* s, size_t len, unsigned base, uint64_t* result) { + uint64_t n = 0; + if (!len) return -1; + for (size_t i = 0; i < len; i++) { + unsigned char c = (unsigned char)s[i]; + unsigned d = c >= '0' && c <= '9' ? c - '0' : c >= 'a' && c <= 'f' ? c - 'a' + 10 : + c >= 'A' && c <= 'F' ? c - 'A' + 10 : base; + if (d >= base || n > (UINT64_MAX - d) / base) return -1; + n = n * base + d; + } + *result = n; return 0; +} + +// Накопление прекращается ровно после CRLFCRLF, независимо от размера recv. +int http_proxy_headers_feed(struct http_proxy_message* m, const uint8_t* data, size_t len, size_t* used) { + *used = 0; + while (*used < len && !m->headers_done) { + if (m->header_len == HTTP_PROXY_HEADER_LIMIT) return http_error(431, "headers limit"); + unsigned char c = data[(*used)++]; + if (!c || c == 127 || (c < 32 && c != '\r' && c != '\n' && c != '\t')) return http_error(400, "header control byte"); + m->headers[m->header_len++] = (char)c; + size_t n = m->header_len; + if (c == '\n' && (n < 2 || m->headers[n-2] != '\r')) return http_error(400, "bare LF"); + if (n > 1 && m->headers[n-2] == '\r' && c != '\n') return http_error(400, "bare CR"); + if (!m->first_line_len) { + if (n > HTTP_PROXY_LINE_LIMIT + 2) return http_error(414, "start line limit"); + if (c == '\n') m->first_line_len = n; + } + if (n >= 4 && memcmp(m->headers + n - 4, "\r\n\r\n", 4) == 0) m->headers_done = 1; + } + m->headers[m->header_len] = 0; + return m->headers_done; +} + +// Список полей хранит ссылки на собственный проверенный буфер. +static int parse_fields(struct http_proxy_message* m) { + char* p = m->headers + m->first_line_len; + m->headers[m->first_line_len - 2] = 0; + while (*p != '\r') { + if (m->field_count == HTTP_PROXY_FIELDS_LIMIT) return http_error(431, "field count limit"); + char* end = strstr(p, "\r\n"); + if (!end) return http_error(400, "field terminator"); + char* colon = strchr(p, ':'); + if (!colon || colon >= end || colon == p) return http_error(400, "field name"); + for (char* q = p; q < colon; q++) if (!token_char((unsigned char)*q)) return http_error(400, "field name token"); + *colon = 0; *end = 0; + char* value = colon + 1; while (*value == ' ' || *value == '\t') value++; + char* tail = end; while (tail > value && (tail[-1] == ' ' || tail[-1] == '\t')) *--tail = 0; + m->fields[m->field_count++] = (struct http_proxy_field){p, value}; + p = end + 2; + } + unsigned lengths = 0, transfers = 0, hosts = 0; + for (unsigned i = 0; i < m->field_count; i++) { + const char* name = m->fields[i].name; const char* value = m->fields[i].value; + if (ascii_equal(name, strlen(name), "Content-Length")) { + uint64_t n; + if (parse_number(value, strlen(value), 10, &n) < 0 || (lengths && n != m->length)) + return http_error(400, "Content-Length"); + m->length = n; m->has_length = 1; lengths++; + } else if (ascii_equal(name, strlen(name), "Transfer-Encoding")) { + if (++transfers != 1 || !ascii_equal(value, strlen(value), "chunked") || !m->minor) + return http_error(400, "unsupported Transfer-Encoding"); + m->chunked = 1; + } else if (ascii_equal(name, strlen(name), "Connection")) { + int close = list_contains(value, "close"), upgrade = list_contains(value, "upgrade"); + if (close < 0 || upgrade < 0) return http_error(400, "Connection list"); + m->close |= close; m->upgrade |= upgrade; + if (list_contains(value, "Content-Length") || list_contains(value, "Transfer-Encoding") || list_contains(value, "Host")) + return http_error(400, "Connection names framing field"); + } else if (ascii_equal(name, strlen(name), "Host")) { + if (++hosts != 1 || !*value || strpbrk(value, " \t/@?#,")) return http_error(400, "Host"); + } + } + if (m->chunked && m->has_length) return http_error(400, "Transfer-Encoding with Content-Length"); + m->body = m->chunked ? HTTP_BODY_CHUNK_SIZE : m->has_length && m->length ? HTTP_BODY_LENGTH : HTTP_BODY_DONE; + m->remaining = m->length; + return 0; +} + +// Определяет поля текущего соединения; framing реконструируется отдельно. +static int remove_field(const struct http_proxy_message* m, const char* name) { + const char* fixed[] = {"Connection", "Proxy-Connection", "Proxy-Authorization", "Proxy-Authenticate", + "Keep-Alive", "TE", "Trailer", "Transfer-Encoding", "Content-Length", "Upgrade"}; + for (size_t i = 0; i < sizeof(fixed) / sizeof(fixed[0]); i++) if (ascii_equal(name, strlen(name), fixed[i])) return 1; + for (unsigned i = 0; i < m->field_count; i++) + if (ascii_equal(m->fields[i].name, strlen(m->fields[i].name), "Connection") && list_contains(m->fields[i].value, name) == 1) + return 1; + return 0; +} + +// Завершает сериализацию после проверенных полей; исходные байты учитываются независимо от результата. +static int emit_headers(struct http_proxy_message* m, struct strbuf* out, int close, const char* upgrade, + http_proxy_emit emit, void* arg) { + for (unsigned i = 0; i < m->field_count; i++) { + const char* name = m->fields[i].name; + if (remove_field(m, name) || (m->status == 0 && ascii_equal(name, strlen(name), "Host"))) continue; + if (strbuf_addf(out, "%s: %s\r\n", name, m->fields[i].value) < 0) goto alloc_error; + } + if (m->has_length && strbuf_addf(out, "Content-Length: %llu\r\n", (unsigned long long)m->length) < 0) goto alloc_error; + if (m->chunked && !m->decode_chunks && strbuf_addf(out, "Transfer-Encoding: chunked\r\n") < 0) goto alloc_error; + if (upgrade && *upgrade) { + if (strbuf_addf(out, "Connection: Upgrade\r\nUpgrade: %s\r\n", upgrade) < 0) goto alloc_error; + } else if (close && strbuf_addf(out, "Connection: close\r\n") < 0) goto alloc_error; + if (strbuf_addf(out, "Via: 1.%u utun\r\n\r\n", m->minor) < 0) goto alloc_error; + if (emit(arg, (const uint8_t*)strbuf_str(out), out->len, m->header_len) < 0) goto alloc_error; + return 0; +alloc_error: + DEBUG_ERROR(DEBUG_CATEGORY_PROXY, "HTTP codec: header output failed bytes=%zu", out->len); + return -500; +} + +// Разбирает authority; IPv6 распознаётся, но общий CONNECT wire содержит только IPv4. +static int parse_authority(const char* s, size_t len, int require_port, struct http_proxy_request* r) { + if (!len || memchr(s, '@', len) || memchr(s, '#', len)) return http_error(400, "authority"); + if (s[0] == '[') return http_error(502, "IPv6 destination unsupported"); + const char* colon = memchr(s, ':', len); + size_t host_len = colon ? (size_t)(colon - s) : len; + if (!host_len || host_len >= sizeof(r->host)) return http_error(400, "host length"); + for (size_t i = 0; i < host_len; i++) { + unsigned char c = (unsigned char)s[i]; + if (!((c >= 'a' && c <= 'z') || (c >= 'A' && c <= 'Z') || (c >= '0' && c <= '9') || c == '.' || c == '-')) + return http_error(400, "host syntax"); + } + r->port = 80; + if (colon) { + uint64_t port; + if (parse_number(colon + 1, len - host_len - 1, 10, &port) < 0 || !port || port > 65535) + return http_error(400, "port syntax/range"); + r->port = (uint16_t)port; + } else if (require_port) return http_error(400, "CONNECT port missing"); + memcpy(r->host, s, host_len); r->host[host_len] = 0; + snprintf(r->authority, sizeof(r->authority), "%s:%u", r->host, r->port); + return 0; +} + +// Строгая request line и origin-form без декодирования path/query. +int http_proxy_request_parse(struct http_proxy_message* m, struct http_proxy_request* r, http_proxy_emit emit, void* arg) { + char* line = m->headers; + char* space = memchr(line, ' ', m->first_line_len - 2); + if (!space || space == line) return http_error(400, "request method"); + for (char* p = line; p < space; p++) if (!token_char((unsigned char)*p)) return http_error(400, "method token"); + char* target = space + 1; char* version = strchr(target, ' '); + if (!version || version == target) return http_error(400, "request target/version"); + if ((size_t)(m->headers + m->first_line_len - 2 - version) != 9 || memcmp(version, " HTTP/1.", 8) || + (version[8] != '0' && version[8] != '1')) return http_error(400, "HTTP version"); + m->minor = version[8] - '0'; *space = 0; *version = 0; + for (char* p = target; *p; p++) if ((unsigned char)*p <= 32 || *p == '#' || (unsigned char)*p >= 127) + return http_error(400, "request target byte"); + int ret = parse_fields(m); if (ret < 0) return ret; + unsigned hosts = 0; + for (unsigned i = 0; i < m->field_count; i++) if (ascii_equal(m->fields[i].name, strlen(m->fields[i].name), "Host")) hosts++; + if (m->minor && hosts != 1) return http_error(400, "HTTP/1.1 Host missing"); + r->connect = strcmp(line, "CONNECT") == 0; r->head = strcmp(line, "HEAD") == 0; + r->minor = m->minor; r->close = m->close || !m->minor; + if (r->connect) { + if (m->chunked || m->length) return http_error(400, "CONNECT content"); + return parse_authority(target, strlen(target), 1, r); + } + if (strlen(target) < 7 || !ascii_equal(target, 7, "http://")) return http_error(400, "absolute http URL required"); + char* authority = target + 7; char* path = authority + strcspn(authority, "/?"); + ret = parse_authority(authority, path - authority, 0, r); if (ret < 0) return ret; + if (m->upgrade) { + if (!m->minor) return http_error(400, "Upgrade requires HTTP/1.1"); + unsigned upgrades = 0; + for (unsigned i = 0; i < m->field_count; i++) if (ascii_equal(m->fields[i].name, strlen(m->fields[i].name), "Upgrade")) { + if (++upgrades > 1 || upgrade_contains(m->fields[i].value, NULL) < 0 || strlen(m->fields[i].value) >= sizeof(r->upgrade)) + return http_error(400, "Upgrade value"); + strcpy(r->upgrade, m->fields[i].value); + } + if (!upgrades) return http_error(400, "Upgrade missing"); + } + struct strbuf out = strbuf_new(); + const char* prefix = !*path || *path == '?' ? "/" : ""; + const char* request_path = strcmp(line, "OPTIONS") == 0 && !*path ? "*" : path; + if (*request_path == '*') prefix = ""; + if (strbuf_addf(&out, "%s %s%s HTTP/1.1\r\nHost: %s\r\n", line, prefix, request_path, r->authority) < 0) ret = -500; + else ret = emit_headers(m, &out, !m->upgrade, r->upgrade, emit, arg); + strbuf_free(&out); + if (ret == -500) DEBUG_ERROR(DEBUG_CATEGORY_PROXY, "HTTP codec: request serialization failed"); + return ret; +} + +// Ответы HEAD/1xx/204/304 не имеют тела; Upgrade подтверждается только согласованным 101. +int http_proxy_response_parse(struct http_proxy_message* m, const struct http_proxy_request* r, + int close_client, http_proxy_emit emit, void* arg) { + char* line = m->headers; + if (m->first_line_len < 15 || memcmp(line, "HTTP/1.", 7) || (line[7] != '0' && line[7] != '1') || line[8] != ' ') + return http_error(502, "response version/status line"); + uint64_t status; + if (parse_number(line + 9, 3, 10, &status) < 0 || status < 100 || status > 599 || line[12] != ' ') + return http_error(502, "response status"); + m->minor = line[7] - '0'; m->status = (unsigned)status; + int ret = parse_fields(m); if (ret < 0) return http_error(502, "response fields/framing"); + const char* upgrade = NULL; + if (m->status == 101) { + if (!r->upgrade[0] || !m->upgrade) return http_error(502, "unsolicited Upgrade"); + for (unsigned i = 0; i < m->field_count; i++) if (ascii_equal(m->fields[i].name, strlen(m->fields[i].name), "Upgrade")) { + if (upgrade || upgrade_contains(r->upgrade, m->fields[i].value) != 1) return http_error(502, "Upgrade mismatch"); + upgrade = m->fields[i].value; + } + if (!upgrade) return http_error(502, "101 Upgrade missing"); + } + if (r->head || m->status < 200 || m->status == 204 || m->status == 304) m->body = HTTP_BODY_DONE; + else if (!m->chunked && !m->has_length) { m->body = HTTP_BODY_EOF; close_client = 1; } + if ((m->status < 200 || m->status == 204) && (m->chunked || m->has_length)) return http_error(502, "body headers on no-content response"); + m->close = close_client; + m->decode_chunks = !r->minor && m->chunked; + if (!r->minor && m->status < 200) { + if (emit(arg, NULL, 0, m->header_len) < 0) return http_error(500, "interim credit output failed"); + return 0; + } + struct strbuf out = strbuf_new(); + if (strbuf_addf(&out, "HTTP/1.%u %s\r\n", r->minor, line + 9) < 0) ret = -500; + else ret = emit_headers(m, &out, close_client, upgrade, emit, arg); + strbuf_free(&out); + if (ret == -500) DEBUG_ERROR(DEBUG_CATEGORY_PROXY, "HTTP codec: response serialization failed"); + return ret; +} + +// Chunk extensions: tokens/quoted-string с ограниченным размером строки. +static int chunk_size(const char* line, size_t len, uint64_t* size) { + size_t n = 0; + while (n < len && ((line[n] >= '0' && line[n] <= '9') || (line[n] >= 'a' && line[n] <= 'f') || + (line[n] >= 'A' && line[n] <= 'F'))) n++; + if (parse_number(line, n, 16, size) < 0) return -1; + while (n < len) { + while (n < len && (line[n] == ' ' || line[n] == '\t')) n++; + if (n == len || line[n++] != ';') return -1; + while (n < len && (line[n] == ' ' || line[n] == '\t')) n++; + size_t start = n; while (n < len && token_char((unsigned char)line[n])) n++; + if (start == n) return -1; + while (n < len && (line[n] == ' ' || line[n] == '\t')) n++; + if (n < len && line[n] == '=') { + n++; while (n < len && (line[n] == ' ' || line[n] == '\t')) n++; + if (n < len && line[n] == '"') { + n++; + while (n < len && line[n] != '"') { + if (line[n] == '\\') { if (++n == len) return -1; } + if ((unsigned char)line[n] < 32 && line[n] != '\t') return -1; + n++; + } + if (n == len) return -1; + n++; + } else { + start = n; while (n < len && token_char((unsigned char)line[n])) n++; + if (start == n) return -1; + } + } + } + return 0; +} + +// Тело не буферизуется целиком; только chunk size/trailer line ожидают полного CRLF. +int http_proxy_body_feed(struct http_proxy_message* m, const uint8_t* data, size_t len, size_t* used, + http_proxy_emit emit, void* arg) { + *used = 0; + while (*used < len && m->body != HTTP_BODY_DONE) { + if (m->body == HTTP_BODY_LENGTH || m->body == HTTP_BODY_CHUNK_DATA || m->body == HTTP_BODY_EOF) { + size_t n = len - *used; + if (m->body != HTTP_BODY_EOF && m->remaining < n) n = (size_t)m->remaining; + if (emit(arg, data + *used, n, n) < 0) return http_error(500, "body output failed"); + *used += n; + if (m->body != HTTP_BODY_EOF) { + m->remaining -= n; + if (!m->remaining) m->body = m->body == HTTP_BODY_LENGTH ? HTTP_BODY_DONE : HTTP_BODY_CHUNK_CRLF; + } + } else if (m->body == HTTP_BODY_CHUNK_CRLF) { + unsigned char c = data[(*used)++]; + if (c != (m->crlf_pos ? '\n' : '\r')) return http_error(400, "chunk data CRLF"); + if (++m->crlf_pos == 2) { + if (emit(arg, (const uint8_t*)"\r\n", m->decode_chunks ? 0 : 2, 2) < 0) + return http_error(500, "chunk CRLF output failed"); + m->crlf_pos = 0; m->body = HTTP_BODY_CHUNK_SIZE; + } + } else { + if (m->line_len == HTTP_PROXY_LINE_LIMIT) return http_error(431, "chunk/trailer line limit"); + unsigned char c = data[(*used)++]; + if (!c || c == 127 || (c < 32 && c != '\r' && c != '\n' && c != '\t')) return http_error(400, "chunk/trailer control byte"); + m->line[m->line_len++] = (char)c; + size_t n = m->line_len; + if (c == '\n' && (n < 2 || m->line[n-2] != '\r')) return http_error(400, "chunk/trailer bare LF"); + if (n > 1 && m->line[n-2] == '\r' && c != '\n') return http_error(400, "chunk/trailer bare CR"); + if (c != '\n') continue; + m->line[n] = 0; + int drop = m->decode_chunks; + if (m->body == HTTP_BODY_CHUNK_SIZE) { + if (chunk_size(m->line, n - 2, &m->remaining) < 0) return http_error(400, "chunk size/extensions"); + m->body = m->remaining ? HTTP_BODY_CHUNK_DATA : HTTP_BODY_TRAILERS; + } else { + m->trailer_len += n; + if (m->trailer_len > HTTP_PROXY_HEADER_LIMIT) return http_error(431, "trailers limit"); + if (n == 2) m->body = HTTP_BODY_DONE; + else { + char* colon = memchr(m->line, ':', n - 2); + if (!colon || colon == m->line) return http_error(400, "trailer name"); + for (char* p = m->line; p < colon; p++) if (!token_char((unsigned char)*p)) return http_error(400, "trailer token"); + *colon = 0; + if (ascii_equal(m->line, strlen(m->line), "Host")) return http_error(400, "Host trailer"); + if (ascii_equal(m->line, strlen(m->line), "Content-Length") || + ascii_equal(m->line, strlen(m->line), "Transfer-Encoding")) return http_error(400, "framing trailer"); + drop |= remove_field(m, m->line); *colon = ':'; + } + } + if (emit(arg, (const uint8_t*)m->line, drop ? 0 : n, n) < 0) return http_error(500, "chunk/trailer output failed"); + m->line_len = 0; + } + } + return m->body == HTTP_BODY_DONE; +} diff --git a/src/proxy/http_proxy.h b/src/proxy/http_proxy.h new file mode 100644 index 00000000..cc16a199 --- /dev/null +++ b/src/proxy/http_proxy.h @@ -0,0 +1,66 @@ +/* Потоковый HTTP/1.x codec локального proxy. Не владеет сокетами/DNS/ETCP. + * Заголовки проверяются целиком до пересылки; тело передаётся ограниченными порциями. + * emit получает исходное число байт для возврата ETCP credit после доставки результата. + * Возврат feed: 0=нужно продолжение, 1=сообщение завершено, <0=HTTP error status. */ +#ifndef HTTP_PROXY_H +#define HTTP_PROXY_H + +#include +#include + +#define HTTP_PROXY_LINE_LIMIT 8192 +#define HTTP_PROXY_HEADER_LIMIT 32768 +#define HTTP_PROXY_FIELDS_LIMIT 512 + +typedef int (*http_proxy_emit)(void* arg, const uint8_t* data, size_t len, size_t source_len); + +struct http_proxy_field { char* name; char* value; }; +enum http_proxy_body_state { HTTP_BODY_NONE, HTTP_BODY_LENGTH, HTTP_BODY_EOF, HTTP_BODY_CHUNK_SIZE, + HTTP_BODY_CHUNK_DATA, HTTP_BODY_CHUNK_CRLF, HTTP_BODY_TRAILERS, HTTP_BODY_DONE }; + +struct http_proxy_message { + char headers[HTTP_PROXY_HEADER_LIMIT + 1]; + size_t header_len; + size_t first_line_len; + struct http_proxy_field fields[HTTP_PROXY_FIELDS_LIMIT]; + unsigned field_count; + uint8_t headers_done; + uint8_t minor; + uint8_t close; + uint8_t chunked; + uint8_t has_length; + uint8_t upgrade; + uint8_t decode_chunks; + unsigned status; + uint64_t length; + enum http_proxy_body_state body; + uint64_t remaining; + char line[HTTP_PROXY_LINE_LIMIT + 1]; + size_t line_len; + size_t trailer_len; + unsigned crlf_pos; +}; + +struct http_proxy_request { + char host[256]; + char authority[264]; + char upgrade[HTTP_PROXY_LINE_LIMIT + 1]; + uint16_t port; + uint8_t connect; + uint8_t head; + uint8_t close; + uint8_t minor; +}; + +// Накапливает только заголовки; хвост входа остаётся у вызывающей стороны. +int http_proxy_headers_feed(struct http_proxy_message* m, const uint8_t* data, size_t len, size_t* used); +// Проверяет запрос и отправляет реконструированные заголовки через emit (кроме CONNECT). +int http_proxy_request_parse(struct http_proxy_message* m, struct http_proxy_request* r, http_proxy_emit emit, void* arg); +// Проверяет ответ, учитывает HEAD/Upgrade и формирует собственные connection headers. +int http_proxy_response_parse(struct http_proxy_message* m, const struct http_proxy_request* r, + int close_client, http_proxy_emit emit, void* arg); +// Определяет точную границу тела, фильтрует trailers; поддерживает произвольные TCP-разбиения. +int http_proxy_body_feed(struct http_proxy_message* m, const uint8_t* data, size_t len, size_t* used, + http_proxy_emit emit, void* arg); + +#endif diff --git a/src/proxy/socks_proxy.c b/src/proxy/socks_proxy.c index 4a92c717..6f1df5fb 100644 --- a/src/proxy/socks_proxy.c +++ b/src/proxy/socks_proxy.c @@ -1,5 +1,6 @@ // socks_proxy.c — SOCKS5 / HTTP CONNECT proxy (client side) #include "socks_proxy.h" +#include "http_proxy.h" #include "tcp_proxy_server.h" #include "etcp.h" #include "etcp_api.h" @@ -40,6 +41,30 @@ static void socks_dns_done_cb(const struct adns_result* res, void* arg); static void socks_dns_error_and_close(struct socks_proxy_conn* c); static void socks_bp_register(struct socks_proxy_conn* c); static void socks_maybe_relay_fin(struct socks_proxy_conn* c); +static void http_finish(struct socks_proxy_conn* c); +static void http_fail(struct socks_proxy_conn* c, int status, const char* reason); +static void http_read(struct socks_proxy_conn* c); +static void http_receive(struct socks_proxy_conn* c, const uint8_t* data, size_t len); +static void http_deadline(void* arg); + +struct socks_http_state { + struct http_proxy_message request; + struct http_proxy_message response; + struct http_proxy_request target; + void* timer; + unsigned busy; + uint32_t rx_forwarded; + uint8_t active; + uint8_t request_done; + uint8_t response_done; + uint8_t response_started; + uint8_t close_client; + uint8_t retire_cmd; + uint8_t input_eof; + uint8_t tunnel; + uint8_t stop_upload; + uint8_t sending; +}; 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); @@ -71,7 +96,7 @@ static int send_connect(struct socks_proxy_conn* c) { memcpy(buf, c->dest_ip, 4); memcpy(buf + 4, &c->dest_port, 2); DEBUG_INFO(DEBUG_CATEGORY_PROXY, "SOCKS proxy: CONNECT sid=%08x to %d.%d.%d.%d:%d via_node=%016llx %s", c->stream_id, c->dest_ip[0], c->dest_ip[1], c->dest_ip[2], c->dest_ip[3], - ntohs(c->dest_port), (unsigned long long)c->via_node_id, c->is_http ? "http" : "socks"); + ntohs(c->dest_port), (unsigned long long)c->via_node_id, c->http ? "http" : "socks"); return send_msg(c->inst, TOPO_GROUP_UTUN, c->via_node_id, TCP_PROXY_SUBCMD_CONNECT, c->stream_id, buf, 6, 1); } @@ -115,6 +140,247 @@ static int write_to_client(struct socks_proxy_conn* c, const uint8_t* data, uint return 0; } +// Deadline относится к фазе HTTP, а не к каждому принятому байту заголовка. +static void http_arm(struct socks_proxy_conn* c, int timeout_tb) { + if (c->http->timer) uasync_cancel_timeout(c->ua, c->http->timer); + c->http->timer = uasync_set_timeout(c->ua, timeout_tb, c, http_deadline, "http_deadline"); + if (!c->http->timer) { + DEBUG_ERROR(DEBUG_CATEGORY_PROXY, "HTTP: deadline allocation failed sid=%08x", c->stream_id); + on_error_cb(c->tc, ENOMEM, c); + } +} + +// Request output ограничен одним входным блоком плюс реконструированные заголовки. +static int http_request_output(void* arg, const uint8_t* data, size_t len, size_t source_len) { + (void)source_len; + struct socks_proxy_conn* c = arg; + if (!len) return 0; + size_t total = c->tx_len + len; + if (total > HTTP_PROXY_HEADER_LIMIT + HTTP_PROXY_LINE_LIMIT + TCP_PROXY_CHUNK) { + DEBUG_ERROR(DEBUG_CATEGORY_PROXY, "HTTP: request output limit sid=%08x bytes=%zu", c->stream_id, total); return -1; + } + uint8_t* buf = u_realloc(c->tx_buf, total); + if (!buf) { DEBUG_ERROR(DEBUG_CATEGORY_PROXY, "HTTP: request allocation failed sid=%08x bytes=%zu", c->stream_id, total); return -1; } + memcpy(buf + c->tx_len, data, len); c->tx_buf = buf; c->tx_len = (uint16_t)total; + return 0; +} + +// Deferred write consumer исключает flush посреди реконструкции одного header block. +static int http_response_output(void* arg, const uint8_t* data, size_t len, size_t source_len) { + struct socks_proxy_conn* c = arg; + struct socks_http_state* h = c->http; + if (h->response.status >= 200 || h->response.status == 101) h->response_started = 1; + for (size_t off = 0; off < len;) { + size_t n = len - off; if (n > TCP_PROXY_CHUNK) n = TCP_PROXY_CHUNK; + if (write_to_client(c, data + off, (uint16_t)n) < 0) return -1; + off += n; + } + h->rx_forwarded += (uint32_t)source_len; + tcp_conn_set_flushed(c->tc, socks_flushed_cb); + return 0; +} + +// Завершение отдельного upstream не закрывает клиентский keep-alive. +static void http_finish(struct socks_proxy_conn* c) { + struct socks_http_state* h = c->http; + if (!h || h->busy || h->tunnel || c->freed || c->free_soon_id) return; + if (h->retire_cmd) { + uint8_t cmd = h->retire_cmd; + h->busy++; + int ret = send_msg(c->inst, TOPO_GROUP_UTUN, c->via_node_id, cmd, c->stream_id, NULL, 0, 1); + h->busy--; + if (ret < 0) { + socks_bp_register(c); return; + } + DEBUG_INFO(DEBUG_CATEGORY_PROXY, "HTTP: upstream retired sid=%08x command=%u", c->stream_id, cmd); + h->retire_cmd = 0; h->active = 0; c->stream_id = 0; + if (!c->rem_closed) { + if (h->timer) uasync_cancel_timeout(c->ua, h->timer); + int eof = h->input_eof; + memset(h, 0, sizeof(*h)); h->input_eof = eof; c->state = HTTP_STATE_REQUEST; + http_arm(c, 600000); + DEBUG_DEBUG(DEBUG_CATEGORY_PROXY, "HTTP: keep-alive ready buffered=%u", c->buf_len); + http_read(c); return; + } + } + if (c->rem_closed) { + if (!c->close_queued) { + c->close_queued = 1; + if (tcp_conn_push_close(c->tc) < 0) on_error_cb(c->tc, ENOMEM, c); + } + return; + } + if (h->input_eof && !h->active && !c->buf_len && !c->tc->read_queue->head) { + if (h->request.header_len) { + http_fail(c, 400, "EOF in request headers"); return; + } + c->rem_closed = 1; http_finish(c); return; + } + if (!h->active) return; + if (h->input_eof && c->flow.ready && !h->request_done && !h->stop_upload && !c->buf_len && !c->tc->read_queue->head) { + http_fail(c, 400, "EOF in request body"); return; + } + if ((h->request_done || h->stop_upload) && !c->tx_buf && !h->target.upgrade[0] && !c->flow.fin_sent && c->flow.ready) { + h->busy++; + c->flow.fin_pending = 1; proxy_flow_flush(&c->flow); + h->busy--; + if (c->rem_closed || c->free_soon_id) { http_finish(c); return; } + } + if (!h->response_done || c->tx_buf || c->tc->write_buf || c->tc->write_queue->head) return; + if (!h->request_done && !h->close_client) return; + h->retire_cmd = h->close_client && !h->request_done ? TCP_PROXY_SUBCMD_ERROR : TCP_PROXY_SUBCMD_CLOSE; + proxy_flow_destroy(&c->flow); c->flow.ready = 0; + etcp_router_cancel_send_ready(c->inst, TOPO_GROUP_UTUN, c->via_node_id, ETCP_RT_ID_TCP_PROXY_SERVER, &c->tx_waiter); + if (c->tx_retry_timer) { uasync_cancel_timeout(c->ua, c->tx_retry_timer); c->tx_retry_timer = NULL; } + int close = h->close_client || h->target.close || (h->input_eof && !c->buf_len && !c->tc->read_queue->head); + if (h->timer) { uasync_cancel_timeout(c->ua, h->timer); h->timer = NULL; } + // Ответ полностью передан в OS TCP: старый credit/FIN не может попасть в новую транзакцию. + if (close) c->rem_closed = 1; + http_finish(c); + if (h->retire_cmd || c->rem_closed) http_arm(c, 600000); +} + +// Ошибка до final headers получает отдельный ответ; после них закрываем прерванное сообщение. +static void http_fail(struct socks_proxy_conn* c, int status, const char* reason) { + struct socks_http_state* h = c->http; + DEBUG_WARN(DEBUG_CATEGORY_PROXY, "HTTP: fail sid=%08x status=%d reason=%s active=%u request_done=%u response_started=%u", + c->stream_id, status, reason, h->active, h->request_done, h->response_started); + if (c->rem_closed) { socks_proxy_conn_free_soon(c); return; } + c->rem_closed = 1; + if (c->dns_q) { adns_cancel(c->dns_q); c->dns_q = NULL; } + proxy_flow_destroy(&c->flow); c->flow.ready = 0; + etcp_router_cancel_send_ready(c->inst, TOPO_GROUP_UTUN, c->via_node_id, ETCP_RT_ID_TCP_PROXY_SERVER, &c->tx_waiter); + h->stop_upload = 1; + if (c->tx_buf && !h->busy) { u_free(c->tx_buf); c->tx_buf = NULL; c->tx_len = 0; } + if (h->active) h->retire_cmd = TCP_PROXY_SUBCMD_ERROR; + if (!h->response_started) { + char reply[160]; + const char* text = status == 431 ? "Request Header Fields Too Large" : status == 414 ? "URI Too Long" : + status == 408 ? "Request Timeout" : status == 504 ? "Gateway Timeout" : + status == 502 ? "Bad Gateway" : status == 500 ? "Internal Server Error" : "Bad Request"; + int n = snprintf(reply, sizeof(reply), "HTTP/1.1 %d %s\r\nContent-Length: 0\r\nConnection: close\r\n\r\n", status, text); + if (write_to_client(c, (const uint8_t*)reply, (uint16_t)n) < 0) { on_error_cb(c->tc, ENOMEM, c); return; } + } + http_arm(c, 600000); http_finish(c); +} + +// Ожидание заголовков/connect ограничено; зависшее закрытие освобождается принудительно. +static void http_deadline(void* arg) { + struct socks_proxy_conn* c = arg; + c->http->timer = NULL; + if (c->http->tunnel) return; + if (c->rem_closed) { + DEBUG_WARN(DEBUG_CATEGORY_PROXY, "HTTP: close drain deadline expired sid=%08x", c->stream_id); + socks_proxy_conn_free_soon(c); return; + } + if (!c->http->active && !c->http->request.header_len && !c->buf_len) { + DEBUG_INFO(DEBUG_CATEGORY_PROXY, "HTTP: idle connection expired"); + c->rem_closed = 1; http_finish(c); return; + } + http_fail(c, c->http->active ? 504 : 408, "deadline expired"); +} + +// Один read block может содержать конец тела и начало следующего запроса. +static void http_read(struct socks_proxy_conn* c) { + struct socks_http_state* h = c->http; + if (c->rem_closed || c->free_soon_id || c->tx_buf || h->stop_upload || (h->active && (!c->flow.ready || h->request_done))) return; + h->busy++; + if (!c->buf_len) { + struct ll_entry* e = queue_data_get(c->tc->read_queue); + if (!e) { h->busy--; http_finish(c); return; } + memcpy(c->buf, e->dgram, e->len); c->buf_len = e->len; + memory_pool_free(c->tc->data_pool, e->dgram); queue_entry_free(e); + } + size_t used = 0; + int ret; + if (!h->active) { + if (!h->request.header_len) http_arm(c, 300000); + ret = http_proxy_headers_feed(&h->request, c->buf, c->buf_len, &used); + consume_header(c, (uint16_t)used); + if (ret > 0) { + ret = http_proxy_request_parse(&h->request, &h->target, http_request_output, c); + if (ret == 0) { + c->stream_id = ++(*c->next_stream_id); + if (!c->stream_id) c->stream_id = ++(*c->next_stream_id); + proxy_flow_init(&c->flow, c->inst, c->ua, c->via_node_id, ETCP_RT_ID_TCP_PROXY_SERVER, + c->stream_id, socks_flow_wake, c); + h->active = 1; h->close_client = h->target.close; h->request_done = h->request.body == HTTP_BODY_DONE; + c->dest_port = htons(h->target.port); http_arm(c, 150000); + DEBUG_INFO(DEBUG_CATEGORY_PROXY, "HTTP: request sid=%08x host=%s port=%u connect=%u framing=%u length=%llu", + c->stream_id, h->target.host, h->target.port, h->target.connect, h->request.body, + (unsigned long long)h->request.length); + socks_issue_dns(c, h->target.host, h->target.connect ? DNSK_HTTP_CONNECT : DNSK_HTTP_PROXY); + } + } + } else { + ret = http_proxy_body_feed(&h->request, c->buf, c->buf_len, &used, http_request_output, c); + consume_header(c, (uint16_t)used); http_arm(c, 600000); + if (ret > 0) h->request_done = 1; + } + h->busy--; + if (ret < 0) { http_fail(c, -ret, "request parse/output"); return; } + if (c->flow.ready && (c->tx_buf || (c->buf_len && !h->request_done))) socks_flow_wake(c); + else if (!h->active || (!h->request_done && !c->tx_buf)) queue_resume_callback(c->tc->read_queue); + http_finish(c); +} + +// Принятые ETCP bytes остаются charged до доставки реконструированного результата. +static void http_receive(struct socks_proxy_conn* c, const uint8_t* data, size_t len) { + struct socks_http_state* h = c->http; + if (!h->active || h->response_done) { http_fail(c, 502, "unsolicited response bytes"); return; } + h->busy++; + int ret = 0; + while (len && !h->response_done) { + size_t used = 0; + if (!h->response.headers_done) { + if (!h->response.header_len) http_arm(c, 300000); + ret = http_proxy_headers_feed(&h->response, data, len, &used); + data += used; len -= used; + if (ret <= 0) break; + int early = h->response.headers[9] >= '2' && (!h->request_done || c->tx_buf); + if (early) { + h->close_client = 1; + h->stop_upload = 1; + DEBUG_INFO(DEBUG_CATEGORY_PROXY, "HTTP: early final response sid=%08x; stop upload", c->stream_id); + } + ret = http_proxy_response_parse(&h->response, &h->target, h->close_client, http_response_output, c); + if (ret < 0) break; + http_arm(c, 600000); + if (h->response.status == 101) { + if (!h->request_done || c->tx_buf) { ret = -502; break; } + h->tunnel = 1; + if (h->timer) { uasync_cancel_timeout(c->ua, h->timer); h->timer = NULL; } + if (len && http_response_output(c, data, len, len) < 0) { ret = -500; break; } + if (c->buf_len) { + c->tx_buf = u_malloc(c->buf_len); + if (!c->tx_buf) { ret = -500; break; } + memcpy(c->tx_buf, c->buf, c->buf_len); c->tx_len = c->buf_len; c->buf_len = 0; + } + len = 0; c->state = SOCKS_STATE_RELAY; socks_flow_wake(c); break; + } + if (h->response.status < 200) { memset(&h->response, 0, sizeof(h->response)); continue; } + h->close_client |= h->response.close; + } + if (h->response.body == HTTP_BODY_DONE) h->response_done = 1; + else { + ret = http_proxy_body_feed(&h->response, data, len, &used, http_response_output, c); + data += used; len -= used; http_arm(c, 600000); + if (ret < 0) break; + if (ret > 0) h->response_done = 1; + } + } + if (len && ret >= 0 && !h->tunnel) ret = -502; + h->busy--; + if (ret < 0) { + if (ret == -500) on_error_cb(c->tc, ENOMEM, c); + else http_fail(c, 502, "response parse/boundary"); + return; + } + if (h->stop_upload && c->tx_buf && !h->busy) socks_flow_wake(c); + if (h->rx_forwarded && !c->tc->write_buf && !c->tc->write_queue->head) socks_flushed_cb(c->tc, c); + http_finish(c); +} + // ==================================================================== // SOCKS5 handshake // ==================================================================== @@ -183,127 +449,6 @@ error: { } } -// ==================================================================== -// HTTP CONNECT / HTTP proxy handshake -// ==================================================================== -static void process_http_request(struct socks_proxy_conn* c) { - char* line_end = memmem(c->buf, c->buf_len, "\r\n", 2); - if (!line_end) return; - - 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)) { 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 ===== - if (strncmp(line, "CONNECT ", 8) == 0) { - char host_port[256]; - if (sscanf(line, "CONNECT %255s HTTP/", host_port) != 1) { - DEBUG_ERROR(DEBUG_CATEGORY_PROXY, "socks_proxy: bad http connect: %s", line); - goto error; - } - - char* colon = strrchr(host_port, ':'); - if (!colon) { DEBUG_ERROR(DEBUG_CATEGORY_PROXY, "socks_proxy: no port in %s", host_port); goto error; } - *colon = '\0'; int port = atoi(colon + 1); - if (port <= 0 || port > 65535) { DEBUG_ERROR(DEBUG_CATEGORY_PROXY, "socks_proxy: bad port %d", port); goto error; } - - 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); - consume_header(c, headers_len); - socks_issue_dns(c, host_port, DNSK_HTTP_CONNECT); - return; - } - - // ===== Не-CONNECT: HTTP-прокси (GET, POST, PUT, HEAD, OPTIONS, ...) ===== - - char method[16] = {0}, url[512] = {0}; - if (sscanf(line, "%15s %511s", method, url) < 2) { - DEBUG_ERROR(DEBUG_CATEGORY_PROXY, "socks_proxy: bad http request line: %s", line); - goto error; - } - - if (strncmp(url, "http://", 7) != 0) { - DEBUG_ERROR(DEBUG_CATEGORY_PROXY, "socks_proxy: not an absolute http URL: %s", url); - goto error; - } - - char host[256]; - int port = 80; - const char* host_start = url + 7; - const char* path_start = strchr(host_start, '/'); - const char* col = NULL; - for (const char* p = host_start; (path_start ? p < path_start : *p); p++) { - if (*p == ':') { col = p; break; } - } - - if (col && (!path_start || col < path_start)) { - size_t hlen = (size_t)(col - host_start); - if (hlen >= sizeof(host)) goto error; - memcpy(host, host_start, hlen); host[hlen] = '\0'; - port = atoi(col + 1); - if (port <= 0 || port > 65535) { DEBUG_ERROR(DEBUG_CATEGORY_PROXY, "socks_proxy: bad port %d in %s", port, url); goto error; } - } else if (path_start) { - size_t hlen = (size_t)(path_start - host_start); - if (hlen >= sizeof(host)) goto error; - memcpy(host, host_start, hlen); host[hlen] = '\0'; - } else { - size_t hlen = strlen(host_start); - if (hlen >= sizeof(host)) goto error; - memcpy(host, host_start, hlen); host[hlen] = '\0'; - } - - if (host[0] == '\0') { DEBUG_ERROR(DEBUG_CATEGORY_PROXY, "socks_proxy: empty host in %s", url); goto error; } - if (!path_start || *path_start == '\0') path_start = "/"; - - c->dest_port = htons((uint16_t)port); - DEBUG_INFO(DEBUG_CATEGORY_PROXY, "socks_proxy: HTTP %s %s:%d%s sid=%08x", - method, host, port, path_start, c->stream_id); - - const char* version = strstr(line, " HTTP/"); - char version_str[16] = " HTTP/1.1"; - if (version) { - size_t vlen = strlen(version); - if (vlen >= sizeof(version_str)) vlen = sizeof(version_str) - 1; - memcpy(version_str, version, vlen); - version_str[vlen] = '\0'; - } - - uint16_t rest_off = line_len + 2; - uint16_t rest_len = headers_len - rest_off; - 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); - goto error; - } - memcpy(hdr_buf + hdr_n, c->buf + rest_off, rest_len); - uint16_t hdr_chunk = (uint16_t)hdr_n + rest_len; - uint16_t extra = (headers_len < c->buf_len) ? (uint16_t)(c->buf_len - headers_len) : 0; - uint16_t total = hdr_chunk + extra; - - uint8_t* pkt = u_malloc(total); - if (!pkt) { DEBUG_ERROR(DEBUG_CATEGORY_PROXY, "socks_proxy: tx pkt malloc=%u failed sid=%08x", total, c->stream_id); goto error; } - memcpy(pkt, hdr_buf, hdr_chunk); - if (extra) memcpy(pkt + hdr_chunk, c->buf + headers_len, extra); - - c->dns_pending = pkt; - c->dns_pending_len = total; - c->buf_len = 0; - socks_issue_dns(c, host, DNSK_HTTP_PROXY); - return; - -error: { - uint8_t resp[] = "HTTP/1.1 400 Bad Request\r\n\r\n"; - write_to_client(c, resp, (uint16_t)strlen((char*)resp)); tcp_conn_push_close(c->tc); -} -} - // ==================================================================== // DNS-резолвинг (неблокирующий) + финализация рукопожатия // ==================================================================== @@ -323,18 +468,24 @@ static void socks_connected(struct socks_proxy_conn* c) { return; } c->flow.ready = 1; + if (c->http && !c->http->target.connect) { + c->state = HTTP_STATE_RELAY; + http_arm(c, 600000); + DEBUG_INFO(DEBUG_CATEGORY_PROXY, "HTTP: upstream connected sid=%08x buffered=%u", c->stream_id, c->tx_len); + socks_flow_wake(c); return; + } 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; } + char reply[64]; + int len = snprintf(reply, sizeof(reply), "HTTP/1.%u 200 Connection Established\r\n\r\n", c->http->target.minor); + if (write_to_client(c, (const uint8_t*)reply, (uint16_t)len) < 0) { on_error_cb(c->tc, ENOMEM, c); return; } + c->http->tunnel = 1; + if (c->http->timer) { uasync_cancel_timeout(c->ua, c->http->timer); c->http->timer = NULL; } } 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; - } else if (c->buf_len) { + 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; @@ -344,13 +495,9 @@ static void socks_connected(struct socks_proxy_conn* c) { } static void socks_dns_error_and_close(struct socks_proxy_conn* c) { - if (c->is_http) { - uint8_t resp[] = "HTTP/1.1 502 Bad Gateway\r\n\r\n"; - write_to_client(c, resp, (uint16_t)strlen((char*)resp)); - } else { - uint8_t err[] = { 0x05, 0x04, 0x00, 0x01, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00 }; - write_to_client(c, err, 10); - } + if (c->http) { http_fail(c, 502, "DNS/connect failed"); return; } + uint8_t err[] = { 0x05, 0x04, 0x00, 0x01, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00 }; + write_to_client(c, err, sizeof(err)); c->rem_closed = 1; tcp_conn_push_close(c->tc); } @@ -362,7 +509,7 @@ static void socks_issue_dns(struct socks_proxy_conn* c, const char* host, uint8_ socks_finalize(c); return; } - c->state = c->is_http ? HTTP_STATE_CONNECTING : SOCKS_STATE_CONNECTING; + c->state = SOCKS_STATE_CONNECTING; c->dns_q = adns_resolve(c->ua, host, NULL, socks_dns_done_cb, c); if (!c->dns_q) { DEBUG_ERROR(DEBUG_CATEGORY_PROXY, "socks_proxy: DNS resolve start failed for %s sid=%08x", host, c->stream_id); @@ -389,6 +536,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->http && !c->http->tunnel) { http_read(c); return; } 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; } @@ -397,8 +545,7 @@ static void on_read_cb(struct ll_queue* q, void* arg) { queue_dgram_free(e); queue_entry_free(e); queue_resume_callback(q); return; } - if (c->state == SOCKS_STATE_GREETING || c->state == SOCKS_STATE_REQUEST || - c->state == HTTP_STATE_REQUEST) { + if (c->state == SOCKS_STATE_GREETING || c->state == SOCKS_STATE_REQUEST) { // накапливаем данные для рукопожатия uint16_t space = sizeof(c->buf) - c->buf_len; if (space < e->len) { @@ -409,18 +556,14 @@ 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); - if (c->is_http) { - process_http_request(c); - } else { - if (c->state == SOCKS_STATE_GREETING) process_socks_greeting(c); - if (!c->rem_closed && c->state == SOCKS_STATE_REQUEST) process_socks_request(c); - } + if (c->state == SOCKS_STATE_GREETING) process_socks_greeting(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; } 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); + DEBUG_TRACE(DEBUG_CATEGORY_PROXY, "SOCKS PROXY SEND sid=%08x len=%u ret=%d http=%d", c->stream_id, e->len, ret, c->http != NULL); if (ret == 0) { memory_pool_free(c->tc->data_pool, e->dgram); queue_entry_free(e); queue_resume_callback(q); @@ -438,17 +581,59 @@ static void on_read_cb(struct ll_queue* q, void* arg) { // Отправка буфера частями сохраняет HTTP headers/body, независимо от размера одного DATA. static void socks_flow_wake(void* arg) { struct socks_proxy_conn* c = arg; + if (c->http && c->http->retire_cmd) { http_finish(c); return; } + if (c->http && c->http->sending) return; if (c->rem_closed || c->freed || !c->flow.ready) return; + if (c->http) c->http->busy++; while (c->tx_buf) { + if (c->http && c->http->stop_upload) { u_free(c->tx_buf); c->tx_buf = NULL; c->tx_len = 0; break; } 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; } + uint8_t chunk[TCP_PROXY_CHUNK]; + const uint8_t* data = c->tx_buf; + // Текущий chunk уже принадлежит send: синхронный 101 должен видеть только ещё не отправленный хвост. + if (c->http) { + memcpy(chunk, c->tx_buf, len); data = chunk; + 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; } + c->http->sending = 1; + } + int ret = send_data(c, data, len, 0); + if (c->http) c->http->sending = 0; + if (c->http && (c->http->stop_upload || c->rem_closed || c->free_soon_id)) { + u_free(c->tx_buf); c->tx_buf = NULL; c->tx_len = 0; + c->http->busy--; http_finish(c); return; + } + if (ret < 0) { + if (c->http) { + uint8_t* restored = u_realloc(c->tx_buf, c->tx_len + len); + if (!restored) { + c->http->busy--; + DEBUG_ERROR(DEBUG_CATEGORY_PROXY, "HTTP: send rollback allocation failed sid=%08x", c->stream_id); + on_error_cb(c->tc, ENOMEM, c); return; + } + memmove(restored + len, restored, c->tx_len); memcpy(restored, chunk, len); + c->tx_buf = restored; c->tx_len += len; + } + socks_bp_register(c); + if (c->http) c->http->busy--; + return; + } + if (!c->http) { + 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; } + } + } + if (c->http) c->http->busy--; + if (c->http && !c->http->tunnel) { + if (c->buf_len && !c->http->busy && !c->http->request_done) http_read(c); + else if (!c->http->request_done) queue_resume_callback(c->tc->read_queue); + http_finish(c); + } else { + queue_resume_callback(c->tc->read_queue); + socks_maybe_relay_fin(c); } - queue_resume_callback(c->tc->read_queue); - socks_maybe_relay_fin(c); } static void tx_waiter_cb(struct ll_queue* q, void* arg) { @@ -482,12 +667,26 @@ static void socks_maybe_relay_fin(struct socks_proxy_conn* c) { } static void on_fin_cb(struct tcp_conn* tc, void* arg) { - (void)tc; + struct socks_proxy_conn* c = arg; + if (c->http && !c->http->tunnel) { + c->http->input_eof = tc->fin_remote; http_finish(c); return; + } socks_maybe_relay_fin(arg); } static void socks_flushed_cb(struct tcp_conn* tc, void* arg) { struct socks_proxy_conn* c = arg; + if (c->http && !c->http->tunnel) { + if (!tc->write_buf && !tc->write_queue->head) { + uint32_t credit = c->http->rx_forwarded; c->http->rx_forwarded = 0; + if (credit && c->http->active && !c->http->retire_cmd && !c->rem_closed) { + c->http->busy++; + proxy_flow_consume(&c->flow, credit); + c->http->busy--; + } + } + http_finish(c); return; + } 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); @@ -500,7 +699,7 @@ static void on_error_cb(struct tcp_conn* tc, int err, void* arg) { "bytes_client=%u bytes_exit=%u wq=%d", c->stream_id, err, c->state, tc->fin_remote, tc->fin_local, c->rem_closed, c->close_sent, c->close_pending, c->bytes_to_client, c->bytes_from_exit, c->tc->write_queue->count); - int notify = !c->close_sent && !c->close_pending; + int notify = !c->close_sent && !c->close_pending && (!c->http || c->http->active); struct UTUN_INSTANCE* inst = c->inst; uint64_t via = c->via_node_id; uint32_t sid = c->stream_id; @@ -513,7 +712,9 @@ 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->flow.fin_sent && !c->rem_closed && !c->close_sent && !c->close_pending) send_close(c); + if (c->http && c->http->active && !c->http->retire_cmd) + c->http->retire_cmd = TCP_PROXY_SUBCMD_ERROR; + if (!c->http && !c->flow.fin_sent && !c->rem_closed && !c->close_sent && !c->close_pending) send_close(c); socks_proxy_conn_free_soon(c); } @@ -530,7 +731,7 @@ static void on_accept_cb(socket_t sock, void* arg) { struct socks_proxy_conn* c = u_calloc(1, sizeof(struct socks_proxy_conn)); if (!c) { socket_close_wrapper(csock); DEBUG_ERROR(DEBUG_CATEGORY_PROXY, "socks_proxy: u_calloc failed"); return; } c->stream_id = ++(*ctx->next_stream_id); - c->is_http = (uint8_t)ctx->is_http; + c->next_stream_id = ctx->next_stream_id; 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; @@ -538,6 +739,14 @@ static void on_accept_cb(socket_t sock, void* arg) { 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; } + if (ctx->is_http) { + c->http = u_calloc(1, sizeof(*c->http)); + if (!c->http) { + DEBUG_ERROR(DEBUG_CATEGORY_PROXY, "HTTP: state allocation failed"); tcp_conn_destroy(c->tc); u_free(c); return; + } + queue_set_callback_defer(c->tc->write_queue, 1); + http_arm(c, 300000); + } 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); @@ -565,6 +774,7 @@ int socks_proxy_handle_etcp(struct socks_proxy_conn** head, int* count, 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 (c->http && !c->http->active) 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); @@ -576,6 +786,7 @@ int socks_proxy_handle_etcp(struct socks_proxy_conn** head, int* count, 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 (c->http && !c->http->tunnel) { http_receive(c, data, data_len); return 1; } if (data_len > 0) { if (data_len > c->tc->data_pool->object_size) { DEBUG_ERROR(DEBUG_CATEGORY_PROXY, "socks_proxy: data_len=%zu > pool_sz=%zu, dropping sid=%08x", @@ -600,6 +811,11 @@ int socks_proxy_handle_etcp(struct socks_proxy_conn** head, int* count, } if (subcmd == TCP_PROXY_SUBCMD_CLOSE) { + if (c->http && !c->http->tunnel) { + if (c->http->response_done) http_finish(c); + else http_fail(c, 502, "upstream CLOSE before completion"); + return 1; + } DEBUG_INFO(DEBUG_CATEGORY_PROXY, "socks_proxy: REM_CLOSED sid=%08x", stream_id); 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); @@ -609,6 +825,7 @@ int socks_proxy_handle_etcp(struct socks_proxy_conn** head, int* count, if (subcmd == TCP_PROXY_SUBCMD_ERROR) { DEBUG_WARN(DEBUG_CATEGORY_PROXY, "socks_proxy: ERROR from exit sid=%08x", stream_id); + if (c->http && !c->http->tunnel) { http_fail(c, 502, "upstream ERROR"); return 1; } 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); @@ -620,6 +837,14 @@ int socks_proxy_handle_etcp(struct socks_proxy_conn** head, int* count, DEBUG_INFO(DEBUG_CATEGORY_PROXY, "socks_proxy: FIN from exit sid=%08x", stream_id); if (c->flow.fin_received) return 1; c->fin_remote = 1; c->flow.fin_received = 1; + if (c->http && !c->http->tunnel) { + struct socks_http_state* h = c->http; + if (!h->response.headers_done) http_fail(c, 502, "EOF in response headers"); + else if (h->response.body == HTTP_BODY_EOF) { h->response_done = 1; http_finish(c); } + else if (!h->response_done) http_fail(c, 502, "EOF in response body"); + else http_finish(c); + return 1; + } if (tcp_conn_push_fin(c->tc) < 0) on_error_cb(c->tc, ENOMEM, c); return 1; } @@ -642,6 +867,15 @@ void socks_proxy_conn_free(struct socks_proxy_conn* c) { if (!c) return; if (c->freed) return; c->freed = 1; + if (c->http) { + if (c->http->timer) uasync_cancel_timeout(c->ua, c->http->timer); + if (c->http->active) { + uint8_t cmd = c->http->retire_cmd ? c->http->retire_cmd : TCP_PROXY_SUBCMD_ERROR; + if (send_msg(c->inst, TOPO_GROUP_UTUN, c->via_node_id, cmd, c->stream_id, NULL, 0, 1) < 0) + DEBUG_WARN(DEBUG_CATEGORY_PROXY, "HTTP: upstream cancellation unavailable sid=%08x", c->stream_id); + } + u_free(c->http); c->http = NULL; + } 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; @@ -653,7 +887,6 @@ void socks_proxy_conn_free(struct socks_proxy_conn* c) { } if (head) { struct socks_proxy_conn** prev = head; while (*prev) { if (*prev == c) { *prev = c->next; if (count) (*count)--; break; } prev = &(*prev)->next; } } if (c->tx_buf) { u_free(c->tx_buf); c->tx_buf = NULL; c->tx_len = 0; } - if (c->dns_pending) { u_free(c->dns_pending); c->dns_pending = NULL; c->dns_pending_len = 0; } if (c->dns_q) { adns_cancel(c->dns_q); c->dns_q = NULL; } if (c->tx_retry_timer) { uasync_cancel_timeout(c->ua, c->tx_retry_timer); c->tx_retry_timer = NULL; } if (c->tc) { tcp_conn_destroy(c->tc); c->tc = NULL; } @@ -671,7 +904,7 @@ void socks_proxy_conn_free_all(struct socks_proxy_conn** head, int* count) { struct listen_ctx* socks_proxy_init_listen(struct UASYNC* ua, const char* addr_str, struct socks_proxy_conn** conns, int* count, uint32_t* next_stream_id, struct UTUN_INSTANCE* inst, uint64_t via_node_id, int is_http) { - char ip[64], *colon = strrchr(addr_str, ':'); + char ip[64]; const char* colon = strrchr(addr_str, ':'); if (!colon) { DEBUG_ERROR(DEBUG_CATEGORY_PROXY, "socks_proxy: bad addr '%s' (need IP:PORT)", addr_str); return NULL; } size_t ip_len = (size_t)(colon - addr_str); if (ip_len > sizeof(ip) - 1) ip_len = sizeof(ip) - 1; diff --git a/src/proxy/socks_proxy.h b/src/proxy/socks_proxy.h index a8246d24..359c5873 100644 --- a/src/proxy/socks_proxy.h +++ b/src/proxy/socks_proxy.h @@ -20,6 +20,7 @@ struct ll_entry; struct ll_queue; struct tcp_conn; struct adns_query; +struct socks_http_state; enum { SOCKS_STATE_GREETING = 0, // ждём SOCKS5 greeting (ver + methods) @@ -28,11 +29,9 @@ enum { SOCKS_STATE_RELAY // релей данных (== HTTP_STATE_RELAY) }; -// HTTP-состояния совпадают по значениям с SOCKS-состояниями (общая state-машина -// в on_read_cb), поэтому явно привязаны к SOCKS_STATE_* — иначе возникает коллизия -// SOCKS_STATE_CONNECTING(2) == HTTP_STATE_RELAY(2) и RELAY-данные дропаются. +// DNS/connect и сырой туннель используют общие состояния; HTTP framing хранится отдельно в http. enum { - HTTP_STATE_REQUEST = SOCKS_STATE_REQUEST, // ждём первую строку HTTP-запроса + HTTP_STATE_REQUEST = SOCKS_STATE_REQUEST, // ждём заголовки очередного HTTP-запроса HTTP_STATE_CONNECTING = SOCKS_STATE_CONNECTING, // DNS в процессе HTTP_STATE_RELAY = SOCKS_STATE_RELAY // релей данных }; @@ -47,13 +46,14 @@ struct socks_proxy_conn { uint64_t via_node_id; struct proxy_flow flow; uint32_t stream_id; + uint32_t* next_stream_id; // HTTP: новый поток для следующей транзакции + struct socks_http_state* http; uint8_t dest_ip[4]; uint16_t dest_port; uint8_t fin_remote; uint8_t rem_closed; uint8_t close_sent; uint8_t close_pending; - uint8_t is_http; uint8_t state; uint8_t freed; uint8_t close_queued; @@ -69,8 +69,6 @@ struct socks_proxy_conn { uint8_t fin_deferred; // FIN отложен до сброса tx_buf struct adns_query* dns_q; // активный DNS-запрос (NULL когда нет) uint8_t dns_kind; // DNSK_SOCKS / DNSK_HTTP_CONNECT / DNSK_HTTP_PROXY - uint8_t* dns_pending; // HTTP proxy: накопленный запрос (headers+body) до DNS - uint16_t dns_pending_len; }; struct listen_ctx; diff --git a/src/proxy/socks_proxy_doc.md b/src/proxy/socks_proxy_doc.md index 8e2b534b..0ef99ccd 100644 --- a/src/proxy/socks_proxy_doc.md +++ b/src/proxy/socks_proxy_doc.md @@ -1,11 +1,12 @@ # socks_proxy ## 1. Назначение -Модуль реализует клиентскую сторону SOCKS5 и HTTP CONNECT прокси. Принимает TCP-соединения от локальных приложений (curl, браузер), выполняет handshake (SOCKS5 или HTTP CONNECT), определяет целевой IP:port (с разрешением доменных имён через DNS) и мультиплексирует трафик через ETCP к exit node. Взаимодействует с `tcp_proxy_server` на стороне exit node. +Модуль реализует локальные SOCKS5 CONNECT, HTTP/1.0–1.1 forward proxy и HTTP CONNECT. Разрешает доменное имя назначения и передаёт TCP через ETCP к `tcp_proxy_server` на exit node. Обычный HTTP разбирается до границы каждого запроса и ответа; CONNECT и согласованный Upgrade переходят в сырой туннель. Файлы: - `src/proxy/socks_proxy.c` — реализация - `src/proxy/socks_proxy.h` — интерфейс +- `src/proxy/http_proxy.c/h` — потоковый HTTP codec, независимый от DNS, TCP и ETCP ## 2. Как пользоваться @@ -26,7 +27,7 @@ struct listen_ctx* listen = socks_proxy_init_listen( int handled = socks_proxy_handle_etcp(&conns, &conn_count, stream_id, subcmd, data, data_len); // handled=1 — найден и обработан, 0 — stream_id не найден ``` -Поддерживаемые subcmd: `DATA`, `CLOSE`, `ERROR`, `FIN`. +Поддерживаемые subcmd: `CONNECTED`, `DATA`, `WINDOW`, `CLOSE`, `ERROR`, `FIN`. ### Поиск соединения ```c @@ -42,7 +43,7 @@ socks_proxy_close_listen(ua, listen, &sock_out); // слушатель ### Ключевые нюансы - **Бэкпрессур (backpressure):** при неудаче `etcp_route_send()` данные сохраняются в `tx_buf/tx_len`, и регистрируется waiter через `etcp_router_on_send_ready()`. Когда ETCP готов к отправке, `tx_waiter_cb` повторяет попытку и возобновляет `read_queue`. -- **Несколько соединений:** управляются через связный список `struct socks_proxy_conn*` + счётчик. Каждое соединение — отдельный TCP-сокет (`tcp_conn`) и независимый ETCP stream. +- **Несколько соединений:** управляются через связный список `struct socks_proxy_conn*` + счётчик. В HTTP keep-alive каждое сообщение получает новый upstream stream; клиентский TCP-сокет сохраняется. Pipelining обрабатывается последовательно, включая запросы к разным адресатам. - **stream_id** генерируется монотонно (`++next_stream_id`), передаётся в заголовке ETCP-пакетов для мультиплексирования. ## 3. API @@ -54,9 +55,10 @@ socks_proxy_close_listen(ua, listen, &sock_out); // слушатель - `stream_id` — уникальный идентификатор потока в ETCP - `dest_ip[4]`, `dest_port` — целевой адрес (куда хочет подключиться клиент) - `state` — текущее состояние (SOCKS_STATE_* или HTTP_STATE_*) -- `is_http` — флаг: 0=SOCKS5, 1=HTTP +- `http` — отдельное состояние HTTP-транзакции и таймер; NULL для SOCKS5 +- `next_stream_id` — заимствованный генератор идентификаторов последующих HTTP-транзакций - `rem_closed`, `fin_remote`, `close_sent`, `close_pending` — флаги состояния закрытия -- `buf[1024]`, `buf_len` — буфер накопления данных для парсинга handshake +- `buf[16384]`, `buf_len` — непрочитанный хвост входного блока; HTTP headers хранятся в codec - `tx_buf`, `tx_len`, `tx_waiter` — бэкпрессур: отложенные данные и waiter handle - `bytes_to_client`, `bytes_from_exit` — счётчики переданных байт (диагностика) - `head`, `count` — указатели на голову списка и счётчик (для самоудаления) @@ -87,7 +89,9 @@ socks_proxy_close_listen(ua, listen, &sock_out); // слушатель | `write_to_client(c, data, len)` | Запись данных в `tc->write_queue` (отправка клиенту через TCP). | | `process_socks_greeting(c)` | Парсинг SOCKS5 GREETING: ver=5, выбор метода 0x00 (no auth), ответ 05 00. | | `process_socks_request(c)` | Парсинг SOCKS5 CONNECT REQUEST (atyp=1/3/4, IPv4/IPv6/domain). DNS-резолвинг для доменов. Ответ 05 00 00 01 + dummy addr. | -| `process_http_request(c)` | Парсинг HTTP: CONNECT host:port (туннель) или GET/POST с абсолютным URL (режим HTTP-прокси с переписыванием URL в относительный). | +| `http_read(c)`, `http_receive(c)` | Потоковый разбор клиентского запроса и upstream ответа через HTTP codec. | +| `http_finish(c)` | Завершение/отмена upstream, ожидание доставки ответа клиенту и переход к следующему запросу. | +| `http_fail(c)`, `http_deadline(c)` | Ошибки и тайм-ауты: собственный HTTP error до final headers, закрытие оборванного сообщения после них. | | `on_accept_cb(sock, arg)` | Callback accept: создаёт `socks_proxy_conn`, `tcp_conn_create`, вешает read callback. | | `on_read_cb(q, arg)` | Callback чтения из TCP: handshake-парсинг или релей данных через ETCP с бэкпрессуром. | | `tx_waiter_cb(q, arg)` | Callback бэкпрессура: повторная отправка отложенных данных при готовности ETCP. | @@ -107,9 +111,10 @@ uTun → exit (ETCP): TCP_PROXY CONNECT[dest_ip+port] ### HTTP CONNECT flow (is_http=1) ``` -Клиент → uTun: CONNECT host:port HTTP/1.1\r\n\r\n -uTun → Клиент: HTTP/1.1 200 Connection Established\r\n\r\n +Клиент → uTun: CONNECT host:port HTTP/1.1\r\nHost: host:port\r\n\r\n uTun → exit (ETCP): TCP_PROXY CONNECT[dest_ip+port] +exit → uTun (ETCP): CONNECTED +uTun → Клиент: HTTP/1.1 200 Connection Established\r\n\r\n Клиент ↔ uTun ↔ exit: RELAY данных через ETCP DATA ``` @@ -117,9 +122,22 @@ uTun → exit (ETCP): TCP_PROXY CONNECT[dest_ip+port] ``` Клиент → uTun: GET http://host:port/path HTTP/1.1\r\nHost: ...\r\n\r\n uTun → exit (ETCP): TCP_PROXY CONNECT[dest_ip+port] + TCP_PROXY DATA[GET /path HTTP/1.1\r\n...] -Клиент ↔ uTun ↔ exit: RELAY данных через ETCP DATA +exit → uTun: HTTP headers/body через ETCP DATA, затем FIN +uTun → Клиент: проверенный и реконструированный HTTP response +uTun → exit (ETCP): CLOSE завершённого stream +Клиент → uTun: следующий запрос → новый CONNECT с новым stream_id ``` -URL переписывается в относительный (`/path`), остальные заголовки передаются как есть. +URL преобразуется в origin-form без потери query; Host строится из absolute URL. Заголовки текущего соединения, включая поля из Connection, и Proxy-Authorization удаляются в обоих направлениях. Framing и Connection формируются заново, добавляется Via. Upstream обычного запроса получает Connection: close; frontend HTTP/1.1 сохраняет keep-alive для ответов с определённой длиной. + +Codec поддерживает Content-Length, chunked с extensions/trailers, промежуточные ответы, HEAD/204/304, тела ответа до EOF и подтверждённый 101. Тело передаётся по мере поступления, без накопления целиком. HTTP/1.0 получает декодированное chunked тело и закрытие соединения. Неоднозначный framing (CL+TE, разные CL), obs-fold, неверные порты и управляющие байты отклоняются до отправки заголовков. + +ETCP credit возвращается за исходные принятые байты после опустошения клиентской write queue и write_buf. Частичный header/chunk marker остаётся charged до проверки и доставки результата; удалённые trailers также учитываются. Переход к новому stream разрешён только после доставки полного ответа и отправки CLOSE/ERROR старого stream. Синхронная доставка через loopback защищена счётчиком busy и флагом sending. + +Лимиты: headers 32 KiB, start line и chunk/trailer line 8 KiB, не более 512 полей. Назначения — IPv4 или DNS-имя с IPv4-результатом; IPv6 отклоняется с 502 из-за формата общего proxy CONNECT. Поддерживаемый transfer coding — chunked. + +Тайм-ауты: request/response headers 30 s, DNS/connect 15 s, отсутствие прогресса передачи и idle keep-alive 60 s. Заголовки имеют общий deadline, который не продлевается каждым байтом. Зависшее закрытие принудительно освобождается через 60 s. CONNECT/Upgrade после установки туннеля используют обычный lifecycle TCP. Ошибки и смена транзакций логируются в категории proxy; credentials не выводятся. + +ERROR на exit немедленно уничтожает сокет и отложенный upload; CLOSE сохраняет обычное завершение очереди. ## 5. Зависимости - `tcp_proxy_server.h` — константы subcmd (`TCP_PROXY_SUBCMD_*`) и размер заголовка (`TCP_PROXY_HDR_SIZE=6`) diff --git a/src/proxy/tcp_proxy_server.c b/src/proxy/tcp_proxy_server.c index 08a399a6..f1209594 100644 --- a/src/proxy/tcp_proxy_server.c +++ b/src/proxy/tcp_proxy_server.c @@ -74,6 +74,7 @@ static void server_maybe_finish(struct tcp_proxy_server_conn* rc) { if (tc->fin_remote && !tc->read_queue->head && !rc->flow.fin_sent) { rc->flow.fin_pending = 1; proxy_flow_flush(&rc->flow); + if (rc->freed) return; } if (rc->flow.fin_sent && rc->flow.fin_received && tc->fin_local) { server_queue_close(rc); @@ -172,9 +173,14 @@ 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 = arg; if (!rc->tc || rc->cli_closed) return; + struct tcp_conn* tc = rc->tc; 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); + // Loopback ERROR может отменить exit внутри send; pools tc освобождаются отложенно. + if (rc->freed) { + memory_pool_free(tc->data_pool, e->dgram); queue_entry_free(e); return; + } if (ret == 0) { rc->bytes_relayed += e->len; rc->drain_count++; memory_pool_free(rc->tc->data_pool, e->dgram); queue_entry_free(e); @@ -213,6 +219,8 @@ static int conn_total(struct tcp_proxy_server_conn* rc) { return n; } +static void server_conn_free_deferred(void* arg) { u_free(arg); } + void tcp_proxy_server_conn_free(struct tcp_proxy_server_conn* rc) { if (!rc) return; DEBUG_DEBUG(DEBUG_CATEGORY_PROXY, "SOCK:FREE enter rc=%p freed=%d sid=%08x", (void*)rc, rc->freed, rc->stream_id); @@ -233,8 +241,9 @@ void tcp_proxy_server_conn_free(struct tcp_proxy_server_conn* rc) { if (rc->retry_timer) { uasync_cancel_timeout(rc->ua, rc->retry_timer); rc->retry_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); - u_free(rc); + DEBUG_DEBUG(DEBUG_CATEGORY_PROXY, "SOCK:FREE deferred rc=%p", (void*)rc); + if (!uasync_call_soon(rc->ua, rc, server_conn_free_deferred)) + DEBUG_ERROR(DEBUG_CATEGORY_PROXY, "proxy exit: deferred free allocation failed sid=%08x", rc->stream_id); } struct tcp_proxy_server_conn* tcp_proxy_server_find_conn(struct tcp_proxy_server* ctx, uint64_t peer, uint32_t stream_id) { @@ -401,7 +410,10 @@ void tcp_proxy_server_handle_error(struct UTUN_INSTANCE* inst, uint64_t peer, ui DEBUG_DEBUG(DEBUG_CATEGORY_PROXY, "SOCK:ERROR_RECV fd=%d sid=%08x total=%d fin=%d write_pend=%d", rc->tc ? (int)rc->tc->sock : -1, stream_id, conn_total(rc), rc->tc ? rc->tc->fin_remote : 0, rc->tc ? write_pending(rc->tc) : 0); - tcp_proxy_server_handle_close(inst, peer, stream_id); + // ERROR отменяет поток: очередь назначения может ждать бесконечно, дренировать её нельзя. + DEBUG_INFO(DEBUG_CATEGORY_PROXY, "proxy exit: abort peer=%016llx sid=%08x queued=%zu", + (unsigned long long)peer, stream_id, rc->tc ? queue_total_bytes(rc->tc->write_queue) : 0); + tcp_proxy_server_conn_free(rc); } void tcp_proxy_server_handle_fin(struct UTUN_INSTANCE* inst, uint64_t peer, uint32_t stream_id) { diff --git a/tests/Makefile.am b/tests/Makefile.am index 2f51ab33..dc10c6cf 100644 --- a/tests/Makefile.am +++ b/tests/Makefile.am @@ -54,6 +54,7 @@ check_PROGRAMS = \ test_icmp_proxy \ test_proxy_packets \ test_proxy_regressions \ + test_http_proxy \ test_tcp_proxy_client \ test_socks_http_proxy \ test_socks_client \ @@ -723,7 +724,10 @@ 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 +test_proxy_regressions_LDFLAGS = -Wl,--wrap=etcp_route_send -Wl,--wrap=uasync_set_timeout + +test_http_proxy_SOURCES = test_http_proxy.c ../src/proxy/http_proxy.c +test_http_proxy_LDADD = $(COMMON_LIBS) check_PROGRAMS += test_tcp_io_flush test_tcp_io_flush_SOURCES = test_tcp_io_flush.c diff --git a/tests/test_http_proxy.c b/tests/test_http_proxy.c new file mode 100644 index 00000000..675c9436 --- /dev/null +++ b/tests/test_http_proxy.c @@ -0,0 +1,160 @@ +// Проверки HTTP codec на каждом возможном разбиении headers и побайтовом chunked body. +#include +#include +#include +#include "proxy/http_proxy.h" +#include "debug_config.h" +#include "mem.h" + +#define CHECK(x) do { if (!(x)) { fprintf(stderr, "FAIL line %d: %s\n", __LINE__, #x); exit(1); } } while (0) +struct output { char bytes[100000]; size_t len, source; }; + +static int collect(void* arg, const uint8_t* data, size_t len, size_t source) { + struct output* out = arg; + CHECK(out->len + len < sizeof(out->bytes)); + if (len) memcpy(out->bytes + out->len, data, len); + out->len += len; out->source += source; out->bytes[out->len] = 0; + return 0; +} + +static int request(const char* text, struct http_proxy_message* m, struct http_proxy_request* r, struct output* out) { + memset(m, 0, sizeof(*m)); memset(r, 0, sizeof(*r)); memset(out, 0, sizeof(*out)); + size_t used; + int ret = http_proxy_headers_feed(m, (const uint8_t*)text, strlen(text), &used); + if (ret <= 0) return ret; + return http_proxy_request_parse(m, r, collect, out); +} + +static int response(const char* text, struct http_proxy_message* m, const struct http_proxy_request* r, struct output* out) { + memset(m, 0, sizeof(*m)); memset(out, 0, sizeof(*out)); + size_t used; + int ret = http_proxy_headers_feed(m, (const uint8_t*)text, strlen(text), &used); + if (ret <= 0) return ret; + return http_proxy_response_parse(m, r, r->close, collect, out); +} + +static void header_tests(void) { + const char text[] = "POST http://example.org:8080?x=1 HTTP/1.1\r\nHost: wrong.example\r\n" + "Proxy-Authorization: Basic secret\r\nConnection: X-Private, keep-alive\r\n" + "X-Private: secret\r\nAuthorization: Bearer origin\r\nContent-Length: 3\r\n\r\nabcNEXT"; + size_t header_len = strstr(text, "\r\n\r\n") + 4 - text; + for (size_t split = 0; split <= header_len; split++) { + struct http_proxy_message m = {0}; struct http_proxy_request r = {0}; struct output out = {0}; size_t used; + CHECK(http_proxy_headers_feed(&m, (const uint8_t*)text, split, &used) == (split == header_len)); + CHECK(used == split && out.source == 0); + CHECK(http_proxy_headers_feed(&m, (const uint8_t*)text + split, sizeof(text)-1-split, &used) == 1); + CHECK(used == header_len - split); + CHECK(http_proxy_request_parse(&m, &r, collect, &out) == 0); + CHECK(!strcmp(r.host, "example.org") && r.port == 8080); + CHECK(strstr(out.bytes, "POST /?x=1 HTTP/1.1\r\nHost: example.org:8080\r\n")); + CHECK(!strstr(out.bytes, "secret") && !strstr(out.bytes, "wrong.example")); + CHECK(strstr(out.bytes, "Authorization: Bearer origin\r\n")); + CHECK(strstr(out.bytes, "Connection: close\r\n") && out.source == header_len); + CHECK(http_proxy_body_feed(&m, (const uint8_t*)text + header_len, 7, &used, collect, &out) == 1); + CHECK(used == 3 && out.source == header_len + 3 && !strcmp(out.bytes + out.len - 3, "abc")); + } + struct http_proxy_message m; struct http_proxy_request r; struct output out; + CHECK(request("OPTIONS http://example.org HTTP/1.1\r\nHost: ignored\r\n\r\n", &m, &r, &out) == 0); + CHECK(strstr(out.bytes, "OPTIONS * HTTP/1.1\r\n")); + CHECK(request("CONNECT example.org:443 HTTP/1.1\r\nHost: ignored\r\n\r\n", &m, &r, &out) == 0); + CHECK(r.connect && r.port == 443 && out.len == 0); + const char* invalid[] = { + "CONNECT example.org:443junk HTTP/1.1\r\nHost: x\r\n\r\n", + "CONNECT example.org:443 GARBAGE\r\nHost: x\r\n\r\n", + "CONNECT example.org HTTP/1.1\r\nHost: x\r\n\r\n", + "CONNECT example.org:443 HTTP/1.1\r\nHost: x\r\nContent-Length: 1\r\n\r\n", + "GET http://x/ HTTP/1.1\r\n\r\n", + "GET http://x/ HTTP/1.1\r\nHost: x\r\nHost: x\r\n\r\n", + "GET http://x/ HTTP/1.1\r\nHost : x\r\n\r\n", + "GET http://x/ HTTP/1.1\nHost: x\n\n", + "GET http://x/ HTTP/1.1\r\nHost: x\r\nX: a\r\n folded\r\n\r\n", + "POST http://x/ HTTP/1.1\r\nHost: x\r\nContent-Length: 1\r\nContent-Length: 2\r\n\r\n", + "POST http://x/ HTTP/1.1\r\nHost: x\r\nContent-Length: 18446744073709551616\r\n\r\n", + "POST http://x/ HTTP/1.1\r\nHost: x\r\nContent-Length: +1\r\n\r\n", + "POST http://x/ HTTP/1.1\r\nHost: x\r\nContent-Length: 1\r\nTransfer-Encoding: chunked\r\n\r\n", + "POST http://x/ HTTP/1.1\r\nHost: x\r\nTransfer-Encoding: chunked, chunked\r\n\r\n", + "POST http://x/ HTTP/1.1\r\nHost: x\r\nTransfer-Encoding: gzip\r\n\r\n", + "GET http://x/ HTTP/1.1\r\nHost: x\r\nConnection: Host\r\n\r\n", + "GET http://x/ HTTP/1.1\r\nHost: x\r\nConnection: Upgrade\r\nUpgrade: bad/version/extra\r\n\r\n", + "GET http://user@x/ HTTP/1.1\r\nHost: x\r\n\r\n", + "GET http://x:65536/ HTTP/1.1\r\nHost: x\r\n\r\n", + "GET http://x:0/ HTTP/1.1\r\nHost: x\r\n\r\n", + "GET http://x/#fragment HTTP/1.1\r\nHost: x\r\n\r\n" + }; + for (size_t i = 0; i < sizeof(invalid)/sizeof(invalid[0]); i++) { + CHECK(request(invalid[i], &m, &r, &out) == -400 && out.len == 0); + } + CHECK(request("GET http://[::1]/ HTTP/1.1\r\nHost: [::1]\r\n\r\n", &m, &r, &out) == -502); + char long_text[40000]; int n = snprintf(long_text, sizeof(long_text), "GET http://x/"); + memset(long_text+n, 'a', 3000); n += 3000; + snprintf(long_text+n, sizeof(long_text)-n, " HTTP/1.1\r\nHost: x\r\n\r\n"); + CHECK(request(long_text, &m, &r, &out) == 0); + n = snprintf(long_text, sizeof(long_text), "POST http://x/ HTTP/1.1\r\nHost: x\r\nX-Pad: "); + memset(long_text+n, 'a', 30000); n += 30000; + snprintf(long_text+n, sizeof(long_text)-n, "\r\nContent-Length: 2000\r\n\r\n"); + CHECK(request(long_text, &m, &r, &out) == 0); + n = snprintf(long_text, sizeof(long_text), "GET http://x/ HTTP/1.1\r\nHost: x\r\nX-Pad: "); + memset(long_text+n, 'a', 33000); long_text[n+33000] = 0; + CHECK(request(long_text, &m, &r, &out) == -431); + puts("[PASS] HTTP strict headers, absolute URL/Host, hop fields, arbitrary splits and limits"); +} + +static void body_tests(void) { + struct http_proxy_message m; struct http_proxy_request r; struct output out; + const char header[] = "POST http://x/ HTTP/1.1\r\nHost: x\r\nTransfer-Encoding: chunked\r\nConnection: X-Hop\r\n\r\n"; + const char body[] = "3 ; foo=\"x\\\"y\"\r\nabc\r\n2\r\nde\r\n0\r\nX-End: yes\r\nX-Hop: secret\r\nProxy-Authorization: secret\r\n\r\n"; + CHECK(request(header, &m, &r, &out) == 0); + size_t head_out = out.len; + for (size_t i = 0; i < sizeof(body)-1; i++) { + size_t used; + CHECK(http_proxy_body_feed(&m, (const uint8_t*)body+i, 1, &used, collect, &out) == (i == sizeof(body)-2)); + CHECK(used == 1 && out.source <= strlen(header)+i+1); + } + CHECK(out.source == strlen(header)+strlen(body)); + CHECK(!strcmp(out.bytes+head_out, "3 ; foo=\"x\\\"y\"\r\nabc\r\n2\r\nde\r\n0\r\nX-End: yes\r\n\r\n")); + const char* bad[] = {"Z\r\n", "10000000000000000\r\n", "1\r\na!", "1\na", "0\r\nContent-Length: 1\r\n\r\n", + "0\r\nHost: x\r\n\r\n", "0\r\n folded\r\n\r\n", "1;foo=\"bad\r\n"}; + for (size_t i = 0; i < sizeof(bad)/sizeof(bad[0]); i++) { + size_t used; CHECK(request(header, &m, &r, &out) == 0); + CHECK(http_proxy_body_feed(&m, (const uint8_t*)bad[i], strlen(bad[i]), &used, collect, &out) == -400); + } + r = (struct http_proxy_request){.minor=0,.close=1}; + const char reply[] = "HTTP/1.1 200 OK\r\nTransfer-Encoding: chunked\r\n\r\n"; + CHECK(response(reply, &m, &r, &out) == 0 && m.decode_chunks); + CHECK(!strstr(out.bytes, "Transfer-Encoding") && strstr(out.bytes, "HTTP/1.0 200 OK\r\n")); + head_out = out.len; + for (size_t i = 0; i < sizeof(body)-1; i++) { + size_t used; CHECK(http_proxy_body_feed(&m, (const uint8_t*)body+i, 1, &used, collect, &out) >= 0); + } + CHECK(!strcmp(out.bytes+head_out, "abcde") && out.source == strlen(reply)+strlen(body)); + puts("[PASS] chunked boundaries, extensions/trailers, source byte credit and HTTP/1.0 decoding"); +} + +static void response_tests(void) { + struct http_proxy_message m; struct output out; struct http_proxy_request r = {.minor=1}; + CHECK(response("HTTP/1.1 100 Continue\r\n\r\n", &m, &r, &out) == 0 && m.body == HTTP_BODY_DONE); + CHECK(response("HTTP/1.1 204 No Content\r\n\r\n", &m, &r, &out) == 0 && m.body == HTTP_BODY_DONE); + CHECK(response("HTTP/1.1 304 Not Modified\r\nContent-Length: 999\r\n\r\n", &m, &r, &out) == 0 && m.body == HTTP_BODY_DONE); + CHECK(response("HTTP/1.1 200 OK\r\n\r\n", &m, &r, &out) == 0 && m.body == HTTP_BODY_EOF && m.close); + CHECK(strstr(out.bytes, "Connection: close\r\n")); + r.head = 1; + CHECK(response("HTTP/1.1 200 OK\r\nContent-Length: 999\r\n\r\n", &m, &r, &out) == 0 && m.body == HTTP_BODY_DONE); + r.head = 0; + CHECK(response("HTTP/1.1 204 No Content\r\nContent-Length: 0\r\n\r\n", &m, &r, &out) == -502 && out.len == 0); + CHECK(response("HTTP/1.1 200 OK\r\nContent-Length: 2\r\nTransfer-Encoding: chunked\r\n\r\n", &m, &r, &out) == -502); + CHECK(response("HTTP/1.1 101 Switching Protocols\r\nConnection: Upgrade\r\nUpgrade: websocket\r\n\r\n", &m, &r, &out) == -502); + strcpy(r.upgrade, "websocket, example/1.2"); + CHECK(response("HTTP/1.1 101 Switching Protocols\r\nConnection: Upgrade\r\nUpgrade: example/1.2\r\n\r\n", &m, &r, &out) == 0); + CHECK(strstr(out.bytes, "Connection: Upgrade\r\nUpgrade: example/1.2\r\n")); + CHECK(response("HTTP/1.1 101 Switching Protocols\r\nConnection: Upgrade\r\nUpgrade: other\r\n\r\n", &m, &r, &out) == -502); + r.minor = 0; r.close = 1; r.upgrade[0] = 0; + CHECK(response("HTTP/1.1 100 Continue\r\n\r\n", &m, &r, &out) == 0 && out.len == 0 && out.source > 0); + puts("[PASS] HEAD/1xx/204/304/EOF response framing and negotiated Upgrade"); +} + +int main(void) { + debug_config_init(); debug_set_level(DEBUG_LEVEL_WARN); + header_tests(); body_tests(); response_tests(); + CHECK(u_get_allocated_count() == 0); + return 0; +} diff --git a/tests/test_proxy_regressions.c b/tests/test_proxy_regressions.c index f1c895a3..cd04bc9b 100644 --- a/tests/test_proxy_regressions.c +++ b/tests/test_proxy_regressions.c @@ -24,6 +24,16 @@ struct message { struct UTUN_INSTANCE* inst; uint64_t peer; uint8_t svc, cmd; ui static struct message messages[512]; static unsigned message_count; static struct UASYNC* ua; +static int short_http_deadline; +static int reject_retire; +static int immediate_response; +static int abort_exit_data; + +void* __real_uasync_set_timeout(struct UASYNC* ua, int time_tb, void* arg, void (*fn)(void*), const char* name); +void* __wrap_uasync_set_timeout(struct UASYNC* async, int time_tb, void* arg, void (*fn)(void*), const char* name) { + if (short_http_deadline && !strcmp(name, "http_deadline")) time_tb = 1000; + return __real_uasync_set_timeout(async, time_tb, arg, fn, name); +} // Перехватывается только транспорт. Парсеры, tcp_io, очереди и lifecycle — настоящие. int __wrap_etcp_route_send(struct UTUN_INSTANCE* inst, uint64_t group, uint64_t peer, @@ -35,6 +45,20 @@ int __wrap_etcp_route_send(struct UTUN_INSTANCE* inst, uint64_t group, uint64_t 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); + if (abort_exit_data && m->svc == ETCP_RT_ID_TCP_PROXY_CLIENT && m->cmd == TCP_PROXY_SUBCMD_DATA) { + abort_exit_data = 0; + tcp_proxy_server_handle_error(inst, m->peer, m->sid); + } + if (reject_retire && (m->cmd == TCP_PROXY_SUBCMD_CLOSE || m->cmd == TCP_PROXY_SUBCMD_ERROR)) return -1; + if (immediate_response && m->cmd == TCP_PROXY_SUBCMD_DATA) { + int response = immediate_response; immediate_response = 0; + const char* text = response == 1 ? "HTTP/1.1 413 Too Large\r\nContent-Length: 0\r\n\r\n" : + response == 3 ? "HTTP/1.1 101 Switching Protocols\r\nConnection: Upgrade\r\nUpgrade: websocket\r\n\r\nRAW" : + response == 4 ? "HTTP/1.1 200 OK\r\nContent-Length: 0\r\n\r\n" : "BROKEN\r\n\r\n"; + struct tcp_proxy_client* p = inst->tcp_proxy_client; + CHECK(socks_proxy_handle_etcp(&p->http_conns, &p->http_conn_count, m->sid, + TCP_PROXY_SUBCMD_DATA, (const uint8_t*)text, strlen(text)) == 1); + } return 0; } @@ -174,15 +198,208 @@ static void parser_tests(struct UTUN_INSTANCE* inst) { 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); + CHECK(total < sizeof(expected)); expected[total] = 0; + CHECK(strstr(expected, "POST /upload HTTP/1.1\r\nHost: 127.0.0.1:80\r\n")); + CHECK(strstr(expected, "Content-Length: 20000\r\nConnection: close\r\n")); + char* body = strstr(expected, "\r\n\r\n"); CHECK(body && total - (body + 4 - expected) == 20000); + CHECK(memcmp(body + 4, upload + header, 20000) == 0 && count_cmd(TCP_PROXY_SUBCMD_FIN) == 1); 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 size_t data_sent(uint32_t sid, char* out, size_t cap) { + size_t len = 0; + for (unsigned i = 0; i < message_count; i++) if (messages[i].sid == sid && messages[i].cmd == TCP_PROXY_SUBCMD_DATA) { + CHECK(len + messages[i].len < cap); + memcpy(out+len, messages[i].data, messages[i].len); len += messages[i].len; + } + out[len] = 0; return len; +} + +static size_t client_reply(socket_t fd, char* out, size_t cap) { + size_t len = 0; + for (;;) { + ssize_t n = recv(fd, out+len, cap-len-1, 0); + if (n <= 0) break; + len += n; CHECK(len+1 < cap); + } + out[len] = 0; return len; +} + +static void upstream(struct UTUN_INSTANCE* inst, uint32_t sid, const char* text) { + for (size_t off = 0; off < strlen(text);) { + size_t n = strlen(text)-off; if (n > TCP_PROXY_CHUNK) n = TCP_PROXY_CHUNK; + deliver(inst, 42, ETCP_RT_ID_TCP_PROXY_CLIENT, TCP_PROXY_SUBCMD_DATA, sid, text+off, n); off += n; + } +} + +static void http_tests(struct UTUN_INSTANCE* inst) { + uint16_t port; socket_t reserved = listen_local(&port); socket_close_wrapper(reserved); + char addr[64]; snprintf(addr, sizeof(addr), "127.0.0.1:%u", port); + struct tcp_proxy_client* p = tcp_proxy_client_create(inst, ua, NULL, NULL, 1500, 1, NULL, 0, 42, 0, NULL, 1, addr); + CHECK(p); inst->tcp_proxy_client = p; + char out[64000]; socket_t fd = connect_local(port); message_count = 0; + const char pipeline[] = "GET http://127.0.0.1:8080?x=1 HTTP/1.1\r\nHost: wrong\r\n" + "Proxy-Authorization: secret\r\nConnection: X-Hop\r\nX-Hop: secret\r\n\r\n" + "GET http://127.0.0.2:8081/b HTTP/1.1\r\nHost: wrong\r\n\r\n"; + send_bytes(fd, pipeline, sizeof(pipeline)-1); + uint32_t first = p->http_conns->stream_id; + CHECK(count_cmd(TCP_PROXY_SUBCMD_CONNECT) == 1); + deliver(inst, 42, ETCP_RT_ID_TCP_PROXY_CLIENT, TCP_PROXY_SUBCMD_CONNECTED, first, NULL, 0); + CHECK(data_sent(first, out, sizeof(out)) > 0 && strstr(out, "GET /?x=1 HTTP/1.1\r\nHost: 127.0.0.1:8080\r\n")); + CHECK(!strstr(out, "secret") && !strstr(out, "/b")); + const char reply[] = "HTTP/1.1 200 OK\r\nContent-Length: 3\r\nConnection: close, X-Hop\r\nX-Hop: secret\r\n\r\none"; + // CLOSE сразу после полного response не должен оборвать ещё не записанные в OS данные. + CHECK(socks_proxy_handle_etcp(&p->http_conns, &p->http_conn_count, first, + TCP_PROXY_SUBCMD_DATA, (const uint8_t*)reply, sizeof(reply)-1) == 1); + CHECK(p->http_conns->stream_id == first && p->http_conns->tc->write_queue->head); + CHECK(socks_proxy_handle_etcp(&p->http_conns, &p->http_conn_count, first, TCP_PROXY_SUBCMD_CLOSE, NULL, 0) == 1); + reject_retire = 1; + pump(); + CHECK(p->http_conns->stream_id == first && count_cmd(TCP_PROXY_SUBCMD_CONNECT) == 1); + reject_retire = 0; + uint32_t credit = 1; + CHECK(socks_proxy_handle_etcp(&p->http_conns, &p->http_conn_count, first, + TCP_PROXY_SUBCMD_WINDOW, (uint8_t*)&credit, 4) == 1); + pump(); + CHECK(p->http_conn_count == 1 && count_cmd(TCP_PROXY_SUBCMD_CONNECT) == 2); + uint32_t second = p->http_conns->stream_id; CHECK(second && first != second); + CHECK(messages[message_count-1].cmd == TCP_PROXY_SUBCMD_CONNECT); + CHECK(messages[message_count-1].data[3] == 2); + CHECK(client_reply(fd, out, sizeof(out)) > 0 && strstr(out, "\r\n\r\none") && !strstr(out, "secret")); + CHECK(!strstr(out, "Connection: close")); + CHECK(socks_proxy_handle_etcp(&p->http_conns, &p->http_conn_count, first, TCP_PROXY_SUBCMD_FIN, NULL, 0) == 0); + CHECK(socks_proxy_handle_etcp(&p->http_conns, &p->http_conn_count, first, TCP_PROXY_SUBCMD_WINDOW, (uint8_t*)&credit, 4) == 0); + deliver(inst, 42, ETCP_RT_ID_TCP_PROXY_CLIENT, TCP_PROXY_SUBCMD_CONNECTED, second, NULL, 0); + CHECK(data_sent(second, out, sizeof(out)) > 0 && strstr(out, "GET /b HTTP/1.1\r\nHost: 127.0.0.2:8081\r\n")); + upstream(inst, second, "HTTP/1.1 200 OK\r\nContent-Length: 3\r\n\r\ntwo"); + CHECK(client_reply(fd, out, sizeof(out)) > 0 && strstr(out, "\r\n\r\ntwo")); + socket_close_wrapper(fd); pump(); CHECK(p->http_conn_count == 0); + + // EOF в незавершённых заголовках и теле освобождает frontend, вместо вечного ожидания. + fd = connect_local(port); message_count = 0; + const char truncated[] = "GET http://127.0.0.1/ HTTP/1.1\r\nHost:"; + send_bytes(fd, truncated, sizeof(truncated)-1); + CHECK(shutdown(fd, SHUT_WR) == 0); pump(); + CHECK(client_reply(fd, out, sizeof(out)) > 0 && strstr(out, "HTTP/1.1 400")); + CHECK(p->http_conn_count == 0 && count_cmd(TCP_PROXY_SUBCMD_CONNECT) == 0); socket_close_wrapper(fd); + + const char post[] = "POST http://127.0.0.1/ HTTP/1.1\r\nHost: x\r\nContent-Length: 9\r\nExpect: 100-continue\r\n\r\n"; + fd = connect_local(port); message_count = 0; send_bytes(fd, post, sizeof(post)-1); + first = p->http_conns->stream_id; + deliver(inst, 42, ETCP_RT_ID_TCP_PROXY_CLIENT, TCP_PROXY_SUBCMD_CONNECTED, first, NULL, 0); + upstream(inst, first, "HTTP/1.1 100 Continue\r\n\r\n"); + CHECK(client_reply(fd, out, sizeof(out)) > 0 && strstr(out, "100 Continue")); + send_bytes(fd, "123456789", 9); + CHECK(data_sent(first, out, sizeof(out)) > 9 && !strcmp(out+strlen(out)-9, "123456789")); + upstream(inst, first, "HTTP/1.1 200 OK\r\nTransfer-Encoding: chunked\r\nConnection: X-Hop\r\n\r\n"); + const char chunks[] = "3\r\nabc\r\n0\r\nX-Hop: secret\r\nX-End: yes\r\n\r\n"; + for (size_t i = 0; i < sizeof(chunks)-1; i++) + deliver(inst, 42, ETCP_RT_ID_TCP_PROXY_CLIENT, TCP_PROXY_SUBCMD_DATA, first, chunks+i, 1); + CHECK(client_reply(fd, out, sizeof(out)) > 0 && strstr(out, "3\r\nabc\r\n0\r\nX-End: yes\r\n\r\n")); + CHECK(!strstr(out, "secret")); + uint32_t returned = 0; + for (unsigned i = 0; i < message_count; i++) if (messages[i].sid == first && messages[i].cmd == TCP_PROXY_SUBCMD_WINDOW) { + uint32_t n; memcpy(&n, messages[i].data, 4); returned += n; + } + CHECK(returned == strlen("HTTP/1.1 100 Continue\r\n\r\n") + + strlen("HTTP/1.1 200 OK\r\nTransfer-Encoding: chunked\r\nConnection: X-Hop\r\n\r\n") + strlen(chunks)); + socket_close_wrapper(fd); pump(); CHECK(p->http_conn_count == 0); + + // Ранний final и malformed response синхронно внутри отправки request headers. + for (int which = 1; which <= 2; which++) { + fd = connect_local(port); message_count = 0; send_bytes(fd, post, sizeof(post)-1); + first = p->http_conns->stream_id; immediate_response = which; + deliver(inst, 42, ETCP_RT_ID_TCP_PROXY_CLIENT, TCP_PROXY_SUBCMD_CONNECTED, first, NULL, 0); + CHECK(client_reply(fd, out, sizeof(out)) > 0 && strstr(out, which == 1 ? "HTTP/1.1 413" : "HTTP/1.1 502")); + CHECK(p->http_conn_count == 0 && count_cmd(TCP_PROXY_SUBCMD_ERROR) >= 1); + socket_close_wrapper(fd); + } + fd = connect_local(port); message_count = 0; send_bytes(fd, post, sizeof(post)-1); + first = p->http_conns->stream_id; + deliver(inst, 42, ETCP_RT_ID_TCP_PROXY_CLIENT, TCP_PROXY_SUBCMD_CONNECTED, first, NULL, 0); + send_bytes(fd, "abc", 3); CHECK(shutdown(fd, SHUT_WR) == 0); pump(); + CHECK(client_reply(fd, out, sizeof(out)) > 0 && strstr(out, "HTTP/1.1 400")); + CHECK(p->http_conn_count == 0); socket_close_wrapper(fd); + + // Close-delimited response завершается только FIN и закрывает frontend. + const char get[] = "GET http://127.0.0.1/ HTTP/1.1\r\nHost: x\r\n\r\n"; + fd = connect_local(port); message_count = 0; send_bytes(fd, get, sizeof(get)-1); + first = p->http_conns->stream_id; + deliver(inst, 42, ETCP_RT_ID_TCP_PROXY_CLIENT, TCP_PROXY_SUBCMD_CONNECTED, first, NULL, 0); + upstream(inst, first, "HTTP/1.1 200 OK\r\n\r\nbody"); CHECK(p->http_conn_count == 1); + deliver(inst, 42, ETCP_RT_ID_TCP_PROXY_CLIENT, TCP_PROXY_SUBCMD_FIN, first, NULL, 0); + CHECK(client_reply(fd, out, sizeof(out)) > 0 && strstr(out, "Connection: close") && strstr(out, "\r\n\r\nbody")); + CHECK(p->http_conn_count == 0); socket_close_wrapper(fd); + + // HEAD не ждёт объявленное Content-Length тело; pipeline сохраняется и после клиентского EOF. + fd = connect_local(port); message_count = 0; + const char heads[] = "HEAD http://127.0.0.1/ HTTP/1.1\r\nHost: x\r\n\r\n" + "GET http://127.0.0.2/ HTTP/1.1\r\nHost: x\r\n\r\n"; + send_bytes(fd, heads, sizeof(heads)-1); CHECK(shutdown(fd, SHUT_WR) == 0); pump(); + first = p->http_conns->stream_id; + deliver(inst, 42, ETCP_RT_ID_TCP_PROXY_CLIENT, TCP_PROXY_SUBCMD_CONNECTED, first, NULL, 0); + upstream(inst, first, "HTTP/1.1 200 OK\r\nContent-Length: 999\r\n\r\n"); + CHECK(p->http_conn_count == 1 && count_cmd(TCP_PROXY_SUBCMD_CONNECT) == 2); + second = p->http_conns->stream_id; + deliver(inst, 42, ETCP_RT_ID_TCP_PROXY_CLIENT, TCP_PROXY_SUBCMD_CONNECTED, second, NULL, 0); + upstream(inst, second, "HTTP/1.1 204 No Content\r\n\r\n"); + CHECK(p->http_conn_count == 0 && client_reply(fd, out, sizeof(out)) > 0 && strstr(out, "204 No Content")); + socket_close_wrapper(fd); + + // HTTP/1.0 получает декодированное chunked тело и close; credits не зависят от размера результата. + fd = connect_local(port); message_count = 0; + const char old_get[] = "GET http://127.0.0.1/ HTTP/1.0\r\n\r\n"; + send_bytes(fd, old_get, sizeof(old_get)-1); first = p->http_conns->stream_id; + deliver(inst, 42, ETCP_RT_ID_TCP_PROXY_CLIENT, TCP_PROXY_SUBCMD_CONNECTED, first, NULL, 0); + upstream(inst, first, "HTTP/1.1 100 Continue\r\n\r\nHTTP/1.1 200 OK\r\nTransfer-Encoding: chunked\r\n\r\n"); + upstream(inst, first, "3\r\nabc\r\n0\r\n\r\n"); + CHECK(p->http_conn_count == 0 && client_reply(fd, out, sizeof(out)) > 0); + CHECK(strstr(out, "HTTP/1.0 200 OK") && strstr(out, "\r\n\r\nabc") && !strstr(out, "100 Continue")); + CHECK(!strstr(out, "Transfer-Encoding")); socket_close_wrapper(fd); + + // Синхронный 101 видит отправленный последний header chunk; ранние raw bytes сохраняются. + fd = connect_local(port); message_count = 0; + const char upgrade[] = "GET http://127.0.0.1/ HTTP/1.1\r\nHost: x\r\nConnection: Upgrade\r\nUpgrade: websocket\r\n\r\nCLIENT"; + send_bytes(fd, upgrade, sizeof(upgrade)-1); first = p->http_conns->stream_id; immediate_response = 3; + deliver(inst, 42, ETCP_RT_ID_TCP_PROXY_CLIENT, TCP_PROXY_SUBCMD_CONNECTED, first, NULL, 0); + CHECK(client_reply(fd, out, sizeof(out)) > 0 && strstr(out, "101 Switching Protocols") && strstr(out, "\r\n\r\nRAW")); + CHECK(data_sent(first, out, sizeof(out)) > 6 && !strcmp(out+strlen(out)-6, "CLIENT")); + send_bytes(fd, "NEXT", 4); upstream(inst, first, "SERVER"); + CHECK(client_reply(fd, out, sizeof(out)) == 6 && !strcmp(out, "SERVER")); + deliver(inst, 42, ETCP_RT_ID_TCP_PROXY_CLIENT, TCP_PROXY_SUBCMD_ERROR, first, NULL, 0); + socket_close_wrapper(fd); pump(); CHECK(p->http_conn_count == 0); + + // Закрытое окно возвращает снятый перед send chunk обратно в tx_buf без потери/дублирования. + fd = connect_local(port); message_count = 0; send_bytes(fd, get, sizeof(get)-1); first = p->http_conns->stream_id; + p->http_conns->flow.tx_credit = 0; + deliver(inst, 42, ETCP_RT_ID_TCP_PROXY_CLIENT, TCP_PROXY_SUBCMD_CONNECTED, first, NULL, 0); + CHECK(count_cmd(TCP_PROXY_SUBCMD_DATA) == 0 && p->http_conns->tx_buf); + credit = TCP_PROXY_WINDOW; + deliver(inst, 42, ETCP_RT_ID_TCP_PROXY_CLIENT, TCP_PROXY_SUBCMD_WINDOW, first, &credit, 4); + CHECK(count_cmd(TCP_PROXY_SUBCMD_DATA) == 1 && data_sent(first, out, sizeof(out)) > 0); + CHECK(strstr(out, "GET / HTTP/1.1\r\nHost: 127.0.0.1:80\r\n") && !p->http_conns->tx_buf); + upstream(inst, first, "HTTP/1.1 200 OK\r\nContent-Length: 0\r\n\r\n"); + socket_close_wrapper(fd); pump(); CHECK(p->http_conn_count == 0); + + // Ограниченные по времени headers/connect и принудительное завершение зависшего drain. + short_http_deadline = 1; + for (int phase = 0; phase < 3; phase++) { + fd = connect_local(port); message_count = 0; + if (phase == 2) queue_set_callback(p->http_conns->tc->write_queue, NULL, NULL); + send_bytes(fd, phase == 1 ? get : "GET", phase == 1 ? sizeof(get)-1 : 3); + for (int i = 0; i < 22; i++) pump(); + CHECK(p->http_conn_count == 0); + if (phase != 2) CHECK(client_reply(fd, out, sizeof(out)) > 0 && strstr(out, phase == 1 ? "504" : "408")); + socket_close_wrapper(fd); + } + short_http_deadline = 0; + tcp_proxy_client_destroy(p); inst->tcp_proxy_client = NULL; pump(); + puts("[PASS] HTTP keep-alive origins, pipeline, EOF, Continue/chunks/credits, reentrancy and deadlines"); +} + 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; @@ -213,6 +430,20 @@ static void server_tests(struct UTUN_INSTANCE* inst) { CHECK(shutdown(c, SHUT_WR) == 0); pump(); CHECK(inst->tcp_proxy_server.conn_count == 0 && count_cmd(TCP_PROXY_SUBCMD_FIN) == 1); socket_close_wrapper(c); + // ERROR уничтожает exit сразу, не доставляя отложенный upload после раннего отказа origin. + deliver(inst, 11, ETCP_RT_ID_TCP_PROXY_SERVER, TCP_PROXY_SUBCMD_CONNECT, 4, connect_data, 6); + c = accept(listener, NULL, NULL); CHECK(c != SOCKET_INVALID); socket_set_nonblocking(c); + struct tcp_proxy_server_conn* rc = tcp_proxy_server_find_conn(&inst->tcp_proxy_server, 11, 4); CHECK(rc); + queue_set_callback(rc->tc->write_queue, NULL, NULL); + deliver(inst, 11, ETCP_RT_ID_TCP_PROXY_SERVER, TCP_PROXY_SUBCMD_DATA, 4, "discard", 7); + CHECK(rc->tc->write_queue->count == 1); + deliver(inst, 11, ETCP_RT_ID_TCP_PROXY_SERVER, TCP_PROXY_SUBCMD_ERROR, 4, NULL, 0); + CHECK(inst->tcp_proxy_server.conn_count == 0 && recv(c, buf, sizeof(buf), 0) == 0); socket_close_wrapper(c); + // Синхронная отмена внутри exit DATA send не освобождает текущий callback/пулы преждевременно. + deliver(inst, 11, ETCP_RT_ID_TCP_PROXY_SERVER, TCP_PROXY_SUBCMD_CONNECT, 5, connect_data, 6); + c = accept(listener, NULL, NULL); CHECK(c != SOCKET_INVALID); socket_set_nonblocking(c); + abort_exit_data = 1; send_bytes(c, "response", 8); + CHECK(!abort_exit_data && inst->tcp_proxy_server.conn_count == 0); socket_close_wrapper(c); 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"); @@ -293,7 +524,7 @@ int main(void) { 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); + instance_tests(inst); parser_tests(inst); http_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();