You can not select more than 25 topics Topics must start with a letter or number, can include dashes ('-') and can be up to 35 characters long.
 
 
 
 
 

621 lines
30 KiB

// socks_proxy.c — SOCKS5 / HTTP CONNECT proxy (client side)
#include "socks_proxy.h"
#include "tcp_proxy_server.h"
#include "etcp.h"
#include "etcp_api.h"
#include "etcp_router.h"
#include "utun_instance.h"
#include "../routing_layer/topo_node.h"
#include "../lib/u_async.h"
#include "../lib/debug_config.h"
#include "../lib/ll_queue.h"
#include "../lib/memory_pool.h"
#include "../lib/tcp_io.h"
#include "../lib/mem.h"
#include <stdlib.h>
#include <string.h>
#include <errno.h>
#ifndef _WIN32
#include <unistd.h>
#include <sys/socket.h>
#include <netinet/in.h>
#include <netdb.h>
#include <arpa/inet.h>
#endif
static void on_accept_cb(socket_t sock, void* arg);
static void on_read_cb(struct ll_queue* q, void* arg);
static void on_fin_cb(struct tcp_conn* tc, void* arg);
static void on_error_cb(struct tcp_conn* tc, int err, void* arg);
static void on_closed_cb(struct tcp_conn* tc, void* arg);
static void tx_waiter_cb(struct ll_queue* q, void* arg);
static int send_msg(struct UTUN_INSTANCE* inst, uint64_t group_id, uint64_t dst, uint8_t subcmd, uint32_t sid, const uint8_t* data, size_t len, int force);
struct listen_ctx {
socket_t listen_sock;
void* socket_id;
struct UASYNC* ua;
struct UTUN_INSTANCE* inst;
uint64_t via_node_id;
int is_http;
struct socks_proxy_conn** conns;
int* conn_count;
uint32_t* next_stream_id;
};
// ====================================================================
// Отправка сообщений через ETCP
// ====================================================================
static int send_msg(struct UTUN_INSTANCE* inst, uint64_t group_id, uint64_t dst, uint8_t subcmd, uint32_t sid, const uint8_t* data, size_t len, int force) {
struct ll_entry* e = queue_entry_new(0);
if (!e) { DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "socks_proxy: queue_entry_new failed subcmd=%02x sid=%08x", subcmd, sid); return -1; }
e->dgram = u_malloc(TCP_PROXY_HDR_SIZE + len);
if (!e->dgram) { DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "socks_proxy: malloc(%zu) failed", TCP_PROXY_HDR_SIZE + len); queue_entry_free(e); return -1; }
e->dgram[0] = ETCP_RT_ID_TCP_PROXY;
e->dgram[1] = subcmd;
memcpy(e->dgram + 2, &sid, 4);
if (len > 0) memcpy(e->dgram + TCP_PROXY_HDR_SIZE, data, len);
if (TCP_PROXY_HDR_SIZE + len > UINT16_MAX) {
DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "socks_proxy: msg too large len=%zu subcmd=%02x", len, subcmd);
queue_dgram_free(e); queue_entry_free(e); return -1;
}
e->len = (uint16_t)(TCP_PROXY_HDR_SIZE + len);
int ret = etcp_route_send(inst, group_id, dst, e, force);
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) {
uint8_t buf[6];
memcpy(buf, c->dest_ip, 4); memcpy(buf + 4, &c->dest_port, 2);
DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "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");
return send_msg(c->inst, TOPO_GROUP_UTUN, c->via_node_id, TCP_PROXY_SUBCMD_CONNECT, c->stream_id, buf, 6, 1);
}
static int send_data(struct socks_proxy_conn* c, const uint8_t* data, uint16_t len) {
return send_msg(c->inst, TOPO_GROUP_UTUN, c->via_node_id, TCP_PROXY_SUBCMD_DATA, c->stream_id, data, len, 0);
}
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, TOPO_GROUP_UTUN, 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, TOPO_GROUP_UTUN, 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, TOPO_GROUP_UTUN, c->via_node_id, TCP_PROXY_SUBCMD_FIN, c->stream_id, NULL, 0, 1);
}
// ====================================================================
// Запись ответа в TCP клиенту через write_queue
// ====================================================================
static int write_to_client(struct socks_proxy_conn* c, const uint8_t* data, uint16_t len) {
struct ll_entry* e = queue_entry_new_from_pool(c->tc->entry_pool);
uint8_t* buf = memory_pool_alloc(c->tc->data_pool);
if (!e || !buf) {
DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "socks_proxy: write_to_client alloc failed sid=%08x", c->stream_id);
if (e) queue_entry_free(e); if (buf) memory_pool_free(c->tc->data_pool, buf);
return -1;
}
if (len > c->tc->data_pool->object_size) {
DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "socks_proxy: write_to_client len=%u > pool_sz=%zu sid=%08x",
len, c->tc->data_pool->object_size, c->stream_id);
queue_entry_free(e); memory_pool_free(c->tc->data_pool, buf);
return -1;
}
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;
}
// ====================================================================
// SOCKS5 handshake
// ====================================================================
static void process_socks_greeting(struct socks_proxy_conn* c) {
if (c->buf_len < 3) return;
uint8_t ver = c->buf[0], nmethods = c->buf[1];
if (ver != 5 || c->buf_len < (uint16_t)(2 + nmethods)) return;
DEBUG_DEBUG(DEBUG_CATEGORY_SOCKET, "socks_proxy: greeting ver=%d nmethods=%d", ver, nmethods);
uint8_t reply[] = { 0x05, 0x00 };
write_to_client(c, reply, 2);
c->buf_len = 0;
c->state = SOCKS_STATE_REQUEST;
}
static void process_socks_request(struct socks_proxy_conn* c) {
if (c->buf_len < 10) return;
uint8_t ver = c->buf[0], cmd = c->buf[1], atyp = c->buf[3];
if (ver != 5) { DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "socks_proxy: bad ver=%d", ver); goto error; }
if (cmd != 1) { DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "socks_proxy: unsupported cmd=%d (only CONNECT supported)", cmd); goto error; }
uint16_t need;
if (atyp == 1) need = 10; // IPv4: 4+2=6 more bytes
else if (atyp == 3) { // domain: 1+len+2
if (c->buf_len < 5) return;
need = (uint16_t)(5 + c->buf[4] + 2);
}
else if (atyp == 4) need = 22; // IPv6: 16+2=18 more
else { DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "socks_proxy: unsupported atyp=%d", atyp); goto error; }
if (c->buf_len < need) return;
if (atyp == 1) {
memcpy(c->dest_ip, c->buf + 4, 4);
memcpy(&c->dest_port, c->buf + 8, 2);
} else if (atyp == 4) {
// Извлекаем первые 4 байта IPv6 в dest_ip (для упрощения: IPv6 mapped IPv4 или реальный IPv6)
memcpy(c->dest_ip, c->buf + 12, 4);
memcpy(&c->dest_port, c->buf + 20, 2);
// TODO: proper IPv6 support
DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "socks_proxy: IPv6 unsupported, using last 4 bytes of addr");
} else { // atyp == 3 (domain)
uint8_t dlen = c->buf[4];
char domain[256]; memcpy(domain, c->buf + 5, dlen); domain[dlen] = '\0';
memcpy(&c->dest_port, c->buf + 5 + dlen, 2);
// resolve domain
struct hostent* he = gethostbyname(domain);
if (!he || he->h_addrtype != AF_INET) {
DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "socks_proxy: DNS failed for %s", domain);
goto error;
}
memcpy(c->dest_ip, he->h_addr, 4);
DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "socks_proxy: resolved %s → %d.%d.%d.%d:%d",
domain, c->dest_ip[0], c->dest_ip[1], c->dest_ip[2], c->dest_ip[3], ntohs(c->dest_port));
}
c->buf_len = 0;
c->state = SOCKS_STATE_RELAY; // relay immediately, no waiting for first DATA
uint8_t reply[] = { 0x05, 0x00, 0x00, 0x01, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00 };
write_to_client(c, reply, 10);
if (send_connect(c) < 0) {
DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "socks_proxy: send_connect failed sid=%08x", c->stream_id);
tcp_conn_push_close(c->tc);
}
return;
error: {
uint8_t err_reply[] = { 0x05, 0x01, 0x00, 0x01, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00 };
write_to_client(c, err_reply, 10);
tcp_conn_push_close(c->tc);
}
}
// ====================================================================
// 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_SOCKET, "socks_proxy: http request too short"); goto error; }
char line[512]; if (line_len > sizeof(line) - 1) line_len = sizeof(line) - 1;
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_SOCKET, "socks_proxy: bad http connect: %s", line);
goto error;
}
char* colon = strrchr(host_port, ':');
if (!colon) { DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "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_SOCKET, "socks_proxy: bad port %d", port); goto error; }
struct hostent* he = gethostbyname(host_port);
if (!he || he->h_addrtype != AF_INET) {
DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "socks_proxy: DNS failed for %s", host_port);
uint8_t resp[] = "HTTP/1.1 502 Bad Gateway\r\n\r\n";
write_to_client(c, resp, (uint16_t)strlen((char*)resp)); tcp_conn_push_close(c->tc); return;
}
memcpy(c->dest_ip, he->h_addr, 4); c->dest_port = htons((uint16_t)port);
DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "socks_proxy: HTTP CONNECT %s → %d.%d.%d.%d:%d",
host_port, c->dest_ip[0], c->dest_ip[1], c->dest_ip[2], c->dest_ip[3], port);
c->buf_len = 0;
c->state = HTTP_STATE_RELAY;
uint8_t resp[] = "HTTP/1.1 200 Connection Established\r\n\r\n";
write_to_client(c, resp, (uint16_t)strlen((char*)resp));
if (send_connect(c) < 0) { tcp_conn_push_close(c->tc); }
return;
}
// ===== Не-CONNECT: HTTP-прокси (GET, POST, PUT, HEAD, OPTIONS, ...) =====
char* hdr_end = memmem(c->buf, c->buf_len, "\r\n\r\n", 4);
if (!hdr_end) return;
uint16_t headers_len = (uint16_t)((uint8_t*)hdr_end - c->buf) + 4;
char method[16] = {0}, url[512] = {0};
if (sscanf(line, "%15s %511s", method, url) < 2) {
DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "socks_proxy: bad http request line: %s", line);
goto error;
}
if (strncmp(url, "http://", 7) != 0) {
DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "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_SOCKET, "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_SOCKET, "socks_proxy: empty host in %s", url); goto error; }
if (!path_start || *path_start == '\0') path_start = "/";
struct hostent* he = gethostbyname(host);
if (!he || he->h_addrtype != AF_INET) {
DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "socks_proxy: DNS failed for %s", host);
uint8_t resp[] = "HTTP/1.1 502 Bad Gateway\r\n\r\n";
write_to_client(c, resp, (uint16_t)strlen((char*)resp)); tcp_conn_push_close(c->tc); return;
}
memcpy(c->dest_ip, he->h_addr, 4); c->dest_port = htons((uint16_t)port);
DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "socks_proxy: HTTP %s %s:%d%s → %d.%d.%d.%d:%d",
method, host, port, path_start,
c->dest_ip[0], c->dest_ip[1], c->dest_ip[2], c->dest_ip[3], port);
if (send_connect(c) < 0) {
DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "socks_proxy: send_connect failed for %s", host);
uint8_t resp[] = "HTTP/1.1 502 Bad Gateway\r\n\r\n";
write_to_client(c, resp, (uint16_t)strlen((char*)resp)); tcp_conn_push_close(c->tc); return;
}
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[2048];
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_SOCKET, "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_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_on_send_ready(c->inst, TOPO_GROUP_UTUN, c->via_node_id, ETCP_RT_ID_TCP_PROXY, &c->tx_waiter, tx_waiter_cb, c);
}
c->buf_len = 0;
c->state = HTTP_STATE_RELAY;
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);
}
}
// ====================================================================
// tcp_io read callback — парсинг рукопожатия + релей данных
// ====================================================================
static void on_read_cb(struct ll_queue* q, void* arg) {
struct socks_proxy_conn* c = (struct socks_proxy_conn*)arg;
struct ll_entry* e = queue_data_get(q);
if (!e) { queue_resume_callback(q); return; }
if (c->rem_closed) {
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) {
// накапливаем данные для рукопожатия
uint16_t space = sizeof(c->buf) - c->buf_len;
if (space < e->len) {
DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "socks_proxy: handshake buffer overflow sid=%08x", c->stream_id);
memory_pool_free(c->tc->data_pool, e->dgram); queue_entry_free(e);
tcp_conn_push_close(c->tc); queue_resume_callback(q); return;
}
memcpy(c->buf + c->buf_len, e->dgram, e->len);
c->buf_len += e->len;
memory_pool_free(c->tc->data_pool, e->dgram); queue_entry_free(e);
queue_resume_callback(q);
if (c->is_http) {
process_http_request(c);
} else {
if (c->state == SOCKS_STATE_GREETING) process_socks_greeting(c);
if (c->state == SOCKS_STATE_REQUEST) process_socks_request(c);
}
return;
}
// 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 ret=%d is_http=%d", c->stream_id, e->len, ret, c->is_http);
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_on_send_ready(c->inst, TOPO_GROUP_UTUN, c->via_node_id, ETCP_RT_ID_TCP_PROXY, &c->tx_waiter, tx_waiter_cb, c);
DEBUG_DEBUG(DEBUG_CATEGORY_SOCKET, "SOCKS PROXY BP: sid=%08x tx_buf=%u waiter_reg", c->stream_id, c->tx_len);
}
}
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->tx_buf || c->rem_closed || c->close_sent) return;
int ret = send_data(c, c->tx_buf, c->tx_len);
DEBUG_DEBUG(DEBUG_CATEGORY_SOCKET, "SOCKS PROXY WAKE: sid=%08x tx_buf=%u ret=%d is_http=%d", c->stream_id, c->tx_len, ret, c->is_http);
if (ret == 0) {
u_free(c->tx_buf); c->tx_buf = NULL; c->tx_len = 0;
queue_resume_callback(c->tc->read_queue);
} else {
etcp_router_on_send_ready(c->inst, TOPO_GROUP_UTUN, c->via_node_id, ETCP_RT_ID_TCP_PROXY, &c->tx_waiter, tx_waiter_cb, c);
}
}
// ====================================================================
// tcp_io: FIN / Error / Closed
// ====================================================================
static void on_fin_cb(struct tcp_conn* tc, void* arg) {
struct socks_proxy_conn* c = (struct socks_proxy_conn*)arg;
if (c->rem_closed || c->close_sent) return;
if (!tc->fin_local && !tc->write_buf && !tc->write_queue->head) {
DEBUG_DEBUG(DEBUG_CATEGORY_SOCKET, "socks_proxy: local FIN → relay FIN sid=%08x", c->stream_id);
send_fin(c);
if (c->fin_remote && !c->close_sent && !c->close_pending) send_close(c);
} else {
// данные ещё в write_queue, отложим FIN
DEBUG_DEBUG(DEBUG_CATEGORY_SOCKET, "socks_proxy: local FIN deferred (wq=%d wbuf=%s) sid=%08x",
tc->write_queue->count, tc->write_buf ? "y" : "n", c->stream_id);
tcp_conn_set_flushed(tc, NULL);
}
}
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 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_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);
}
// ====================================================================
// Accept callback для слушающего сокета
// ====================================================================
static void on_accept_cb(socket_t sock, void* arg) {
struct listen_ctx* ctx = (struct listen_ctx*)arg;
struct sockaddr_in addr; socklen_t alen = sizeof(addr);
socket_t csock = accept(sock, (struct sockaddr*)&addr, &alen);
if (csock == SOCKET_INVALID) { DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "socks_proxy: accept failed errno=%d", errno); return; }
socket_set_nonblocking(csock);
struct socks_proxy_conn* c = u_calloc(1, sizeof(struct socks_proxy_conn));
if (!c) { socket_close_wrapper(csock); DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "socks_proxy: u_calloc failed"); return; }
c->stream_id = ++(*ctx->next_stream_id);
c->is_http = (uint8_t)ctx->is_http;
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;
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_SOCKET, "socks_proxy: tcp_conn_create failed"); u_free(c); socket_close_wrapper(csock); return; }
c->tc->on_closed = on_closed_cb;
queue_set_callback(c->tc->read_queue, on_read_cb, c);
queue_set_waiter_defer(c->tc->read_queue, 1);
c->next = *ctx->conns; *ctx->conns = c; (*ctx->conn_count)++;
char ip[INET_ADDRSTRLEN]; inet_ntop(AF_INET, &addr.sin_addr, ip, sizeof(ip));
DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "socks_proxy: accepted %s conn from %s:%d sid=%08x total=%d",
ctx->is_http ? "http" : "socks", ip, ntohs(addr.sin_port), c->stream_id, *ctx->conn_count);
}
// ====================================================================
// Публичное API для tcp_proxy_client (etcp dispatch)
// ====================================================================
struct socks_proxy_conn* socks_proxy_find_conn(struct socks_proxy_conn* head, uint32_t stream_id) {
struct socks_proxy_conn* c;
for (c = head; c; c = c->next) if (c->stream_id == stream_id) return c;
return NULL;
}
int socks_proxy_handle_etcp(struct socks_proxy_conn** head, int* count,
uint32_t stream_id, uint8_t subcmd,
const uint8_t* data, size_t data_len) {
struct socks_proxy_conn* c = socks_proxy_find_conn(*head, stream_id);
if (!c) return 0;
if (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",
data_len, c->tc->data_pool->object_size, stream_id);
return 1;
}
struct ll_entry* e = queue_entry_new_from_pool(c->tc->entry_pool);
uint8_t* buf = memory_pool_alloc(c->tc->data_pool);
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 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);
}
}
return 1;
}
if (subcmd == TCP_PROXY_SUBCMD_CLOSE) {
DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "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, &c->tx_waiter);
tcp_conn_push_close(c->tc);
return 1;
}
if (subcmd == TCP_PROXY_SUBCMD_ERROR) {
DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "socks_proxy: ERROR from exit 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, &c->tx_waiter);
tcp_conn_push_close(c->tc);
return 1;
}
if (subcmd == TCP_PROXY_SUBCMD_FIN) {
DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "socks_proxy: FIN from exit sid=%08x", stream_id);
c->fin_remote = 1;
tcp_conn_push_fin(c->tc);
if (c->tc->fin_local && !c->close_sent && !c->close_pending) send_close(c);
return 1;
}
return 1;
}
void socks_proxy_conn_free(struct socks_proxy_conn* c) {
if (!c) return;
if (c->freed) return;
c->freed = 1;
struct socks_proxy_conn** head = c->head; int* count = c->count;
DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "socks_proxy: FREE sid=%08x state=%d total=%d", c->stream_id, c->state, count ? *count : 0);
if (c->close_pending && c->inst) {
c->close_pending = 0;
uint8_t subcmd = c->tc && c->tc->error ? TCP_PROXY_SUBCMD_ERROR : TCP_PROXY_SUBCMD_CLOSE;
send_msg(c->inst, TOPO_GROUP_UTUN, 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, TOPO_GROUP_UTUN, c->via_node_id, &c->tx_waiter);
u_free(c);
}
void socks_proxy_conn_free_all(struct socks_proxy_conn** head, int* count) {
while (*head) { struct socks_proxy_conn* next = (*head)->next; socks_proxy_conn_free(*head); }
}
// ====================================================================
// Слушающий сокет
// ====================================================================
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, ':');
if (!colon) { DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "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;
memcpy(ip, addr_str, ip_len); ip[ip_len] = '\0';
int port = atoi(colon + 1);
if (port <= 0 || port > 65535) { DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "socks_proxy: bad port %d", port); return NULL; }
struct listen_ctx* ctx = u_calloc(1, sizeof(struct listen_ctx));
if (!ctx) { DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "socks_proxy: u_calloc failed"); return NULL; }
ctx->ua = ua; ctx->inst = inst; ctx->via_node_id = via_node_id;
ctx->is_http = is_http; ctx->conns = conns; ctx->conn_count = count;
ctx->next_stream_id = next_stream_id;
socket_t sock = socket(AF_INET, SOCK_STREAM, 0);
if (sock == SOCKET_INVALID) { DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "socks_proxy: socket() failed errno=%d", errno); u_free(ctx); return NULL; }
socket_set_nonblocking(sock);
socket_set_reuseaddr(sock, 1);
struct sockaddr_in addr; memset(&addr, 0, sizeof(addr));
addr.sin_family = AF_INET; addr.sin_port = htons((uint16_t)port);
if (inet_pton(AF_INET, ip, &addr.sin_addr) != 1) { DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "socks_proxy: inet_pton(%s) failed", ip); socket_close_wrapper(sock); u_free(ctx); return NULL; }
if (bind(sock, (struct sockaddr*)&addr, sizeof(addr)) < 0) { DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "socks_proxy: bind(%s:%d) failed errno=%d", ip, port, errno); socket_close_wrapper(sock); u_free(ctx); return NULL; }
if (listen(sock, 32) < 0) { DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "socks_proxy: listen() failed errno=%d", errno); socket_close_wrapper(sock); u_free(ctx); return NULL; }
ctx->listen_sock = sock;
ctx->socket_id = uasync_add_socket_t(ua, sock, on_accept_cb, NULL, NULL, ctx);
if (!ctx->socket_id) { DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "socks_proxy: uasync_add_socket_t failed"); socket_close_wrapper(sock); u_free(ctx); return NULL; }
DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "socks_proxy: %s listening on %s:%d sock=%d",
is_http ? "HTTP" : "SOCKS", ip, port, (int)sock);
return ctx;
}
void socks_proxy_close_listen(struct UASYNC* ua, struct listen_ctx* ctx, socket_t* sock_out) {
if (!ctx) return;
if (ctx->socket_id) uasync_remove_socket_t(ua, ctx->listen_sock);
socket_close_wrapper(ctx->listen_sock);
if (sock_out) *sock_out = SOCKET_INVALID;
DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "socks_proxy: listener closed sock=%d", (int)ctx->listen_sock);
u_free(ctx);
}