如何在同一个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和TaskSpawnertrait,让用户可以传入tokio/async-std等不同运行时的实现。
内容的提问来源于stack exchange,提问作者al3x
相关产品推荐
相关产品推荐

