Browse Source

1

etcp-inflight-fix
Evgeny 4 months ago
parent
commit
57281a975e
  1. 20
      lib/tcp_io.c
  2. 8
      lib/tcp_io.h
  3. 49
      src/proxy/tcp_proxy_client.c
  4. 3
      src/proxy/tcp_proxy_client.h
  5. 33
      src/proxy/tcp_proxy_server.c
  6. 4
      tests/test_tcp_io.c

20
lib/tcp_io.c

@ -97,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 fin_sent=%d closed=%d",
(int)tc->sock, tc->connected, tc->error, tc->fin, tc->fin_sent, tc->closed);
DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "tcp_conn_destroy: fd=%d connected=%d error=%d fin_remote=%d fin_local=%d closed=%d",
(int)tc->sock, tc->connected, tc->error, tc->fin_remote, tc->fin_local, tc->closed);
if (tc->socket_id) {
uasync_remove_socket_t(tc->ua, tc->sock);
@ -150,12 +150,6 @@ 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);
@ -167,7 +161,7 @@ static void read_cb(socket_t sock, void* arg) {
} else if (n == 0) {
memory_pool_free(tc->data_pool, buf);
queue_entry_free(e);
tc->fin = 1;
tc->fin_remote = 1;
DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "tcp_io: FIN fd=%d", (int)tc->sock);
uasync_set_socket_read(tc->ua, tc->socket_id, 0);
if (tc->read_queue->count == 0) {
@ -264,9 +258,7 @@ static void write_queue_fetch_cb(struct ll_queue* q, void* arg) {
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);
tc->fin_local = 1;
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);
@ -360,7 +352,7 @@ static void error_cb(socket_t sock, void* arg) {
struct tcp_conn* tc = (struct tcp_conn*)arg;
if (!tc || tc->sock == SOCKET_INVALID) return;
tc->error = 1;
DEBUG_TRACE(DEBUG_CATEGORY_SOCKET, "tcp_io: error_cb fd=%d connected=%d fin=%d wmon=%d", (int)tc->sock, tc->connected, tc->fin, tc->write_monitor);
DEBUG_TRACE(DEBUG_CATEGORY_SOCKET, "tcp_io: error_cb fd=%d connected=%d fin_remote=%d wmon=%d", (int)tc->sock, tc->connected, tc->fin_remote, tc->write_monitor);
DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "tcp_io: async error fd=%d", (int)tc->sock);
if (tc->on_error) tc->on_error(tc, -1, tc->arg);
}
@ -371,7 +363,7 @@ static void error_cb(socket_t sock, void* arg) {
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->fin_local) { DEBUG_WARN(DEBUG_CATEGORY_SOCKET, "tcp_io: push_fin already fin_local 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; }

8
lib/tcp_io.h

@ -59,8 +59,8 @@ struct tcp_conn {
uint8_t write_monitor; // 1 = EPOLLOUT активен
uint8_t connected;
uint8_t error;
uint8_t fin; // FIN получен от сокета (recv == 0)
uint8_t fin_sent; // FIN отправлен (shutdown SHUT_WR), входящие данные отбрасываем
uint8_t fin_remote; // FIN получен от удалённой стороны (recv == 0)
uint8_t fin_local; // FIN отправлен удалённой стороне (shutdown SHUT_WR)
uint8_t closed; // сокет полностью закрыт (close)
// Частичная отправка (из write_pool, не в очереди — досылается первой)
@ -77,8 +77,8 @@ struct tcp_conn {
struct queue_waiter_handle read_waiter;
// Коллбэки
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_fin)(struct tcp_conn* tc, void* arg); // FIN получен от удалённой стороны (fin_remote=1)
void (*on_fin_sent)(struct tcp_conn* tc, void* arg); // FIN отправлен удалённой стороне (fin_local=1)
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)

49
src/proxy/tcp_proxy_client.c

@ -173,7 +173,8 @@ static int tcp_proxy_client_handle_non_tcp(struct tcp_proxy_client* p, uint8_t*
// Помощник: передача данных из очереди to_lwip в lwIP TCP
// ====================================================================
static void tcp_proxy_client_feed_from_transport(struct tcp_proxy_client_conn *pc) {
if (!pc->to_lwip || !pc->pcb || pc->tun_closed) return;
if (!pc->to_lwip || !pc->pcb) return;
if (pc->rem_closed) return;
int sent_any = 0;
uint32_t q_pre = queue_entry_count(pc->to_lwip);
while (1) {
@ -254,8 +255,12 @@ static void tcp_proxy_client_tx_waiter_cb(struct ll_queue* q, void* arg) {
if (ret == 0) {
tcp_recved(pc->pcb, pc->tx_len);
u_free(pc->tx_buf); pc->tx_buf = NULL; pc->tx_len = 0;
if (pc->tun_closed && !pc->close_sent && !pc->close_pending)
tcp_proxy_client_send_fin(pc);
if (pc->fin_local && !pc->close_sent && !pc->close_pending) {
if (pc->fin_remote || pc->rem_closed)
tcp_proxy_client_send_close(pc);
else
tcp_proxy_client_send_fin(pc);
}
}
}
@ -264,13 +269,24 @@ static err_t tcp_proxy_client_recv_cb(void *arg, struct tcp_pcb *pcb, struct pbu
if (!pc) { if (p) pbuf_free(p); return LERR_OK; }
if (p == NULL || err != LERR_OK) {
DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "PROXY FIN sid=%08x pcb_state=%u sndbuf=%u cwnd=%u unsent=%p unacked=%p",
pc->stream_id, pcb->state, pcb->snd_buf, pcb->cwnd,
(void*)pcb->unsent, (void*)pcb->unacked);
{ uint16_t wnd_gap = TCP_WND_MAX(pcb) - pcb->rcv_wnd;
if (wnd_gap > 0) tcp_recved(pcb, wnd_gap); }
pc->tun_closed = 1;
if (!pc->tx_buf) tcp_proxy_client_send_fin(pc);
pc->fin_local = 1;
if (p == NULL && err == ERR_OK) {
DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "PROXY FIN sid=%08x pcb_state=%u sndbuf=%u cwnd=%u unsent=%p unacked=%p",
pc->stream_id, pcb->state, pcb->snd_buf, pcb->cwnd,
(void*)pcb->unsent, (void*)pcb->unacked);
{ uint16_t wnd_gap = TCP_WND_MAX(pcb) - pcb->rcv_wnd;
if (wnd_gap > 0) tcp_recved(pcb, wnd_gap); }
if (!pc->tx_buf) {
tcp_proxy_client_send_fin(pc);
if (pc->fin_remote && !pc->close_sent && !pc->close_pending)
tcp_proxy_client_send_close(pc);
}
} else {
DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "PROXY ERROR recv sid=%08x pcb_state=%u err=%d",
pc->stream_id, pcb->state, err);
pc->error = 1;
tcp_proxy_client_send_error(pc);
}
return LERR_OK;
}
@ -308,8 +324,8 @@ static err_t tcp_proxy_client_sent_cb(void *arg, struct tcp_pcb *pcb, uint16_t l
static void tcp_proxy_client_err_cb(void *arg, err_t err) {
struct tcp_proxy_client_conn *pc = (struct tcp_proxy_client_conn *)arg;
if (!pc) { DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "PROXY tcp_proxy_client_err_cb: pc=NULL err=%d", err); return; }
DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "TCP proxy client: error %d sid=%08x pcb_state=%u tun_closed=%d rem_closed=%d",
err, pc->stream_id, pc->pcb ? pc->pcb->state : 0, pc->tun_closed, pc->rem_closed);
DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "TCP proxy client: error %d sid=%08x pcb_state=%u fin_local=%d rem_closed=%d",
err, pc->stream_id, pc->pcb ? pc->pcb->state : 0, pc->fin_local, pc->rem_closed);
pc->error = 1;
if (tcp_proxy_client_send_error(pc) < 0) DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "PROXY send_error failed sid=%08x", pc->stream_id);
}
@ -417,7 +433,7 @@ static struct tcp_proxy_client_conn* tcp_proxy_client_find_conn(struct tcp_proxy
static void tcp_proxy_client_handle_data(struct tcp_proxy_client* p, struct ETCP_CONN* conn, uint32_t stream_id, struct ll_entry* entry) {
struct tcp_proxy_client_conn* pc = tcp_proxy_client_find_conn(p, stream_id);
if (!pc || pc->rem_closed || pc->error || pc->tun_closed) {
if (!pc || pc->rem_closed || pc->error) {
DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "PROXY DATA sid=%08x — no conn/closed, dropping", stream_id);
queue_dgram_free(entry); queue_entry_free(entry); return;
}
@ -434,7 +450,7 @@ static void tcp_proxy_client_handle_data(struct tcp_proxy_client* p, struct ETCP
static void tcp_proxy_client_handle_close(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) { DEBUG_WARN(DEBUG_CATEGORY_SOCKET, "PROXY CLOSE sid=%08x — no conn, dropping", stream_id); return; }
DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "PROXY REM_CLOSED sid=%08x tun_closed=%d", stream_id, pc->tun_closed);
DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "PROXY REM_CLOSED sid=%08x fin_local=%d", stream_id, pc->fin_local);
pc->rem_closed = 1;
if (pc->tx_buf) { u_free(pc->tx_buf); pc->tx_buf = NULL; pc->tx_len = 0; }
etcp_router_waiter_cancel(pc->proxy->inst, pc->proxy->via_node_id, &pc->tx_waiter);
@ -445,7 +461,7 @@ static void tcp_proxy_client_handle_close(struct tcp_proxy_client* p, uint32_t s
static void tcp_proxy_client_handle_error(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) { DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "PROXY ERROR sid=%08x — no conn, dropping", stream_id); return; }
DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "PROXY ERROR from exit sid=%08x tun_closed=%d", stream_id, pc->tun_closed);
DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "PROXY ERROR from exit sid=%08x fin_local=%d", stream_id, pc->fin_local);
tcp_proxy_client_handle_close(p, stream_id);
}
@ -453,7 +469,10 @@ static void tcp_proxy_client_handle_fin(struct tcp_proxy_client* p, uint32_t str
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);
pc->fin_remote = 1;
tcp_shutdown(pc->pcb, 0, 1);
if (pc->fin_local && !pc->close_sent && !pc->close_pending)
tcp_proxy_client_send_close(pc);
}
// ====================================================================

3
src/proxy/tcp_proxy_client.h

@ -29,7 +29,8 @@ struct tcp_proxy_client_conn {
struct ll_queue* to_lwip; // DATA от exit → lwIP (управление потоком)
uint8_t tun_closed; // сторона lwIP отправила FIN
uint8_t fin_local; // локальная сторона отправила FIN (lwIP)
uint8_t fin_remote; // exit отправил FIN (назначение закрылось)
uint8_t rem_closed; // exit отправил CLOSE
uint8_t error; // ошибка, немедленная очистка
uint8_t close_sent; // отправили CLOSE/ERROR в exit

33
src/proxy/tcp_proxy_server.c

@ -111,18 +111,31 @@ 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);
if (pend)
if (tc->fin_local) {
DEBUG_INFO(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) {
tcp_conn_set_flushed(tc, on_flushed_cb);
else
DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "SOCK:DST_FIN fd=%d sid=%08x write_pend=%d → deferred",
(int)tc->sock, rc->stream_id, pend);
} 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");
DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "SOCK:DST_FIN fd=%d sid=%08x → relay now",
(int)tc->sock, rc->stream_id);
}
}
static void on_flushed_cb(struct tcp_conn* tc, void* arg) {
struct tcp_proxy_server_conn* rc = (struct tcp_proxy_server_conn*)arg;
if (rc->cli_closed) { shutdown(tc->sock, SHUT_WR); tcp_proxy_server_conn_free(rc); return; }
if (rc->cli_closed) return;
if (tc->fin_local) {
DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "SOCK:FLUSHED 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); return;
}
send_fin(rc);
}
@ -158,7 +171,7 @@ static void diag_timer_cb(void* arg) {
"SOCK:DIAG fd=%d sid=%08x conn=%d fin=%d "
"rq=%d(%zub) wq=%d(%zub) wbuf=%s "
"ent_f=%d dat_f=%d tcp_rb=%zu tcp_sb=%zu",
(int)tc->sock, rc->stream_id, tc->connected, tc->fin,
(int)tc->sock, rc->stream_id, tc->connected, tc->fin_remote,
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",
@ -177,7 +190,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 || rc->tc->fin_sent) {
if (rc->cli_closed) {
do {
memory_pool_free(rc->tc->data_pool, e->dgram);
queue_entry_free(e);
@ -229,7 +242,7 @@ void tcp_proxy_server_conn_free(struct tcp_proxy_server_conn* rc) {
rc->freed = 1;
int total = conn_total(rc);
DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "SOCK:FREE fd=%d sid=%08x total=%d cli_closed=%d fin=%d error=%d",
rc->tc ? (int)rc->tc->sock : -1, rc->stream_id, total, rc->cli_closed, rc->tc ? rc->tc->fin : 0, rc->tc ? rc->tc->error : 0);
rc->tc ? (int)rc->tc->sock : -1, rc->stream_id, total, rc->cli_closed, rc->tc ? rc->tc->fin_remote : 0, rc->tc ? rc->tc->error : 0);
if (rc->close_pending && rc->ctx && rc->ctx->inst) {
rc->close_pending = 0;
uint8_t subcmd = (rc->tc && rc->tc->error) ? TCP_PROXY_SUBCMD_ERROR : TCP_PROXY_SUBCMD_CLOSE;
@ -335,7 +348,7 @@ void tcp_proxy_server_handle_close(struct UTUN_INSTANCE* inst, uint32_t stream_i
if (!rc) { DEBUG_WARN(DEBUG_CATEGORY_SOCKET, "TCP proxy server: CLOSE sid=%08x — no conn", stream_id); return; }
DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "SOCK:CLOSE_RECV fd=%d sid=%08x total=%d fin=%d write_pend=%d",
rc->tc ? (int)rc->tc->sock : -1, stream_id, conn_total(rc),
rc->tc ? rc->tc->fin : 0, rc->tc ? write_pending(rc->tc) : 0);
rc->tc ? rc->tc->fin_remote : 0, rc->tc ? write_pending(rc->tc) : 0);
rc->cli_closed = 1;
if (!rc->tc) { tcp_proxy_server_conn_free(rc); return; }
tcp_conn_pause_read(rc->tc);
@ -357,7 +370,7 @@ void tcp_proxy_server_handle_error(struct UTUN_INSTANCE* inst, uint32_t stream_i
if (!rc) { DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "TCP proxy server: ERROR sid=%08x — no conn", stream_id); return; }
DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "SOCK:ERROR_RECV fd=%d sid=%08x total=%d fin=%d write_pend=%d",
rc->tc ? (int)rc->tc->sock : -1, stream_id, conn_total(rc),
rc->tc ? rc->tc->fin : 0, rc->tc ? write_pending(rc->tc) : 0);
rc->tc ? rc->tc->fin_remote : 0, rc->tc ? write_pending(rc->tc) : 0);
tcp_proxy_server_handle_close(inst, stream_id);
}

4
tests/test_tcp_io.c

@ -129,7 +129,7 @@ static void test_basic_send_recv(void) {
close(sv[1]);
uasync_poll(ua, 10);
ASSERT_EQ(g_fin_count, 1, "on_fin not called");
ASSERT_EQ(tc->fin, 1, "tc->fin not set");
ASSERT_EQ(tc->fin_remote, 1, "tc->fin_remote not set");
tcp_conn_destroy(tc);
close(sv[0]);
@ -348,7 +348,7 @@ static void test_push_fin(void) {
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(tc->fin_local, 1, "tc->fin_local 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");

Loading…
Cancel
Save