Browse Source

Fix HTTP proxy framing, keep-alive and stream lifecycle

master
evgeny 6 days ago
parent
commit
ba7146f90a
  1. 2
      cross-build-win.sh
  2. 4
      src/Makefile.am
  3. 380
      src/proxy/http_proxy.c
  4. 66
      src/proxy/http_proxy.h
  5. 549
      src/proxy/socks_proxy.c
  6. 12
      src/proxy/socks_proxy.h
  7. 38
      src/proxy/socks_proxy_doc.md
  8. 18
      src/proxy/tcp_proxy_server.c
  9. 6
      tests/Makefile.am
  10. 160
      tests/test_http_proxy.c
  11. 239
      tests/test_proxy_regressions.c

2
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
)

4
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 \

380
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 <string.h>
#include <stdio.h>
#include <limits.h>
// Все протокольные ошибки имеют причину; значения 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;
}

66
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 <stddef.h>
#include <stdint.h>
#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

549
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;

12
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;

38
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`)

18
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) {

6
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

160
tests/test_http_proxy.c

@ -0,0 +1,160 @@
// Проверки HTTP codec на каждом возможном разбиении headers и побайтовом chunked body.
#include <stdio.h>
#include <stdlib.h>
#include <string.h>
#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;
}

239
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();

Loading…
Cancel
Save