调用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
相关产品推荐
相关产品推荐

