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

Rust中Rocket+Capnproto出现future跨线程安全发送问题求助

问题

我用Rocket框架结合Capnproto开发,Rocket接收Web API请求后,需通过Capnproto的RPC转发到另一服务器。单独运行Capnproto示例代码正常,但集成到Rocket项目后出现future cannot be sent between threads safely错误。

相关代码

Rocket主程序

#[rocket::main]
async fn main() {   
    rocket::build()
    .mount("/API",routes![invoke])    
    .launch()
    .await.ok();
}

路由方法

#[post("/API", data = "<request>")]
async fn invoke(request: api_request::APIRequest<'_>) -> Result<Json<api_response::ApiResponse>, Json<api_response::ApiResponseError>>{
    let result = capnp_rpc::client::run_client(String::from("Hello World")).await;
    // 后续处理逻辑...
}

Capnproto客户端代码

pub async fn run_client( message: String ) -> Result<String, Box<dyn std::error::Error>> {
    let server_addr : String = "127.0.0.1:4000".to_string();
    let addr = server_addr
        .to_socket_addrs().unwrap()
        .next()
        .expect("could not parse address");       
        
    rocket::tokio::task::LocalSet::new()
        .run_until( async move {
            let stream = rocket::tokio::net::TcpStream::connect(&addr).await?;
            stream.set_nodelay(true).unwrap();
            let (reader, writer) =
                tokio_util::compat::TokioAsyncReadCompatExt::compat(stream).split();
            let rpc_network = Box::new(twoparty::VatNetwork::new(
                futures::io::BufReader::new(reader),
                futures::io::BufWriter::new(writer),
                rpc_twoparty_capnp::Side::Client,
                Default::default(),
            ));
            let mut rpc_system = RpcSystem::new(rpc_network, None);
            let hello_world: hello_world::Client =
                rpc_system.bootstrap(rpc_twoparty_capnp::Side::Server);

            rocket::tokio::task::spawn_local(rpc_system);           
           
            let mut request = hello_world.say_hello_request();
            request.get().init_request().set_name(&message[..]);

            let reply = request.send().promise.await?;           
            let reply_message  = reply.get()?.get_reply()?.get_message()?.to_str()?;
            println!("received: {}", reply_message);
            Ok(reply_message.to_string())
        }).await
}

错误信息

error: future cannot be sent between threads safely
   --> src/main.rs:92:1
    |
92  | #[post("/API", data = "<request>")]
    | ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^ future created by async block is not `Send`
    |
    = help: within `{async block@src/main.rs:92:1: 92:36}`, the trait `std::marker::Send` is not implemented for `Rc<tokio::task::local::Context>`, which is required by `{async block@src/main.rs:92:1: 92:36}: std::marker::Send`
note: future is not `Send` as this value is used across an await
   --> src\infrastructure\capnp_rpc\client.rs:284:16
    |
257 |          rocket::tokio::task::LocalSet::new()
    |          ------------------------------------ has type `LocalSet` which is not `Send`
...
284 |             }).await.unwrap()
    |                ^^^^^ await occurs here, with `rocket::tokio::task::LocalSet::new()` maybe used later
    = note: required for the cast from `Pin<Box<{async block@src/main.rs:92:1: 92:36}>>` to `Pin<Box<dyn futures::Future<Output = Outcome<rocket::Response<'_>, Status, (rocket::Data<'_>, Status)>> + std::marker::Send>>`
    = note: this error originates in the attribute macro `post` (in Nightly builds, run with -Z macro-backtrace for more info)

error: future cannot be sent between threads safely
   --> src/main.rs:92:1
    |
92  | #[post("/API", data = "<request>")]
    | ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^ future created by async block is not `Send`
    |
    = help: within `{async block@src/main.rs:92:1: 92:36}`, the trait `std::marker::Send` is not implemented for `*const ()`, which is required by `{async block@src/main.rs:92:1: 92:36}: std::marker::Send`
note: future is not `Send` as this value is used across an await
   --> src\infrastructure\capnp_rpc\client.rs:284:16
    |
257 |          rocket::tokio::task::LocalSet::new()
    |          ------------------------------------ has type `LocalSet` which is not `Send`
...
284 |             }).await.unwrap()
    |                ^^^^^ await occurs here, with `rocket::tokio::task::LocalSet::new()` maybe used later
    = note: required for the cast from `Pin<Box<{async block@src/main.rs:92:1: 92:36}>>` to `Pin<Box<dyn futures::Future<Output = Outcome<rocket::Response<'_>, Status, (rocket::Data<'_>, Status)>> + std::marker::Send>>`
    = note: this error originates in the attribute macro `post` (in Nightly builds, run with -Z macro-backtrace for more info)

error: future cannot be sent between threads safely
   --> src/main.rs:92:1
    |
92  | #[post("/API", data = "<request>")]
    | ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^ future created by async block is not `Send`
    |
    = help: the trait `std::marker::Send` is not implemented for `dyn StdError`, which is required by `{async block@src/main.rs:92:1: 92:36}: std::marker::Send`
note: future is not `Send` as it awaits another future which is not `Send`
   --> src\infrastructure\capnp_rpc\client.rs:257:10
    |
257 | /          rocket::tokio::task::LocalSet::new()
258 | |             .spawn_local( async move {
259 | |                 let stream = rocket::tokio::net::TcpStream::connect(&addr).await?;
260 | |                 stream.set_nodelay(true).unwrap();
...   |
283 | |                 Ok(reply_message.to_string())
284 | |             }).await.unwrap()
    | |______________^ await occurs here on type `tokio::task::JoinHandle<Result<std::string::String, Box<dyn StdError>>>`, which is not `Send`
    = note: required for the cast from `Pin<Box<{async block@src/main.rs:92:1: 92:36}>>` to `Pin<Box<dyn futures::Future<Output = Outcome<rocket::Response<'_>, Status, (rocket::Data<'_>, Status)>> + std::marker::Send>>`
    = note: this error originates in the attribute macro `post` (in Nightly builds, run with -Z macro-backtrace for more info)

推测是Rocket的工作线程模型和当前代码启动新线程的逻辑冲突,试过Mutex未解决,需要适配Rocket主工作线程的方案。


解决方案

核心问题是LocalSet和spawn_local创建的任务不满足Send trait,而Rocket的路由处理future必须是Send的(Rocket会在多线程环境调度这些任务)。解决思路是移除非Send相关的代码,改用Tokio原生的异步任务调度,适配Rocket的线程池。

修改后的Capnproto客户端代码

pub async fn run_client(message: String) -> Result<String, Box<dyn std::error::Error>> {
    let server_addr: String = "127.0.0.1:4000".to_string();
    let addr = server_addr
        .to_socket_addrs()?
        .next()
        .expect("could not parse address");

    // 直接建立TCP连接,无需LocalSet包装
    let stream = rocket::tokio::net::TcpStream::connect(&addr).await?;
    stream.set_nodelay(true)?;
    
    let (reader, writer) = tokio_util::compat::TokioAsyncReadCompatExt::compat(stream).split();
    let rpc_network = Box::new(twoparty::VatNetwork::new(
        futures::io::BufReader::new(reader),
        futures::io::BufWriter::new(writer),
        rpc_twoparty_capnp::Side::Client,
        Default::default(),
    ));

    let mut rpc_system = RpcSystem::new(rpc_network, None);
    let hello_world: hello_world::Client = rpc_system.bootstrap(rpc_twoparty_capnp::Side::Server);

    // 使用spawn而非spawn_local,确保任务满足Send trait
    rocket::tokio::spawn(async move {
        if let Err(e) = rpc_system.run().await {
            eprintln!("RPC system error: {}", e);
        }
    });

    let mut request = hello_world.say_hello_request();
    request.get().init_request().set_name(&message[..]);

    let reply = request.send().promise.await?;
    let reply_message = reply.get()?.get_reply()?.get_message()?.to_str()?;
    println!("received: {}", reply_message);

    Ok(reply_message.to_string())
}

关键修改点

  • 移除LocalSet:Rocket本身运行在Tokio多线程runtime上,无需额外创建LocalSet处理非Send任务
  • 替换spawn_local为spawn:spawn创建的任务默认满足Send,能在Rocket的线程池中正常调度
  • 优化错误处理:将unwrap()替换为?,避免程序panic,同时让错误正确向上传递

路由方法补充(完善返回逻辑)

#[post("/API", data = "<request>")]
async fn invoke(request: api_request::APIRequest<'_>) -> Result<Json<api_response::ApiResponse>, Json<api_response::ApiResponseError>>{
    match capnp_rpc::client::run_client(String::from("Hello World")).await {
        Ok(msg) => Ok(Json(api_response::ApiResponse {
            message: msg,
            // 填充其他响应字段
        })),
        Err(e) => Err(Json(api_response::ApiResponseError {
            error: format!("RPC调用失败: {}", e),
            // 填充其他错误字段
        }))
    }
}

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.20 23:24:52