agent_client_protocol/mcp_server/
service.rs1use std::sync::{
4 Arc, Mutex,
5 atomic::{AtomicBool, Ordering},
6};
7
8use futures::{
9 channel::oneshot,
10 future::{BoxFuture, FutureExt, Shared},
11};
12use serde_json::{Map, Value};
13
14use super::McpConnectionTo;
15use crate::{
16 Error, RequestCancellation, Role,
17 schema::v1::{McpError, McpRequestId, McpServerAcpId},
18};
19
20#[derive(Debug)]
22pub struct McpRequest {
23 pub method: String,
25 pub params: Option<Map<String, Value>>,
27}
28
29#[derive(Debug)]
31pub enum McpOutcome {
32 Result(Value),
34 Error(McpError),
36}
37
38type Notify = dyn Fn(String, Option<Map<String, Value>>) -> BoxFuture<'static, Result<(), Error>>
39 + Send
40 + Sync;
41
42#[derive(Clone)]
44pub struct McpOperationCancellation {
45 state: Arc<CancellationState>,
46}
47
48struct CancellationState {
49 cancelled: AtomicBool,
50 sender: Mutex<Option<oneshot::Sender<()>>>,
51 signal: Shared<BoxFuture<'static, ()>>,
52}
53
54impl std::fmt::Debug for McpOperationCancellation {
55 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
56 f.debug_struct("McpOperationCancellation")
57 .field("cancelled", &self.is_cancelled())
58 .finish()
59 }
60}
61
62impl McpOperationCancellation {
63 pub(crate) fn new() -> Self {
64 let (sender, receiver) = oneshot::channel();
65 Self {
66 state: Arc::new(CancellationState {
67 cancelled: AtomicBool::new(false),
68 sender: Mutex::new(Some(sender)),
69 signal: receiver.map(|_| ()).boxed().shared(),
70 }),
71 }
72 }
73
74 pub(crate) fn cancel(&self) {
75 self.state.cancelled.store(true, Ordering::Release);
76 drop(
77 self.state
78 .sender
79 .lock()
80 .expect("MCP cancellation poisoned")
81 .take(),
82 );
83 }
84
85 pub async fn cancelled(&self) {
87 self.state.signal.clone().await;
88 }
89
90 #[must_use]
92 pub fn is_cancelled(&self) -> bool {
93 self.state.cancelled.load(Ordering::Acquire)
94 }
95}
96
97#[derive(Clone)]
99pub struct McpRequestContext<Counterpart: Role> {
100 server_id: McpServerAcpId,
101 request_id: McpRequestId,
102 connection: McpConnectionTo<Counterpart>,
103 metadata: Map<String, Value>,
104 cancellation: RequestCancellation,
105 operation_cancellation: McpOperationCancellation,
106 notify: Arc<Notify>,
107}
108
109impl<Counterpart: Role> std::fmt::Debug for McpRequestContext<Counterpart> {
110 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
111 f.debug_struct("McpRequestContext")
112 .field("server_id", &self.server_id)
113 .field("request_id", &self.request_id)
114 .field("metadata", &self.metadata)
115 .field("operation_cancellation", &self.operation_cancellation)
116 .finish_non_exhaustive()
117 }
118}
119
120impl<Counterpart: Role> McpRequestContext<Counterpart> {
121 pub(crate) fn new(
122 server_id: McpServerAcpId,
123 request_id: McpRequestId,
124 connection: McpConnectionTo<Counterpart>,
125 metadata: Map<String, Value>,
126 cancellation: RequestCancellation,
127 operation_cancellation: McpOperationCancellation,
128 notify: Arc<Notify>,
129 ) -> Self {
130 Self {
131 server_id,
132 request_id,
133 connection,
134 metadata,
135 cancellation,
136 operation_cancellation,
137 notify,
138 }
139 }
140
141 pub fn server_id(&self) -> &McpServerAcpId {
143 &self.server_id
144 }
145 pub fn request_id(&self) -> &McpRequestId {
147 &self.request_id
148 }
149 pub fn connection(&self) -> &McpConnectionTo<Counterpart> {
151 &self.connection
152 }
153 pub fn metadata(&self) -> &Map<String, Value> {
155 &self.metadata
156 }
157 pub fn cancellation(&self) -> &RequestCancellation {
159 &self.cancellation
160 }
161 pub fn operation_cancellation(&self) -> &McpOperationCancellation {
163 &self.operation_cancellation
164 }
165
166 pub async fn send_notification(
168 &self,
169 method: impl Into<String>,
170 params: Option<Map<String, Value>>,
171 ) -> Result<(), Error> {
172 if self.cancellation.is_cancelled() || self.operation_cancellation.is_cancelled() {
173 return Err(Error::request_cancelled());
174 }
175 (self.notify)(method.into(), params).await
176 }
177}
178
179pub trait McpService<Counterpart: Role>: Send + Sync + 'static {
181 fn execute(
187 &self,
188 request: McpRequest,
189 context: McpRequestContext<Counterpart>,
190 ) -> BoxFuture<'static, Result<McpOutcome, Error>>;
191}