Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
8 changes: 0 additions & 8 deletions crates/rmcp/src/handler/client.rs
Original file line number Diff line number Diff line change
Expand Up @@ -53,14 +53,6 @@ impl<H: ClientHandler> Service<RoleClient> for H {
Ok(())
}

fn get_peer(&self) -> Option<Peer<RoleClient>> {
self.get_peer()
}

fn set_peer(&mut self, peer: Peer<RoleClient>) {
self.set_peer(peer);
}

fn get_info(&self) -> <RoleClient as ServiceRole>::Info {
self.get_info()
}
Expand Down
8 changes: 0 additions & 8 deletions crates/rmcp/src/handler/server.rs
Original file line number Diff line number Diff line change
Expand Up @@ -89,14 +89,6 @@ impl<H: ServerHandler> Service<RoleServer> for H {
Ok(())
}

fn get_peer(&self) -> Option<Peer<RoleServer>> {
self.get_peer()
}

fn set_peer(&mut self, peer: Peer<RoleServer>) {
self.set_peer(peer);
}

fn get_info(&self) -> <RoleServer as ServiceRole>::Info {
self.get_info()
}
Expand Down
54 changes: 16 additions & 38 deletions crates/rmcp/src/service.rs
Original file line number Diff line number Diff line change
Expand Up @@ -36,8 +36,10 @@ use tracing::instrument;
pub enum ServiceError {
#[error("Mcp error: {0}")]
McpError(McpError),
#[error("Transport error: {0}")]
Transport(std::io::Error),
#[error("Transport send error: {0}")]
TransportSend(Box<dyn std::error::Error + Send + Sync>),
#[error("Transport closed")]
TransportClosed,
#[error("Unexpected response type")]
UnexpectedResponse,
#[error("task cancelled for reason {}", reason.as_deref().unwrap_or("<unknown>"))]
Expand Down Expand Up @@ -98,8 +100,6 @@ pub trait Service<R: ServiceRole>: Send + Sync + 'static {
&self,
notification: R::PeerNot,
) -> impl Future<Output = Result<(), McpError>> + Send + '_;
fn get_peer(&self) -> Option<Peer<R>>;
fn set_peer(&mut self, peer: Peer<R>);
fn get_info(&self) -> R::Info;
}

Expand Down Expand Up @@ -148,14 +148,6 @@ impl<R: ServiceRole> Service<R> for Box<dyn DynService<R>> {
DynService::handle_notification(self.as_ref(), notification)
}

fn get_peer(&self) -> Option<Peer<R>> {
DynService::get_peer(self.as_ref())
}

fn set_peer(&mut self, peer: Peer<R>) {
DynService::set_peer(self.as_mut(), peer)
}

fn get_info(&self) -> R::Info {
DynService::get_info(self.as_ref())
}
Expand All @@ -168,8 +160,6 @@ pub trait DynService<R: ServiceRole>: Send + Sync {
context: RequestContext<R>,
) -> BoxFuture<Result<R::Resp, McpError>>;
fn handle_notification(&self, notification: R::PeerNot) -> BoxFuture<Result<(), McpError>>;
fn get_peer(&self) -> Option<Peer<R>>;
fn set_peer(&mut self, peer: Peer<R>);
fn get_info(&self) -> R::Info;
}

Expand All @@ -184,12 +174,6 @@ impl<R: ServiceRole, S: Service<R>> DynService<R> for S {
fn handle_notification(&self, notification: R::PeerNot) -> BoxFuture<Result<(), McpError>> {
Box::pin(self.handle_notification(notification))
}
fn get_peer(&self) -> Option<Peer<R>> {
self.get_peer()
}
fn set_peer(&mut self, peer: Peer<R>) {
self.set_peer(peer)
}
fn get_info(&self) -> R::Info {
self.get_info()
}
Expand Down Expand Up @@ -255,9 +239,7 @@ impl<R: ServiceRole> RequestHandle<R> {
pub async fn await_response(self) -> Result<R::PeerResp, ServiceError> {
if let Some(timeout) = self.options.timeout {
let timeout_result = tokio::time::timeout(timeout, async move {
self.rx
.await
.map_err(|_e| ServiceError::Transport(std::io::Error::other("disconnected")))?
self.rx.await.map_err(|_e| ServiceError::TransportClosed)?
})
.await;
match timeout_result {
Expand All @@ -278,9 +260,7 @@ impl<R: ServiceRole> RequestHandle<R> {
}
}
} else {
self.rx
.await
.map_err(|_e| ServiceError::Transport(std::io::Error::other("disconnected")))?
self.rx.await.map_err(|_e| ServiceError::TransportClosed)?
}
}

Expand Down Expand Up @@ -373,12 +353,8 @@ impl<R: ServiceRole> Peer<R> {
responder,
})
.await
.map_err(|_m| {
ServiceError::Transport(std::io::Error::other("disconnected: receiver dropped"))
})?;
receiver.await.map_err(|_e| {
ServiceError::Transport(std::io::Error::other("disconnected: responder dropped"))
})?
.map_err(|_m| ServiceError::TransportClosed)?;
receiver.await.map_err(|_e| ServiceError::TransportClosed)?
}
pub async fn send_request(&self, request: R::Req) -> Result<R::PeerResp, ServiceError> {
self.send_request_with_option(request, PeerRequestOptions::no_options())
Expand Down Expand Up @@ -416,7 +392,7 @@ impl<R: ServiceRole> Peer<R> {
responder,
})
.await
.map_err(|_m| ServiceError::Transport(std::io::Error::other("disconnected")))?;
.map_err(|_m| ServiceError::TransportClosed)?;
Ok(RequestHandle {
id,
rx: receiver,
Expand All @@ -428,6 +404,10 @@ impl<R: ServiceRole> Peer<R> {
pub fn peer_info(&self) -> &R::PeerInfo {
&self.info
}

pub fn is_transport_closed(&self) -> bool {
self.tx.is_closed()
}
}

#[derive(Debug)]
Expand Down Expand Up @@ -518,7 +498,7 @@ where

#[instrument(skip_all)]
async fn serve_inner<R, S, T, E, A>(
mut service: S,
service: S,
transport: T,
peer: Peer<R>,
mut peer_rx: tokio::sync::mpsc::Receiver<PeerSinkMessage<R>>,
Expand All @@ -540,7 +520,6 @@ where
tracing::info!(?peer_info, "Service initialized as server");
}

service.set_peer(peer.clone());
let mut local_responder_pool =
HashMap::<RequestId, Responder<Result<R::PeerResp, ServiceError>>>::new();
let mut local_ct_pool = HashMap::<RequestId, CancellationToken>::new();
Expand Down Expand Up @@ -631,8 +610,7 @@ where
Event::SendTaskResult(SendTaskResult::Request { id, result }) => {
if let Err(e) = result {
if let Some(responder) = local_responder_pool.remove(&id) {
let _ = responder
.send(Err(ServiceError::Transport(std::io::Error::other(e))));
let _ = responder.send(Err(ServiceError::TransportSend(Box::new(e))));
}
}
}
Expand All @@ -642,7 +620,7 @@ where
cancellation_param,
}) => {
let response = if let Err(e) = result {
Err(ServiceError::Transport(std::io::Error::other(e)))
Err(ServiceError::TransportSend(Box::new(e)))
} else {
Ok(())
};
Expand Down
12 changes: 1 addition & 11 deletions crates/rmcp/src/service/tower.rs
Original file line number Diff line number Diff line change
Expand Up @@ -2,12 +2,11 @@ use std::{future::poll_fn, marker::PhantomData};

use tower_service::Service as TowerService;

use crate::service::{Peer, RequestContext, Service, ServiceRole};
use crate::service::{RequestContext, Service, ServiceRole};

pub struct TowerHandler<S, R: ServiceRole> {
pub service: S,
pub info: R::Info,
pub peer: Option<Peer<R>>,
role: PhantomData<R>,
}

Expand All @@ -17,7 +16,6 @@ impl<S, R: ServiceRole> TowerHandler<S, R> {
service,
role: PhantomData,
info,
peer: None,
}
}
}
Expand Down Expand Up @@ -48,14 +46,6 @@ where
std::future::ready(Ok(()))
}

fn get_peer(&self) -> Option<Peer<R>> {
self.peer.clone()
}

fn set_peer(&mut self, peer: Peer<R>) {
self.peer = Some(peer);
}

fn get_info(&self) -> R::Info {
self.info.clone()
}
Expand Down
Loading