From 3258a64df48e06b0ace1a897767976994a2acf59 Mon Sep 17 00:00:00 2001 From: Evgeny Date: Wed, 13 May 2026 09:02:24 +0300 Subject: [PATCH] add periodic cleanup timers for ICMP pending requests and UDP flows --- src/icmp_proxy.c | 18 ++++++++++++++++++ src/icmp_proxy.h | 1 + src/udp_proxy.c | 15 +++++++++++++++ src/udp_proxy.h | 1 + 4 files changed, 35 insertions(+) diff --git a/src/icmp_proxy.c b/src/icmp_proxy.c index 82884a16..9cab6b1a 100644 --- a/src/icmp_proxy.c +++ b/src/icmp_proxy.c @@ -26,6 +26,9 @@ static struct icmp_proxy_ctx* g_icmp_ctx = NULL; +static void req_expire_timer_cb(void* arg); +static void req_expire(struct icmp_proxy_ctx* ctx); + static struct icmp_request* req_find_by_id(struct icmp_request* head, uint16_t echo_id, uint16_t echo_seq) { struct icmp_request* r; for (r = head; r; r = r->next) if (r->echo_id == echo_id && r->echo_seq == echo_seq) return r; @@ -70,6 +73,8 @@ static int exit_send_echo(struct UTUN_INSTANCE* inst, uint64_t client_node_id, if (payload_len > 0) memcpy(r->payload, payload, r->payload_len); r->sent_tb = get_time_tb(); r->next = g_icmp_ctx->pending; g_icmp_ctx->pending = r; + if (!g_icmp_ctx->expire_timer) + g_icmp_ctx->expire_timer = uasync_set_timeout(g_icmp_ctx->ua, g_icmp_ctx->request_timeout_tb, g_icmp_ctx, req_expire_timer_cb, "icmp_expire"); return 0; } @@ -93,6 +98,8 @@ static void raw_read_cb(socket_t sock, void* arg) { struct icmp_request* r = req_find_by_id(g_icmp_ctx->pending, icmp_hdr->icmp_id, icmp_hdr->icmp_seq); if (!r) { DEBUG_WARN(DEBUG_CATEGORY_SOCKET, "icmp_proxy: unclaimed echo reply id=0x%04x seq=%u from=0x%08x", icmp_hdr->icmp_id, icmp_hdr->icmp_seq, from.sin_addr.s_addr); return; } + { struct icmp_request** rp = &g_icmp_ctx->pending; + while (*rp) { if (*rp == r) { *rp = r->next; break; } rp = &(*rp)->next; } } DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "icmp_proxy: echo reply id=0x%04x seq=%u from=0x%08x", icmp_hdr->icmp_id, icmp_hdr->icmp_seq, from.sin_addr.s_addr); @@ -116,6 +123,7 @@ static void raw_read_cb(socket_t sock, void* arg) { int ret = etcp_route_send(g_icmp_ctx->inst, r->client_node_id, e); if (ret != 0) DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "icmp_proxy: etcp_route_send reply failed: %d", ret); else DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "icmp_proxy: reply forwarded to client %016llx", (unsigned long long)r->client_node_id); + u_free(r); } // ==================================================================== @@ -272,6 +280,14 @@ void icmp_proxy_set_test_loopback(struct UTUN_INSTANCE* inst, int enabled) { if (g_icmp_ctx) g_icmp_ctx->test_loopback = enabled; } +static void req_expire_timer_cb(void* arg) { + struct icmp_proxy_ctx* ctx = (struct icmp_proxy_ctx*)arg; + req_expire(ctx); + ctx->expire_timer = NULL; + if (ctx->pending) + ctx->expire_timer = uasync_set_timeout(ctx->ua, ctx->request_timeout_tb, ctx, req_expire_timer_cb, "icmp_expire"); +} + static void req_expire(struct icmp_proxy_ctx* ctx) { uint64_t now = get_time_tb(); struct icmp_request** prev = &ctx->pending; while (*prev) { @@ -292,6 +308,7 @@ int icmp_proxy_init(struct UTUN_INSTANCE* inst, struct UASYNC* ua) { ctx->request_timeout_tb = ICMP_TIMEOUT_TB; ctx->is_exit = inst->remote_proxy.enabled; ctx->test_loopback = 0; + ctx->expire_timer = NULL; g_icmp_ctx = ctx; if (ctx->is_exit) { @@ -314,6 +331,7 @@ int icmp_proxy_init(struct UTUN_INSTANCE* inst, struct UASYNC* ua) { void icmp_proxy_destroy(struct UTUN_INSTANCE* inst) { if (!inst || !g_icmp_ctx) return; etcp_router_unbind(inst, ETCP_ID_ICMP_PROXY); + if (g_icmp_ctx->expire_timer) { uasync_cancel_timeout(g_icmp_ctx->ua, g_icmp_ctx->expire_timer); g_icmp_ctx->expire_timer = NULL; } if (g_icmp_ctx->raw_sock != SOCKET_INVALID) { if (g_icmp_ctx->raw_read_id) { uasync_remove_socket_t(g_icmp_ctx->ua, g_icmp_ctx->raw_sock); g_icmp_ctx->raw_read_id = NULL; } socket_close_wrapper(g_icmp_ctx->raw_sock); diff --git a/src/icmp_proxy.h b/src/icmp_proxy.h index e23d0478..f5704f1e 100644 --- a/src/icmp_proxy.h +++ b/src/icmp_proxy.h @@ -42,6 +42,7 @@ struct icmp_proxy_ctx { struct icmp_request* pending; uint64_t request_timeout_tb; int test_loopback; // 1 = virtual loopback when no raw socket (test only) + void* expire_timer; // periodic cleanup timer for pending requests }; int icmp_proxy_init(struct UTUN_INSTANCE* inst, struct UASYNC* ua); diff --git a/src/udp_proxy.c b/src/udp_proxy.c index 7d84743b..686d5afc 100644 --- a/src/udp_proxy.c +++ b/src/udp_proxy.c @@ -24,6 +24,9 @@ static struct udp_proxy_ctx* g_udp_ctx = NULL; +static void flow_expire_timer_cb(void* arg); +static void flow_expire(struct udp_proxy_ctx* ctx); + static struct udp_flow* flow_find(struct udp_flow* head, uint64_t client_node_id, uint32_t src_ip, uint16_t src_port, uint32_t dst_ip, uint16_t dst_port) { @@ -95,6 +98,8 @@ static void exit_handle_data(struct ETCP_CONN* conn, struct ll_entry* entry) { f->read_id = uasync_add_socket(g_udp_ctx->ua, f->sock, flow_read_cb, NULL, NULL, f); if (!f->read_id) { socket_close_wrapper(f->sock); u_free(f); goto drop; } f->next = g_udp_ctx->flows; g_udp_ctx->flows = f; g_udp_ctx->flow_count++; + if (!g_udp_ctx->expire_timer) + g_udp_ctx->expire_timer = uasync_set_timeout(g_udp_ctx->ua, g_udp_ctx->flow_timeout_tb, g_udp_ctx, flow_expire_timer_cb, "udp_expire"); } f->last_activity_tb = get_time_tb(); @@ -207,6 +212,14 @@ int udp_proxy_deliver_reply(struct UTUN_INSTANCE* inst, return 0; } +static void flow_expire_timer_cb(void* arg) { + struct udp_proxy_ctx* ctx = (struct udp_proxy_ctx*)arg; + flow_expire(ctx); + ctx->expire_timer = NULL; + if (ctx->flows) + ctx->expire_timer = uasync_set_timeout(ctx->ua, ctx->flow_timeout_tb, ctx, flow_expire_timer_cb, "udp_expire"); +} + static void flow_expire(struct udp_proxy_ctx* ctx) { uint64_t now = get_time_tb(); struct udp_flow** prev = &ctx->flows; while (*prev) { @@ -228,6 +241,7 @@ int udp_proxy_init(struct UTUN_INSTANCE* inst, struct UASYNC* ua) { if (!ctx) return -1; ctx->inst = inst; ctx->ua = ua; ctx->flow_timeout_tb = UDP_FLOW_TIMEOUT_TB; + ctx->expire_timer = NULL; ctx->is_exit = inst->remote_proxy.enabled; g_udp_ctx = ctx; @@ -240,6 +254,7 @@ int udp_proxy_init(struct UTUN_INSTANCE* inst, struct UASYNC* ua) { void udp_proxy_destroy(struct UTUN_INSTANCE* inst) { if (!inst || !g_udp_ctx) return; etcp_router_unbind(inst, ETCP_ID_UDP_PROXY); + if (g_udp_ctx->expire_timer) { uasync_cancel_timeout(g_udp_ctx->ua, g_udp_ctx->expire_timer); g_udp_ctx->expire_timer = NULL; } struct udp_flow* f = g_udp_ctx->flows; while (f) { struct udp_flow* n = f->next; if (f->read_id) { uasync_remove_socket_t(g_udp_ctx->ua, f->sock); f->read_id = NULL; } diff --git a/src/udp_proxy.h b/src/udp_proxy.h index d4cf821d..c9bdf2b2 100644 --- a/src/udp_proxy.h +++ b/src/udp_proxy.h @@ -40,6 +40,7 @@ struct udp_proxy_ctx { struct udp_flow* flows; int flow_count; uint64_t flow_timeout_tb; + void* expire_timer; // periodic cleanup timer for expired flows }; int udp_proxy_init(struct UTUN_INSTANCE* inst, struct UASYNC* ua);