Rust循环内futures select未按预期工作的问题排查
问题
实现了一个接收UDP消息的循环,通过broadcast通道结合select!宏处理Ctrl+C触发的停机事件,预期发送停机事件时能终止循环,但实际接收第一条UDP消息后,Ctrl+C事件无法被选中,循环无法终止。
复现步骤
- 运行
cargo run启动服务端 - 在另一个终端运行
cargo run --release -- c发送UDP消息 - 在服务端终端按Ctrl+C,服务未终止
- 再次运行
cargo run --release -- c发送消息,服务才终止
相关代码
// [dependencies] // tokio = { version = "1.23", features = ["full"] } // tokio-stream = { version = "0.1" , features = ["sync"]} // socket2 = "0.4" // stream-cancel = "0.8" // async-stream = "0.3" // futures = "0.3" use std::env::args; use std::error::Error; use std::net::{Ipv4Addr, SocketAddr, SocketAddrV4}; use std::sync::Arc; use socket2::{Domain, Protocol, SockAddr, Socket, Type}; use tokio::net::UdpSocket; use tokio::sync::broadcast; use tokio::sync::broadcast::{Receiver, Sender}; const DEFAULT_PORT: u16 = 10020; const DEFAULT_MULTICAST: &str = "224.0.1.1"; #[tokio::main] async fn main() -> Result<(), Box<dyn Error>> { let mode = args().nth(1).unwrap_or_else(|| "server".to_string()); let (shutdown_tx, shutdown_rx) = broadcast::channel(1); shutdown_watcher(shutdown_tx); return if mode == "server" { run_server(shutdown_rx).await } else { run_client().await }; } async fn run_server(rx: Receiver<()>) -> Result<(), Box<dyn Error>> { let multicast = create_multicast()?; let multicast = UdpSocket::from_std(multicast)?; let multicast = Arc::new(multicast); println!("Start Discovery Server"); start_discovery_server(multicast, rx).await; Ok(()) } async fn run_client() -> Result<(), Box<dyn Error>> { let multicast = create_multicast()?; let multicast = UdpSocket::from_std(multicast)?; let multicast = Arc::new(multicast); let msg = "My message".to_string(); let addr = SocketAddrV4::new(DEFAULT_MULTICAST.parse::<Ipv4Addr>()?, DEFAULT_PORT); let len = multicast.send_to(msg.as_bytes(), &addr).await?; println!("Client Sent {len} bytes."); Ok(()) } fn create_multicast() -> Result<std::net::UdpSocket, Box<dyn Error>> { let socket = Socket::new(Domain::IPV4, Type::DGRAM, Some(Protocol::UDP)).unwrap(); socket.set_reuse_address(true)?; let multiaddr = Ipv4Addr::new(224, 0, 1, 1); socket.set_multicast_loop_v4(true)?; let addr = SocketAddr::new("0.0.0.0".parse()?, DEFAULT_PORT); socket.bind(&SockAddr::from(addr))?; let interface = socket2::InterfaceIndexOrAddress::Index(0); socket.join_multicast_v4_n(&multiaddr, &interface)?; Ok(socket.into()) } async fn read_new_message(udp: Arc<UdpSocket>) -> String { let mut buf = vec![0u8; 1024]; let result = udp.recv_from(&mut buf).await; match result { Ok((len, _addr)) => { let msg = String::from_utf8_lossy(&buf[..len]); msg.to_string() } Err(_) => "Error!".to_string(), } } async fn start_discovery_server(udp: Arc<UdpSocket>, mut shutdown: Receiver<()>) { loop { tokio::select! { msg = read_new_message(udp.clone()) => { println!("Got: {msg:?}"); } res = shutdown.recv() => { println!("Got {res:?} for shutdown"); break } else => { println!("Both channels closed"); break } } println!("loop"); } } fn shutdown_watcher(tx: Sender<()>) { tokio::spawn(async move { println!("watcher started"); let _ = tokio::signal::ctrl_c().await; let r = tx.send(()); println!("got ctrl+C {r:?}"); }); }
问题原因与修复方案
原因分析
问题核心在于broadcast::Receiver的错误处理逻辑缺失:
- broadcast通道的消息是一次性的,当接收者错过消息后,后续调用
recv()会返回Err(broadcast::RecvError::Lagged(_)),原代码未处理该错误,直接进入下一次循环,但此时Receiver的内部状态已经处于"滞后"状态,无法正确监听新的停机信号。 - 第一次UDP消息处理完成后,循环再次进入
select!,shutdown.recv()因滞后状态无法响应新的Ctrl+C信号,直到下一次UDP消息触发select!才会检查通道状态。 read_new_message忽略了recv_from的错误,可能导致socket异常后仍无限等待。
修复方案
- 正确处理
broadcast::RecvError:区分Lagged、Closed两种错误类型,遇到Lagged时继续监听新信号,而非直接退出。 - 内联UDP接收逻辑:去掉
read_new_message函数,直接在select!中处理recv_from的结果,减少异步嵌套。 - 处理socket错误:当
recv_from返回错误时,直接退出循环,避免无效等待。
修改后的核心代码
async fn start_discovery_server(udp: Arc<UdpSocket>, mut shutdown: Receiver<()>) { loop { let mut buf = vec![0u8; 1024]; tokio::select! { result = udp.recv_from(&mut buf) => { match result { Ok((len, _addr)) => { let msg = String::from_utf8_lossy(&buf[..len]); println!("Got: {msg:?}"); } Err(e) => { println!("Socket error: {e}"); break; } } } res = shutdown.recv() => { match res { Ok(_) => { println!("Got shutdown signal"); break; } Err(broadcast::RecvError::Lagged(_)) => { // 错过历史消息,继续监听新的停机信号 continue; } Err(broadcast::RecvError::Closed) => { println!("Shutdown channel closed"); break; } } } else => { println!("Both tasks completed"); break; } } } }
内容的提问来源于stack exchange,提问作者mojtab23
相关产品推荐
相关产品推荐

