Browse Source

tcp_io: FIN/close через write_queue, отладка всех этапов закрытия

- tcp_io.h/c: push_fin/push_close через очередь отправки
  Кодирование: dgram=NULL→close, len=0+dgram!=NULL→FIN (sentinel)
  FIN: shutdown(SHUT_WR), стоп чтения, on_fin_sent
  CLOSE: close сокета, on_closed
  read_cb: отбрасывание входящих после fin_sent
  guard wq_fetch: error||closed (не fin)
- tcp_proxy_server: адаптация handle_fin→push_fin, handle_close→push_close
  read_queue_drain_cb: discard при fin_sent
  DEBUG_INFO на всех точках закрытия
- tcp_proxy_client: send_fin, handle_fin — DEBUG_INFO
- test_tcp_io: +6 тестов (push_fin, push_close, ordering, double-push, fin+close)
etcp-inflight-fix
Evgeny 4 months ago
parent
commit
b953a6daf0
  1. 66
      lib/tcp_io.c
  2. 33
      lib/tcp_io.h
  3. 3
      src/etcp_api.h
  4. 50
      src/proxy/tcp_proxy_client.c
  5. 2
      src/proxy/tcp_proxy_client.h
  6. 89
      src/proxy/tcp_proxy_server.c
  7. 2
      src/proxy/tcp_proxy_server.h
  8. 251
      tests/test_tcp_io.c

66
lib/tcp_io.c

@ -21,6 +21,7 @@ static void resume_read_cb(struct ll_queue* q, void* arg);
static void fin_deferred_cb(struct ll_queue* q, void* arg);
static void write_queue_fetch_cb(struct ll_queue* q, void* arg);
static void flush_write_buf(struct tcp_conn* tc);
static const uint8_t tcp_fin_sentinel;
struct tcp_conn* tcp_conn_create(
struct UASYNC* ua, socket_t sock,
@ -96,8 +97,8 @@ struct tcp_conn* tcp_conn_create(
void tcp_conn_destroy(struct tcp_conn* tc) {
if (!tc) return;
DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "tcp_conn_destroy: fd=%d connected=%d error=%d fin=%d",
(int)tc->sock, tc->connected, tc->error, tc->fin);
DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "tcp_conn_destroy: fd=%d connected=%d error=%d fin=%d fin_sent=%d closed=%d",
(int)tc->sock, tc->connected, tc->error, tc->fin, tc->fin_sent, tc->closed);
if (tc->socket_id) {
uasync_remove_socket_t(tc->ua, tc->sock);
@ -114,13 +115,14 @@ void tcp_conn_destroy(struct tcp_conn* tc) {
queue_entry_free(e); queue_resume_callback(tc->read_queue);
}
while ((e = queue_data_get(tc->write_queue)) != NULL) {
if (e->dgram) memory_pool_free(tc->data_pool, e->dgram);
if (e->dgram && e->dgram != &tcp_fin_sentinel) memory_pool_free(tc->data_pool, e->dgram);
queue_entry_free(e); queue_resume_callback(tc->write_queue);
}
queue_free(tc->read_queue);
queue_free(tc->write_queue);
if (tc->write_buf) memory_pool_free(tc->data_pool, tc->write_buf);
if (tc->sock != SOCKET_INVALID) socket_close_wrapper(tc->sock);
memory_pool_destroy(tc->entry_pool);
memory_pool_destroy(tc->data_pool);
u_free(tc);
@ -148,6 +150,12 @@ static void read_cb(socket_t sock, void* arg) {
ssize_t n = recv(tc->sock, buf, tc->entry_data_size, 0);
if (n > 0) {
if (tc->fin_sent || tc->closed) {
memory_pool_free(tc->data_pool, buf);
queue_entry_free(e);
DEBUG_TRACE(DEBUG_CATEGORY_SOCKET, "tcp_io: read_cb discard %zd bytes fd=%d (fin_sent=%d closed=%d)", n, (int)tc->sock, tc->fin_sent, tc->closed);
return;
}
e->dgram = buf;
e->len = (uint16_t)n;
queue_data_put(tc->read_queue, e);
@ -221,7 +229,7 @@ static void flush_write_buf(struct tcp_conn* tc) {
static void write_queue_fetch_cb(struct ll_queue* q, void* arg) {
struct tcp_conn* tc = (struct tcp_conn*)arg;
if (tc->error || tc->fin) { DEBUG_TRACE(DEBUG_CATEGORY_SOCKET, "tcp_io: wq_fetch skip fd=%d error=%d fin=%d", (int)tc->sock, tc->error, tc->fin); return; }
if (tc->error || tc->closed) { DEBUG_TRACE(DEBUG_CATEGORY_SOCKET, "tcp_io: wq_fetch skip fd=%d error=%d closed=%d", (int)tc->sock, tc->error, tc->closed); return; }
if (tc->write_buf) {
flush_write_buf(tc);
@ -244,6 +252,29 @@ static void write_queue_fetch_cb(struct ll_queue* q, void* arg) {
return;
}
if (e->len == 0) {
if (!e->dgram) {
DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "tcp_io: CLOSE fd=%d", (int)tc->sock);
queue_entry_free(e); queue_resume_callback(q);
uasync_remove_socket_t(tc->ua, tc->sock);
socket_close_wrapper(tc->sock);
tc->socket_id = NULL; tc->sock = SOCKET_INVALID; tc->closed = 1;
if (tc->on_closed) tc->on_closed(tc, tc->arg);
} else if (e->dgram == &tcp_fin_sentinel) {
DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "tcp_io: FIN fd=%d", (int)tc->sock);
queue_entry_free(e); queue_resume_callback(q);
shutdown(tc->sock, SHUT_WR);
tc->fin_sent = 1;
uasync_set_socket_read(tc->ua, tc->socket_id, 0); tc->read_paused = 1;
queue_waiter_cancel(tc->read_queue, &tc->read_waiter);
if (tc->on_fin_sent) tc->on_fin_sent(tc, tc->arg);
} else {
DEBUG_WARN(DEBUG_CATEGORY_SOCKET, "tcp_io: wq_fetch zero-len entry with unknown dgram=%p fd=%d", (void*)e->dgram, (int)tc->sock);
queue_entry_free(e); queue_resume_callback(q);
}
return;
}
ssize_t n = send(tc->sock, e->dgram, e->len, MSG_NOSIGNAL);
if (n == (ssize_t)e->len) {
memory_pool_free(tc->data_pool, e->dgram);
@ -334,6 +365,33 @@ static void error_cb(socket_t sock, void* arg) {
if (tc->on_error) tc->on_error(tc, -1, tc->arg);
}
// ====================================================================
// FIN / Close через очередь отправки
// ====================================================================
int tcp_conn_push_fin(struct tcp_conn* tc) {
if (!tc || tc->sock == SOCKET_INVALID) { DEBUG_WARN(DEBUG_CATEGORY_SOCKET, "tcp_io: push_fin invalid tc"); return -1; }
if (tc->fin_sent) { DEBUG_WARN(DEBUG_CATEGORY_SOCKET, "tcp_io: push_fin already fin_sent fd=%d", (int)tc->sock); return -1; }
if (tc->closed) { DEBUG_WARN(DEBUG_CATEGORY_SOCKET, "tcp_io: push_fin already closed fd=%d", (int)tc->sock); return -1; }
struct ll_entry* e = queue_entry_new_from_pool(tc->entry_pool);
if (!e) { DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "tcp_io: push_fin entry_pool exhausted fd=%d", (int)tc->sock); return -1; }
e->dgram = (uint8_t*)&tcp_fin_sentinel; e->len = 0;
DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "tcp_io: push_fin fd=%d wq=%d", (int)tc->sock, tc->write_queue->count + 1);
queue_data_put(tc->write_queue, e);
return 0;
}
int tcp_conn_push_close(struct tcp_conn* tc) {
if (!tc || tc->sock == SOCKET_INVALID) { DEBUG_WARN(DEBUG_CATEGORY_SOCKET, "tcp_io: push_close invalid tc"); return -1; }
if (tc->closed) { DEBUG_WARN(DEBUG_CATEGORY_SOCKET, "tcp_io: push_close already closed fd=%d", (int)tc->sock); return -1; }
struct ll_entry* e = queue_entry_new_from_pool(tc->entry_pool);
if (!e) { DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "tcp_io: push_close entry_pool exhausted fd=%d", (int)tc->sock); return -1; }
e->dgram = NULL; e->len = 0;
DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "tcp_io: push_close fd=%d wq=%d", (int)tc->sock, tc->write_queue->count + 1);
queue_data_put(tc->write_queue, e);
return 0;
}
void tcp_conn_set_flushed(struct tcp_conn* tc, void (*on_flushed)(struct tcp_conn* tc, void* arg)) {
if (!tc) return;
tc->on_flushed = on_flushed;

33
lib/tcp_io.h

@ -20,6 +20,15 @@
// EAGAIN → write_buf (из data_pool) + EPOLLOUT ON → write_cb досылает → resume автозабора
// Всё отправлено → EPOLLOUT OFF + on_flushed (если установлен)
//
// FIN / Close через очередь отправки:
// tcp_conn_push_fin(tc) — ставит FIN в write_queue (dgram=&sentinel, len=0)
// tcp_conn_push_close(tc) — ставит CLOSE в write_queue (dgram=NULL, len=0)
// Оба сигнала обрабатываются после всех предшествующих данных в очереди:
// FIN → shutdown(SHUT_WR), стоп чтения, on_fin_sent
// CLOSE → close сокета, on_closed
// После FIN входящие данные отбрасываются (read_cb → discard).
// После CLOSE сокет закрыт, tc->closed=1, дальнейшая отправка невозможна.
//
// Connect:
// tcp_conn_create регистрирует сокет с read_cb + write_cb в uasync
// После create вызывается connect() (неблокирующий, EINPROGRESS)
@ -50,7 +59,9 @@ struct tcp_conn {
uint8_t write_monitor; // 1 = EPOLLOUT активен
uint8_t connected;
uint8_t error;
uint8_t fin; // FIN от сокета
uint8_t fin; // FIN получен от сокета (recv == 0)
uint8_t fin_sent; // FIN отправлен (shutdown SHUT_WR), входящие данные отбрасываем
uint8_t closed; // сокет полностью закрыт (close)
// Частичная отправка (из write_pool, не в очереди — досылается первой)
uint8_t* write_buf;
@ -66,9 +77,11 @@ struct tcp_conn {
struct queue_waiter_handle read_waiter;
// Коллбэки
void (*on_fin)(struct tcp_conn* tc, void* arg);
void (*on_error)(struct tcp_conn* tc, int err, void* arg);
void (*on_flushed)(struct tcp_conn* tc, void* arg); // все данные записи отправлены
void (*on_fin)(struct tcp_conn* tc, void* arg); // FIN получен от удалённой стороны (recv == 0)
void (*on_fin_sent)(struct tcp_conn* tc, void* arg); // FIN отправлен через очередь (после shutdown SHUT_WR)
void (*on_error)(struct tcp_conn* tc, int err, void* arg); // ошибка сокета
void (*on_flushed)(struct tcp_conn* tc, void* arg); // все данные записи отправлены (write_queue + write_buf пусты)
void (*on_closed)(struct tcp_conn* tc, void* arg); // сокет закрыт через очередь (после close)
void* arg;
};
@ -90,6 +103,18 @@ void tcp_conn_destroy(struct tcp_conn* tc);
// Перед push проверять порог: queue_set_threshold в tcp_conn_create (32 entries).
// При заполнении — queue_waiter_wait на освобождение.
// Поставить сигналы в очередь отправки. Все данные в очереди перед сигналом
// будут отправлены до его обработки. Кодирование: dgram=NULL → close, len=0+dgram!=NULL → FIN.
// tcp_conn_push_fin(tc) — FIN: после отправки предшествующих данных вызывает
// shutdown(SHUT_WR), останавливает чтение, вызывает on_fin_sent.
// Повторный вызов игнорируется (fin_sent уже установлен).
// tcp_conn_push_close(tc) — CLOSE: после отправки предшествующих данных закрывает
// сокет (close), вызывает on_closed.
// После close сокет удалён из uasync, tc->closed=1.
// Повторный вызов игнорируется.
int tcp_conn_push_fin(struct tcp_conn* tc);
int tcp_conn_push_close(struct tcp_conn* tc);
// Одноразовый коллбэк: вызывается когда write_queue + write_buf полностью опустели.
// После вызова сбрасывается. Установить повторно можно в любой момент.
void tcp_conn_set_flushed(struct tcp_conn* tc, void (*on_flushed)(struct tcp_conn* tc, void* arg));

3
src/etcp_api.h

@ -25,9 +25,10 @@
#define ETCP_ID_ROUTE_ENTRY 0x01 // Элемент роутинг-таблицы
#define ETCP_ID_NAT 0x02 // NAT трафик между узлами
#define ETCP_ID_SVC_ROUTE 0x03 // Маршрутизируемые сервисные пакеты (etcp_router)
#define ETCP_ID_TCP_PROXY 0x04 // TCP proxy через удаленный узел (tcp_proxy_server)
#define ETCP_ID_TCP_PROXY 0x04 // TCP proxy exit node (server)
#define ETCP_ID_UDP_PROXY 0x05 // UDP datagram прокси (client ↔ exit)
#define ETCP_ID_ICMP_PROXY 0x06 // ICMP echo прокси (ping через exit)
#define ETCP_ID_TCP_PROXY_CLIENT 0x07 // TCP proxy client (клиент)
// Forward declarations
struct ETCP_CONN;

50
src/proxy/tcp_proxy_client.c

@ -113,7 +113,7 @@ static int tcp_proxy_client_send_error(struct tcp_proxy_client_conn* pc) {
}
static int tcp_proxy_client_send_fin(struct tcp_proxy_client_conn* pc) {
DEBUG_TRACE(DEBUG_CATEGORY_TRAFFIC, "PROXY FIN RELAY sid=%08x", pc->stream_id);
DEBUG_INFO(DEBUG_CATEGORY_TRAFFIC, "PROXY FIN RELAY sid=%08x", pc->stream_id);
return tcp_proxy_client_send_msg(pc->proxy->inst, pc->proxy->via_node_id,
TCP_PROXY_SUBCMD_FIN, pc->stream_id, NULL, 0, 1);
}
@ -452,35 +452,27 @@ static void tcp_proxy_client_handle_error(struct tcp_proxy_client* p, uint32_t s
static void tcp_proxy_client_handle_fin(struct tcp_proxy_client* p, uint32_t stream_id) {
struct tcp_proxy_client_conn* pc = tcp_proxy_client_find_conn(p, stream_id);
if (!pc || !pc->pcb) return;
DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "PROXY FIN FROM exit sid=%08x — shutdown read", stream_id);
tcp_shutdown(pc->pcb, 0, 1);
}
// ====================================================================
// Единый обработчик etcp_router (диспетчеризация в tcp_proxy_client или tcp_proxy_server)
// ETCP коллбэк клиентской стороны (принимает DATA/CLOSE/ERROR/FIN от сервера)
// ====================================================================
void tcp_proxy_client_etcp_recv_cb(struct ETCP_CONN* conn, struct ll_entry* entry) {
void tcp_proxy_client_router_recv_cb(struct ETCP_CONN* conn, struct ll_entry* entry) {
struct UTUN_INSTANCE* inst = conn ? conn->instance : NULL;
struct tcp_proxy_client* proxy = inst ? inst->tcp_proxy_client : NULL;
if (!entry || !entry->dgram || entry->len < TCP_PROXY_HDR_SIZE) {
// conn==NULL + len==9: close/restart уведомление от роутера
if (entry && entry->len == 9 && !conn && entry->dgram[0] == ETCP_ID_TCP_PROXY) {
if (entry && entry->len == 9 && !conn && entry->dgram[0] == ETCP_ID_TCP_PROXY_CLIENT) {
uint64_t peer_id;
memcpy(&peer_id, entry->dgram + 1, 8);
DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "PROXY CLOSE_ALL from %016llx — clearing all conns for peer",
DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "PROXY CLOSE_ALL from %016llx — clearing client conns for peer",
(unsigned long long)peer_id);
if (proxy) {
struct tcp_proxy_client_conn *pc, *next;
for (pc = proxy->conns; pc; pc = next) { next = pc->next; tcp_proxy_client_conn_free(pc); }
proxy->conns = NULL; proxy->conn_count = 0;
inst = proxy->inst;
}
if (inst && inst->tcp_proxy_server.enabled) {
struct tcp_proxy_server_conn *rc, *next;
for (rc = inst->tcp_proxy_server.conns; rc; rc = next) {
next = rc->next;
if (rc->peer_node_id == peer_id) tcp_proxy_server_conn_free(rc);
}
}
}
if (entry) { queue_dgram_free(entry); queue_entry_free(entry); }
@ -489,28 +481,14 @@ void tcp_proxy_client_etcp_recv_cb(struct ETCP_CONN* conn, struct ll_entry* entr
uint8_t subcmd = entry->dgram[1];
uint32_t stream_id; memcpy(&stream_id, entry->dgram + 2, 4);
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);
return;
}
if (inst && inst->tcp_proxy_server.enabled) {
struct tcp_proxy_server_conn* rc = tcp_proxy_server_find_conn(&inst->tcp_proxy_server, stream_id);
if (rc) {
if (subcmd == TCP_PROXY_SUBCMD_DATA) { tcp_proxy_server_handle_data(inst, conn, entry, stream_id); return; }
if (subcmd == TCP_PROXY_SUBCMD_CLOSE) { tcp_proxy_server_handle_close(inst, stream_id); queue_dgram_free(entry); queue_entry_free(entry); return; }
if (subcmd == TCP_PROXY_SUBCMD_ERROR) { tcp_proxy_server_handle_error(inst, stream_id); queue_dgram_free(entry); queue_entry_free(entry); return; }
if (subcmd == TCP_PROXY_SUBCMD_FIN) { tcp_proxy_server_handle_fin(inst, stream_id); queue_dgram_free(entry); queue_entry_free(entry); return; }
}
}
if (proxy) {
if (subcmd == TCP_PROXY_SUBCMD_DATA) { tcp_proxy_client_handle_data(proxy, conn, stream_id, entry); return; }
if (subcmd == TCP_PROXY_SUBCMD_CLOSE) { tcp_proxy_client_handle_close(proxy, stream_id); queue_dgram_free(entry); queue_entry_free(entry); return; }
if (subcmd == TCP_PROXY_SUBCMD_ERROR) { tcp_proxy_client_handle_error(proxy, stream_id); queue_dgram_free(entry); queue_entry_free(entry); return; }
if (subcmd == TCP_PROXY_SUBCMD_FIN) { tcp_proxy_client_handle_fin(proxy, stream_id); queue_dgram_free(entry); queue_entry_free(entry); return; }
if (subcmd == TCP_PROXY_SUBCMD_DATA) { tcp_proxy_client_handle_data(proxy, conn, stream_id, entry); return; }
if (subcmd == TCP_PROXY_SUBCMD_CLOSE) { tcp_proxy_client_handle_close(proxy, stream_id); queue_dgram_free(entry); queue_entry_free(entry); return; }
if (subcmd == TCP_PROXY_SUBCMD_ERROR) { tcp_proxy_client_handle_error(proxy, stream_id); queue_dgram_free(entry); queue_entry_free(entry); return; }
if (subcmd == TCP_PROXY_SUBCMD_FIN) { tcp_proxy_client_handle_fin(proxy, stream_id); queue_dgram_free(entry); queue_entry_free(entry); return; }
}
if (subcmd == TCP_PROXY_SUBCMD_ERROR) {
DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "PROXY ERROR sid=%08x — conn not found, silent drop", stream_id);
DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "PROXY ERROR sid=%08x — client conn not found, silent drop", stream_id);
queue_dgram_free(entry); queue_entry_free(entry); return;
}
DEBUG_WARN(DEBUG_CATEGORY_SOCKET, "TCP proxy client: unhandled subcmd=%02x sid=%08x", subcmd, stream_id);
@ -563,10 +541,10 @@ struct tcp_proxy_client* tcp_proxy_client_create(struct UTUN_INSTANCE* inst, str
}
if (inst) {
if (etcp_router_bind(inst, ETCP_ID_TCP_PROXY, tcp_proxy_client_etcp_recv_cb) != 0) {
if (etcp_router_bind(inst, ETCP_ID_TCP_PROXY_CLIENT, tcp_proxy_client_router_recv_cb) != 0) {
DEBUG_WARN(DEBUG_CATEGORY_SOCKET, "TCP proxy client: etcp_router_bind failed");
} else {
DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "TCP proxy client: etcp_router bind registered for ID=0x%02x", ETCP_ID_TCP_PROXY);
DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "TCP proxy client: etcp_router bind registered for ID=0x%02x", ETCP_ID_TCP_PROXY_CLIENT);
udp_proxy_init(inst, ua);
icmp_proxy_init(inst, ua);
}
@ -579,7 +557,7 @@ struct tcp_proxy_client* tcp_proxy_client_create(struct UTUN_INSTANCE* inst, str
void tcp_proxy_client_destroy(struct tcp_proxy_client* p) {
if (!p) return;
DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "TCP proxy client destroying: conns=%d", p->conn_count);
if (p->inst) etcp_router_unbind(p->inst, ETCP_ID_TCP_PROXY);
if (p->inst) etcp_router_unbind(p->inst, ETCP_ID_TCP_PROXY_CLIENT);
udp_proxy_destroy(p->inst);
icmp_proxy_destroy(p->inst);

2
src/proxy/tcp_proxy_client.h

@ -63,6 +63,6 @@ struct tcp_proxy_client* tcp_proxy_client_create(struct UTUN_INSTANCE* inst, str
struct tcp_proxy_client_mapping_config* mappings, int mapping_count,
uint64_t via_node_id);
void tcp_proxy_client_destroy(struct tcp_proxy_client* p);
void tcp_proxy_client_etcp_recv_cb(struct ETCP_CONN* conn, struct ll_entry* entry);
void tcp_proxy_client_router_recv_cb(struct ETCP_CONN* conn, struct ll_entry* entry);
#endif // TCP_PROXY_CLIENT_H

89
src/proxy/tcp_proxy_server.c

@ -1,6 +1,5 @@
// tcp_proxy_server.c — TCP прокси-сервер (exit node)
#include "tcp_proxy_server.h"
#include "tcp_proxy_client.h"
#include "udp_proxy.h"
#include "icmp_proxy.h"
#include "etcp.h"
@ -29,6 +28,7 @@ static struct tcp_proxy_server* g_tcp_proxy_server_ctx = NULL;
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_flushed_cb(struct tcp_conn* tc, void* arg);
static void on_closed_cb(struct tcp_conn* tc, void* arg);
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);
@ -53,7 +53,7 @@ static int send_msg(struct UTUN_INSTANCE* inst, uint64_t dst, uint8_t subcmd,
if (!e) { DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "tcp_proxy_server: 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, "tcp_proxy_server: malloc(%zu) failed", TCP_PROXY_HDR_SIZE + len); queue_entry_free(e); return -1; }
e->dgram[0] = ETCP_ID_TCP_PROXY;
e->dgram[0] = ETCP_ID_TCP_PROXY_CLIENT;
e->dgram[1] = subcmd;
memcpy(e->dgram + 2, &sid, 4);
if (len > 0) memcpy(e->dgram + TCP_PROXY_HDR_SIZE, data, len);
@ -111,10 +111,13 @@ 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 (write_pending(tc))
int pend = write_pending(tc);
if (pend)
tcp_conn_set_flushed(tc, on_flushed_cb);
else
send_fin(rc);
DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "SOCK:DST_FIN fd=%d sid=%08x write_pend=%d → %s",
(int)tc->sock, rc->stream_id, pend, pend ? "deferred" : "relay now");
}
static void on_flushed_cb(struct tcp_conn* tc, void* arg) {
@ -123,8 +126,16 @@ static void on_flushed_cb(struct tcp_conn* tc, void* arg) {
send_fin(rc);
}
static void on_closed_cb(struct tcp_conn* tc, void* arg) {
struct tcp_proxy_server_conn* rc = (struct tcp_proxy_server_conn*)arg;
DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "SOCK:CLOSED fd=%d sid=%08x — freeing",
tc ? (int)tc->sock : -1, rc->stream_id);
tcp_proxy_server_conn_free(rc);
}
static void on_error_cb(struct tcp_conn* tc, int err, void* arg) {
(void)tc; (void)err;
(void)err;
if (tc->closed) return;
struct tcp_proxy_server_conn* rc = (struct tcp_proxy_server_conn*)arg;
send_error(rc);
tcp_proxy_server_conn_free(rc);
@ -166,7 +177,7 @@ static void read_queue_drain_cb(struct ll_queue* q, void* 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; }
if (rc->cli_closed) {
if (rc->cli_closed || rc->tc->fin_sent) {
do {
memory_pool_free(rc->tc->data_pool, e->dgram);
queue_entry_free(e);
@ -185,8 +196,8 @@ static void read_queue_drain_cb(struct ll_queue* q, void* arg) {
} else {
DEBUG_INFO(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,
&rc->pause_waiter, pause_resume_cb, rc);
etcp_router_on_send_ready(inst, rc->peer_node_id, ETCP_ID_TCP_PROXY_CLIENT,
&rc->pause_waiter, pause_resume_cb, rc);
}
}
@ -231,7 +242,7 @@ void tcp_proxy_server_conn_free(struct tcp_proxy_server_conn* rc) {
if (rc->tc) { tcp_conn_destroy(rc->tc); rc->tc = 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->ctx && rc->ctx->inst) etcp_router_cancel_send_ready(rc->ctx->inst, rc->peer_node_id, ETCP_ID_TCP_PROXY, &rc->pause_waiter);
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);
DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "SOCK:FREE u_free rc=%p", (void*)rc);
u_free(rc);
}
@ -268,6 +279,7 @@ int tcp_proxy_server_handle_connect(struct UTUN_INSTANCE* inst, struct ll_entry*
rc->tc = tcp_conn_create(inst->ua, sock, 1500, 8192, 32, 8, on_fin_cb, on_error_cb, rc);
if (!rc->tc) { DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "TCP proxy server: tcp_conn_create failed"); socket_close_wrapper(sock); ctx->conn_count--; u_free(rc); queue_dgram_free(entry); queue_entry_free(entry); return -1; }
rc->tc->on_closed = on_closed_cb;
queue_set_callback(rc->tc->read_queue, read_queue_drain_cb, rc);
queue_set_waiter_defer(rc->tc->read_queue, 1);
@ -335,12 +347,7 @@ void tcp_proxy_server_handle_close(struct UTUN_INSTANCE* inst, uint32_t stream_i
}
queue_resume_callback(rc->tc->read_queue);
}
if (write_pending(rc->tc)) {
tcp_conn_set_flushed(rc->tc, on_flushed_cb);
} else {
if (rc->tc->connected) shutdown(rc->tc->sock, SHUT_WR);
tcp_proxy_server_conn_free(rc);
}
if (rc->tc->connected) tcp_conn_push_close(rc->tc); else tcp_proxy_server_conn_free(rc);
}
void tcp_proxy_server_handle_error(struct UTUN_INSTANCE* inst, uint32_t stream_id) {
@ -359,7 +366,57 @@ void tcp_proxy_server_handle_fin(struct UTUN_INSTANCE* inst, uint32_t stream_id)
struct tcp_proxy_server* ctx = &inst->tcp_proxy_server;
struct tcp_proxy_server_conn* rc = tcp_proxy_server_find_conn(ctx, stream_id);
if (!rc || !rc->tc || !rc->tc->connected) return;
shutdown(rc->tc->sock, SHUT_WR);
DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "SOCK:FIN_RECV fd=%d sid=%08x — pushing FIN to wq",
(int)rc->tc->sock, stream_id);
tcp_conn_push_fin(rc->tc);
}
// ====================================================================
// ETCP коллбэк серверной стороны (принимает CONNECT/DATA/CLOSE/ERROR/FIN от клиента)
// ====================================================================
void tcp_proxy_server_recv_cb(struct ETCP_CONN* conn, struct ll_entry* entry) {
struct UTUN_INSTANCE* inst = conn ? conn->instance : NULL;
if (!entry || !entry->dgram || entry->len < TCP_PROXY_HDR_SIZE) {
if (entry && entry->len == 9 && !conn && entry->dgram[0] == ETCP_ID_TCP_PROXY) {
uint64_t peer_id;
memcpy(&peer_id, entry->dgram + 1, 8);
DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "PROXY CLOSE_ALL from %016llx — clearing server conns for peer",
(unsigned long long)peer_id);
if (inst && inst->tcp_proxy_server.enabled) {
struct tcp_proxy_server_conn *rc, *next;
for (rc = inst->tcp_proxy_server.conns; rc; rc = next) {
next = rc->next;
if (rc->peer_node_id == peer_id) tcp_proxy_server_conn_free(rc);
}
}
}
if (entry) { queue_dgram_free(entry); queue_entry_free(entry); }
return;
}
uint8_t subcmd = entry->dgram[1];
uint32_t stream_id; memcpy(&stream_id, entry->dgram + 2, 4);
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);
return;
}
if (inst && inst->tcp_proxy_server.enabled) {
struct tcp_proxy_server_conn* rc = tcp_proxy_server_find_conn(&inst->tcp_proxy_server, stream_id);
if (rc) {
if (subcmd == TCP_PROXY_SUBCMD_DATA) { tcp_proxy_server_handle_data(inst, conn, entry, stream_id); return; }
if (subcmd == TCP_PROXY_SUBCMD_CLOSE) { tcp_proxy_server_handle_close(inst, stream_id); queue_dgram_free(entry); queue_entry_free(entry); return; }
if (subcmd == TCP_PROXY_SUBCMD_ERROR) { tcp_proxy_server_handle_error(inst, stream_id); queue_dgram_free(entry); queue_entry_free(entry); return; }
if (subcmd == TCP_PROXY_SUBCMD_FIN) { tcp_proxy_server_handle_fin(inst, stream_id); queue_dgram_free(entry); queue_entry_free(entry); return; }
}
}
if (subcmd == TCP_PROXY_SUBCMD_ERROR) {
DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "PROXY ERROR sid=%08x — server conn not found, silent drop", stream_id);
queue_dgram_free(entry); queue_entry_free(entry); return;
}
DEBUG_WARN(DEBUG_CATEGORY_SOCKET, "TCP proxy server: unhandled subcmd=%02x sid=%08x", subcmd, stream_id);
queue_dgram_free(entry); queue_entry_free(entry);
}
// ====================================================================
@ -375,7 +432,7 @@ int tcp_proxy_server_init(struct UTUN_INSTANCE* inst) {
ctx->inst = inst;
if (!ctx->enabled) return 0;
g_tcp_proxy_server_ctx = ctx;
etcp_router_bind(inst, ETCP_ID_TCP_PROXY, tcp_proxy_client_etcp_recv_cb);
etcp_router_bind(inst, ETCP_ID_TCP_PROXY, tcp_proxy_server_recv_cb);
if (!inst->config->global.tcp_proxy_client_enabled) {
udp_proxy_init(inst, inst->ua);
icmp_proxy_init(inst, inst->ua);

2
src/proxy/tcp_proxy_server.h

@ -55,6 +55,8 @@ void tcp_proxy_server_destroy(struct UTUN_INSTANCE* inst);
struct tcp_proxy_server_conn* tcp_proxy_server_find_conn(struct tcp_proxy_server* ctx, uint32_t stream_id);
void tcp_proxy_server_conn_free(struct tcp_proxy_server_conn* rc);
void tcp_proxy_server_recv_cb(struct ETCP_CONN* conn, struct ll_entry* entry);
int tcp_proxy_server_handle_connect(struct UTUN_INSTANCE* inst, struct ll_entry* entry,
uint32_t stream_id, uint64_t src_node_id);
int tcp_proxy_server_handle_data(struct UTUN_INSTANCE* inst, struct ETCP_CONN* conn, struct ll_entry* entry, uint32_t stream_id);

251
tests/test_tcp_io.c

@ -28,6 +28,9 @@ static int g_fin_count = 0;
static int g_error_count = 0;
static int g_last_error = 0;
static struct tcp_conn* g_last_fin_tc = NULL;
static int g_fin_sent_count = 0;
static int g_closed_count = 0;
static int g_last_event = 0; // 1=on_fin_sent, 2=on_closed — проверка порядка
static void on_fin(struct tcp_conn* tc, void* arg) {
(void)arg;
@ -41,11 +44,24 @@ static void on_error(struct tcp_conn* tc, int err, void* arg) {
g_last_error = err;
}
static void on_fin_sent_cb(struct tcp_conn* tc, void* arg) {
(void)tc; (void)arg;
g_fin_sent_count++; g_last_event = 1;
}
static void on_closed_cb(struct tcp_conn* tc, void* arg) {
(void)tc; (void)arg;
g_closed_count++; g_last_event = 2;
}
static void reset_counters(void) {
g_fin_count = 0;
g_error_count = 0;
g_last_error = 0;
g_last_fin_tc = NULL;
g_fin_sent_count = 0;
g_closed_count = 0;
g_last_event = 0;
}
static int push_write(struct tcp_conn* tc, const uint8_t* data, size_t len) {
@ -292,6 +308,235 @@ static void test_error_callback(void) {
TEST_PASS();
}
// ====================================================================
// FIN / Close через очередь — новые тесты
// ====================================================================
static void test_push_fin(void) {
TEST_START("push_fin via write_queue");
uasync_t* ua = uasync_create();
ASSERT_TRUE(ua != NULL, "uasync_create failed");
reset_counters();
int sv[2];
ASSERT_EQ(socketpair(AF_UNIX, SOCK_STREAM, 0, sv), 0, "socketpair failed");
for (int i = 0; i < 2; i++) fcntl(sv[i], F_SETFL, fcntl(sv[i], F_GETFL, 0) | O_NONBLOCK);
struct tcp_conn* tc = tcp_conn_create(ua, sv[0], 1500, 8192, 32, 8, on_fin, on_error, NULL);
ASSERT_TRUE(tc != NULL, "tcp_conn_create failed");
tc->on_fin_sent = on_fin_sent_cb;
uint8_t data[64];
memset(data, 'D', sizeof(data));
ASSERT_EQ(push_write(tc, data, sizeof(data)), 0, "push_write failed");
ASSERT_EQ(tcp_conn_push_fin(tc), 0, "push_fin failed");
uasync_poll(ua, 10);
uint8_t recv_buf[128] = {0};
ssize_t total = 0;
while (total < (ssize_t)sizeof(data)) {
ssize_t n = recv(sv[1], recv_buf + total, sizeof(recv_buf) - total, 0);
ASSERT_TRUE(n >= 0, "recv failed on peer");
total += n;
}
ASSERT_EQ(total, (ssize_t)sizeof(data), "should receive all data before EOF");
ASSERT_TRUE(memcmp(data, recv_buf, sizeof(data)) == 0, "data mismatch");
// После всех данных — EOF (FIN)
ssize_t n = recv(sv[1], recv_buf, sizeof(recv_buf), 0);
ASSERT_EQ(n, 0, "should get EOF after FIN");
ASSERT_EQ(tc->fin_sent, 1, "tc->fin_sent not set");
ASSERT_EQ(g_fin_sent_count, 1, "on_fin_sent not called");
ASSERT_EQ(g_closed_count, 0, "on_closed should not be called");
tcp_conn_destroy(tc);
close(sv[1]);
uasync_destroy(ua, 0);
TEST_PASS();
}
static void test_push_close(void) {
TEST_START("push_close via write_queue");
uasync_t* ua = uasync_create();
ASSERT_TRUE(ua != NULL, "uasync_create failed");
reset_counters();
int sv[2];
ASSERT_EQ(socketpair(AF_UNIX, SOCK_STREAM, 0, sv), 0, "socketpair failed");
for (int i = 0; i < 2; i++) fcntl(sv[i], F_SETFL, fcntl(sv[i], F_GETFL, 0) | O_NONBLOCK);
struct tcp_conn* tc = tcp_conn_create(ua, sv[0], 1500, 8192, 32, 8, on_fin, on_error, NULL);
ASSERT_TRUE(tc != NULL, "tcp_conn_create failed");
tc->on_closed = on_closed_cb;
uint8_t data[64];
memset(data, 'E', sizeof(data));
ASSERT_EQ(push_write(tc, data, sizeof(data)), 0, "push_write failed");
ASSERT_EQ(tcp_conn_push_close(tc), 0, "push_close failed");
uasync_poll(ua, 10);
// Peer должен получить данные перед закрытием
uint8_t recv_buf[128] = {0};
ssize_t total = recv(sv[1], recv_buf, sizeof(recv_buf), 0);
ASSERT_TRUE(total == sizeof(data), "peer should receive data");
ASSERT_TRUE(memcmp(data, recv_buf, sizeof(data)) == 0, "data mismatch");
ASSERT_EQ(tc->closed, 1, "tc->closed not set");
ASSERT_EQ(tc->sock, SOCKET_INVALID, "sock should be invalid after close");
ASSERT_EQ(g_closed_count, 1, "on_closed not called");
tcp_conn_destroy(tc);
close(sv[1]);
uasync_destroy(ua, 0);
TEST_PASS();
}
static void test_fin_data_ordering(void) {
TEST_START("Data before FIN ordering");
uasync_t* ua = uasync_create();
ASSERT_TRUE(ua != NULL, "uasync_create failed");
reset_counters();
int sv[2];
ASSERT_EQ(socketpair(AF_UNIX, SOCK_STREAM, 0, sv), 0, "socketpair failed");
for (int i = 0; i < 2; i++) fcntl(sv[i], F_SETFL, fcntl(sv[i], F_GETFL, 0) | O_NONBLOCK);
struct tcp_conn* tc = tcp_conn_create(ua, sv[0], 1500, 8192, 32, 8, on_fin, on_error, NULL);
ASSERT_TRUE(tc != NULL, "tcp_conn_create failed");
tc->on_fin_sent = on_fin_sent_cb;
uint8_t data1[64], data2[32], data3[48];
memset(data1, 0x01, sizeof(data1)); memset(data2, 0x02, sizeof(data2)); memset(data3, 0x03, sizeof(data3));
ASSERT_EQ(push_write(tc, data1, sizeof(data1)), 0, "push_write data1");
ASSERT_EQ(push_write(tc, data2, sizeof(data2)), 0, "push_write data2");
ASSERT_EQ(push_write(tc, data3, sizeof(data3)), 0, "push_write data3");
ASSERT_EQ(tcp_conn_push_fin(tc), 0, "push_fin");
uasync_poll(ua, 10);
// TCP — поток, блоки могут склеиться. Читаем всё и сверяем суммарно.
uint8_t buf[256], ref[sizeof(data1)+sizeof(data2)+sizeof(data3)];
memcpy(ref, data1, sizeof(data1)); memcpy(ref + sizeof(data1), data2, sizeof(data2)); memcpy(ref + sizeof(data1) + sizeof(data2), data3, sizeof(data3));
ssize_t total = 0;
while (total < (ssize_t)sizeof(ref)) {
ssize_t n = recv(sv[1], buf + total, sizeof(buf) - total, 0);
ASSERT_TRUE(n > 0, "recv failed on peer");
total += n;
}
ASSERT_EQ(total, (ssize_t)sizeof(ref), "total received mismatch");
ASSERT_TRUE(memcmp(buf, ref, sizeof(ref)) == 0, "data content mismatch");
// EOF
ssize_t n = recv(sv[1], buf, sizeof(buf), 0);
ASSERT_EQ(n, 0, "should get EOF after data");
ASSERT_EQ(g_fin_sent_count, 1, "on_fin_sent should be called");
tcp_conn_destroy(tc);
close(sv[1]);
uasync_destroy(ua, 0);
TEST_PASS();
}
static void test_double_push_fin(void) {
TEST_START("Double push_fin ignored");
uasync_t* ua = uasync_create();
ASSERT_TRUE(ua != NULL, "uasync_create failed");
reset_counters();
int sv[2];
ASSERT_EQ(socketpair(AF_UNIX, SOCK_STREAM, 0, sv), 0, "socketpair failed");
for (int i = 0; i < 2; i++) fcntl(sv[i], F_SETFL, fcntl(sv[i], F_GETFL, 0) | O_NONBLOCK);
struct tcp_conn* tc = tcp_conn_create(ua, sv[0], 1500, 8192, 32, 8, on_fin, on_error, NULL);
ASSERT_TRUE(tc != NULL, "tcp_conn_create failed");
tc->on_fin_sent = on_fin_sent_cb;
ASSERT_EQ(tcp_conn_push_fin(tc), 0, "first push_fin");
ASSERT_EQ(tcp_conn_push_fin(tc), -1, "second push_fin should return -1");
uasync_poll(ua, 10);
ASSERT_EQ(g_fin_sent_count, 1, "on_fin_sent called more than once");
tcp_conn_destroy(tc);
close(sv[1]);
uasync_destroy(ua, 0);
TEST_PASS();
}
static void test_double_push_close(void) {
TEST_START("Double push_close ignored");
uasync_t* ua = uasync_create();
ASSERT_TRUE(ua != NULL, "uasync_create failed");
reset_counters();
int sv[2];
ASSERT_EQ(socketpair(AF_UNIX, SOCK_STREAM, 0, sv), 0, "socketpair failed");
for (int i = 0; i < 2; i++) fcntl(sv[i], F_SETFL, fcntl(sv[i], F_GETFL, 0) | O_NONBLOCK);
struct tcp_conn* tc = tcp_conn_create(ua, sv[0], 1500, 8192, 32, 8, on_fin, on_error, NULL);
ASSERT_TRUE(tc != NULL, "tcp_conn_create failed");
tc->on_closed = on_closed_cb;
ASSERT_EQ(tcp_conn_push_close(tc), 0, "first push_close");
uasync_poll(ua, 10);
ASSERT_EQ(g_closed_count, 1, "on_closed not called after first close");
ASSERT_EQ(tcp_conn_push_close(tc), -1, "second push_close should return -1");
ASSERT_EQ(g_closed_count, 1, "on_closed called more than once");
tcp_conn_destroy(tc);
close(sv[1]);
uasync_destroy(ua, 0);
TEST_PASS();
}
static void test_fin_before_close(void) {
TEST_START("FIN before close ordering");
uasync_t* ua = uasync_create();
ASSERT_TRUE(ua != NULL, "uasync_create failed");
reset_counters();
int sv[2];
ASSERT_EQ(socketpair(AF_UNIX, SOCK_STREAM, 0, sv), 0, "socketpair failed");
for (int i = 0; i < 2; i++) fcntl(sv[i], F_SETFL, fcntl(sv[i], F_GETFL, 0) | O_NONBLOCK);
struct tcp_conn* tc = tcp_conn_create(ua, sv[0], 1500, 8192, 32, 8, on_fin, on_error, NULL);
ASSERT_TRUE(tc != NULL, "tcp_conn_create failed");
tc->on_fin_sent = on_fin_sent_cb;
tc->on_closed = on_closed_cb;
uint8_t data[32];
memset(data, 'F', sizeof(data));
ASSERT_EQ(push_write(tc, data, sizeof(data)), 0, "push_write");
ASSERT_EQ(tcp_conn_push_fin(tc), 0, "push_fin");
ASSERT_EQ(tcp_conn_push_close(tc), 0, "push_close");
uasync_poll(ua, 10);
// Peer получает данные
uint8_t buf[64];
ssize_t n = recv(sv[1], buf, sizeof(buf), 0);
ASSERT_EQ(n, (ssize_t)sizeof(data), "peer should get data");
ASSERT_EQ(g_fin_sent_count, 1, "on_fin_sent not called");
ASSERT_EQ(g_closed_count, 1, "on_closed not called");
tcp_conn_destroy(tc);
close(sv[1]);
uasync_destroy(ua, 0);
TEST_PASS();
}
int main(void) {
debug_config_init();
debug_set_level(DEBUG_LEVEL_INFO);
@ -304,6 +549,12 @@ int main(void) {
test_high_water_pause();
test_connect_detection();
test_error_callback();
test_push_fin();
test_push_close();
test_fin_data_ordering();
test_double_push_fin();
test_double_push_close();
test_fin_before_close();
DEBUG_INFO(DEBUG_CATEGORY_UASYNC, "=== Test Statistics ===");
DEBUG_INFO(DEBUG_CATEGORY_UASYNC, "Tests run: %d", tests_run);

Loading…
Cancel
Save