You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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的错误处理逻辑缺失:

  1. broadcast通道的消息是一次性的,当接收者错过消息后,后续调用recv()会返回Err(broadcast::RecvError::Lagged(_)),原代码未处理该错误,直接进入下一次循环,但此时Receiver的内部状态已经处于"滞后"状态,无法正确监听新的停机信号。
  2. 第一次UDP消息处理完成后,循环再次进入select!,shutdown.recv()因滞后状态无法响应新的Ctrl+C信号,直到下一次UDP消息触发select!才会检查通道状态。
  3. read_new_message忽略了recv_from的错误,可能导致socket异常后仍无限等待。

修复方案

  1. 正确处理broadcast::RecvError:区分Lagged、Closed两种错误类型,遇到Lagged时继续监听新信号,而非直接退出。
  2. 内联UDP接收逻辑:去掉read_new_message函数,直接在select!中处理recv_from的结果,减少异步嵌套。
  3. 处理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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.03 09:11:31