Browse Source

add periodic cleanup timers for ICMP pending requests and UDP flows

congestion
Evgeny 5 months ago
parent
commit
3258a64df4
  1. 18
      src/icmp_proxy.c
  2. 1
      src/icmp_proxy.h
  3. 15
      src/udp_proxy.c
  4. 1
      src/udp_proxy.h

18
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);

1
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);

15
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; }

1
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);

Loading…
Cancel
Save