Run a UDP port forwarder: listen on `listen_addr`, forward to `target`.
(
listen_addr: SocketAddr,
target: SocketAddr,
cancel: CancellationToken,
)
| 25 | |
| 26 | /// Run a UDP port forwarder: listen on `listen_addr`, forward to `target`. |
| 27 | pub async fn run_udp_forwarder( |
| 28 | listen_addr: SocketAddr, |
| 29 | target: SocketAddr, |
| 30 | cancel: CancellationToken, |
| 31 | ) { |
| 32 | let host_socket = match UdpSocket::bind(listen_addr).await { |
| 33 | Ok(s) => s, |
| 34 | Err(e) => { |
| 35 | tracing::error!("udp bind {listen_addr} failed: {e}"); |
| 36 | return; |
| 37 | } |
| 38 | }; |
| 39 | tracing::info!("udp forwarding {listen_addr} -> {target}"); |
| 40 | |
| 41 | let host_socket = Arc::new(host_socket); |
| 42 | let clients: Arc<Mutex<HashMap<SocketAddr, ClientState>>> = |
| 43 | Arc::new(Mutex::new(HashMap::new())); |
| 44 | |
| 45 | // Periodic cleanup of idle client entries |
| 46 | let cleanup_clients = clients.clone(); |
| 47 | let cleanup_cancel = cancel.child_token(); |
| 48 | tokio::spawn(async move { |
| 49 | loop { |
| 50 | tokio::select! { |
| 51 | _ = cleanup_cancel.cancelled() => break, |
| 52 | _ = tokio::time::sleep(CLEANUP_INTERVAL) => { |
| 53 | let mut map = cleanup_clients.lock().await; |
| 54 | let now = Instant::now(); |
| 55 | map.retain(|addr, entry| { |
| 56 | let alive = now.duration_since(entry.last_active) < UDP_IDLE_TIMEOUT; |
| 57 | if !alive { |
| 58 | tracing::debug!("udp client {addr} idle timeout"); |
| 59 | } |
| 60 | alive |
| 61 | // Dropping ClientState cancels the return-path task via _cancel |
| 62 | }); |
| 63 | } |
| 64 | } |
| 65 | } |
| 66 | }); |
| 67 | |
| 68 | let mut buf = vec![0u8; UDP_BUF_SIZE]; |
| 69 | |
| 70 | loop { |
| 71 | tokio::select! { |
| 72 | _ = cancel.cancelled() => break, |
| 73 | result = host_socket.recv_from(&mut buf) => { |
| 74 | let (n, client_addr) = match result { |
| 75 | Ok(v) => v, |
| 76 | Err(e) => { |
| 77 | tracing::warn!("udp recv on {listen_addr}: {e}"); |
| 78 | continue; |
| 79 | } |
| 80 | }; |
| 81 | |
| 82 | let data = &buf[..n]; |
| 83 | let mut map = clients.lock().await; |
| 84 |