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

Rust中如何存储返回Future的闭包列表并在线程间共享?

基于Tokio的多线程TCP事件驱动服务器类型问题解决

问题背景

尝试用Tokio编写多线程TCP服务器,采用事件驱动架构,支持为新建连接、收到消息等事件绑定多个闭包。期望的调用方式如下:

server.on_message(|msg: String, stream: &mut TcpStream| {
   async move {
        println!("Recieved {:?}", msg);
        stream.write_all(b"Hello\n").await;
   }
}).await;

服务器核心运行逻辑为每个新连接分配独立线程,并传递回调列表:

pub async fn run(&mut self) {
    let listener = TcpListener::bind("127.0.0.1:9090").await.unwrap();

    loop {
        let (mut socket, _) = listener.accept().await.unwrap();

        let cb = self.on_message.clone();
        tokio::spawn(async move {
            Self::process(socket, cb).await;
        });
    }
}

但在定义回调存储类型和函数参数类型时遇到编译错误,核心问题是跨线程安全的类型约束不满足,原实现代码如下:

type Callback<T> = dyn Fn(T, &mut TcpStream) -> Pin<Box<dyn Future<Output=()> + Send>> + Send + 'static;

unsafe impl<T> Send for TcpStreamCallbackList<T> {}
unsafe impl<T> Sync for TcpStreamCallbackList<T> {}

impl<T> TcpStreamCallbackList<T> {

    pub fn new() -> Self {
        Self { callbacks: Vec::new() }
    }

    pub fn push<G: Send + 'static>(&mut self, mut fun: impl Fn(T, &mut TcpStream) -> G + Send + 'static) where G: Future<Output=()> {
        self.callbacks.push(Arc::new(Box::new(move |val:T, stream: &mut TcpStream| Box::pin(fun(val, stream)))));
    }

    pub async fn call(&self, val: T, stream: &mut TcpStream) where T: Clone {
        for cb in self.callbacks.iter() {
            let _cb = cb.clone();
            _cb(val.clone(), stream).await; // 编译错误点
        }
    }
}

编译错误信息:

error[E0277]: `dyn for<'a> Fn(String, &'a mut tokio::net::TcpStream) -> Pin<Box<dyn futures::Future<Output = ()> + std::marker::Send>> + std::marker::Send` cannot be shared between threads safely
   --> src/main.rs:94:26

关联提示:

note: required by a bound in `tokio::spawn`
   --> /Users/lukasz/.cargo/registry/src/github.com-1ecc6299db9ec823/tokio-1.25.0/src/task/spawn.rs:163:21
    |
163 |         T: Future + Send + 'static,
    |                     ^^^^ required by this bound in `tokio::spawn`

问题分析

  1. 缺少Sync约束:Arc包裹的回调类型需要满足Sync才能跨线程安全共享,原Callback<T>仅标注了Send,未标注Sync。
  2. 生命周期不匹配:返回的Future未绑定闭包参数的生命周期,导致编译器无法确认其安全性。
  3. 不安全的Send/Sync实现:手动实现Send和Sync绕过了编译器检查,反而隐藏了真实的类型问题。

解决方案

修正类型约束和实现逻辑,让编译器自动推导线程安全性,代码如下:

use std::sync::Arc;
use tokio::net::TcpStream;
use futures::Future;
use std::pin::Pin;

// 修正回调类型:添加Sync约束,绑定Future生命周期到参数
type Callback<T> = dyn Fn(T, &mut TcpStream) -> Pin<Box<dyn Future<Output = ()> + Send + '_>> + Send + Sync + 'static;

#[derive(Clone)]
pub struct TcpStreamCallbackList<T> {
    callbacks: Vec<Arc<Callback<T>>>,
}

impl<T> TcpStreamCallbackList<T> {
    pub fn new() -> Self {
        Self { callbacks: Vec::new() }
    }

    // 明确闭包和返回Future的线程安全约束
    pub fn push<F, Fut>(&mut self, fun: F)
    where
        F: Fn(T, &mut TcpStream) -> Fut + Send + Sync + 'static,
        Fut: Future<Output = ()> + Send + 'static,
    {
        self.callbacks.push(Arc::new(move |val, stream| Box::pin(fun(val, stream))));
    }

    pub async fn call(&self, val: T, stream: &mut TcpStream)
    where
        T: Clone,
    {
        for cb in &self.callbacks {
            cb(val.clone(), stream).await;
        }
    }
}

关键调整说明

  • 添加Sync约束:Callback<T>新增Sync,确保Arc<Callback<T>>可以安全跨线程共享,满足tokio::spawn的Send要求。
  • 生命周期绑定:给返回的Future添加'_生命周期,绑定到闭包的&mut TcpStream参数,消除生命周期不匹配问题。
  • 移除不安全实现:删除手动的Send/Sync实现,通过类型约束让编译器自动验证线程安全性,避免潜在的未定义行为。
  • 优化泛型约束:在push方法中明确闭包F和返回FutureFut的Send/'static约束,确保回调可以在Tokio任务中安全执行。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.31 06:06:45