返回 CodeWhale
mcp.rs
1 //! The pinned SDK owns MCP framing; Rust owns every process/network operation.
2 //! The selected Host backend immediately adopts proc/* and FetchProxy sessions.
3 use std::collections::HashMap;
4 use std::sync::{Arc, Mutex};
5 use std::time::{Duration, Instant};
6
7 use anyhow::{Context, Result, bail};
8 use serde_json::{Value, json};
9 use tokio::io::BufReader;
10 use tokio::sync::{Mutex as AsyncMutex, mpsc};
11 use tokio_util::sync::CancellationToken;
12
13 use super::protocol::*;
14 use super::registry::OwnerState;
15 use super::supervisor::{HostProcess, HostRequestContext};
16 use super::ticket::{Grant, Presented, Ticket, TicketKind};
17 use super::tier::HostTier;
18 use super::{ExtensionHostManager, ManagerShared};
19 use crate::core::engine::HumanDecision;
20 use crate::mcp::process_broker::BrokerSession;
21 use crate::mcp::{MAX_MCP_RESPONSE_BYTES, McpPool, McpServerConfig};
22
23 const MAX_SESSIONS: usize = 64;
24 const MAX_PENDING: usize = 256;
25 const MAX_DEADLINE: Duration = Duration::from_secs(86_400);
26 const SUPPORTED_METHODS: &[&str] = &[
27 "initialize",
28 "notifications/initialized",
29 "tools/list",
30 "resources/list",
31 "resources/templates/list",
32 "prompts/list",
33 "tools/call",
34 "resources/read",
35 "prompts/get",
36 ];
37
38 #[derive(Clone)]
39 struct Operation {
40 target: Value,
41 ticket: Ticket,
42 expires: Instant,
43 decision: Option<HumanDecision>,
44 }
45 struct BoundRequest {
46 operation_id: String,
47 }
48
49 /// No launch command, environment mapping, credential or decision key is wire data.
50 struct Session {
51 session_id: String,
52 owner: OwnerRef,
53 host_generation: u64,
54 name: String,
55 config: McpServerConfig,
56 cancel: CancellationToken,
57 broker: AsyncMutex<Option<BrokerSession>>,
58 frames: AsyncMutex<mpsc::Receiver<Value>>,
59 frame_sender: mpsc::Sender<Value>,
60 operations: Mutex<HashMap<String, Operation>>,
61 pending: Mutex<HashMap<String, BoundRequest>>,
62 replies: Mutex<HashMap<String, Value>>,
63 server_requests: Mutex<HashMap<String, String>>,
64 decision_key: Mutex<Option<[u8; 32]>>,
65 launched: Mutex<bool>,
66 launch_ticket: Mutex<Option<Ticket>>,
67 http: Option<http::HttpSession>,
68 }
69 impl Session {
70 fn validate(&self, shared: &ManagerShared, generation: u64, owner: &OwnerRef) -> Result<()> {
71 if self.cancel.is_cancelled() || generation != self.host_generation || owner != &self.owner
72 {
73 bail!("MCP broker session is stale or closed");
74 }
75 if shared
76 .builtin
77 .host_generation
78 .load(std::sync::atomic::Ordering::SeqCst)
79 != generation
80 {
81 bail!("MCP broker host generation is stale");
82 }
83 shared
84 .live_owner_authority(HostTier::Builtin, |registry| {
85 registry
86 .owner(&owner.plugin_id)
87 .filter(|entry| entry.owner == *owner && entry.state == OwnerState::Active)
88 .map(|entry| entry.owner.clone())
89 .ok_or_else(|| "MCP builtin owner is no longer live".to_string())
90 })
91 .map_err(anyhow::Error::msg)?;
92 if let Some(source) = self.config.reviewed_plugin.as_ref() {
93 source.validate_before_use(&self.name, "broker operation")?;
94 }
95 Ok(())
96 }
97 fn target(
98 &self,
99 operation_id: &str,
100 method: &str,
101 params: &Value,
102 wire_id: Option<&str>,
103 ) -> Value {
104 json!({"session_id": self.id(), "operation_id": operation_id, "method": method, "params": params, "wire_id": wire_id})
105 }
106 fn observe(&self, frame: &Value) -> Result<()> {
107 if frame.get("jsonrpc").and_then(Value::as_str) != Some("2.0") {
108 bail!("invalid MCP server frame");
109 }
110 if let Some(id) = frame.get("id") {
111 let id = wire_id(id)?;
112 if frame.get("method").is_some() {
113 let mut requests = self
114 .server_requests
115 .lock()
116 .expect("MCP server request lock");
117 if requests.len() >= MAX_PENDING || requests.contains_key(&id) {
118 bail!("MCP server request bound exceeded");
119 }
120 requests.insert(
121 id,
122 frame["method"]
123 .as_str()
124 .context("invalid MCP method")?
125 .to_string(),
126 );
127 } else if let Some(binding) = self.pending.lock().expect("MCP pending lock").remove(&id)
128 {
129 let mut replies = self.replies.lock().expect("MCP reply lock");
130 if replies.len() >= MAX_PENDING || replies.contains_key(&binding.operation_id) {
131 bail!("MCP reply bound exceeded");
132 }
133 replies.insert(binding.operation_id, frame.clone());
134 }
135 }
136 Ok(())
137 }
138 fn cancel(&self) {
139 self.cancel.cancel();
140 if let Some(http) = &self.http {
141 http.close();
142 }
143 }
144 fn id(&self) -> &str {
145 &self.session_id
146 }
147 }
148
149 #[derive(Default)]
150 pub(super) struct Broker {
151 sessions: Mutex<HashMap<String, Arc<Session>>>,
152 }
153 impl Drop for Broker {
154 fn drop(&mut self) {
155 for session in self.sessions.get_mut().expect("MCP broker lock").values() {
156 session.cancel();
157 }
158 }
159 }
160 impl Broker {
161 pub(super) fn revoke_owner(&self, owner: &str) -> u64 {
162 let mut sessions = self.sessions.lock().expect("MCP broker lock");
163 let before = sessions.len();
164 sessions.retain(|_, session| {
165 if session.owner.plugin_id == owner {
166 session.cancel();
167 false
168 } else {
169 true
170 }
171 });
172 (before - sessions.len()) as u64
173 }
174 pub(super) fn revoke_host(&self, tier: HostTier, generation: u64) -> u64 {
175 if tier != HostTier::Builtin {
176 return 0;
177 }
178 let mut sessions = self.sessions.lock().expect("MCP broker lock");
179 let before = sessions.len();
180 sessions.retain(|_, session| {
181 if session.host_generation == generation {
182 session.cancel();
183 false
184 } else {
185 true
186 }
187 });
188 (before - sessions.len()) as u64
189 }
190 fn session(&self, id: &str) -> Result<Arc<Session>> {
191 self.sessions
192 .lock()
193 .expect("MCP broker lock")
194 .get(id)
195 .cloned()
196 .context("MCP broker session is closed")
197 }
198 fn remove(&self, shared: &ManagerShared, id: &str) -> Option<Arc<Session>> {
199 let session = self.sessions.lock().expect("MCP broker lock").remove(id)?;
200 session.cancel();
201 for operation in session
202 .operations
203 .lock()
204 .expect("MCP operation lock")
205 .drain()
206 .map(|(_, v)| v)
207 {
208 shared.core_calls.tickets.revoke(&operation.ticket);
209 }
210 if let Some(ticket) = session
211 .launch_ticket
212 .lock()
213 .expect("MCP launch ticket lock")
214 .take()
215 {
216 shared.core_calls.tickets.revoke(&ticket);
217 }
218 shared
219 .mcp_users
220 .fetch_sub(1, std::sync::atomic::Ordering::SeqCst);
221 Some(session)
222 }
223 pub(super) async fn serve(
224 &self,
225 shared: &ManagerShared,
226 generation: u64,
227 request: HostRequest,
228 cx: HostRequestContext,
229 ) -> Result<Value, RpcErrorWire> {
230 let run = async {
231 match request {
232 request @ (HostRequest::NetStart(_)
233 | HostRequest::NetFetch(_)
234 | HostRequest::NetRead(_)
235 | HostRequest::NetRelease(_)
236 | HostRequest::NetClose(_)) => {
237 http::serve(self, shared, generation, request, &cx).await
238 }
239 HostRequest::ProcLaunch(params) => {
240 self.launch(shared, generation, params, &cx).await
241 }
242 HostRequest::ProcRead(params) => {
243 let session = self.session(&params.session_id)?;
244 session.validate(shared, generation, &params.owner)?;
245 let mut receiver = session.frames.lock().await;
246 let frame = tokio::select! { biased;
247 _ = cx.cancel.cancelled() => bail!("MCP broker read cancelled"),
248 _ = session.cancel.cancelled() => bail!("MCP broker session closed"),
249 frame = receiver.recv() => frame,
250 };
251 session.validate(shared, generation, &params.owner)?;
252 Ok(frame
253 .map_or_else(|| json!({"closed": true}), |frame| json!({"frame": frame})))
254 }
255 HostRequest::ProcWrite(params) => self.write(shared, generation, params, &cx).await,
256 HostRequest::ProcClose(params) => {
257 let session = self.session(&params.session_id)?;
258 // Closing needs exact ownership even after cancellation, but
259 // never grants a stale owner access to a replacement session.
260 if session.owner != params.owner || session.host_generation != generation {
261 bail!("MCP broker close has the wrong owner");
262 }
263 if let Some(session) = self.remove(shared, &params.session_id)
264 && let Some(mut broker) = session.broker.lock().await.take()
265 {
266 broker.shutdown().await;
267 }
268 Ok(json!({}))
269 }
270 _ => bail!("not an MCP broker request"),
271 }
272 };
273 run.await.map_err(|_| RpcErrorWire {
274 code: error_code::REFUSED,
275 message: "MCP broker request refused or closed".to_string(),
276 data: None,
277 })
278 }
279 async fn launch(
280 &self,
281 shared: &ManagerShared,
282 generation: u64,
283 params: ProcLaunchParams,
284 cx: &HostRequestContext,
285 ) -> Result<Value> {
286 let session = self.session(&params.session_id)?;
287 session.validate(shared, generation, &params.owner)?;
288 let mut launch_guard = CancelGuard {
289 cancel: session.cancel.clone(),
290 armed: true,
291 };
292 if session.http.is_some() {
293 bail!("HTTP session cannot launch a process");
294 }
295 let target = json!({"session_id": params.session_id});
296 shared
297 .core_calls
298 .tickets
299 .redeem(&Presented {
300 ticket: &params.ticket,
301 kind: TicketKind::McpLaunch,
302 tier: HostTier::Builtin,
303 host_generation: generation,
304 owner: &params.owner,
305 method: "proc/launch",
306 target: Some(&target),
307 })
308 .map_err(|refused| {
309 if refused.violation {
310 cx.violation("too many invalid MCP broker tickets".to_string());
311 }
312 anyhow::anyhow!("MCP launch ticket refused")
313 })?;
314 session
315 .launch_ticket
316 .lock()
317 .expect("MCP launch ticket lock")
318 .take();
319 {
320 let mut launched = session.launched.lock().expect("MCP launch lock");
321 if *launched {
322 bail!("MCP session already launched");
323 }
324 *launched = true;
325 }
326 if cx.cancel.is_cancelled() {
327 session.cancel();
328 bail!("MCP launch cancelled");
329 }
330 let command = session
331 .config
332 .command
333 .as_deref()
334 .context("MCP host stdio requires a command")?;
335 let (mut broker, stdout) = BrokerSession::spawn(
336 &session.name,
337 command,
338 &session.config,
339 session.cancel.clone(),
340 )?;
341 session.validate(shared, generation, &params.owner)?;
342 if session
343 .config
344 .reviewed_plugin
345 .as_ref()
346 .is_some_and(|source| source.plugin_name() == crate::mcp::COMPUTER_USE_PLUGIN_NAME)
347 {
348 let key = crate::mcp::random_key()?;
349 let keys = json!({"jsonrpc": "2.0", "method": crate::mcp::COMPUTER_USE_HOST_KEYS_METHOD, "params": {"decision_key": crate::mcp::hex_encode(&key), "ledger_key": crate::mcp::computer_use_ledger_key().await}});
350 let mut bytes = serde_json::to_vec(&keys)?;
351 bytes.push(b'\n');
352 tokio::select! { biased;
353 _ = cx.cancel.cancelled() => { session.cancel(); bail!("MCP launch cancelled"); },
354 _ = session.cancel.cancelled() => bail!("MCP launch revoked"),
355 result = broker.write(&bytes) => result?,
356 }
357 *session.decision_key.lock().expect("MCP decision key lock") = Some(key);
358 }
359 *session.broker.lock().await = Some(broker);
360 let reading = Arc::clone(&session);
361 tokio::spawn(async move {
362 let mut reader = BufReader::new(stdout);
363 let run = async {
364 loop {
365 let mut bytes = Vec::new();
366 let count = tokio::select! { biased;
367 _ = reading.cancel.cancelled() => break,
368 result = crate::mcp::read_line_capped(&mut reader, &mut bytes, MAX_MCP_RESPONSE_BYTES) => result?,
369 };
370 if count == 0 {
371 break;
372 }
373 if bytes.iter().all(u8::is_ascii_whitespace) {
374 continue;
375 }
376 let frame: Value = serde_json::from_slice(&bytes)?;
377 reading.observe(&frame)?;
378 tokio::select! { biased;
379 _ = reading.cancel.cancelled() => break,
380 result = reading.frame_sender.send(frame) => result.map_err(|_| anyhow::anyhow!("MCP frame consumer closed"))?,
381 }
382 }
383 Ok::<_, anyhow::Error>(())
384 };
385 let _ = run.await;
386 reading.cancel();
387 if let Some(mut broker) = reading.broker.lock().await.take() {
388 broker.shutdown().await;
389 }
390 });
391 session.validate(shared, generation, &params.owner)?;
392 if cx.cancel.is_cancelled() {
393 bail!("MCP launch cancelled");
394 }
395 launch_guard.armed = false;
396 Ok(json!({}))
397 }
398 async fn write(
399 &self,
400 shared: &ManagerShared,
401 generation: u64,
402 params: ProcWriteParams,
403 cx: &HostRequestContext,
404 ) -> Result<Value> {
405 let session = self.session(&params.session_id)?;
406 session.validate(shared, generation, &params.owner)?;
407 let mut write_guard = CancelGuard {
408 cancel: session.cancel.clone(),
409 armed: true,
410 };
411 if session.http.is_some() {
412 bail!("HTTP session cannot write a pipe");
413 }
414 let owner = params.owner.clone();
415 let (frame, operation_deadline) =
416 self.authorize_frame(shared, generation, &session, params, cx, "proc/write")?;
417 let mut bytes = serde_json::to_vec(&frame)?;
418 bytes.push(b'\n');
419 session.validate(shared, generation, &owner)?;
420 let result = async {
421 let mut slot = session.broker.lock().await;
422 let broker = slot.as_mut().context("MCP process not launched")?;
423 tokio::select! { biased;
424 _ = cx.cancel.cancelled() => bail!("MCP write cancelled"),
425 _ = session.cancel.cancelled() => bail!("MCP write revoked"),
426 result = tokio::time::timeout_at(tokio::time::Instant::from_std(operation_deadline), broker.write(&bytes)) => result.context("MCP pipe write expired")??,
427 }
428 Ok::<_, anyhow::Error>(())
429 }.await;
430 if result.is_err() {
431 session.cancel();
432 }
433 result?;
434 session.validate(shared, generation, &owner)?;
435 write_guard.armed = false;
436 Ok(json!({}))
437 }
438 fn authorize_frame(
439 &self,
440 shared: &ManagerShared,
441 generation: u64,
442 session: &Arc<Session>,
443 params: ProcWriteParams,
444 cx: &HostRequestContext,
445 ticket_method: &'static str,
446 ) -> Result<(Value, Instant)> {
447 let mut frame = params.frame;
448 let object = frame.as_object().context("MCP frame must be an object")?;
449 if object.keys().any(|key| {
450 !["jsonrpc", "id", "method", "params", "error", "result"].contains(&key.as_str())
451 }) || frame["jsonrpc"] != "2.0"
452 {
453 bail!("invalid MCP outbound frame");
454 }
455 let method = frame
456 .get("method")
457 .and_then(Value::as_str)
458 .map(str::to_owned);
459 if method.is_some() && (frame.get("error").is_some() || frame.get("result").is_some()) {
460 bail!("MCP request cannot contain an error");
461 }
462 let mut operation_deadline = Instant::now() + Duration::from_secs(2);
463 match method.as_deref() {
464 Some("notifications/cancelled") => {
465 if params.ticket.is_some() || frame.get("id").is_some() {
466 bail!("invalid MCP cancellation");
467 }
468 let id = wire_id(&frame["params"]["requestId"])?;
469 let pending = session.pending.lock().expect("MCP pending lock");
470 let binding = pending
471 .get(&id)
472 .context("MCP cancellation has no active request")?;
473 if params.operation_id.as_deref() != Some(binding.operation_id.as_str()) {
474 bail!("MCP cancellation target mismatch");
475 }
476 let fields = frame["params"]
477 .as_object()
478 .context("invalid MCP cancellation params")?;
479 if fields
480 .keys()
481 .any(|key| !["requestId", "reason"].contains(&key.as_str()))
482 {
483 bail!("invalid MCP cancellation params");
484 }
485 }
486 Some(method) => {
487 if !SUPPORTED_METHODS.contains(&method) {
488 bail!("MCP method has no Rust grant");
489 }
490 let operation_id = params
491 .operation_id
492 .as_deref()
493 .context("MCP operation id absent")?;
494 let body = frame.get("params").cloned().unwrap_or_else(|| json!({}));
495 if !body.is_object() {
496 bail!("MCP named params required");
497 }
498 let admitted_id = frame
499 .get("id")
500 .map(|id| {
501 id.as_str()
502 .context("MCP wire request ID must be a Rust string")
503 })
504 .transpose()?;
505 let target = session.target(operation_id, method, &body, admitted_id);
506 let ticket = params
507 .ticket
508 .as_deref()
509 .context("MCP operation ticket absent")?;
510 shared
511 .core_calls
512 .tickets
513 .redeem(&Presented {
514 ticket,
515 kind: TicketKind::McpOperation,
516 tier: HostTier::Builtin,
517 host_generation: generation,
518 owner: &params.owner,
519 method: ticket_method,
520 target: Some(&target),
521 })
522 .map_err(|refused| {
523 if refused.violation {
524 cx.violation("too many invalid MCP broker tickets".to_string());
525 }
526 anyhow::anyhow!("MCP operation ticket refused")
527 })?;
528 let operation = session
529 .operations
530 .lock()
531 .expect("MCP operation lock")
532 .remove(operation_id)
533 .context("MCP operation absent")?;
534 if operation.target != target || ticket != operation.ticket.expose() {
535 bail!("MCP operation does not match");
536 }
537 if operation.expires <= Instant::now() {
538 bail!("MCP operation expired");
539 }
540 operation_deadline = operation.expires;
541 if method == "tools/call" {
542 let tool = body["name"].as_str().context("MCP tool name absent")?;
543 if !session.config.is_tool_enabled(tool) {
544 bail!("MCP tool disabled");
545 }
546 if body.get("_meta").is_some() {
547 bail!("host-supplied MCP decision metadata refused");
548 }
549 if let Some(decision) = operation.decision.as_ref() {
550 if !decision.authorizes(
551 &McpPool::mcp_model_tool_name(&session.name, tool),
552 &body["arguments"],
553 ) {
554 bail!("MCP decision does not match");
555 }
556 if let Some(key) = session
557 .decision_key
558 .lock()
559 .expect("MCP decision key lock")
560 .as_ref()
561 {
562 frame["params"]["_meta"] = json!({});
563 frame["params"]["_meta"][crate::mcp::COMPUTER_USE_DECISION_META] =
564 crate::mcp::attest_decision(key, tool, &body["arguments"])?;
565 }
566 }
567 }
568 if method == "notifications/initialized" {
569 if frame.get("id").is_some() {
570 bail!("initialized must be a notification");
571 }
572 } else {
573 let id = wire_id(frame.get("id").context("MCP request id absent")?)?;
574 let mut pending = session.pending.lock().expect("MCP pending lock");
575 if pending.len() >= MAX_PENDING || pending.contains_key(&id) {
576 bail!("MCP request id unavailable");
577 }
578 pending.insert(
579 id,
580 BoundRequest {
581 operation_id: operation_id.to_string(),
582 },
583 );
584 }
585 }
586 None => {
587 if params.ticket.is_some()
588 || params.operation_id.is_some()
589 || frame.get("params").is_some()
590 {
591 bail!("invalid MCP server request answer");
592 }
593 let id = wire_id(frame.get("id").context("MCP response id absent")?)?;
594 let mut requests = session
595 .server_requests
596 .lock()
597 .expect("MCP server request lock");
598 let method = requests
599 .get(&id)
600 .context("MCP response has no observed server request")?;
601 let refusal = frame.get("result").is_none()
602 && frame["error"]["code"] == -32601
603 && frame["error"]["message"].is_string();
604 let ping = method == "ping"
605 && frame.get("error").is_none()
606 && frame
607 .get("result")
608 .and_then(Value::as_object)
609 .is_some_and(|result| result.is_empty());
610 if !refusal && !ping {
611 bail!("MCP server request success requires Rust authority");
612 }
613 requests.remove(&id);
614 }
615 }
616 if serde_json::to_vec(&frame)?.len() > MAX_MCP_RESPONSE_BYTES {
617 bail!("MCP outbound frame exceeds bound");
618 }
619 Ok((frame, operation_deadline))
620 }
621 }
622 struct CancelGuard {
623 cancel: CancellationToken,
624 armed: bool,
625 }
626 impl Drop for CancelGuard {
627 fn drop(&mut self) {
628 if self.armed {
629 self.cancel.cancel();
630 }
631 }
632 }
633 fn wire_id(id: &Value) -> Result<String> {
634 if !(id.is_string() || id.as_i64().is_some() || id.as_u64().is_some()) {
635 bail!("invalid MCP request id");
636 }
637 Ok(serde_json::to_string(id)?)
638 }
639
640 /// Semantic adapter: the existing Rust MCP connection still consumes raw
641 /// server replies for policy/catalog admission; it never writes MCP bytes.
642 pub(crate) struct SdkTransport {
643 manager: Arc<ExtensionHostManager>,
644 host: Arc<HostProcess>,
645 session: Arc<Session>,
646 session_id: String,
647 replies: std::collections::VecDeque<Vec<u8>>,
648 discovery_timeout: Duration,
649 }
650 impl SdkTransport {
651 #[cfg(test)]
652 pub(crate) async fn connect(
653 name: &str,
654 config: &McpServerConfig,
655 cancel: CancellationToken,
656 timeout: Duration,
657 ) -> Result<Self> {
658 Self::connect_with_http(name, config, cancel, timeout, None).await
659 }
660 pub(crate) async fn connect_with_http(
661 name: &str,
662 config: &McpServerConfig,
663 cancel: CancellationToken,
664 timeout: Duration,
665 client: Option<crate::mcp::http_client::McpHttpClient>,
666 ) -> Result<Self> {
667 if (config.url.is_some() != client.is_some())
668 || (config.url.is_none() && config.command.is_none())
669 {
670 bail!(
671 "Host MCP backend requires a Rust-prepared HTTP client or a stdio command; no Rust fallback is allowed"
672 );
673 }
674 let manager = super::manager();
675 manager
676 .shared
677 .mcp_users
678 .fetch_add(1, std::sync::atomic::Ordering::SeqCst);
679 let ready = manager.ensure_mcp_builtin().await;
680 let (host, owner, host_generation) = match ready {
681 Ok(ready) => ready,
682 Err(error) => {
683 manager
684 .shared
685 .mcp_users
686 .fetch_sub(1, std::sync::atomic::Ordering::SeqCst);
687 return Err(anyhow::Error::msg(error));
688 }
689 };
690 let session_id = uuid::Uuid::new_v4().to_string();
691 let (frame_sender, frames) = mpsc::channel(1);
692 let session = Arc::new(Session {
693 session_id: session_id.clone(),
694 owner,
695 host_generation,
696 name: name.to_string(),
697 config: config.clone(),
698 // The facade owns caller cancellation. Retiring this transport
699 // must release its resources without cancelling that facade and
700 // masking the concrete reply or failure it is about to receive.
701 cancel: cancel.child_token(),
702 broker: AsyncMutex::new(None),
703 frames: AsyncMutex::new(frames),
704 frame_sender,
705 operations: Mutex::new(HashMap::new()),
706 pending: Mutex::new(HashMap::new()),
707 replies: Mutex::new(HashMap::new()),
708 server_requests: Mutex::new(HashMap::new()),
709 decision_key: Mutex::new(None),
710 launched: Mutex::new(false),
711 launch_ticket: Mutex::new(None),
712 http: client.map(|client| {
713 http::HttpSession::new(
714 client,
715 config.url.as_deref().expect("HTTP config URL"),
716 crate::mcp::is_legacy_sse_transport(config),
717 )
718 }),
719 });
720 if session.http.is_some() {
721 // Handler cancellation can drop a body future outside an exchange.
722 // The weak watch releases retained network bodies on every token
723 // cancellation without extending the broker session's lifetime.
724 let weak = Arc::downgrade(&session);
725 let cancel = session.cancel.clone();
726 tokio::spawn(async move {
727 cancel.cancelled().await;
728 if let Some(session) = weak.upgrade()
729 && let Some(http) = &session.http
730 {
731 http.close();
732 }
733 });
734 }
735 {
736 let mut sessions = manager
737 .shared
738 .mcp_broker
739 .sessions
740 .lock()
741 .expect("MCP broker lock");
742 if sessions.len() >= MAX_SESSIONS {
743 manager
744 .shared
745 .mcp_users
746 .fetch_sub(1, std::sync::atomic::Ordering::SeqCst);
747 bail!("MCP host session limit exceeded");
748 }
749 sessions.insert(session_id.clone(), Arc::clone(&session));
750 }
751 let transport = Self {
752 manager,
753 host,
754 session,
755 session_id,
756 replies: Default::default(),
757 discovery_timeout: timeout,
758 };
759 if let Some(http) = &transport.session.http {
760 http.preflight(&transport.session, &transport.manager.shared)
761 .await?;
762 }
763 Ok(transport)
764 }
765 fn grant(
766 &self,
767 method: &str,
768 params: Value,
769 timeout: Duration,
770 decision: Option<&HumanDecision>,
771 wire_id: Option<&str>,
772 ) -> Result<McpOperationGrant> {
773 if timeout.is_zero()
774 || timeout > MAX_DEADLINE
775 || !params.is_object()
776 || serde_json::to_vec(&params)?.len() > MAX_MCP_RESPONSE_BYTES
777 {
778 bail!("MCP grant exceeds bound");
779 }
780 self.session.validate(
781 &self.manager.shared,
782 self.session.host_generation,
783 &self.session.owner,
784 )?;
785 let operation_id = uuid::Uuid::new_v4().to_string();
786 let target = self.session.target(&operation_id, method, &params, wire_id);
787 let mut operations = self.session.operations.lock().expect("MCP operation lock");
788 if operations.len() >= MAX_PENDING {
789 bail!("MCP grant bound exceeded");
790 }
791 let ticket = self.manager.shared.core_calls.tickets.mint(Grant {
792 kind: TicketKind::McpOperation,
793 tier: HostTier::Builtin,
794 host_generation: self.session.host_generation,
795 owner: self.session.owner.clone(),
796 method: if self.session.http.is_some() {
797 "net/fetch"
798 } else {
799 "proc/write"
800 },
801 target: target.clone(),
802 ttl: timeout,
803 uses: 1,
804 });
805 let wire_ticket = ticket.expose().to_string();
806 operations.insert(
807 operation_id.clone(),
808 Operation {
809 target,
810 ticket,
811 expires: Instant::now() + timeout,
812 decision: decision.cloned(),
813 },
814 );
815 Ok(McpOperationGrant {
816 ticket: wire_ticket,
817 operation_id,
818 method: method.to_string(),
819 wire_id: wire_id.map(str::to_owned),
820 params,
821 })
822 }
823 async fn exchange(
824 &mut self,
825 bytes: Vec<u8>,
826 decision: Option<&HumanDecision>,
827 timeout: Duration,
828 ) -> Result<()> {
829 let value: Value = serde_json::from_slice(&bytes)?;
830 let method = value["method"]
831 .as_str()
832 .context("MCP semantic method absent")?;
833 // The SDK has already sent this notification during connect, under
834 // its separate exact Rust grant; the old connection facade sees one
835 // successful initialized exchange, not a second pipe write.
836 if method == "notifications/initialized" {
837 return Ok(());
838 }
839 if !SUPPORTED_METHODS.contains(&method) {
840 bail!("MCP semantic operation unavailable");
841 }
842 let params = value.get("params").cloned().unwrap_or_else(|| json!({}));
843 let wire_id = value
844 .get("id")
845 .map(|id| {
846 id.as_str()
847 .context("Rust MCP facade request ID must be a string")
848 })
849 .transpose()?;
850 let grant = self.grant(method, params, timeout, decision, wire_id)?;
851 let operation_id = grant.operation_id.clone();
852 let mut guard = ExchangeGuard {
853 transport: self,
854 armed: true,
855 };
856 let request = if method == "initialize" {
857 let initialized = guard.transport.grant(
858 "notifications/initialized",
859 json!({}),
860 timeout,
861 None,
862 None,
863 )?;
864 let target = json!({"session_id": guard.transport.session_id});
865 let launch = guard
866 .transport
867 .manager
868 .shared
869 .core_calls
870 .tickets
871 .mint(Grant {
872 kind: TicketKind::McpLaunch,
873 tier: HostTier::Builtin,
874 host_generation: guard.transport.session.host_generation,
875 owner: guard.transport.session.owner.clone(),
876 method: if guard.transport.session.http.is_some() {
877 "net/start"
878 } else {
879 "proc/launch"
880 },
881 target,
882 ttl: timeout,
883 uses: 1,
884 });
885 *guard
886 .transport
887 .session
888 .launch_ticket
889 .lock()
890 .expect("MCP launch ticket lock") = Some(launch.clone());
891 CoreRequest::McpOpen(Box::new(McpOpenParams {
892 owner: guard.transport.session.owner.clone(),
893 session_id: guard.transport.session_id.clone(),
894 launch_ticket: launch.expose().to_string(),
895 transport: if guard.transport.session.http.is_none() {
896 "stdio"
897 } else if crate::mcp::is_legacy_sse_transport(&guard.transport.session.config) {
898 "sse"
899 } else {
900 "http"
901 }
902 .to_string(),
903 initialize_grant: grant,
904 initialized_grant: initialized,
905 client_version: env!("CARGO_PKG_VERSION").to_string(),
906 deadline_ms: timeout.as_millis() as u64,
907 }))
908 } else {
909 CoreRequest::McpRequest(McpRequestParams {
910 owner: guard.transport.session.owner.clone(),
911 session_id: guard.transport.session_id.clone(),
912 grant,
913 deadline_ms: timeout.as_millis() as u64,
914 })
915 };
916 let outcome = guard
917 .transport
918 .host
919 .call(request, Some("host:mcp".to_string()))
920 .await;
921 let raw = guard
922 .transport
923 .session
924 .replies
925 .lock()
926 .expect("MCP reply lock")
927 .remove(&operation_id);
928 if let Some(mut raw) = raw {
929 // Preserve the server's exact result/error (including extra MCP
930 // fields), changing only correlation back to the Rust facade id.
931 raw["id"] = value["id"].clone();
932 // Rust's existing initialize validation owns the diagnostic and
933 // accepted-revision contract, including when the SDK refused it.
934 guard.transport.replies.push_back(serde_json::to_vec(&raw)?);
935 guard.armed = false;
936 return Ok(());
937 }
938 if let Some(error) = guard
939 .transport
940 .session
941 .http
942 .as_ref()
943 .and_then(http::HttpSession::failure)
944 {
945 return Err(error);
946 }
947 outcome
948 .map_err(anyhow::Error::msg)
949 .context("MCP SDK request failed before an admitted reply")?;
950 bail!("MCP SDK connection closed before an admitted reply")
951 }
952 }
953 struct ExchangeGuard<'a> {
954 transport: &'a mut SdkTransport,
955 armed: bool,
956 }
957 impl Drop for ExchangeGuard<'_> {
958 fn drop(&mut self) {
959 if self.armed {
960 self.transport.session.cancel();
961 }
962 }
963 }
964 impl Drop for SdkTransport {
965 fn drop(&mut self) {
966 self.manager
967 .shared
968 .mcp_broker
969 .remove(&self.manager.shared, &self.session_id);
970 }
971 }
972 #[async_trait::async_trait]
973 impl crate::mcp::McpTransport for SdkTransport {
974 async fn send(&mut self, bytes: Vec<u8>) -> Result<()> {
975 self.exchange(bytes, None, self.discovery_timeout).await
976 }
977 async fn send_decided(
978 &mut self,
979 bytes: Vec<u8>,
980 decision: Option<&HumanDecision>,
981 timeout: Duration,
982 ) -> Result<()> {
983 self.exchange(bytes, decision, timeout).await
984 }
985 async fn recv(&mut self) -> Result<Vec<u8>> {
986 self.replies
987 .pop_front()
988 .context("MCP SDK connection closed without a reply")
989 }
990 async fn last_stderr_line(&self) -> Option<String> {
991 let slot = self.session.broker.lock().await;
992 slot.as_ref()?.last_stderr_line().await
993 }
994 fn probe_dead(&self) -> bool {
995 self.session.cancel.is_cancelled() || self.host.has_exited() || self.host.is_retiring()
996 }
997 async fn shutdown(&mut self) {
998 if let Some(session) = self
999 .manager
1000 .shared
1001 .mcp_broker
1002 .remove(&self.manager.shared, &self.session_id)
1003 && let Some(mut broker) = session.broker.lock().await.take()
1004 {
1005 broker.shutdown().await;
1006 }
1007 let _ = self
1008 .host
1009 .call(
1010 CoreRequest::McpClose(McpCloseParams {
1011 owner: self.session.owner.clone(),
1012 session_id: self.session_id.clone(),
1013 deadline_ms: 2500,
1014 }),
1015 Some("host:mcp".to_string()),
1016 )
1017 .await;
1018 }
1019 }
1020
1021 #[cfg(test)]
1022 #[path = "mcp_tests.rs"]
1023 mod tests;
1024
1025 #[path = "mcp_http.rs"]
1026 mod http;
1027
1027 lines RUST