Browse Source

fix: memory leaks in TCP proxy server/client

- tcp_proxy_server: close socket + conn_count-- on uasync_add_socket_t failure
- tcp_proxy_server: conn_count-- on connect() failure
- tcp_proxy_client: queue_free(to_lwip) + tcp_abort on send_connect failure
- tcp_proxy_client: free tx_buf + cancel tx_waiter in destroy loop
etcp-inflight-fix
Evgeny 4 months ago
parent
commit
3533ec3122
  1. 4
      src/proxy/tcp_proxy_client.c
  2. 41
      src/proxy/tcp_proxy_server.c

4
src/proxy/tcp_proxy_client.c

@ -229,7 +229,7 @@ static err_t tcp_proxy_client_accept_cb(void *arg, struct tcp_pcb *newpcb, err_t
pc->to_lwip = queue_new(p->ua, 0, 0, 0, "to_lwip");
if (tcp_proxy_client_send_connect(pc) != 0) { DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "TCP proxy client: send_connect failed sid=%08x to %d.%d.%d.%d:%d", pc->stream_id, pc->dest_ip[0], pc->dest_ip[1], pc->dest_ip[2], pc->dest_ip[3], ntohs(pc->dest_port)); u_free(pc); return LERR_MEM; }
if (tcp_proxy_client_send_connect(pc) != 0) { DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "TCP proxy client: send_connect failed sid=%08x to %d.%d.%d.%d:%d", pc->stream_id, pc->dest_ip[0], pc->dest_ip[1], pc->dest_ip[2], pc->dest_ip[3], ntohs(pc->dest_port)); tcp_arg(newpcb, NULL); tcp_abort(newpcb); queue_free(pc->to_lwip); u_free(pc); return LERR_MEM; }
pc->next = p->conns; p->conns = pc; p->conn_count++;
DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "TCP proxy client: new conn sid=%08x local_port=%u -> %d.%d.%d.%d:%d total_conns=%d snd_wnd=%u mss=%u",
@ -579,6 +579,8 @@ void tcp_proxy_client_destroy(struct tcp_proxy_client* p) {
while ((e = queue_data_get(pc->to_lwip))) { queue_dgram_free(e); queue_entry_free(e); }
queue_free(pc->to_lwip); pc->to_lwip = NULL;
}
if (pc->tx_buf) { u_free(pc->tx_buf); pc->tx_buf = NULL; pc->tx_len = 0; }
if (p->inst) etcp_router_waiter_cancel(p->inst, p->via_node_id, &pc->tx_waiter);
u_free(pc); pc = next;
}
p->conns = NULL;

41
src/proxy/tcp_proxy_server.c

@ -37,6 +37,8 @@ static void tcp_proxy_server_close_retry_cb(void* arg);
static void tcp_proxy_server_sock_send(struct tcp_proxy_server_conn* rc, const uint8_t* data, size_t len);
void tcp_proxy_server_conn_free(struct tcp_proxy_server_conn* rc);
// Отправка сообщения через ETCP-маршрут: выделяет queue_entry, заполняет заголовок TCP_PROXY_HDR_SIZE, вызывает etcp_route_send.
// Вызывается из всех send-функций (close/error/data).
static int tcp_proxy_server_send_msg(struct UTUN_INSTANCE* inst, uint64_t dst, uint8_t subcmd,
uint32_t sid, const uint8_t* data, size_t len) {
struct ll_entry* e = queue_entry_new(0);
@ -51,6 +53,7 @@ static int tcp_proxy_server_send_msg(struct UTUN_INSTANCE* inst, uint64_t dst, u
return etcp_route_send(inst, dst, e);
}
// Таймер-коллбэк повтора CLOSE/ERROR с экспоненциальным backoff (50..5000 tb). Запускается из send_close/send_error при неудаче etcp_route_send.
static void tcp_proxy_server_close_retry_cb(void* arg) {
struct tcp_proxy_server_conn* rc = (struct tcp_proxy_server_conn*)arg;
if (!rc) return;
@ -67,6 +70,7 @@ static void tcp_proxy_server_close_retry_cb(void* arg) {
rc->close_pending = 0;
}
// Отправка TCP_PROXY_SUBCMD_CLOSE пиру; при неудаче — close_pending=1 + retry-таймер. Вызывается из sock_read_cb (EOF), sock_retry_cb (завершение send).
static void tcp_proxy_server_send_close(struct tcp_proxy_server_conn* rc) {
struct UTUN_INSTANCE* inst = rc->ctx ? rc->ctx->inst : NULL;
if (!inst) return;
@ -78,6 +82,7 @@ static void tcp_proxy_server_send_close(struct tcp_proxy_server_conn* rc) {
}
}
// Отправка TCP_PROXY_SUBCMD_ERROR пиру; при неудаче — close_pending=1 + retry-таймер. Вызывается из sock_write_cb (ошибка connect), sock_error_cb.
static void tcp_proxy_server_send_error(struct tcp_proxy_server_conn* rc) {
struct UTUN_INSTANCE* inst = rc->ctx ? rc->ctx->inst : NULL;
if (!inst) return;
@ -89,6 +94,7 @@ static void tcp_proxy_server_send_error(struct tcp_proxy_server_conn* rc) {
}
}
// Подсчёт всех коннектов в ctx->conns (для диагностики).
static int tcp_proxy_server_conn_total(struct tcp_proxy_server_conn* rc) {
if (!rc || !rc->ctx) return 0;
int n = 0; struct tcp_proxy_server_conn* c;
@ -96,6 +102,8 @@ static int tcp_proxy_server_conn_total(struct tcp_proxy_server_conn* rc) {
return n;
}
// Неблокирующий send в TCP-сокет с обработкой EAGAIN/EWOULDBLOCK/EINTR. Возвращает: 0=done, 1=need_retry, -1=error.
// Вызывается из sock_send и sock_retry_cb.
// ====================================================================
// Неблокирующий send с EAGAIN и таймером повтора
// ====================================================================
@ -118,6 +126,7 @@ static int tcp_proxy_server_sock_try_send(struct tcp_proxy_server_conn* rc) {
return 0;
}
// Таймер-коллбэк повтора sock_try_send с backoff (10..5000 tb). При успешном завершении — shutdown(SHUT_WR) если cli_closed, send_close если sock_closed, conn_free если обе стороны закрыты.
static void tcp_proxy_server_sock_retry_cb(void* arg) {
struct tcp_proxy_server_conn* rc = (struct tcp_proxy_server_conn*)arg;
if (!rc || !rc->out_buf) { rc->out_timer = NULL; return; }
@ -137,6 +146,8 @@ static void tcp_proxy_server_sock_retry_cb(void* arg) {
}
}
// Буферизованная неблокирующая отправка данных в TCP-сокет. Если уже есть out_buf — дописывает (realloc). Иначе sock_try_send, при EAGAIN ставит retry-таймер.
// Вызывается из sock_write_cb (сброс после connect) и handle_data (приём DATA от клиента).
static void tcp_proxy_server_sock_send(struct tcp_proxy_server_conn* rc, const uint8_t* data, size_t len) {
if (!rc || !data || len == 0) return;
if (rc->out_buf) {
@ -160,6 +171,9 @@ static void tcp_proxy_server_sock_send(struct tcp_proxy_server_conn* rc, const u
}
}
// [queue] Коллбэк освобождения очереди ETCP (backpressure relief). Срабатывает когда etcp_router_waiter сигнализирует что очередь освободилась.
// Отправляет сохранённый в pause_buf буфер через ETCP, затем восстанавливает uasync-чтение из сокета.
// Цепочка: queue waiter (threshold) → etcp_router → эта функция.
static void tcp_proxy_server_pause_waiter_cb(struct ll_queue* q, void* arg) {
(void)q;
struct tcp_proxy_server_conn* rc = (struct tcp_proxy_server_conn*)arg;
@ -179,8 +193,12 @@ static void tcp_proxy_server_pause_waiter_cb(struct ll_queue* q, void* arg) {
}
// ====================================================================
// Socket callbacks
// Socket callbacks — вызываются из uasync при событиях на TCP-сокете
// ====================================================================
// Сокет-коллбэк чтения. Читает данные из TCP-сокета, отправляет через ETCP пиру (DATA).
// При ошибке ETCP-send (backpressure) — буферизует в pause_buf, снимает сокет с uasync, регистрирует waiter в очереди.
// При EOF (n=0) — sock_closed=1, send_close.
static void tcp_proxy_server_sock_read_cb(socket_t sock, void* arg) {
(void)sock; struct tcp_proxy_server_conn* rc = (struct tcp_proxy_server_conn*)arg;
if (!rc || rc->sock == SOCKET_INVALID) return;
@ -220,6 +238,7 @@ static void tcp_proxy_server_sock_read_cb(socket_t sock, void* arg) {
}
}
// Сокет-коллбэк записи. Срабатывает при завершении неблокирующего connect(). Проверяет SO_ERROR, при успехе — connected=1, сбрасывает pending_buf в сокет.
static void tcp_proxy_server_sock_write_cb(socket_t sock, void* arg) {
(void)sock; struct tcp_proxy_server_conn* rc = (struct tcp_proxy_server_conn*)arg;
if (!rc || rc->sock == SOCKET_INVALID || rc->connected) return;
@ -246,6 +265,7 @@ static void tcp_proxy_server_sock_write_cb(socket_t sock, void* arg) {
}
}
// Сокет-коллбэк ошибки. Асинхронная ошибка сокета — отправляет ERROR пиру и освобождает коннект через conn_free.
static void tcp_proxy_server_sock_error_cb(socket_t sock, void* arg) {
(void)sock; struct tcp_proxy_server_conn* rc = (struct tcp_proxy_server_conn*)arg;
if (!rc || rc->sock == SOCKET_INVALID) return;
@ -255,6 +275,8 @@ static void tcp_proxy_server_sock_error_cb(socket_t sock, void* arg) {
tcp_proxy_server_conn_free(rc);
}
// Полное освобождение коннекта: удаление из ctx->conns, закрытие сокета, отмена таймеров (out/close), освобождение буферов (pending/out/pause), отмена queue waiter.
// Вызывается из всех точек завершения/очистки коннекта.
void tcp_proxy_server_conn_free(struct tcp_proxy_server_conn* rc) {
if (!rc) return;
int total = tcp_proxy_server_conn_total(rc);
@ -278,14 +300,18 @@ void tcp_proxy_server_conn_free(struct tcp_proxy_server_conn* rc) {
}
// ====================================================================
// Public API
// Public API — входные точки из ETCP-диспетчера (tcp_proxy_client_etcp_recv_cb)
// ====================================================================
// Поиск коннекта по stream_id в цепочке ctx->conns.
struct tcp_proxy_server_conn* tcp_proxy_server_find_conn(struct tcp_proxy_server* ctx, uint32_t stream_id) {
struct tcp_proxy_server_conn* c;
for (c = ctx->conns; c; c = c->next) if (c->stream_id == stream_id) return c;
return NULL;
}
// Входная точка: обработка CONNECT от клиента. Создаёт коннект, инициирует неблокирующий connect() к dest_ip:dest_port, регистрирует сокет в uasync.
// Вызывается из ETCP-диспетчера при приёме пакета с subcmd=CONNECT.
int tcp_proxy_server_handle_connect(struct UTUN_INSTANCE* inst, struct ll_entry* entry,
uint32_t stream_id, uint64_t src_node_id) {
if (!inst || !inst->tcp_proxy_server.enabled) { if (entry) { queue_dgram_free(entry); queue_entry_free(entry); } return -1; }
@ -307,7 +333,7 @@ int tcp_proxy_server_handle_connect(struct UTUN_INSTANCE* inst, struct ll_entry*
(int)rc->sock, stream_id, dest_ip[0],dest_ip[1],dest_ip[2],dest_ip[3],ntohs(dest_port), ctx->conn_count);
socket_set_nonblocking(rc->sock);
rc->read_id = uasync_add_socket_t(rc->ua, rc->sock, tcp_proxy_server_sock_read_cb, tcp_proxy_server_sock_write_cb, tcp_proxy_server_sock_error_cb, rc);
if (!rc->read_id) { DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "TCP proxy server: uasync_add_socket_t failed"); u_free(rc); queue_dgram_free(entry); queue_entry_free(entry); return -1; }
if (!rc->read_id) { DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "TCP proxy server: uasync_add_socket_t failed"); ctx->conn_count--; socket_close_wrapper(rc->sock); u_free(rc); queue_dgram_free(entry); queue_entry_free(entry); return -1; }
struct sockaddr_in addr; memset(&addr, 0, sizeof(addr));
addr.sin_family = AF_INET; memcpy(&addr.sin_addr.s_addr, dest_ip, 4); addr.sin_port = dest_port;
@ -316,7 +342,7 @@ int tcp_proxy_server_handle_connect(struct UTUN_INSTANCE* inst, struct ll_entry*
DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "TCP proxy server: connect() to %d.%d.%d.%d:%d failed: %s",
dest_ip[0], dest_ip[1], dest_ip[2], dest_ip[3], ntohs(dest_port), strerror(errno));
tcp_proxy_server_send_msg(inst, src_node_id, TCP_PROXY_SUBCMD_ERROR, stream_id, NULL, 0);
uasync_remove_socket_t(rc->ua, rc->sock); socket_close_wrapper(rc->sock); u_free(rc);
uasync_remove_socket_t(rc->ua, rc->sock); socket_close_wrapper(rc->sock); ctx->conn_count--; u_free(rc);
queue_dgram_free(entry); queue_entry_free(entry); return -1;
}
@ -325,6 +351,8 @@ int tcp_proxy_server_handle_connect(struct UTUN_INSTANCE* inst, struct ll_entry*
return 0;
}
// Входная точка: обработка DATA от клиента. Пересылает данные в целевой TCP-сокет через sock_send; если connect не завершён — буферизует в pending_buf.
// Вызывается из ETCP-диспетчера при приёме DATA.
int tcp_proxy_server_handle_data(struct UTUN_INSTANCE* inst, struct ETCP_CONN* conn, struct ll_entry* entry, uint32_t stream_id) {
if (!inst) { queue_dgram_free(entry); queue_entry_free(entry); return -1; }
struct tcp_proxy_server* ctx = &inst->tcp_proxy_server;
@ -352,6 +380,8 @@ int tcp_proxy_server_handle_data(struct UTUN_INSTANCE* inst, struct ETCP_CONN* c
return 0;
}
// Входная точка: обработка CLOSE от клиента. cli_closed=1, shutdown(SHUT_WR) сокета, при sock_closed && !out_buf — conn_free.
// Вызывается из ETCP-диспетчера при приёме CLOSE.
void tcp_proxy_server_handle_close(struct UTUN_INSTANCE* inst, uint32_t stream_id) {
if (!inst) return;
struct tcp_proxy_server* ctx = &inst->tcp_proxy_server;
@ -372,6 +402,8 @@ void tcp_proxy_server_handle_close(struct UTUN_INSTANCE* inst, uint32_t stream_i
DEBUG_WARN(DEBUG_CATEGORY_SOCKET, "TCP proxy server: CLOSE sid=%08x — no conn, dropping", stream_id);
}
// Инициализация TCP proxy server: читает enabled из конфига, биндит tcp_proxy_client_etcp_recv_cb на ETCP_ID_TCP_PROXY.
// Если клиент не включён — инициализирует также UDP/ICMP proxy. Вызывается из utun_instance_init.
int tcp_proxy_server_init(struct UTUN_INSTANCE* inst) {
if (!inst) return -1;
struct tcp_proxy_server* ctx = &inst->tcp_proxy_server;
@ -390,6 +422,7 @@ int tcp_proxy_server_init(struct UTUN_INSTANCE* inst) {
return 0;
}
// Деинициализация: освобождает все коннекты через conn_free, сбрасывает ctx, вызывает udp/icmp proxy destroy. Вызывается из utun_instance_destroy.
void tcp_proxy_server_destroy(struct UTUN_INSTANCE* inst) {
if (!inst) return;
struct tcp_proxy_server* ctx = &inst->tcp_proxy_server;

Loading…
Cancel
Save