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

如何在同一个Socket中处理多个Cap'n Proto RPC服务?

在Cap'n Proto中实现多服务部署(类似Tonic的高阶库)

首先明确核心问题:Cap'n Proto的RpcSystem初始化时仅接受一个根对象,但你可以通过定义一个根服务容器来统一暴露所有子服务,而非尝试把多个RpcSystem绑定到同一个流(这会导致第一个服务独占流处理,后续服务无法响应)。以下是具体实现方案:

1. 理解Cap'n Proto RPC的核心模型

Cap'n Proto RPC基于对象引用机制:客户端先通过连接获取根对象的引用,再通过根对象的方法获取其他子服务的引用。这和HTTP/REST的路由模式完全不同,不能直接用"路由匹配"的思路处理多服务。

2. 定义根服务的Cap'n Proto Schema

先编写schema,定义根服务和所有子服务的接口,根服务提供获取子服务的方法:

@0xd1d7b9f0b5c4a7e1;

// 子服务1接口
interface SomeService {
  sayHello @0 (name :Text) -> (response :Text);
}

// 子服务2接口
interface SomeService2 {
  doSomething @0 (input :UInt32) -> (output :UInt32);
}

// 根服务:统一暴露所有子服务
interface RootService {
  getSomeService @0 () -> (service :SomeService);
  getSomeService2 @1 () -> (service :SomeService2);
}

3. 实现子服务与根服务

先完成子服务的业务逻辑实现:

use async_trait::async_trait;
use capnp::Result;

// SomeService的具体实现
#[derive(Clone)]
struct SomeServiceImpl;

#[async_trait]
impl some_service::Server for SomeServiceImpl {
    async fn say_hello(
        &self,
        params: some_service::SayHelloParams,
        mut results: some_service::SayHelloResults,
    ) -> Result<()> {
        let name = params.get()?.get_name()?;
        results.get().set_response(&format!("Hello, {}!", name));
        Ok(())
    }
}

// SomeService2的具体实现
#[derive(Clone)]
struct SomeService2Impl;

#[async_trait]
impl some_service2::Server for SomeService2Impl {
    async fn do_something(
        &self,
        params: some_service2::DoSomethingParams,
        mut results: some_service2::DoSomethingResults,
    ) -> Result<()> {
        let input = params.get()?.get_input();
        results.get().set_output(input * 2);
        Ok(())
    }
}

然后实现根服务,它持有所有子服务的实例,负责向客户端提供子服务的引用:

#[derive(Clone)]
struct RootServiceImpl {
    some_service: SomeServiceImpl,
    some_service2: SomeService2Impl,
}

#[async_trait]
impl root_service::Server for RootServiceImpl {
    async fn get_some_service(
        &self,
        _params: root_service::GetSomeServiceParams,
        mut results: root_service::GetSomeServiceResults,
    ) -> Result<()> {
        let client = capnp_rpc::new_client(self.some_service.clone());
        results.get().set_service(client.client);
        Ok(())
    }

    async fn get_some_service2(
        &self,
        _params: root_service::GetSomeService2Params,
        mut results: root_service::GetSomeService2Results,
    ) -> Result<()> {
        let client = capnp_rpc::new_client(self.some_service2.clone());
        results.get().set_service(client.client);
        Ok(())
    }
}

4. 封装类似Tonic的Server结构体

现在可以封装你的高阶Server,让它支持链式添加服务,并自动生成根服务:

use capnp_rpc::{RpcSystem, twoparty};
use std::net::SocketAddr;

struct Server {
    some_service: Option<SomeServiceImpl>,
    some_service2: Option<SomeService2Impl>,
}

impl Server {
    // 创建空服务器实例
    pub fn new() -> Self {
        Server {
            some_service: None,
            some_service2: None,
        }
    }

    // 添加SomeService
    pub fn add_service(mut self, service: SomeServiceImpl) -> Self {
        self.some_service = Some(service);
        self
    }

    // 添加SomeService2
    pub fn add_service2(mut self, service: SomeService2Impl) -> Self {
        self.some_service2 = Some(service);
        self
    }

    // 启动服务
    pub async fn serve(self, addr: SocketAddr) -> std::io::Result<()> {
        // 检查所有必要服务是否已添加
        let root_service = RootServiceImpl {
            some_service: self.some_service.expect("SomeService未添加"),
            some_service2: self.some_service2.expect("SomeService2未添加"),
        };

        let listener = async_std::net::TcpListener::bind(addr).await?;
        let mut incoming = listener.incoming();

        while let Some(stream) = incoming.next().await {
            let stream = stream?;
            let (reader, writer) = stream.split();

            // 创建网络层
            let network = twoparty::VatNetwork::new(
                reader,
                writer,
                rpc_twoparty_capnp::Side::Server,
                Default::default(),
            );

            // 以根服务为核心启动RPC系统
            let rpc_system = RpcSystem::new(
                Box::new(network),
                Some(capnp_rpc::new_client(root_service.clone()).client)
            );

            // 异步处理每个连接
            async_std::task::spawn_local(rpc_system).await?;
        }

        Ok(())
    }
}

5. 服务器使用示例

现在可以按照你期望的方式启动多服务:

use my_lib::server::Server;
use crate::services::{SomeServiceImpl, SomeService2Impl};

#[async_std::main]
async fn main() -> async_std::io::Result<()> {
    let service_1 = SomeServiceImpl;
    let service_2 = SomeService2Impl;

    Server::new()
        .add_service(service_1)
        .add_service2(service_2)
        .serve("0.0.0.0:4567".parse()?)
        .await?;

    Ok(())
}

6. 客户端调用示例

客户端需要先获取根服务,再通过根服务拿到目标子服务:

#[async_std::main]
async fn main() -> capnp::Result<()> {
    let stream = async_std::net::TcpStream::connect("127.0.0.1:4567").await?;
    let (reader, writer) = stream.split();

    let network = twoparty::VatNetwork::new(
        reader,
        writer,
        rpc_twoparty_capnp::Side::Client,
        Default::default(),
    );

    let mut rpc_system = RpcSystem::new(Box::new(network), None);
    // 获取根服务引用
    let root: root_service::Client = rpc_system.bootstrap(rpc_twoparty_capnp::Side::Server);

    async_std::task::spawn_local(rpc_system);

    // 调用SomeService的方法
    let some_service = root.get_some_service().send().await?.get()?.get_service()?;
    let response = some_service.say_hello().send().await?.get()?.get_response()?;
    println!("{}", response);

    // 调用SomeService2的方法
    let some_service2 = root.get_some_service2().send().await?.get()?.get_service()?;
    let output = some_service2.do_something().send().await?.get()?.get_output();
    println!("输出: {}", output);

    Ok(())
}

关键优化点

  • 如果要支持任意数量的子服务,可以用HashMap<TypeId, Box<dyn CapnpService>>存储服务实例,配合动态分发实现更通用的根服务。
  • 要实现运行时无关,可以抽象Listener和TaskSpawner trait,让用户可以传入tokio/async-std等不同运行时的实现。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.18 08:57:55