diff --git a/lib/tcp_io.c b/lib/tcp_io.c index a97f36be..5a14f822 100644 --- a/lib/tcp_io.c +++ b/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; } diff --git a/lib/tcp_io.h b/lib/tcp_io.h index 45fbd6e1..19820684 100644 --- a/lib/tcp_io.h +++ b/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) diff --git a/src/proxy/tcp_proxy_client.c b/src/proxy/tcp_proxy_client.c index ee0a85e8..472c38e5 100644 --- a/src/proxy/tcp_proxy_client.c +++ b/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); } // ==================================================================== diff --git a/src/proxy/tcp_proxy_client.h b/src/proxy/tcp_proxy_client.h index 09a30001..8c034774 100644 --- a/src/proxy/tcp_proxy_client.h +++ b/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 diff --git a/src/proxy/tcp_proxy_server.c b/src/proxy/tcp_proxy_server.c index 6ff69530..929b3f05 100644 --- a/src/proxy/tcp_proxy_server.c +++ b/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); } diff --git a/tests/test_tcp_io.c b/tests/test_tcp_io.c index 7c389cf3..09703606 100644 --- a/tests/test_tcp_io.c +++ b/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");