Browse Source

fix: backpressure data loss in socks/http proxy + comprehensive test

socks_proxy (client side):
- tx_buf buffering on send_data failure in RELAY path, retry in tx_waiter_cb
- process_http_request: version_str null-term fix (memcpy(vlen)+explicit \0)
- process_http_request: assemble hdr+body+extra into heap pkt, atomic send
- socks_proxy_handle_etcp: alloc failure logs write_queue state + byte counters

tcp_proxy_server (exit node):
- read_queue_drain_cb: tx_buf buffering on backpressure, retry_timer (500ms, force=1)
- on_fin_cb: defer FIN when read_queue+tx_buf has pending data (dst_fin_deferred)
- handle_data + drain_cb: per-stream byte counters (sent/recv/relayed/drain_count)
- on_fin_cb/on_error_cb: detailed state logging (counters, queues, buffers)

test_socks_http_proxy:
- 4 workers: socks_ipv4, socks_domain (DNS resolution via SOCKS), http_proxy, https_connect
- 4 file sizes: 1KB, 100KB, 1MB, 10MB — deterministic content, byte-by-byte cmp -s
- HTTPS server with self-signed cert, POST/echo, HEAD, OPTIONS
- 51/51 tests PASS
chatgui
Evgeny 3 months ago
parent
commit
2fba5c31c4
  1. 62
      src/proxy/socks_proxy.c
  2. 4
      src/proxy/socks_proxy.h
  3. 2
      src/proxy/tcp_proxy_client.c
  4. 115
      src/proxy/tcp_proxy_server.c
  5. 9
      src/proxy/tcp_proxy_server.h
  6. 263
      tests/test_socks_http_proxy.c

62
src/proxy/socks_proxy.c

@ -60,7 +60,10 @@ static int send_msg(struct UTUN_INSTANCE* inst, uint64_t dst, uint8_t subcmd, ui
queue_dgram_free(e); queue_entry_free(e); return -1; queue_dgram_free(e); queue_entry_free(e); return -1;
} }
e->len = (uint16_t)(TCP_PROXY_HDR_SIZE + len); e->len = (uint16_t)(TCP_PROXY_HDR_SIZE + len);
return etcp_route_send(inst, dst, e, force); int ret = etcp_route_send(inst, dst, e, force);
DEBUG_TRACE(DEBUG_CATEGORY_TRAFFIC, "SOCKS send_msg subcmd=%02x sid=%08x len=%zu force=%d → ret=%d",
subcmd, sid, len, force, ret);
return ret;
} }
static int send_connect(struct socks_proxy_conn* c) { static int send_connect(struct socks_proxy_conn* c) {
@ -77,16 +80,19 @@ static int send_data(struct socks_proxy_conn* c, const uint8_t* data, uint16_t l
} }
static void send_close(struct socks_proxy_conn* c) { static void send_close(struct socks_proxy_conn* c) {
DEBUG_TRACE(DEBUG_CATEGORY_TRAFFIC, "SOCKS send CLOSE sid=%08x", c->stream_id);
if (send_msg(c->inst, c->via_node_id, TCP_PROXY_SUBCMD_CLOSE, c->stream_id, NULL, 0, 1) < 0) c->close_pending = 1; if (send_msg(c->inst, c->via_node_id, TCP_PROXY_SUBCMD_CLOSE, c->stream_id, NULL, 0, 1) < 0) c->close_pending = 1;
else { c->close_pending = 0; c->close_sent = 1; } else { c->close_pending = 0; c->close_sent = 1; }
} }
static void send_error(struct socks_proxy_conn* c) { static void send_error(struct socks_proxy_conn* c) {
DEBUG_TRACE(DEBUG_CATEGORY_TRAFFIC, "SOCKS send ERROR sid=%08x", c->stream_id);
if (send_msg(c->inst, c->via_node_id, TCP_PROXY_SUBCMD_ERROR, c->stream_id, NULL, 0, 1) < 0) c->close_pending = 1; if (send_msg(c->inst, c->via_node_id, TCP_PROXY_SUBCMD_ERROR, c->stream_id, NULL, 0, 1) < 0) c->close_pending = 1;
else { c->close_pending = 0; c->close_sent = 1; } else { c->close_pending = 0; c->close_sent = 1; }
} }
static void send_fin(struct socks_proxy_conn* c) { static void send_fin(struct socks_proxy_conn* c) {
DEBUG_TRACE(DEBUG_CATEGORY_TRAFFIC, "SOCKS send FIN sid=%08x", c->stream_id);
send_msg(c->inst, c->via_node_id, TCP_PROXY_SUBCMD_FIN, c->stream_id, NULL, 0, 1); send_msg(c->inst, c->via_node_id, TCP_PROXY_SUBCMD_FIN, c->stream_id, NULL, 0, 1);
} }
@ -110,6 +116,7 @@ static int write_to_client(struct socks_proxy_conn* c, const uint8_t* data, uint
memcpy(buf, data, len); memcpy(buf, data, len);
e->dgram = buf; e->len = len; e->dgram = buf; e->len = len;
queue_data_put(c->tc->write_queue, e); queue_data_put(c->tc->write_queue, e);
c->bytes_to_client += len;
return 0; return 0;
} }
@ -297,7 +304,8 @@ static void process_http_request(struct socks_proxy_conn* c) {
if (version) { if (version) {
size_t vlen = strlen(version); size_t vlen = strlen(version);
if (vlen >= sizeof(version_str)) vlen = sizeof(version_str) - 1; if (vlen >= sizeof(version_str)) vlen = sizeof(version_str) - 1;
memcpy(version_str, version, vlen + 1); memcpy(version_str, version, vlen);
version_str[vlen] = '\0';
} }
uint16_t rest_off = line_len + 2; uint16_t rest_off = line_len + 2;
@ -309,10 +317,21 @@ static void process_http_request(struct socks_proxy_conn* c) {
goto error; goto error;
} }
memcpy(hdr_buf + hdr_n, c->buf + rest_off, rest_len); memcpy(hdr_buf + hdr_n, c->buf + rest_off, rest_len);
send_data(c, hdr_buf, (uint16_t)(hdr_n + 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;
if (headers_len < c->buf_len) { uint8_t* pkt = u_malloc(total);
send_data(c, c->buf + headers_len, c->buf_len - headers_len); if (!pkt) { DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "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);
int ret = send_data(c, pkt, total);
if (ret == 0) {
u_free(pkt);
} else {
c->tx_buf = pkt; c->tx_len = total;
etcp_router_waiter_register(c->inst, c->via_node_id, &c->tx_waiter, tx_waiter_cb, c);
} }
c->buf_len = 0; c->buf_len = 0;
@ -361,12 +380,20 @@ static void on_read_cb(struct ll_queue* q, void* arg) {
} }
// RELAY = релей данных в ETCP // RELAY = релей данных в ETCP
if (c->tx_buf) {
DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "socks_proxy: RELAY with pending tx_buf sid=%08x — drop new len=%u", c->stream_id, e->len);
memory_pool_free(c->tc->data_pool, e->dgram); queue_entry_free(e);
return;
}
int ret = send_data(c, e->dgram, e->len); int ret = send_data(c, e->dgram, e->len);
DEBUG_TRACE(DEBUG_CATEGORY_TRAFFIC, "SOCKS PROXY SEND sid=%08x len=%u", c->stream_id, e->len); DEBUG_TRACE(DEBUG_CATEGORY_TRAFFIC, "SOCKS PROXY SEND sid=%08x len=%u", c->stream_id, e->len);
if (ret == 0) { if (ret == 0) {
memory_pool_free(c->tc->data_pool, e->dgram); queue_entry_free(e); memory_pool_free(c->tc->data_pool, e->dgram); queue_entry_free(e);
queue_resume_callback(q); queue_resume_callback(q);
} else { } else {
c->tx_buf = u_malloc(e->len);
if (c->tx_buf) { memcpy(c->tx_buf, e->dgram, e->len); c->tx_len = e->len; }
else { DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "socks_proxy: tx_buf malloc=%u failed sid=%08x — drop", e->len, c->stream_id); }
memory_pool_free(c->tc->data_pool, e->dgram); queue_entry_free(e); memory_pool_free(c->tc->data_pool, e->dgram); queue_entry_free(e);
etcp_router_waiter_register(c->inst, c->via_node_id, &c->tx_waiter, tx_waiter_cb, c); etcp_router_waiter_register(c->inst, c->via_node_id, &c->tx_waiter, tx_waiter_cb, c);
} }
@ -375,8 +402,14 @@ static void on_read_cb(struct ll_queue* q, void* arg) {
static void tx_waiter_cb(struct ll_queue* q, void* arg) { static void tx_waiter_cb(struct ll_queue* q, void* arg) {
(void)q; (void)q;
struct socks_proxy_conn* c = (struct socks_proxy_conn*)arg; struct socks_proxy_conn* c = (struct socks_proxy_conn*)arg;
if (c->rem_closed || c->close_sent) return; if (!c->tx_buf || c->rem_closed || c->close_sent) return;
queue_resume_callback(c->tc->read_queue); int ret = send_data(c, c->tx_buf, c->tx_len);
if (ret == 0) {
u_free(c->tx_buf); c->tx_buf = NULL; c->tx_len = 0;
queue_resume_callback(c->tc->read_queue);
} else {
etcp_router_waiter_register(c->inst, c->via_node_id, &c->tx_waiter, tx_waiter_cb, c);
}
} }
// ==================================================================== // ====================================================================
@ -400,14 +433,18 @@ static void on_fin_cb(struct tcp_conn* tc, void* arg) {
static void on_error_cb(struct tcp_conn* tc, int err, void* arg) { static void on_error_cb(struct tcp_conn* tc, int err, void* arg) {
(void)err; (void)err;
struct socks_proxy_conn* c = (struct socks_proxy_conn*)arg; struct socks_proxy_conn* c = (struct socks_proxy_conn*)arg;
DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "socks_proxy: tcp error sid=%08x err=%d", c->stream_id, err); DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "socks_proxy: tcp error sid=%08x err=%d state=%d fm_r=%d fm_l=%d rem_cl=%d cs=%d cp=%d "
"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);
if (!c->close_sent && !c->close_pending) send_error(c); if (!c->close_sent && !c->close_pending) send_error(c);
socks_proxy_conn_free(c); socks_proxy_conn_free(c);
} }
static void on_closed_cb(struct tcp_conn* tc, void* arg) { static void on_closed_cb(struct tcp_conn* tc, void* arg) {
struct socks_proxy_conn* c = (struct socks_proxy_conn*)arg; struct socks_proxy_conn* c = (struct socks_proxy_conn*)arg;
DEBUG_DEBUG(DEBUG_CATEGORY_SOCKET, "socks_proxy: tcp closed sid=%08x", c->stream_id); DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "socks_proxy: tcp closed sid=%08x state=%d fm_r=%d fm_l=%d rem_cl=%d bytes_client=%u bytes_exit=%u",
c->stream_id, c->state, tc->fin_remote, tc->fin_local, c->rem_closed, c->bytes_to_client, c->bytes_from_exit);
if (!c->close_sent && !c->close_pending) send_close(c); if (!c->close_sent && !c->close_pending) send_close(c);
socks_proxy_conn_free(c); socks_proxy_conn_free(c);
} }
@ -459,6 +496,8 @@ int socks_proxy_handle_etcp(struct socks_proxy_conn** head, int* count,
if (!c) return 0; if (!c) return 0;
if (subcmd == TCP_PROXY_SUBCMD_DATA) { if (subcmd == TCP_PROXY_SUBCMD_DATA) {
DEBUG_TRACE(DEBUG_CATEGORY_TRAFFIC, "SOCKS DATA <- sid=%08x len=%zu", stream_id, data_len);
c->bytes_from_exit += (uint32_t)data_len;
if (data_len > 0) { if (data_len > 0) {
if (data_len > c->tc->data_pool->object_size) { if (data_len > c->tc->data_pool->object_size) {
DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "socks_proxy: data_len=%zu > pool_sz=%zu, dropping sid=%08x", DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "socks_proxy: data_len=%zu > pool_sz=%zu, dropping sid=%08x",
@ -470,8 +509,10 @@ int socks_proxy_handle_etcp(struct socks_proxy_conn** head, int* count,
if (e && buf) { if (e && buf) {
memcpy(buf, data, data_len); e->dgram = buf; e->len = (uint16_t)data_len; memcpy(buf, data, data_len); e->dgram = buf; e->len = (uint16_t)data_len;
queue_data_put(c->tc->write_queue, e); queue_data_put(c->tc->write_queue, e);
c->bytes_to_client += (uint32_t)data_len;
} else { } else {
DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "socks_proxy: handle_data alloc failed sid=%08x", stream_id); DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "socks_proxy: handle_data alloc failed sid=%08x wq=%d bytes_in=%u bytes_out=%u",
stream_id, c->tc->write_queue->count, c->bytes_from_exit, c->bytes_to_client);
if (e) queue_entry_free(e); if (buf) memory_pool_free(c->tc->data_pool, buf); if (e) queue_entry_free(e); if (buf) memory_pool_free(c->tc->data_pool, buf);
} }
} }
@ -517,6 +558,7 @@ void socks_proxy_conn_free(struct socks_proxy_conn* c) {
send_msg(c->inst, c->via_node_id, subcmd, c->stream_id, NULL, 0, 1); send_msg(c->inst, c->via_node_id, subcmd, c->stream_id, NULL, 0, 1);
} }
if (head) { struct socks_proxy_conn** prev = head; while (*prev) { if (*prev == c) { *prev = c->next; if (count) (*count)--; break; } prev = &(*prev)->next; } } 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->tc) { tcp_conn_destroy(c->tc); c->tc = NULL; } if (c->tc) { tcp_conn_destroy(c->tc); c->tc = NULL; }
if (c->inst) etcp_router_waiter_cancel(c->inst, c->via_node_id, &c->tx_waiter); if (c->inst) etcp_router_waiter_cancel(c->inst, c->via_node_id, &c->tx_waiter);
u_free(c); u_free(c);

4
src/proxy/socks_proxy.h

@ -45,6 +45,10 @@ struct socks_proxy_conn {
uint8_t freed; uint8_t freed;
uint8_t buf[1024]; uint8_t buf[1024];
uint16_t buf_len; uint16_t buf_len;
uint8_t* tx_buf; // буфер при backpressure (retry в tx_waiter_cb)
uint16_t tx_len;
uint32_t bytes_to_client; // байт записано клиенту (curl) через write_queue
uint32_t bytes_from_exit; // байт получено от exit node через ETCP DATA
struct queue_waiter_handle tx_waiter; struct queue_waiter_handle tx_waiter;
}; };

2
src/proxy/tcp_proxy_client.c

@ -536,6 +536,8 @@ void tcp_proxy_client_router_recv_cb(struct ETCP_CONN* conn, struct ll_entry* en
uint8_t subcmd = entry->dgram[1]; uint8_t subcmd = entry->dgram[1];
uint32_t stream_id; memcpy(&stream_id, entry->dgram + 2, 4); uint32_t stream_id; memcpy(&stream_id, entry->dgram + 2, 4);
DEBUG_TRACE(DEBUG_CATEGORY_TRAFFIC, "PROXY CLIENT RECV subcmd=%02x sid=%08x pkt_len=%u", subcmd, stream_id, entry->len);
if (proxy) { if (proxy) {
size_t data_len = entry->len > TCP_PROXY_HDR_SIZE ? entry->len - TCP_PROXY_HDR_SIZE : 0; size_t data_len = entry->len > TCP_PROXY_HDR_SIZE ? entry->len - TCP_PROXY_HDR_SIZE : 0;
// Пробуем SOCKS/HTTP first (у них приоритет — могут быть без TUN) // Пробуем SOCKS/HTTP first (у них приоритет — могут быть без TUN)

115
src/proxy/tcp_proxy_server.c

@ -33,6 +33,7 @@ static void read_queue_drain_cb(struct ll_queue* q, void* arg);
static void pause_resume_cb(struct ll_queue* q, void* arg); static void pause_resume_cb(struct ll_queue* q, void* arg);
static void close_retry_cb(void* arg); static void close_retry_cb(void* arg);
static void diag_timer_cb(void* arg); static void diag_timer_cb(void* arg);
static void retry_timer_cb(void* arg);
static int send_msg(struct UTUN_INSTANCE* inst, uint64_t dst, uint8_t subcmd, static int send_msg(struct UTUN_INSTANCE* inst, uint64_t dst, uint8_t subcmd,
uint32_t sid, const uint8_t* data, size_t len, int force); uint32_t sid, const uint8_t* data, size_t len, int force);
static void send_close(struct tcp_proxy_server_conn* rc); static void send_close(struct tcp_proxy_server_conn* rc);
@ -116,19 +117,26 @@ static void send_fin(struct tcp_proxy_server_conn* rc) {
static void on_fin_cb(struct tcp_conn* tc, void* arg) { static void on_fin_cb(struct tcp_conn* tc, void* arg) {
struct tcp_proxy_server_conn* rc = (struct tcp_proxy_server_conn*)arg; struct tcp_proxy_server_conn* rc = (struct tcp_proxy_server_conn*)arg;
if (rc->cli_closed) return; if (rc->cli_closed) return;
int pend = write_pending(tc); int pend_w = write_pending(tc);
int pend_r = (rc->tx_buf != NULL) || (tc->read_queue && tc->read_queue->head);
DEBUG_INFO(DEBUG_CATEGORY_SOCKET,
"SOCK:DST_FIN fd=%d sid=%08x fin_l=%d pend_w=%d pend_r=%d "
"sent=%u recv=%u relayed=%u drain#=%u data#=%u "
"rq=%d(%zub) wq=%d(%zub) wbuf=%s err=%d tx_buf=%s",
(int)tc->sock, rc->stream_id, tc->fin_local, pend_w, pend_r,
rc->bytes_sent, rc->bytes_recv, rc->bytes_relayed,
rc->drain_count, rc->data_count,
tc->read_queue->count, queue_total_bytes(tc->read_queue),
tc->write_queue->count, queue_total_bytes(tc->write_queue),
tc->write_buf ? "y" : "n", tc->error, rc->tx_buf ? "y" : "n");
if (tc->fin_local) { if (tc->fin_local) {
DEBUG_DEBUG(DEBUG_CATEGORY_SOCKET, "SOCK:DST_FIN fd=%d sid=%08x fin_local=%d — both FINs, closing",
(int)tc->sock, rc->stream_id, tc->fin_local);
send_close(rc); tcp_conn_push_close(tc); send_close(rc); tcp_conn_push_close(tc);
} else if (pend) { } else if (pend_w) {
tcp_conn_set_flushed(tc, on_flushed_cb); tcp_conn_set_flushed(tc, on_flushed_cb);
DEBUG_DEBUG(DEBUG_CATEGORY_SOCKET, "SOCK:DST_FIN fd=%d sid=%08x write_pend=%d → deferred", } else if (pend_r) {
(int)tc->sock, rc->stream_id, pend); rc->dst_fin_deferred = 1;
} else { } else {
send_fin(rc); send_fin(rc);
DEBUG_DEBUG(DEBUG_CATEGORY_SOCKET, "SOCK:DST_FIN fd=%d sid=%08x → relay now",
(int)tc->sock, rc->stream_id);
} }
} }
@ -154,6 +162,16 @@ static void on_error_cb(struct tcp_conn* tc, int err, void* arg) {
(void)err; (void)err;
if (tc->closed) return; if (tc->closed) return;
struct tcp_proxy_server_conn* rc = (struct tcp_proxy_server_conn*)arg; struct tcp_proxy_server_conn* rc = (struct tcp_proxy_server_conn*)arg;
DEBUG_ERROR(DEBUG_CATEGORY_SOCKET,
"SOCK:ERR fd=%d sid=%08x err=%d fin_r=%d fin_l=%d "
"sent=%u recv=%u relayed=%u drain#=%u data#=%u "
"rq=%d wq=%d wbuf=%s conn=%d",
tc ? (int)tc->sock : -1, rc->stream_id, err,
tc ? tc->fin_remote : 0, tc ? tc->fin_local : 0,
rc->bytes_sent, rc->bytes_recv, rc->bytes_relayed,
rc->drain_count, rc->data_count,
tc ? tc->read_queue->count : 0, tc ? tc->write_queue->count : 0,
tc && tc->write_buf ? "y" : "n", tc ? tc->connected : 0);
send_error(rc); send_error(rc);
tcp_proxy_server_conn_free(rc); tcp_proxy_server_conn_free(rc);
} }
@ -193,8 +211,33 @@ static void diag_timer_cb(void* arg) {
static void read_queue_drain_cb(struct ll_queue* q, void* arg) { static void read_queue_drain_cb(struct ll_queue* q, void* arg) {
struct tcp_proxy_server_conn* rc = (struct tcp_proxy_server_conn*)arg; struct tcp_proxy_server_conn* rc = (struct tcp_proxy_server_conn*)arg;
struct UTUN_INSTANCE* inst = rc->ctx ? rc->ctx->inst : NULL; struct UTUN_INSTANCE* inst = rc->ctx ? rc->ctx->inst : NULL;
struct ll_entry* e = queue_data_get(q); struct ll_entry* e;
if (!e) { queue_resume_callback(q); return; }
// Сначала добиваем отложенный tx_buf — force=1 чтобы обойти backpressure
if (rc->tx_buf) {
int ret = send_msg(inst, rc->peer_node_id, TCP_PROXY_SUBCMD_DATA, rc->stream_id, rc->tx_buf, rc->tx_len, 0);
if (ret == 0) {
u_free(rc->tx_buf); rc->tx_buf = NULL; rc->tx_len = 0;
if (rc->retry_timer) { uasync_cancel_timeout(rc->ua, rc->retry_timer); rc->retry_timer = NULL; }
if (rc->dst_fin_deferred && !q->head) {
rc->dst_fin_deferred = 0;
send_fin(rc);
return;
}
} else {
if (!rc->retry_timer) rc->retry_timer = uasync_set_timeout(rc->ua, 5000, rc, retry_timer_cb, "tps_retry");
return;
}
}
e = queue_data_get(q);
if (!e) {
if (rc->dst_fin_deferred) {
rc->dst_fin_deferred = 0;
send_fin(rc);
}
queue_resume_callback(q); return;
}
if (rc->cli_closed) { if (rc->cli_closed) {
do { do {
memory_pool_free(rc->tc->data_pool, e->dgram); memory_pool_free(rc->tc->data_pool, e->dgram);
@ -205,27 +248,50 @@ static void read_queue_drain_cb(struct ll_queue* q, void* arg) {
return; return;
} }
int ret = send_msg(inst, rc->peer_node_id, TCP_PROXY_SUBCMD_DATA, rc->stream_id, e->dgram, e->len, 0); int ret = send_msg(inst, rc->peer_node_id, TCP_PROXY_SUBCMD_DATA, rc->stream_id, e->dgram, e->len, 0);
DEBUG_TRACE(DEBUG_CATEGORY_SOCKET, "SOCK:RECV fd=%d sid=%08x len=%u total=%d", rc->bytes_relayed += e->len; rc->drain_count++;
(int)rc->tc->sock, rc->stream_id, e->len, rc->ctx ? rc->ctx->conn_count : 0); if ((rc->drain_count & 63) == 0 || rc->drain_count <= 4)
memory_pool_free(rc->tc->data_pool, e->dgram); DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "TPS DRAIN #%u sid=%08x len=%u → ret=%d sent=%u recv=%u relayed=%u",
queue_entry_free(e); rc->drain_count, rc->stream_id, e->len, ret, rc->bytes_sent, rc->bytes_recv, rc->bytes_relayed);
if (ret == 0) { if (ret == 0) {
memory_pool_free(rc->tc->data_pool, e->dgram); queue_entry_free(e);
queue_resume_callback(q); queue_resume_callback(q);
} else { } else {
DEBUG_DEBUG(DEBUG_CATEGORY_SOCKET, "SOCK:BACKPRESSURE fd=%d sid=%08x — pausing drain", rc->tx_buf = u_malloc(e->len);
(int)rc->tc->sock, rc->stream_id); if (rc->tx_buf) { memcpy(rc->tx_buf, e->dgram, e->len); rc->tx_len = e->len; }
etcp_router_on_send_ready(inst, rc->peer_node_id, ETCP_ID_TCP_PROXY_CLIENT, else { DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "TPS tx_buf malloc=%u failed sid=%08x — drop", e->len, rc->stream_id); }
&rc->pause_waiter, pause_resume_cb, rc); memory_pool_free(rc->tc->data_pool, e->dgram); queue_entry_free(e);
DEBUG_DEBUG(DEBUG_CATEGORY_SOCKET, "SOCK:BACKPRESSURE fd=%d sid=%08x — buffered %u for retry",
(int)rc->tc->sock, rc->stream_id, rc->tx_len);
if (!rc->retry_timer) rc->retry_timer = uasync_set_timeout(rc->ua, 5000, rc, retry_timer_cb, "tps_retry");
} }
} }
static void pause_resume_cb(struct ll_queue* q, void* arg) { static void pause_resume_cb(struct ll_queue* q, void* arg) {
(void)q; (void)q;
struct tcp_proxy_server_conn* rc = (struct tcp_proxy_server_conn*)arg; struct tcp_proxy_server_conn* rc = (struct tcp_proxy_server_conn*)arg;
if (!rc->tc || rc->tc->sock == SOCKET_INVALID) return; if (!rc->tc || rc->tc->sock == SOCKET_INVALID || rc->cli_closed) return;
queue_resume_callback(rc->tc->read_queue); queue_resume_callback(rc->tc->read_queue);
} }
static void retry_timer_cb(void* arg) {
struct tcp_proxy_server_conn* rc = (struct tcp_proxy_server_conn*)arg;
rc->retry_timer = NULL;
if (!rc->tx_buf || !rc->tc || rc->tc->sock == SOCKET_INVALID || rc->cli_closed) return;
struct UTUN_INSTANCE* inst = rc->ctx ? rc->ctx->inst : NULL;
int ret = send_msg(inst, rc->peer_node_id, TCP_PROXY_SUBCMD_DATA, rc->stream_id, rc->tx_buf, rc->tx_len, 1);
if (ret == 0) {
u_free(rc->tx_buf); rc->tx_buf = NULL; rc->tx_len = 0;
if (rc->dst_fin_deferred && rc->tc && rc->tc->read_queue && !rc->tc->read_queue->head) {
rc->dst_fin_deferred = 0;
send_fin(rc);
return;
}
queue_resume_callback(rc->tc->read_queue);
} else {
rc->retry_timer = uasync_set_timeout(rc->ua, 5000, rc, retry_timer_cb, "tps_retry");
}
}
// ==================================================================== // ====================================================================
// Управление жизненным циклом коннекта // Управление жизненным циклом коннекта
// ==================================================================== // ====================================================================
@ -257,7 +323,8 @@ void tcp_proxy_server_conn_free(struct tcp_proxy_server_conn* rc) {
struct tcp_proxy_server_conn** prev = &rc->ctx->conns; struct tcp_proxy_server_conn** prev = &rc->ctx->conns;
while (*prev) { if (*prev == rc) { *prev = rc->next; rc->ctx->conn_count--; break; } prev = &(*prev)->next; } while (*prev) { if (*prev == rc) { *prev = rc->next; rc->ctx->conn_count--; break; } prev = &(*prev)->next; }
} }
if (rc->ctx && rc->ctx->inst) etcp_router_cancel_send_ready(rc->ctx->inst, rc->peer_node_id, ETCP_ID_TCP_PROXY_CLIENT, &rc->pause_waiter); if (rc->tx_buf) { u_free(rc->tx_buf); rc->tx_buf = NULL; rc->tx_len = 0; }
if (rc->retry_timer) { uasync_cancel_timeout(rc->ua, rc->retry_timer); rc->retry_timer = NULL; }
if (rc->close_timer) { uasync_cancel_timeout(rc->ua, rc->close_timer); rc->close_timer = NULL; } if (rc->close_timer) { uasync_cancel_timeout(rc->ua, rc->close_timer); rc->close_timer = NULL; }
if (rc->diag_timer) { uasync_cancel_timeout(rc->ua, rc->diag_timer); rc->diag_timer = NULL; } if (rc->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; } if (rc->tc) { tcp_conn_destroy(rc->tc); rc->tc = NULL; }
@ -291,7 +358,7 @@ int tcp_proxy_server_handle_connect(struct UTUN_INSTANCE* inst, struct ll_entry*
socket_t sock = socket(AF_INET, SOCK_STREAM, 0); socket_t sock = socket(AF_INET, SOCK_STREAM, 0);
if (sock == SOCKET_INVALID) { DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "TCP proxy server: socket() failed"); u_free(rc); queue_dgram_free(entry); queue_entry_free(entry); return -1; } if (sock == SOCKET_INVALID) { DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "TCP proxy server: socket() failed"); u_free(rc); queue_dgram_free(entry); queue_entry_free(entry); return -1; }
ctx->conn_count++; ctx->conn_count++;
DEBUG_DEBUG(DEBUG_CATEGORY_SOCKET, "SOCK:NEW fd=%d sid=%08x dest=%d.%d.%d.%d:%d total=%d", DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "SOCK:NEW fd=%d sid=%08x dest=%d.%d.%d.%d:%d total=%d",
(int)sock, stream_id, dest_ip[0],dest_ip[1],dest_ip[2],dest_ip[3],ntohs(dest_port), ctx->conn_count); (int)sock, stream_id, dest_ip[0],dest_ip[1],dest_ip[2],dest_ip[3],ntohs(dest_port), ctx->conn_count);
socket_set_nonblocking(sock); socket_set_nonblocking(sock);
@ -329,6 +396,10 @@ int tcp_proxy_server_handle_data(struct UTUN_INSTANCE* inst, struct ETCP_CONN* c
queue_dgram_free(entry); queue_entry_free(entry); return -1; queue_dgram_free(entry); queue_entry_free(entry); return -1;
} }
size_t data_len = entry->len - TCP_PROXY_HDR_SIZE; size_t data_len = entry->len - TCP_PROXY_HDR_SIZE;
rc->bytes_sent += (uint32_t)data_len; rc->data_count++;
if ((rc->data_count & 63) == 0 || rc->data_count <= 4)
DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "TPS DATA #%u sid=%08x len=%zu sent=%u recv=%u relayed=%u",
rc->data_count, stream_id, data_len, rc->bytes_sent, rc->bytes_recv, rc->bytes_relayed);
if (data_len > 0) { if (data_len > 0) {
if (data_len > rc->tc->data_pool->object_size) { if (data_len > rc->tc->data_pool->object_size) {
DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "tcp_proxy_server: data_len=%zu > pool_sz=%zu, dropping sid=%08x", DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "tcp_proxy_server: data_len=%zu > pool_sz=%zu, dropping sid=%08x",
@ -420,6 +491,8 @@ void tcp_proxy_server_recv_cb(struct ETCP_CONN* conn, struct ll_entry* entry) {
uint8_t subcmd = entry->dgram[1]; uint8_t subcmd = entry->dgram[1];
uint32_t stream_id; memcpy(&stream_id, entry->dgram + 2, 4); uint32_t stream_id; memcpy(&stream_id, entry->dgram + 2, 4);
DEBUG_TRACE(DEBUG_CATEGORY_TRAFFIC, "PROXY SERVER RECV subcmd=%02x sid=%08x pkt_len=%u", subcmd, stream_id, entry->len);
if (subcmd == TCP_PROXY_SUBCMD_CONNECT) { if (subcmd == TCP_PROXY_SUBCMD_CONNECT) {
uint64_t src_node_id = conn ? conn->peer_node_id : (inst ? inst->node_id : 0); uint64_t src_node_id = conn ? conn->peer_node_id : (inst ? inst->node_id : 0);
tcp_proxy_server_handle_connect(inst, entry, stream_id, src_node_id); tcp_proxy_server_handle_connect(inst, entry, stream_id, src_node_id);

9
src/proxy/tcp_proxy_server.h

@ -36,9 +36,18 @@ struct tcp_proxy_server_conn {
void* close_timer; // таймер повтора CLOSE/ERROR void* close_timer; // таймер повтора CLOSE/ERROR
int close_backoff; // backoff: 50..5000 tb (5ms..500ms) int close_backoff; // backoff: 50..5000 tb (5ms..500ms)
void* diag_timer; // 1-секундный таймер диагностики void* diag_timer; // 1-секундный таймер диагностики
void* retry_timer; // таймер повтора tx_buf при backpressure (500ms)
uint8_t dst_fin_deferred; // FIN от destination отложен (ждём drain read_queue + tx_buf)
uint8_t freed; // 1 = уже освобождён, защита от double-free uint8_t freed; // 1 = уже освобождён, защита от double-free
struct queue_waiter_handle pause_waiter; struct queue_waiter_handle pause_waiter;
uint8_t* tx_buf; // буфер при backpressure (retry в pause_resume_cb)
uint16_t tx_len;
uint32_t bytes_sent; // байт записано в destination (реальный TCP)
uint32_t bytes_recv; // байт прочитано от destination
uint32_t bytes_relayed; // байт отправлено клиенту через ETCP
uint32_t drain_count; // счётчик вызовов read_queue_drain_cb
uint32_t data_count; // счётчик входящих DATA от клиента
struct UASYNC* ua; struct UASYNC* ua;
}; };

263
tests/test_socks_http_proxy.c

@ -1,5 +1,7 @@
// test_socks_http_proxy.c — SOCKS5 + HTTP CONNECT proxy integration test // test_socks_http_proxy.c — Comprehensive SOCKS5 + HTTP proxy integration test
// 2 workers (1 SOCKS + 1 HTTP), 25 reqs each, files up to 100KB // 4 workers: socks_ipv4, socks_domain, http_proxy, https_connect
// Files: 1KB, 100KB, 1MB, 10MB — deterministic content, byte-by-byte verification
// HTTP + HTTPS (self-signed cert) + POST/echo + HEAD + OPTIONS
#include <stdio.h> #include <stdio.h>
#include <stdlib.h> #include <stdlib.h>
#include <string.h> #include <string.h>
@ -28,25 +30,27 @@
#include "../src/routing.h" #include "../src/routing.h"
#include "../src/tun_if.h" #include "../src/tun_if.h"
#define TIMEOUT_MS 120000 #define TIMEOUT_MS 300000
#define WORKERS 3 #define WORKERS 4
#define REQS_PER_WORKER 25
static const int file_sizes[] = { 1024, 10*1024, 50*1024, 100*1024 }; enum { W_SOCKS_IPV4, W_SOCKS_DOMAIN, W_HTTP_PROXY, W_HTTPS_CONNECT };
static const int file_counts[] = { 8, 7, 5, 5 }; // sum = 25
static const char* file_names[] = { "f_1k.bin", "f_10k.bin", "f_50k.bin", "f_100k.bin" }; static const int file_sizes[] = { 1024, 100*1024, 1024*1024, 10*1024*1024 };
static const char* file_names[] = { "f_1k.bin", "f_100k.bin", "f_1m.bin", "f_10m.bin" };
static const int num_files = 4; static const int num_files = 4;
static struct UTUN_INSTANCE* g_cli = NULL; static struct UTUN_INSTANCE* g_cli = NULL;
static struct UTUN_INSTANCE* g_srv = NULL; static struct UTUN_INSTANCE* g_srv = NULL;
static struct UASYNC* g_ua = NULL; static struct UASYNC* g_ua = NULL;
static int g_ok = 0, g_done = 0, g_phase = 0; static int g_ok = 0, g_done = 0, g_phase = 0;
static pid_t g_http_pid = 0; static pid_t g_http_pid = 0, g_https_pid = 0;
static pid_t g_workers[WORKERS]; static pid_t g_workers[WORKERS];
static int g_http_port = 0, g_cli_port = 0, g_srv_port = 0; static int g_http_port = 0, g_https_port = 0, g_cli_port = 0, g_srv_port = 0;
static int g_socks_port = 0, g_http_proxy_port = 0; static int g_socks_port = 0, g_http_proxy_port = 0;
static char g_tmpdir[256]; static char g_tmpdir[256];
static char g_cert[512], g_key[512];
static void* g_mon_id = NULL; static void* g_mon_id = NULL;
static int g_file_seed = 0;
static int alloc_port(void) { static int alloc_port(void) {
int s = socket(AF_INET, SOCK_STREAM, 0); int s = socket(AF_INET, SOCK_STREAM, 0);
@ -62,8 +66,13 @@ static int alloc_port(void) {
static void gen_file(const char* path, int size, int seed) { static void gen_file(const char* path, int size, int seed) {
FILE* f = fopen(path, "wb"); FILE* f = fopen(path, "wb");
if (!f) return; if (!f) { fprintf(stderr, "gen_file: fopen(%s) failed\n", path); return; }
for (int i = 0; i < size; i++) { uint8_t b = (uint8_t)((i ^ seed) & 0xFF); fwrite(&b, 1, 1, f); } uint8_t buf[4096];
for (int off = 0; off < size; off += (int)sizeof(buf)) {
int chunk = size - off; if (chunk > (int)sizeof(buf)) chunk = (int)sizeof(buf);
for (int i = 0; i < chunk; i++) buf[i] = (uint8_t)(((off + i) ^ seed) & 0xFF);
fwrite(buf, 1, (size_t)chunk, f);
}
fclose(f); fclose(f);
} }
@ -72,7 +81,7 @@ static int start_http_server(int port, const char* dir) {
FILE* sf = fopen(script, "w"); FILE* sf = fopen(script, "w");
if (!sf) return -1; if (!sf) return -1;
fprintf(sf, fprintf(sf,
"import sys,os\n" "import sys,os,ssl,threading\n"
"from http.server import HTTPServer,SimpleHTTPRequestHandler\n" "from http.server import HTTPServer,SimpleHTTPRequestHandler\n"
"class H(SimpleHTTPRequestHandler):\n" "class H(SimpleHTTPRequestHandler):\n"
" def do_POST(self):\n" " def do_POST(self):\n"
@ -91,7 +100,7 @@ static int start_http_server(int port, const char* dir) {
" self.end_headers()\n" " self.end_headers()\n"
"port=int(sys.argv[1])\n" "port=int(sys.argv[1])\n"
"os.chdir(sys.argv[2])\n" "os.chdir(sys.argv[2])\n"
"HTTPServer(('',port),H).serve_forever()\n" "HTTPServer(('127.0.0.1',port),H).serve_forever()\n"
); );
fclose(sf); fclose(sf);
@ -107,6 +116,66 @@ static int start_http_server(int port, const char* dir) {
return 0; return 0;
} }
static int start_https_server(int port, const char* dir, const char* cert, const char* key) {
char script[512]; snprintf(script, sizeof(script), "%s/_srv_https.py", dir);
FILE* sf = fopen(script, "w");
if (!sf) return -1;
fprintf(sf,
"import sys,os,ssl,threading,traceback\n"
"from http.server import HTTPServer,SimpleHTTPRequestHandler\n"
"class H(SimpleHTTPRequestHandler):\n"
" def do_GET(self):\n"
" try:\n super().do_GET()\n"
" except Exception as e:\n"
" sys.stderr.write(f'HTTPS_ERR GET {self.path}: {e}\\n{traceback.format_exc()}\\n')\n"
" def do_POST(self):\n"
" if self.path=='/echo':\n"
" n=int(self.headers.get('Content-Length',0))\n"
" b=self.rfile.read(n)\n"
" self.send_response(200)\n"
" self.send_header('Content-Type','application/octet-stream')\n"
" self.send_header('Content-Length',str(len(b)))\n"
" self.end_headers();self.wfile.write(b)\n"
" else:self.send_response(405);self.end_headers()\n"
" def do_OPTIONS(self):\n"
" self.send_response(200)\n"
" self.send_header('Allow','GET,HEAD,POST,OPTIONS')\n"
" self.send_header('Content-Length','0')\n"
" self.end_headers()\n"
"os.chdir(sys.argv[3])\n"
"ctx=ssl.SSLContext(ssl.PROTOCOL_TLS_SERVER)\n"
"ctx.load_cert_chain(sys.argv[1],sys.argv[2])\n"
"httpd=HTTPServer(('127.0.0.1',int(sys.argv[4])),H)\n"
"httpd.socket=ctx.wrap_socket(httpd.socket,server_side=True)\n"
"httpd.timeout=None\n"
"httpd.serve_forever()\n"
);
fclose(sf);
pid_t pid = fork();
if (pid < 0) return -1;
if (pid == 0) {
char port_str[16]; snprintf(port_str, sizeof(port_str), "%d", port);
execlp("python3", "python3", script, cert, key, dir, port_str, (char*)NULL);
_exit(1);
}
g_https_pid = pid;
usleep(300000);
return 0;
}
static int gen_self_signed_cert(const char* tmpdir, char* cert_out, size_t cert_sz, char* key_out, size_t key_sz) {
snprintf(cert_out, cert_sz, "%s/cert.pem", tmpdir);
snprintf(key_out, key_sz, "%s/key.pem", tmpdir);
char cmd[512];
snprintf(cmd, sizeof(cmd),
"openssl req -x509 -newkey rsa:2048 -keyout %s -out %s -days 1 -nodes -subj '/CN=127.0.0.1' 2>/dev/null",
key_out, cert_out);
int rc = system(cmd);
if (rc != 0) { fprintf(stderr, "[FAIL] openssl cert generation\n"); return -1; }
return 0;
}
static char* make_cfg_server(void) { static char* make_cfg_server(void) {
static char buf[1024]; static char buf[1024];
snprintf(buf, sizeof(buf), snprintf(buf, sizeof(buf),
@ -149,64 +218,84 @@ static int conn_ready(struct UTUN_INSTANCE* inst) {
return 0; return 0;
} }
static int worker_main(const char* type, const char* proxy, const char* url_base, const char* tmpdir) { static int worker_main(int wtype, const char* socks_proxy, const char* http_proxy,
const char* url_http, const char* url_https, const char* tmpdir) {
const char* type_names[] = { "socks_ipv4", "socks_domain", "http_proxy", "https_connect" };
const char* tname = type_names[wtype];
for (int fi = 0; fi < num_files; fi++) { for (int fi = 0; fi < num_files; fi++) {
int count = file_counts[fi]; char expected[512], out[512], url[512], cmd[2048];
char expected[512];
snprintf(expected, sizeof(expected), "%s/%s", tmpdir, file_names[fi]); snprintf(expected, sizeof(expected), "%s/%s", tmpdir, file_names[fi]);
snprintf(out, sizeof(out), "%s/out_%d_%d_%s", tmpdir, getpid(), fi, file_names[fi]);
const char* base = (wtype == W_HTTPS_CONNECT) ? url_https : url_http;
snprintf(url, sizeof(url), "%s/%s", base, file_names[fi]);
switch (wtype) {
case W_SOCKS_IPV4:
snprintf(cmd, sizeof(cmd), "curl -s --connect-timeout 15 --max-time 60 --socks5 %s -o %s %s 2>/dev/null",
socks_proxy, out, url);
break;
case W_SOCKS_DOMAIN:
snprintf(cmd, sizeof(cmd), "curl -s --connect-timeout 15 --max-time 60 --socks5-hostname %s -o %s %s 2>/dev/null",
socks_proxy, out, url);
break;
case W_HTTP_PROXY:
snprintf(cmd, sizeof(cmd), "curl -s --connect-timeout 15 --max-time 60 -x http://%s -o %s %s 2>/dev/null",
http_proxy, out, url);
break;
case W_HTTPS_CONNECT:
snprintf(cmd, sizeof(cmd), "curl -s --connect-timeout 15 --max-time 60 --proxytunnel -x %s -k -o %s %s 2>/dev/null",
http_proxy, out, url);
break;
}
for (int i = 0; i < count; i++) { int rc = system(cmd);
char cmd[1024], out[512], url[512]; if (rc != 0) { fprintf(stderr, "[FAIL] %s %s curl exit=%d\n", tname, file_names[fi], WEXITSTATUS(rc)); return 1; }
snprintf(out, sizeof(out), "%s/out_%d_%s_%d.bin", tmpdir, (int)getpid(), file_names[fi], i);
snprintf(url, sizeof(url), "%s/%s", url_base, file_names[fi]); char cmp_cmd[1024];
if (strcmp(type, "socks") == 0) snprintf(cmp_cmd, sizeof(cmp_cmd), "cmp -s %s %s", out, expected);
snprintf(cmd, sizeof(cmd), "curl -s --connect-timeout 10 --max-time 30 --socks5-hostname %s -o %s %s 2>/dev/null", proxy, out, url); if (system(cmp_cmd) != 0) {
else if (strcmp(type, "http_connect") == 0) fprintf(stderr, "[FAIL] %s %s byte mismatch\n", tname, file_names[fi]);
snprintf(cmd, sizeof(cmd), "curl -s --connect-timeout 10 --max-time 30 --proxytunnel -x %s -o %s %s 2>/dev/null", proxy, out, url);
else char sizecmd[1024]; snprintf(sizecmd, sizeof(sizecmd), "wc -c <%s", out);
snprintf(cmd, sizeof(cmd), "curl -s --connect-timeout 10 --max-time 30 -x %s -o %s %s 2>/dev/null", proxy, out, url); FILE* pf = popen(sizecmd, "r"); if (pf) { int sz = 0; fscanf(pf, "%d", &sz); pclose(pf);
fprintf(stderr, " downloaded: %d bytes, expected: %d bytes\n", sz, file_sizes[fi]); }
int rc = system(cmd); unlink(out); return 1;
if (rc != 0) { fprintf(stderr, "[FAIL] worker %s %s iter %d curl exit %d\n", type, file_names[fi], i, WEXITSTATUS(rc)); return 1; }
char cmp_cmd[1024];
snprintf(cmp_cmd, sizeof(cmp_cmd), "cmp -s %s %s", out, expected);
if (system(cmp_cmd) != 0) { fprintf(stderr, "[FAIL] worker %s %s iter %d mismatch\n", type, file_names[fi], i); return 1; }
unlink(out);
} }
unlink(out);
} }
if (strcmp(type, "http_proxy") == 0) {
char cmd[1024], out[512], f1[512]; if (wtype == W_HTTP_PROXY) {
char cmd[2048], out[512], f1[512];
int seed = (int)getpid() ^ 0x55; int seed = (int)getpid() ^ 0x55;
// POST /echo — echo body back, compare snprintf(f1, sizeof(f1), "%s/post_%d.bin", tmpdir, getpid());
snprintf(f1, sizeof(f1), "%s/post_%d.bin", tmpdir, (int)getpid()); gen_file(f1, 256, seed);
gen_file(f1, 512, seed); snprintf(out, sizeof(out), "%s/post_out_%d.bin", tmpdir, getpid());
snprintf(out, sizeof(out), "%s/post_out_%d.bin", tmpdir, (int)getpid());
snprintf(cmd, sizeof(cmd), snprintf(cmd, sizeof(cmd),
"curl -s --connect-timeout 10 --max-time 10 -x %s --data-binary @%s -o %s %s/echo 2>/dev/null", "curl -s --connect-timeout 10 --max-time 30 -x http://%s --data-binary @%s -o %s %s/echo 2>/dev/null",
proxy, f1, out, url_base); http_proxy, f1, out, url_http);
if (system(cmd) != 0) { fprintf(stderr, "[FAIL] http_proxy POST curl\n"); unlink(f1); return 1; } if (system(cmd) != 0) { fprintf(stderr, "[FAIL] %s POST/echo curl\n", tname); unlink(f1); return 1; }
snprintf(cmd, sizeof(cmd), "cmp -s %s %s", out, f1); snprintf(cmd, sizeof(cmd), "cmp -s %s %s", out, f1);
if (system(cmd) != 0) { fprintf(stderr, "[FAIL] http_proxy POST cmp\n"); unlink(f1); unlink(out); return 1; } if (system(cmd) != 0) { fprintf(stderr, "[FAIL] %s POST/echo cmp\n", tname); unlink(f1); unlink(out); return 1; }
unlink(f1); unlink(out); unlink(f1); unlink(out);
// HEAD /f_1k.bin — check status 200
snprintf(cmd, sizeof(cmd), snprintf(cmd, sizeof(cmd),
"curl -s -I --connect-timeout 10 --max-time 10 -x %s %s/f_1k.bin 2>/dev/null | head -1 | grep -q '200 OK'", "curl -s -I --connect-timeout 10 --max-time 30 -x http://%s %s/f_1k.bin 2>/dev/null | head -1 | grep -q '200 OK'",
proxy, url_base); http_proxy, url_http);
if (system(cmd) != 0) { fprintf(stderr, "[FAIL] http_proxy HEAD 200\n"); return 1; } if (system(cmd) != 0) { fprintf(stderr, "[FAIL] %s HEAD 200\n", tname); return 1; }
// OPTIONS /f_1k.bin — check Allow header snprintf(out, sizeof(out), "%s/opt_%d.txt", tmpdir, getpid());
snprintf(out, sizeof(out), "%s/opt_%d.txt", tmpdir, (int)getpid());
snprintf(cmd, sizeof(cmd), snprintf(cmd, sizeof(cmd),
"curl -s -i --connect-timeout 10 --max-time 10 -x %s -X OPTIONS -o %s %s/f_1k.bin 2>/dev/null", "curl -s -i --connect-timeout 10 --max-time 30 -x http://%s -X OPTIONS -o %s %s/f_1k.bin 2>/dev/null",
proxy, out, url_base); http_proxy, out, url_http);
if (system(cmd) != 0) { fprintf(stderr, "[FAIL] http_proxy OPTIONS curl\n"); return 1; } if (system(cmd) != 0) { fprintf(stderr, "[FAIL] %s OPTIONS curl\n", tname); return 1; }
snprintf(cmd, sizeof(cmd), "grep -q 'Allow:' %s", out); snprintf(cmd, sizeof(cmd), "grep -q 'Allow:' %s", out);
if (system(cmd) != 0) { fprintf(stderr, "[FAIL] http_proxy OPTIONS Allow\n"); unlink(out); return 1; } if (system(cmd) != 0) { fprintf(stderr, "[FAIL] %s OPTIONS Allow\n", tname); unlink(out); return 1; }
unlink(out); unlink(out);
} }
return 0; return 0;
} }
@ -218,17 +307,35 @@ static void monitor(void* arg) {
if (conn_ready(g_cli) && conn_ready(g_srv)) { if (conn_ready(g_cli) && conn_ready(g_srv)) {
printf(" ETCP connected, launching %d workers...\n", WORKERS); fflush(stdout); printf(" ETCP connected, launching %d workers...\n", WORKERS); fflush(stdout);
g_phase = 1; g_phase = 1;
char url[128]; snprintf(url, sizeof(url), "http://127.0.0.1:%d", g_http_port);
char socks_proxy[64]; snprintf(socks_proxy, sizeof(socks_proxy), "127.0.0.1:%d", g_socks_port);
char http_proxy[64]; snprintf(http_proxy, sizeof(http_proxy), "http://127.0.0.1:%d", g_http_proxy_port);
const char* worker_types[] = { "socks", "http_connect", "http_proxy" }; char url_http[128]; snprintf(url_http, sizeof(url_http), "http://127.0.0.1:%d", g_http_port);
const char* worker_proxies[] = { socks_proxy, http_proxy, http_proxy }; char url_https[128]; snprintf(url_https, sizeof(url_https), "https://127.0.0.1:%d", g_https_port);
char socks[64]; snprintf(socks, sizeof(socks), "127.0.0.1:%d", g_socks_port);
char hproxy[64]; snprintf(hproxy, sizeof(hproxy), "127.0.0.1:%d", g_http_proxy_port);
// Для socks_domain worker'а: используем localhost вместо 127.0.0.1 чтобы тестировать atyp=3 (DNS через SOCKS)
char url_http_local[128];
snprintf(url_http_local, sizeof(url_http_local), "http://localhost:%d", g_http_port);
struct {
int type;
const char* socks;
const char* hproxy;
const char* url_h;
const char* url_s;
} wcfg[WORKERS] = {
{ W_SOCKS_IPV4, socks, NULL, url_http, NULL },
{ W_SOCKS_DOMAIN,socks, NULL, url_http_local, NULL },
{ W_HTTP_PROXY, NULL, hproxy, url_http, NULL },
{ W_HTTPS_CONNECT,NULL, hproxy, NULL, url_https },
};
for (int i = 0; i < WORKERS; i++) { for (int i = 0; i < WORKERS; i++) {
pid_t pid = fork(); pid_t pid = fork();
if (pid < 0) { printf("[FAIL] fork worker %d: %s\n", i, strerror(errno)); g_done = -1; return; } if (pid < 0) { printf("[FAIL] fork worker %d: %s\n", i, strerror(errno)); g_done = -1; return; }
if (pid == 0) { if (pid == 0) {
int rc = worker_main(worker_types[i], worker_proxies[i], url, g_tmpdir); int rc = worker_main(wcfg[i].type, wcfg[i].socks, wcfg[i].hproxy,
wcfg[i].url_h, wcfg[i].url_s, g_tmpdir);
_exit(rc); _exit(rc);
} }
g_workers[i] = pid; g_workers[i] = pid;
@ -251,42 +358,48 @@ static void monitor(void* arg) {
} }
} }
if (!g_done) g_mon_id = uasync_set_timeout(g_ua, 500, NULL, monitor, "mon"); if (!g_done) g_mon_id = uasync_set_timeout(g_ua, 1000, NULL, monitor, "mon");
} }
static void test_timeout(void* arg) { static void test_timeout(void* arg) {
(void)arg; (void)arg;
if (!g_done) { printf("[FAIL] timeout\n"); g_done = -1; } if (!g_done) { printf("[FAIL] timeout (%ds)\n", TIMEOUT_MS/1000); g_done = -1; }
} }
int main(void) { int main(void) {
printf("=== test_socks_http_proxy ===\n"); fflush(stdout); printf("=== test_socks_http_proxy ===\n"); fflush(stdout);
// Резервируем fd 0
{ int fd = open("/dev/null", O_RDONLY); if (fd > 0 && fd != 0) { dup2(fd, 0); close(fd); } else if (fd < 0) { close(0); open("/dev/null", O_RDONLY); } } { int fd = open("/dev/null", O_RDONLY); if (fd > 0 && fd != 0) { dup2(fd, 0); close(fd); } else if (fd < 0) { close(0); open("/dev/null", O_RDONLY); } }
debug_config_init(); debug_set_level(DEBUG_LEVEL_ERROR); debug_set_categories(DEBUG_CATEGORY_ALL); debug_config_init(); debug_set_level(DEBUG_LEVEL_ERROR); debug_set_categories(DEBUG_CATEGORY_ALL);
utun_instance_set_tun_init_enabled(0); utun_instance_set_tun_init_enabled(0);
srand((unsigned)time(NULL)); srand((unsigned)time(NULL));
g_http_port = alloc_port(); g_cli_port = alloc_port(); g_srv_port = alloc_port(); g_http_port = alloc_port(); g_https_port = alloc_port();
g_cli_port = alloc_port(); g_srv_port = alloc_port();
g_socks_port = alloc_port(); g_http_proxy_port = alloc_port(); g_socks_port = alloc_port(); g_http_proxy_port = alloc_port();
if (g_http_port < 0 || g_cli_port < 0 || g_srv_port < 0 || g_socks_port < 0 || g_http_proxy_port < 0) { if (g_http_port < 0 || g_https_port < 0 || g_cli_port < 0 || g_srv_port < 0 || g_socks_port < 0 || g_http_proxy_port < 0) {
printf("[FAIL] port allocation\n"); return 1; printf("[FAIL] port allocation\n"); return 1;
} }
printf(" ports: http=%d socks=%d http_proxy=%d\n", g_http_port, g_socks_port, g_http_proxy_port); printf(" ports: http=%d https=%d socks=%d http_proxy=%d\n", g_http_port, g_https_port, g_socks_port, g_http_proxy_port);
strcpy(g_tmpdir, "/tmp/utun_test_XXXXXX"); strcpy(g_tmpdir, "/tmp/utun_test_XXXXXX");
if (test_mkdtemp(g_tmpdir) < 0) { printf("[FAIL] mkdtemp: %s\n", strerror(errno)); return 1; } if (test_mkdtemp(g_tmpdir) < 0) { printf("[FAIL] mkdtemp: %s\n", strerror(errno)); return 1; }
int seed = (int)time(NULL) ^ (int)getpid(); if (gen_self_signed_cert(g_tmpdir, g_cert, sizeof(g_cert), g_key, sizeof(g_key)) < 0) goto cleanup;
g_file_seed = (int)time(NULL) ^ (int)getpid();
int total_mb = 0;
for (int fi = 0; fi < num_files; fi++) { for (int fi = 0; fi < num_files; fi++) {
char path[512]; snprintf(path, sizeof(path), "%s/%s", g_tmpdir, file_names[fi]); char path[512]; snprintf(path, sizeof(path), "%s/%s", g_tmpdir, file_names[fi]);
gen_file(path, file_sizes[fi], seed + fi); gen_file(path, file_sizes[fi], g_file_seed + fi);
total_mb += file_sizes[fi];
} }
printf(" test files: 1k, 10k, 50k, 100k\n"); fflush(stdout); printf(" test files: 1KB, 100KB, 1MB, 10MB (total %.1fMB)\n", total_mb / (1024.0*1024.0)); fflush(stdout);
if (start_http_server(g_http_port, g_tmpdir) < 0) { printf("[FAIL] start http server\n"); goto cleanup; } if (start_http_server(g_http_port, g_tmpdir) < 0) { printf("[FAIL] start http server\n"); goto cleanup; }
if (start_https_server(g_https_port, g_tmpdir, g_cert, g_key) < 0) { printf("[FAIL] start https server\n"); goto cleanup; }
printf(" HTTP/HTTPS servers started\n");
g_ua = uasync_create(); g_ua = uasync_create();
if (!g_ua) { printf("[FAIL] uasync_create\n"); goto cleanup; } if (!g_ua) { printf("[FAIL] uasync_create\n"); goto cleanup; }
@ -305,13 +418,19 @@ int main(void) {
if (to_id) uasync_cancel_timeout(g_ua, to_id); if (to_id) uasync_cancel_timeout(g_ua, to_id);
if (g_ok) printf("[PASS] test_socks_http_proxy — %d reqs, %d workers (socks/http_connect/http_proxy), up to 100KB\n", WORKERS * REQS_PER_WORKER, WORKERS); if (g_ok) {
else printf("[FAIL] test_socks_http_proxy\n"); int total_files = WORKERS * num_files;
printf("[PASS] test_socks_http_proxy — %d workers × %d files (%.1fMB each) = %d downloads verified byte-by-byte\n",
WORKERS, num_files, total_mb / (1024.0*1024.0), total_files);
} else {
printf("[FAIL] test_socks_http_proxy\n");
}
cleanup: cleanup:
for (int i = 0; i < WORKERS; i++) if (g_workers[i] > 0) { kill(g_workers[i], SIGKILL); waitpid(g_workers[i], NULL, 0); } for (int i = 0; i < WORKERS; i++) if (g_workers[i] > 0) { kill(g_workers[i], SIGKILL); waitpid(g_workers[i], NULL, 0); }
if (g_mon_id) uasync_cancel_timeout(g_ua, g_mon_id); if (g_mon_id) uasync_cancel_timeout(g_ua, g_mon_id);
if (g_http_pid > 0) { kill(g_http_pid, SIGKILL); waitpid(g_http_pid, NULL, 0); } if (g_http_pid > 0) { kill(g_http_pid, SIGKILL); waitpid(g_http_pid, NULL, 0); }
if (g_https_pid > 0) { kill(g_https_pid, SIGKILL); waitpid(g_https_pid, NULL, 0); }
if (g_cli) { g_cli->running = 0; utun_instance_destroy(g_cli); } if (g_cli) { g_cli->running = 0; utun_instance_destroy(g_cli); }
if (g_srv) { g_srv->running = 0; utun_instance_destroy(g_srv); } if (g_srv) { g_srv->running = 0; utun_instance_destroy(g_srv); }
if (g_ua) uasync_destroy(g_ua, 0); if (g_ua) uasync_destroy(g_ua, 0);

Loading…
Cancel
Save