diff --git a/src/proxy/socks_proxy.c b/src/proxy/socks_proxy.c index 6bac606a..a9afc902 100644 --- a/src/proxy/socks_proxy.c +++ b/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; } 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) { @@ -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) { + 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; else { c->close_pending = 0; c->close_sent = 1; } } 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; else { c->close_pending = 0; c->close_sent = 1; } } 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); } @@ -110,6 +116,7 @@ static int write_to_client(struct socks_proxy_conn* c, const uint8_t* data, uint memcpy(buf, data, len); e->dgram = buf; e->len = len; queue_data_put(c->tc->write_queue, e); + c->bytes_to_client += len; return 0; } @@ -297,7 +304,8 @@ static void process_http_request(struct socks_proxy_conn* c) { if (version) { size_t vlen = strlen(version); 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; @@ -309,10 +317,21 @@ static void process_http_request(struct socks_proxy_conn* c) { goto error; } 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) { - send_data(c, c->buf + headers_len, c->buf_len - headers_len); + uint8_t* pkt = u_malloc(total); + 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; @@ -361,12 +380,20 @@ static void on_read_cb(struct ll_queue* q, void* arg) { } // 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); DEBUG_TRACE(DEBUG_CATEGORY_TRAFFIC, "SOCKS PROXY SEND sid=%08x len=%u", c->stream_id, e->len); if (ret == 0) { memory_pool_free(c->tc->data_pool, e->dgram); queue_entry_free(e); queue_resume_callback(q); } 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); 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) { (void)q; struct socks_proxy_conn* c = (struct socks_proxy_conn*)arg; - if (c->rem_closed || c->close_sent) return; - queue_resume_callback(c->tc->read_queue); + if (!c->tx_buf || c->rem_closed || c->close_sent) return; + 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) { (void)err; 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); socks_proxy_conn_free(c); } static void on_closed_cb(struct tcp_conn* tc, void* 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); 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 (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 > c->tc->data_pool->object_size) { 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) { memcpy(buf, data, data_len); e->dgram = buf; e->len = (uint16_t)data_len; queue_data_put(c->tc->write_queue, e); + c->bytes_to_client += (uint32_t)data_len; } 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); } } @@ -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); } 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->inst) etcp_router_waiter_cancel(c->inst, c->via_node_id, &c->tx_waiter); u_free(c); diff --git a/src/proxy/socks_proxy.h b/src/proxy/socks_proxy.h index 33acded8..ce92a2e1 100644 --- a/src/proxy/socks_proxy.h +++ b/src/proxy/socks_proxy.h @@ -45,6 +45,10 @@ struct socks_proxy_conn { uint8_t freed; uint8_t buf[1024]; 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; }; diff --git a/src/proxy/tcp_proxy_client.c b/src/proxy/tcp_proxy_client.c index fa62e3c1..f3d9abaf 100644 --- a/src/proxy/tcp_proxy_client.c +++ b/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]; 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) { size_t data_len = entry->len > TCP_PROXY_HDR_SIZE ? entry->len - TCP_PROXY_HDR_SIZE : 0; // Пробуем SOCKS/HTTP first (у них приоритет — могут быть без TUN) diff --git a/src/proxy/tcp_proxy_server.c b/src/proxy/tcp_proxy_server.c index 0b0d61f6..5a2f194c 100644 --- a/src/proxy/tcp_proxy_server.c +++ b/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 close_retry_cb(void* arg); static void diag_timer_cb(void* arg); +static void retry_timer_cb(void* arg); static int send_msg(struct UTUN_INSTANCE* inst, uint64_t dst, uint8_t subcmd, uint32_t sid, const uint8_t* data, size_t len, int force); static void send_close(struct tcp_proxy_server_conn* rc); @@ -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) { struct tcp_proxy_server_conn* rc = (struct tcp_proxy_server_conn*)arg; 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) { - 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); - } else if (pend) { + } else if (pend_w) { tcp_conn_set_flushed(tc, on_flushed_cb); - DEBUG_DEBUG(DEBUG_CATEGORY_SOCKET, "SOCK:DST_FIN fd=%d sid=%08x write_pend=%d → deferred", - (int)tc->sock, rc->stream_id, pend); + } else if (pend_r) { + rc->dst_fin_deferred = 1; } else { 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; if (tc->closed) return; 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); 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) { struct tcp_proxy_server_conn* rc = (struct tcp_proxy_server_conn*)arg; struct UTUN_INSTANCE* inst = rc->ctx ? rc->ctx->inst : NULL; - struct ll_entry* e = queue_data_get(q); - if (!e) { queue_resume_callback(q); return; } + struct ll_entry* e; + + // Сначала добиваем отложенный 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) { do { 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; } 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", - (int)rc->tc->sock, rc->stream_id, e->len, rc->ctx ? rc->ctx->conn_count : 0); - memory_pool_free(rc->tc->data_pool, e->dgram); - queue_entry_free(e); + rc->bytes_relayed += e->len; rc->drain_count++; + if ((rc->drain_count & 63) == 0 || rc->drain_count <= 4) + DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "TPS DRAIN #%u sid=%08x len=%u → ret=%d sent=%u recv=%u relayed=%u", + rc->drain_count, rc->stream_id, e->len, ret, rc->bytes_sent, rc->bytes_recv, rc->bytes_relayed); if (ret == 0) { + memory_pool_free(rc->tc->data_pool, e->dgram); queue_entry_free(e); queue_resume_callback(q); } else { - DEBUG_DEBUG(DEBUG_CATEGORY_SOCKET, "SOCK:BACKPRESSURE fd=%d sid=%08x — pausing drain", - (int)rc->tc->sock, rc->stream_id); - etcp_router_on_send_ready(inst, rc->peer_node_id, ETCP_ID_TCP_PROXY_CLIENT, - &rc->pause_waiter, pause_resume_cb, rc); + rc->tx_buf = u_malloc(e->len); + if (rc->tx_buf) { memcpy(rc->tx_buf, e->dgram, e->len); rc->tx_len = e->len; } + else { DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "TPS tx_buf malloc=%u failed sid=%08x — drop", e->len, rc->stream_id); } + 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) { (void)q; 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); } +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; 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->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; } @@ -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); 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++; - 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); 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; } 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 > rc->tc->data_pool->object_size) { 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]; 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) { 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); diff --git a/src/proxy/tcp_proxy_server.h b/src/proxy/tcp_proxy_server.h index fceebd37..7a65a674 100644 --- a/src/proxy/tcp_proxy_server.h +++ b/src/proxy/tcp_proxy_server.h @@ -36,9 +36,18 @@ struct tcp_proxy_server_conn { void* close_timer; // таймер повтора CLOSE/ERROR int close_backoff; // backoff: 50..5000 tb (5ms..500ms) 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 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; }; diff --git a/tests/test_socks_http_proxy.c b/tests/test_socks_http_proxy.c index 6ee3a493..e9352eac 100644 --- a/tests/test_socks_http_proxy.c +++ b/tests/test_socks_http_proxy.c @@ -1,5 +1,7 @@ -// test_socks_http_proxy.c — SOCKS5 + HTTP CONNECT proxy integration test -// 2 workers (1 SOCKS + 1 HTTP), 25 reqs each, files up to 100KB +// test_socks_http_proxy.c — Comprehensive SOCKS5 + HTTP proxy integration test +// 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 #include #include @@ -28,25 +30,27 @@ #include "../src/routing.h" #include "../src/tun_if.h" -#define TIMEOUT_MS 120000 -#define WORKERS 3 -#define REQS_PER_WORKER 25 +#define TIMEOUT_MS 300000 +#define WORKERS 4 -static const int file_sizes[] = { 1024, 10*1024, 50*1024, 100*1024 }; -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" }; +enum { W_SOCKS_IPV4, W_SOCKS_DOMAIN, W_HTTP_PROXY, W_HTTPS_CONNECT }; + +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 struct UTUN_INSTANCE* g_cli = NULL; static struct UTUN_INSTANCE* g_srv = NULL; static struct UASYNC* g_ua = NULL; 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 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 char g_tmpdir[256]; +static char g_cert[512], g_key[512]; static void* g_mon_id = NULL; +static int g_file_seed = 0; static int alloc_port(void) { 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) { FILE* f = fopen(path, "wb"); - if (!f) return; - for (int i = 0; i < size; i++) { uint8_t b = (uint8_t)((i ^ seed) & 0xFF); fwrite(&b, 1, 1, f); } + if (!f) { fprintf(stderr, "gen_file: fopen(%s) failed\n", path); return; } + 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); } @@ -72,7 +81,7 @@ static int start_http_server(int port, const char* dir) { FILE* sf = fopen(script, "w"); if (!sf) return -1; fprintf(sf, - "import sys,os\n" + "import sys,os,ssl,threading\n" "from http.server import HTTPServer,SimpleHTTPRequestHandler\n" "class H(SimpleHTTPRequestHandler):\n" " def do_POST(self):\n" @@ -91,7 +100,7 @@ static int start_http_server(int port, const char* dir) { " self.end_headers()\n" "port=int(sys.argv[1])\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); @@ -107,6 +116,66 @@ static int start_http_server(int port, const char* dir) { 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 buf[1024]; snprintf(buf, sizeof(buf), @@ -149,64 +218,84 @@ static int conn_ready(struct UTUN_INSTANCE* inst) { 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++) { - int count = file_counts[fi]; - char expected[512]; + char expected[512], out[512], url[512], cmd[2048]; 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++) { - char cmd[1024], out[512], url[512]; - 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]); - if (strcmp(type, "socks") == 0) - snprintf(cmd, sizeof(cmd), "curl -s --connect-timeout 10 --max-time 30 --socks5-hostname %s -o %s %s 2>/dev/null", proxy, out, url); - else if (strcmp(type, "http_connect") == 0) - 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 - snprintf(cmd, sizeof(cmd), "curl -s --connect-timeout 10 --max-time 30 -x %s -o %s %s 2>/dev/null", proxy, out, url); - - int rc = system(cmd); - 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); + int rc = system(cmd); + if (rc != 0) { fprintf(stderr, "[FAIL] %s %s curl exit=%d\n", tname, file_names[fi], 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] %s %s byte mismatch\n", tname, file_names[fi]); + + char sizecmd[1024]; snprintf(sizecmd, sizeof(sizecmd), "wc -c <%s", out); + 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]); } + unlink(out); return 1; } + 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; - // POST /echo — echo body back, compare - snprintf(f1, sizeof(f1), "%s/post_%d.bin", tmpdir, (int)getpid()); - gen_file(f1, 512, seed); - snprintf(out, sizeof(out), "%s/post_out_%d.bin", tmpdir, (int)getpid()); + snprintf(f1, sizeof(f1), "%s/post_%d.bin", tmpdir, getpid()); + gen_file(f1, 256, seed); + snprintf(out, sizeof(out), "%s/post_out_%d.bin", tmpdir, getpid()); snprintf(cmd, sizeof(cmd), - "curl -s --connect-timeout 10 --max-time 10 -x %s --data-binary @%s -o %s %s/echo 2>/dev/null", - proxy, f1, out, url_base); - if (system(cmd) != 0) { fprintf(stderr, "[FAIL] http_proxy POST curl\n"); unlink(f1); return 1; } + "curl -s --connect-timeout 10 --max-time 30 -x http://%s --data-binary @%s -o %s %s/echo 2>/dev/null", + http_proxy, f1, out, url_http); + 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); - 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); - // HEAD /f_1k.bin — check status 200 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'", - proxy, url_base); - if (system(cmd) != 0) { fprintf(stderr, "[FAIL] http_proxy HEAD 200\n"); return 1; } + "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'", + http_proxy, url_http); + 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, (int)getpid()); + snprintf(out, sizeof(out), "%s/opt_%d.txt", tmpdir, getpid()); 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", - proxy, out, url_base); - if (system(cmd) != 0) { fprintf(stderr, "[FAIL] http_proxy OPTIONS curl\n"); return 1; } + "curl -s -i --connect-timeout 10 --max-time 30 -x http://%s -X OPTIONS -o %s %s/f_1k.bin 2>/dev/null", + http_proxy, out, url_http); + if (system(cmd) != 0) { fprintf(stderr, "[FAIL] %s OPTIONS curl\n", tname); return 1; } 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); } + return 0; } @@ -218,17 +307,35 @@ static void monitor(void* arg) { if (conn_ready(g_cli) && conn_ready(g_srv)) { printf(" ETCP connected, launching %d workers...\n", WORKERS); fflush(stdout); 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" }; - const char* worker_proxies[] = { socks_proxy, http_proxy, http_proxy }; + char url_http[128]; snprintf(url_http, sizeof(url_http), "http://127.0.0.1:%d", g_http_port); + 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++) { pid_t pid = fork(); if (pid < 0) { printf("[FAIL] fork worker %d: %s\n", i, strerror(errno)); g_done = -1; return; } 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); } 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) { (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) { 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); } } debug_config_init(); debug_set_level(DEBUG_LEVEL_ERROR); debug_set_categories(DEBUG_CATEGORY_ALL); utun_instance_set_tun_init_enabled(0); 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(); - 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(" 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"); 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++) { 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_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(); 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 (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); - else printf("[FAIL] test_socks_http_proxy\n"); + if (g_ok) { + 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: 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_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_srv) { g_srv->running = 0; utun_instance_destroy(g_srv); } if (g_ua) uasync_destroy(g_ua, 0);