Skip to main content

agent_client_protocol/mcp_server/
service.rs

1//! Reusable application services with request-scoped execution authority.
2
3use 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/// One MCP invocation against a reusable service.
21#[derive(Debug)]
22pub struct McpRequest {
23    /// MCP method, without an implicit initialize or discovery exchange.
24    pub method: String,
25    /// Named parameters, including validated request metadata.
26    pub params: Option<Map<String, Value>>,
27}
28
29/// An MCP outcome, distinct from failure of the ACP binding.
30#[derive(Debug)]
31pub enum McpOutcome {
32    /// Opaque successful MCP result.
33    Result(Value),
34    /// Unmodified MCP error, including absent/null data and extension fields.
35    Error(McpError),
36}
37
38type Notify = dyn Fn(String, Option<Map<String, Value>>) -> BoxFuture<'static, Result<(), Error>>
39    + Send
40    + Sync;
41
42/// Cancellation from the caller, provider removal, or connection shutdown.
43#[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    /// Wait until the operation loses its output authority.
86    pub async fn cancelled(&self) {
87        self.state.signal.clone().await;
88    }
89
90    /// Whether cancellation has been requested.
91    #[must_use]
92    pub fn is_cancelled(&self) -> bool {
93        self.state.cancelled.load(Ordering::Acquire)
94    }
95}
96
97/// Authority for a single operation, not for the lifetime of its service.
98#[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    /// Server identifier bound to this invocation.
142    pub fn server_id(&self) -> &McpServerAcpId {
143        &self.server_id
144    }
145    /// Logical request identifier, held until cleanup and terminal output finish.
146    pub fn request_id(&self) -> &McpRequestId {
147        &self.request_id
148    }
149    /// Host connection for application tools.
150    pub fn connection(&self) -> &McpConnectionTo<Counterpart> {
151        &self.connection
152    }
153    /// Validated protocol version and client capability metadata.
154    pub fn metadata(&self) -> &Map<String, Value> {
155        &self.metadata
156    }
157    /// Cancellation of the outer ACP request.
158    pub fn cancellation(&self) -> &RequestCancellation {
159        &self.cancellation
160    }
161    /// Cancellation including provider removal and transport EOF.
162    pub fn operation_cancellation(&self) -> &McpOperationCancellation {
163        &self.operation_cancellation
164    }
165
166    /// Send an operation-scoped notification while this invocation is live.
167    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
179/// Reusable MCP application state with owned invocation futures.
180pub trait McpService<Counterpart: Role>: Send + Sync + 'static {
181    /// Execute one invocation, including backend teardown.
182    ///
183    /// On operation cancellation, stop user work and complete owned cleanup
184    /// before returning. The binding drives this future through cancellation;
185    /// dropping it is not used as a substitute for joining cleanup.
186    fn execute(
187        &self,
188        request: McpRequest,
189        context: McpRequestContext<Counterpart>,
190    ) -> BoxFuture<'static, Result<McpOutcome, Error>>;
191}