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

调用tokio watch Receiver的changed()方法触发DerefMut trait报错

问题:Tokio Watch Receiver的changed()方法报错需要可变引用

我写了下面的Rust函数,用来启动Warp服务器并监听停止信号:

pub async fn start_file_server(
    stop_signal: Arc<watch::Receiver<bool>>, // 用于监听停止信号的Receiver
) -> Result<(), Box<dyn std::error::Error>> {
    // 定义Warp的路由过滤器
    let filter = warp::path("ws")
        .and(warp::ws())
        .map(|ws: warp::ws::Ws| ws.on_upgrade(|socket| handle_ws_connection(socket)));

    // 启动Warp服务器
    let server = warp::serve(filter).run(([127, 0, 0, 1], 3030));

    // 生成任务监听停止信号
    let stop_signal_task = tokio::spawn(async move {
        // 等待停止信号变更
        let _ = stop_signal.changed().await;
        println!("Stop signal received, shutting down...");
    });

    // 使用tokio::select!等待服务器停止或收到停止信号
    tokio::select! {
        _ = server => {
            println!("Server stopped.");
        },
        _ = stop_signal_task => {
            // 信号变更时触发
            println!("Received stop signal, exiting...");
        }
    }

    Ok(())
}

调用这个函数的简化逻辑是:

let (stop_tx, stop_rx) = watch::channel(false);
let stop_signal = Arc::new(stop_rx);

let file_server_task = {
        let stop_signal = Arc::clone(&stop_signal);

        tokio::spawn(async move {
            if let Err(e) = start_file_server(stop_signal).await {
                eprintln!("Failed to start file server: {}", e);
            }
        })
};

执行时,let _ = stop_signal.changed().await;这一行报错:

cannot borrow data in an `Arc` as mutable
trait `DerefMut` is required to modify through a dereference, but it is not implemented for `Arc<tokio::sync::watch::Receiver<bool>>` 

我觉得changed()并没有修改数据,不需要可变引用,而且就算不用Arc直接传stop_rx,错误依然存在,这让我很困惑。


解决方案

问题核心在于:Tokio的watch::Receiver的changed()方法内部需要更新自身状态(比如记录最新的版本号),所以它的签名实际上要求可变引用:

pub async fn changed(&mut self) -> Result<(), RecvError>

不管有没有用Arc,调用changed()都需要拿到Receiver的可变引用,下面是两种可行的解决方式:

方式一:用Arc<Mutex<Receiver<bool>>>包裹

通过Mutex实现对Receiver的可变访问:

// 修改函数参数
pub async fn start_file_server(
    stop_signal: Arc<tokio::sync::Mutex<watch::Receiver<bool>>>,
) -> Result<(), Box<dyn std::error::Error>> {
    // ... 其他代码保持不变 ...

    let stop_signal_task = tokio::spawn(async move {
        // 锁定Mutex获取可变引用
        let mut rx = stop_signal.lock().await;
        let _ = rx.changed().await;
        println!("Stop signal received, shutting down...");
    });

    // ... 其他代码保持不变 ...
}

// 调用处修改
let (stop_tx, stop_rx) = watch::channel(false);
let stop_signal = Arc::new(tokio::sync::Mutex::new(stop_rx));

let file_server_task = {
        let stop_signal = Arc::clone(&stop_signal);

        tokio::spawn(async move {
            if let Err(e) = start_file_server(stop_signal).await {
                eprintln!("Failed to start file server: {}", e);
            }
        })
};

方式二:使用SharedReceiver(推荐)

Tokio为多线程共享场景提供了SharedReceiver,它的changed()方法只需要共享引用,直接通过Receiver的shared()方法转换即可:

// 修改函数参数
pub async fn start_file_server(
    stop_signal: watch::SharedReceiver<bool>,
) -> Result<(), Box<dyn std::error::Error>> {
    // ... 其他代码保持不变 ...

    let stop_signal_task = tokio::spawn(async move {
        // 直接调用changed(),无需可变引用
        let _ = stop_signal.changed().await;
        println!("Stop signal received, shutting down...");
    });

    // ... 其他代码保持不变 ...
}

// 调用处修改
let (stop_tx, stop_rx) = watch::channel(false);
let stop_signal = stop_rx.shared(); // 生成可共享的Receiver实例

let file_server_task = {
        let stop_signal = stop_signal.clone(); // SharedReceiver原生支持clone
        tokio::spawn(async move {
            if let Err(e) = start_file_server(stop_signal).await {
                eprintln!("Failed to start file server: {}", e);
            }
        })
};

这种方式更简洁,是官方针对多线程共享场景推荐的用法,无需额外的Arc和Mutex包装。


内容的提问来源于stack exchange,提问作者Grayza

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.19 15:13:13