1use std::collections::{HashMap, HashSet};
6use std::sync::Arc;
7
8use serde_json::{Map, Value};
9use tokio::sync::RwLock;
10
11use crate::binary::ArgType;
12use crate::dispatch::{BinaryTrieRouter, SharedArgs, ToolResult};
13use crate::protocol::{FieldDef, InputSchema};
14use crate::resource::{ResourceContent, ResourceError, ResourceRegistry};
15use crate::security::{sanitize_text, SecurityAuditAction, SecurityAuditEvent, SecurityAuditLog};
16use crate::{CapabilityManifest, DCPError, SecurityError};
17
18use super::json_rpc::{
19 JsonRpcError, JsonRpcParseError, JsonRpcParser, JsonRpcRequest, JsonRpcResponse, RequestId,
20 DEFAULT_MAX_JSONRPC_REQUEST_SIZE,
21};
22use super::request_replay::{replay_key, RequestReplayGuard};
23
24#[derive(Debug, Clone, thiserror::Error)]
26pub enum CompleteAdapterError {
27 #[error("JSON-RPC parse error: {0}")]
28 ParseError(#[from] JsonRpcParseError),
29 #[error("unknown tool: {0}")]
30 UnknownTool(String),
31 #[error("DCP error: {0}")]
32 DcpError(#[from] DCPError),
33 #[error("resource error: {0}")]
34 ResourceError(String),
35 #[error("prompt error: {0}")]
36 PromptError(String),
37 #[error("serialization error: {0}")]
38 SerializationError(String),
39 #[error("invalid params: {0}")]
40 InvalidParams(String),
41 #[error("invalid request: {0}")]
42 InvalidRequest(String),
43 #[error("lifecycle not initialized")]
44 LifecycleNotInitialized,
45 #[error("lifecycle already initialized")]
46 LifecycleAlreadyInitialized,
47 #[error("capability denied")]
48 CapabilityDenied,
49 #[error("{kind} capacity exceeded")]
50 CapacityExceeded { kind: &'static str, max: usize },
51}
52
53impl From<ResourceError> for CompleteAdapterError {
54 fn from(e: ResourceError) -> Self {
55 Self::ResourceError(e.to_string())
56 }
57}
58
59#[derive(Debug, Clone)]
61pub struct PromptTemplate {
62 pub name: String,
64 pub description: String,
66 pub arguments: Vec<PromptArgument>,
68 pub template: String,
70}
71
72#[derive(Debug, Clone)]
74pub struct PromptArgument {
75 pub name: String,
77 pub description: String,
79 pub required: bool,
81}
82
83impl PromptTemplate {
84 pub fn new(
86 name: impl Into<String>,
87 description: impl Into<String>,
88 template: impl Into<String>,
89 ) -> Self {
90 Self {
91 name: name.into(),
92 description: description.into(),
93 arguments: Vec::new(),
94 template: template.into(),
95 }
96 }
97
98 pub fn with_argument(
100 mut self,
101 name: impl Into<String>,
102 description: impl Into<String>,
103 required: bool,
104 ) -> Self {
105 self.arguments.push(PromptArgument {
106 name: name.into(),
107 description: description.into(),
108 required,
109 });
110 self
111 }
112
113 pub fn render(&self, args: &HashMap<String, String>) -> Result<String, CompleteAdapterError> {
115 for argument_name in args.keys() {
116 if !self.arguments.iter().any(|arg| arg.name == *argument_name) {
117 return Err(CompleteAdapterError::InvalidParams(
118 "unknown prompt argument".to_string(),
119 ));
120 }
121 }
122
123 for arg in &self.arguments {
125 if arg.required && !args.contains_key(&arg.name) {
126 return Err(CompleteAdapterError::PromptError(format!(
127 "missing required argument: {}",
128 arg.name
129 )));
130 }
131 }
132
133 let mut result = self.template.clone();
135 for (key, value) in args {
136 let placeholder = format!("{{{{{}}}}}", key);
137 result = result.replace(&placeholder, value);
138 }
139
140 Ok(result)
141 }
142
143 fn has_argument(&self, name: &str) -> bool {
144 self.arguments.iter().any(|arg| arg.name == name)
145 }
146}
147
148#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
150pub enum LogLevel {
151 Debug,
152 #[default]
153 Info,
154 Warning,
155 Error,
156}
157
158impl LogLevel {
159 fn from_str(s: &str) -> Option<Self> {
160 match s.to_lowercase().as_str() {
161 "debug" => Some(Self::Debug),
162 "info" => Some(Self::Info),
163 "warning" | "warn" => Some(Self::Warning),
164 "error" => Some(Self::Error),
165 _ => None,
166 }
167 }
168}
169
170fn set_manifest_ids(
171 value: Option<&Value>,
172 field: &'static str,
173 max_id_exclusive: usize,
174 mut set: impl FnMut(u16),
175) -> Result<(), CompleteAdapterError> {
176 let Some(value) = value else {
177 return Ok(());
178 };
179 let ids = value.as_array().ok_or_else(|| {
180 CompleteAdapterError::InvalidParams("capability ids must be array".into())
181 })?;
182
183 for id in ids {
184 let id = id.as_u64().ok_or_else(|| {
185 CompleteAdapterError::InvalidParams("capability id must be integer".into())
186 })?;
187 if id >= max_id_exclusive as u64 {
188 return Err(CompleteAdapterError::InvalidParams(format!(
189 "{field} capability id out of range"
190 )));
191 }
192 set(id as u16);
193 }
194
195 Ok(())
196}
197
198fn set_extension_ids(
199 value: Option<&Value>,
200 manifest: &mut CapabilityManifest,
201) -> Result<(), CompleteAdapterError> {
202 let Some(value) = value else {
203 return Ok(());
204 };
205 let ids = value.as_array().ok_or_else(|| {
206 CompleteAdapterError::InvalidParams("capability ids must be array".into())
207 })?;
208
209 for id in ids {
210 let id = id.as_u64().ok_or_else(|| {
211 CompleteAdapterError::InvalidParams("capability id must be integer".into())
212 })?;
213 if id >= 64 {
214 return Err(CompleteAdapterError::InvalidParams(
215 "extension capability id out of range".into(),
216 ));
217 }
218 manifest.set_extension(id as u8);
219 }
220
221 Ok(())
222}
223
224fn optional_object<'a>(
225 value: Option<&'a Value>,
226 field: &'static str,
227) -> Result<Option<&'a Map<String, Value>>, CompleteAdapterError> {
228 match value {
229 None => Ok(None),
230 Some(Value::Object(object)) => Ok(Some(object)),
231 Some(_) => Err(CompleteAdapterError::InvalidParams(format!(
232 "{field} must be object"
233 ))),
234 }
235}
236
237fn tool_input_schema_json(schema: &InputSchema) -> Value {
238 let mut properties = Map::new();
239 let mut required = Vec::new();
240
241 for (idx, field) in schema.fields.iter().take(16).enumerate() {
242 let mut property = Map::new();
243 property.insert(
244 "type".to_string(),
245 Value::String(
246 match field.field_type {
247 ArgType::Null => "null",
248 ArgType::Bool => "boolean",
249 ArgType::I32 | ArgType::I64 => "integer",
250 ArgType::F64 => "number",
251 ArgType::String | ArgType::Bytes => "string",
252 ArgType::Array => "array",
253 ArgType::Object => "object",
254 }
255 .to_string(),
256 ),
257 );
258
259 if matches!(field.field_type, ArgType::String | ArgType::Bytes) && field.size > 0 {
260 property.insert(
261 "maxLength".to_string(),
262 Value::Number(serde_json::Number::from(field.size)),
263 );
264 }
265
266 if let Some((minimum, maximum)) = field.enum_range {
267 property.insert(
268 "minimum".to_string(),
269 Value::Number(serde_json::Number::from(minimum)),
270 );
271 property.insert(
272 "maximum".to_string(),
273 Value::Number(serde_json::Number::from(maximum)),
274 );
275 }
276
277 if schema.is_required(idx) {
278 required.push(Value::String(field.name.to_string()));
279 }
280
281 properties.insert(field.name.to_string(), Value::Object(property));
282 }
283
284 let mut root = Map::new();
285 root.insert("type".to_string(), Value::String("object".to_string()));
286 root.insert("properties".to_string(), Value::Object(properties));
287 root.insert("additionalProperties".to_string(), Value::Bool(false));
288 if !required.is_empty() {
289 root.insert("required".to_string(), Value::Array(required));
290 }
291
292 Value::Object(root)
293}
294
295fn validate_field_capacity(field: &FieldDef, minimum: usize) -> Result<(), CompleteAdapterError> {
296 if field.size as usize >= minimum {
297 Ok(())
298 } else {
299 Err(CompleteAdapterError::DcpError(DCPError::ValidationFailed))
300 }
301}
302
303fn write_field(
304 buffer: &mut [u8],
305 field: &FieldDef,
306 bytes: &[u8],
307) -> Result<(), CompleteAdapterError> {
308 let offset = field.offset as usize;
309 let end = offset
310 .checked_add(bytes.len())
311 .ok_or(CompleteAdapterError::DcpError(DCPError::OutOfBounds))?;
312 if bytes.len() > field.size as usize || end > buffer.len() {
313 return Err(CompleteAdapterError::DcpError(DCPError::OutOfBounds));
314 }
315
316 buffer[offset..end].copy_from_slice(bytes);
317 Ok(())
318}
319
320fn encode_json_argument(
321 buffer: &mut [u8],
322 field: &FieldDef,
323 value: &Value,
324) -> Result<(), CompleteAdapterError> {
325 match field.field_type {
326 ArgType::Null => {
327 if value.is_null() {
328 Ok(())
329 } else {
330 Err(CompleteAdapterError::InvalidParams(
331 "argument type mismatch".to_string(),
332 ))
333 }
334 }
335 ArgType::Bool => {
336 validate_field_capacity(field, 1)?;
337 let value = value.as_bool().ok_or_else(|| {
338 CompleteAdapterError::InvalidParams("argument type mismatch".to_string())
339 })?;
340 write_field(buffer, field, &[u8::from(value)])
341 }
342 ArgType::I32 => {
343 if let Some((minimum, maximum)) = field.enum_range {
344 validate_field_capacity(field, 1)?;
345 let value = value.as_u64().ok_or_else(|| {
346 CompleteAdapterError::InvalidParams("argument type mismatch".to_string())
347 })?;
348 if value < minimum as u64 || value > maximum as u64 {
349 return Err(CompleteAdapterError::InvalidParams(
350 "argument out of range".to_string(),
351 ));
352 }
353 write_field(buffer, field, &[value as u8])
354 } else {
355 validate_field_capacity(field, 4)?;
356 let value = value.as_i64().ok_or_else(|| {
357 CompleteAdapterError::InvalidParams("argument type mismatch".to_string())
358 })?;
359 let value = i32::try_from(value).map_err(|_| {
360 CompleteAdapterError::InvalidParams("argument out of range".to_string())
361 })?;
362 write_field(buffer, field, &value.to_le_bytes())
363 }
364 }
365 ArgType::I64 => {
366 validate_field_capacity(field, 8)?;
367 let value = value.as_i64().ok_or_else(|| {
368 CompleteAdapterError::InvalidParams("argument type mismatch".to_string())
369 })?;
370 write_field(buffer, field, &value.to_le_bytes())
371 }
372 ArgType::F64 => {
373 validate_field_capacity(field, 8)?;
374 let value = value.as_f64().ok_or_else(|| {
375 CompleteAdapterError::InvalidParams("argument type mismatch".to_string())
376 })?;
377 write_field(buffer, field, &value.to_le_bytes())
378 }
379 ArgType::String | ArgType::Bytes => {
380 let value = value.as_str().ok_or_else(|| {
381 CompleteAdapterError::InvalidParams("argument type mismatch".to_string())
382 })?;
383 if value.len() > field.size as usize {
384 return Err(CompleteAdapterError::InvalidParams(
385 "argument too large".to_string(),
386 ));
387 }
388 write_field(buffer, field, value.as_bytes())
389 }
390 ArgType::Array => {
391 if !value.is_array() {
392 return Err(CompleteAdapterError::InvalidParams(
393 "argument type mismatch".to_string(),
394 ));
395 }
396 let bytes = serde_json::to_vec(value)
397 .map_err(|e| CompleteAdapterError::SerializationError(e.to_string()))?;
398 if bytes.len() > field.size as usize {
399 return Err(CompleteAdapterError::InvalidParams(
400 "argument too large".to_string(),
401 ));
402 }
403 write_field(buffer, field, &bytes)
404 }
405 ArgType::Object => {
406 if !value.is_object() {
407 return Err(CompleteAdapterError::InvalidParams(
408 "argument type mismatch".to_string(),
409 ));
410 }
411 let bytes = serde_json::to_vec(value)
412 .map_err(|e| CompleteAdapterError::SerializationError(e.to_string()))?;
413 if bytes.len() > field.size as usize {
414 return Err(CompleteAdapterError::InvalidParams(
415 "argument too large".to_string(),
416 ));
417 }
418 write_field(buffer, field, &bytes)
419 }
420 }
421}
422
423fn encode_json_arguments(
424 schema: &InputSchema,
425 arguments: Option<&Value>,
426) -> Result<(Vec<u8>, u64), CompleteAdapterError> {
427 let empty_arguments = Map::new();
428 let arguments = match arguments {
429 Some(Value::Object(arguments)) => arguments,
430 Some(_) => {
431 return Err(CompleteAdapterError::InvalidParams(
432 "arguments must be object".to_string(),
433 ));
434 }
435 None => &empty_arguments,
436 };
437
438 let max_len = schema
439 .fields
440 .iter()
441 .try_fold(0usize, |max_len, field| {
442 let end = (field.offset as usize)
443 .checked_add(field.size as usize)
444 .ok_or(DCPError::OutOfBounds)?;
445 Ok::<usize, DCPError>(max_len.max(end))
446 })
447 .map_err(CompleteAdapterError::DcpError)?;
448 let mut buffer = vec![0u8; max_len];
449 let mut layout = 0u64;
450
451 for argument_name in arguments.keys() {
452 if !schema
453 .fields
454 .iter()
455 .any(|field| field.name == argument_name.as_str())
456 {
457 return Err(CompleteAdapterError::InvalidParams(
458 "unknown argument".to_string(),
459 ));
460 }
461 }
462
463 for (idx, field) in schema.fields.iter().enumerate() {
464 if idx >= 16 {
465 return Err(CompleteAdapterError::DcpError(DCPError::ValidationFailed));
466 }
467
468 match arguments.get(field.name) {
469 Some(value) => {
470 encode_json_argument(&mut buffer, field, value)?;
471 let shift = idx * 4;
472 layout |= (field.field_type as u64) << shift;
473 }
474 None if schema.is_required(idx) => {
475 return Err(CompleteAdapterError::InvalidParams(
476 "missing required argument".to_string(),
477 ));
478 }
479 None => {}
480 }
481 }
482
483 Ok((buffer, layout))
484}
485
486fn tools_call_params_object(
487 params: Option<&Value>,
488) -> Result<&Map<String, Value>, CompleteAdapterError> {
489 let Some(Value::Object(params)) = params else {
490 return Err(CompleteAdapterError::InvalidParams(
491 "tools/call params must be an object".into(),
492 ));
493 };
494
495 if params
496 .keys()
497 .any(|key| key != "name" && key != "arguments" && key != "_meta")
498 {
499 return Err(CompleteAdapterError::InvalidParams(
500 "tools/call params contain unsupported fields".into(),
501 ));
502 }
503
504 if let Some(meta) = params.get("_meta") {
505 if !meta.is_object() {
506 return Err(CompleteAdapterError::InvalidParams(
507 "tools/call _meta must be an object".into(),
508 ));
509 }
510 }
511
512 Ok(params)
513}
514
515pub struct CompleteMcpAdapter {
517 tool_cache: HashMap<String, u16>,
519 id_to_name: HashMap<u16, String>,
521 resources: Arc<RwLock<ResourceRegistry>>,
523 prompts: HashMap<String, PromptTemplate>,
525 log_level: RwLock<LogLevel>,
527 server_name: String,
529 server_version: String,
531 protocol_version: RwLock<super::mcp2025::ProtocolVersion>,
533 version_negotiator: super::mcp2025::VersionNegotiator,
535 roots: Arc<super::mcp2025::RootsRegistry>,
537 subscriptions: Arc<super::mcp2025::SubscriptionTracker>,
539 elicitation: Arc<super::mcp2025::ElicitationHandler>,
541 resource_templates: Arc<super::mcp2025::ResourceTemplateRegistry>,
543 notifications: Arc<super::mcp2025::NotificationManager>,
545 cancellation: Arc<super::mcp2025::CancellationManager>,
547 progress: Arc<super::mcp2025::ProgressTracker>,
549 prompt_ids: HashMap<String, u16>,
551 server_manifest: CapabilityManifest,
553 authorization_policy: Option<CapabilityManifest>,
555 negotiated_manifest: RwLock<CapabilityManifest>,
557 initialize_seen: RwLock<bool>,
559 initialized_seen: RwLock<bool>,
561 security_audit: SecurityAuditLog,
563 tool_call_replay_guard: RwLock<RequestReplayGuard>,
565 max_request_size: usize,
567}
568
569impl CompleteMcpAdapter {
570 pub fn new() -> Self {
572 Self {
573 tool_cache: HashMap::new(),
574 id_to_name: HashMap::new(),
575 resources: Arc::new(RwLock::new(ResourceRegistry::new())),
576 prompts: HashMap::new(),
577 log_level: RwLock::new(LogLevel::default()),
578 server_name: "dcp-server".to_string(),
579 server_version: env!("CARGO_PKG_VERSION").to_string(),
580 protocol_version: RwLock::new(super::mcp2025::ProtocolVersion::default()),
581 version_negotiator: super::mcp2025::VersionNegotiator::new(),
582 roots: Arc::new(super::mcp2025::RootsRegistry::new()),
583 subscriptions: Arc::new(super::mcp2025::SubscriptionTracker::new()),
584 elicitation: Arc::new(super::mcp2025::ElicitationHandler::new()),
585 resource_templates: Arc::new(super::mcp2025::ResourceTemplateRegistry::new()),
586 notifications: Arc::new(super::mcp2025::NotificationManager::new()),
587 cancellation: Arc::new(super::mcp2025::CancellationManager::new()),
588 progress: Arc::new(super::mcp2025::ProgressTracker::new()),
589 prompt_ids: HashMap::new(),
590 server_manifest: CapabilityManifest::new(1),
591 authorization_policy: None,
592 negotiated_manifest: RwLock::new(CapabilityManifest::new(1)),
593 initialize_seen: RwLock::new(false),
594 initialized_seen: RwLock::new(false),
595 security_audit: SecurityAuditLog::new(),
596 tool_call_replay_guard: RwLock::new(RequestReplayGuard::default()),
597 max_request_size: DEFAULT_MAX_JSONRPC_REQUEST_SIZE,
598 }
599 }
600
601 pub fn with_server_info(mut self, name: impl Into<String>, version: impl Into<String>) -> Self {
603 self.server_name = name.into();
604 self.server_version = version.into();
605 self
606 }
607
608 pub fn with_max_request_size(mut self, max_request_size: usize) -> Self {
610 self.max_request_size = max_request_size;
611 self
612 }
613
614 pub fn with_authorization_policy(mut self, authorization_policy: CapabilityManifest) -> Self {
616 self.authorization_policy = Some(authorization_policy);
617 self
618 }
619
620 pub fn register_tool(
622 &mut self,
623 name: impl Into<String>,
624 tool_id: u16,
625 ) -> Result<u16, CompleteAdapterError> {
626 let name = name.into();
627 if (tool_id as usize) >= CapabilityManifest::MAX_TOOLS {
628 self.audit_tool_registration_failure(
629 "tool_registration_capacity_exceeded",
630 tool_id,
631 &name,
632 );
633 return Err(CompleteAdapterError::CapacityExceeded {
634 kind: "tool",
635 max: CapabilityManifest::MAX_TOOLS,
636 });
637 }
638 if self.tool_cache.contains_key(&name) {
639 self.audit_tool_registration_failure(
640 "tool_registration_duplicate_name",
641 tool_id,
642 &name,
643 );
644 return Err(CompleteAdapterError::InvalidRequest(
645 "duplicate tool name".into(),
646 ));
647 }
648 if self.id_to_name.contains_key(&tool_id) {
649 self.audit_tool_registration_failure("tool_registration_duplicate_id", tool_id, &name);
650 return Err(CompleteAdapterError::InvalidRequest(
651 "duplicate tool id".into(),
652 ));
653 }
654
655 self.tool_cache.insert(name.clone(), tool_id);
656 self.id_to_name.insert(tool_id, name);
657 self.server_manifest.set_tool(tool_id);
658 Ok(tool_id)
659 }
660
661 pub fn resources(&self) -> Arc<RwLock<ResourceRegistry>> {
663 Arc::clone(&self.resources)
664 }
665
666 pub fn register_prompt(
668 &mut self,
669 template: PromptTemplate,
670 ) -> Result<u16, CompleteAdapterError> {
671 if self.prompt_ids.contains_key(&template.name) {
672 return Err(CompleteAdapterError::InvalidRequest(
673 "duplicate prompt name".into(),
674 ));
675 }
676 if self.prompts.len() >= CapabilityManifest::MAX_PROMPTS {
677 return Err(CompleteAdapterError::CapacityExceeded {
678 kind: "prompt",
679 max: CapabilityManifest::MAX_PROMPTS,
680 });
681 }
682
683 let prompt_id = self.prompts.len() as u16;
684 self.prompt_ids.insert(template.name.clone(), prompt_id);
685 self.prompts.insert(template.name.clone(), template);
686 self.server_manifest.set_prompt(prompt_id);
687 Ok(prompt_id)
688 }
689
690 pub fn security_audit(&self) -> SecurityAuditLog {
692 self.security_audit.clone()
693 }
694
695 fn audit_tool_registration_failure(&self, reason: &'static str, tool_id: u16, name: &str) {
696 self.security_audit.record(
697 SecurityAuditEvent::new(SecurityAuditAction::ValidationRejected, reason)
698 .with_field("adapter", "complete_mcp")
699 .with_field("operation", "tool_registration")
700 .with_field("tool_id", tool_id.to_string())
701 .with_field("tool_name", name),
702 );
703 }
704
705 fn client_manifest_from_initialize(
706 params: Option<&Value>,
707 ) -> Result<CapabilityManifest, CompleteAdapterError> {
708 let mut manifest = CapabilityManifest::new(1);
709 let params = optional_object(params, "initialize params")?;
710 let capabilities = if let Some(params) = params {
711 let legacy_dcp_capabilities = params.get("dcpCapabilities");
712 let standard_dcp_capabilities = if let Some(capabilities) = params.get("capabilities") {
713 optional_object(Some(capabilities), "capabilities")?
714 .and_then(|capabilities| capabilities.get("dcp"))
715 } else {
716 None
717 };
718 if legacy_dcp_capabilities.is_some()
719 && standard_dcp_capabilities.is_some()
720 && legacy_dcp_capabilities != standard_dcp_capabilities
721 {
722 return Err(CompleteAdapterError::InvalidParams(
723 "conflicting dcp capabilities".into(),
724 ));
725 }
726 standard_dcp_capabilities.or(legacy_dcp_capabilities)
727 } else {
728 None
729 };
730 let capabilities = optional_object(capabilities, "dcp capabilities")?;
731
732 if let Some(capabilities) = capabilities {
733 set_manifest_ids(
734 capabilities.get("tools"),
735 "tool",
736 CapabilityManifest::MAX_TOOLS,
737 |id| manifest.set_tool(id),
738 )?;
739 set_manifest_ids(
740 capabilities.get("resources"),
741 "resource",
742 CapabilityManifest::MAX_RESOURCES,
743 |id| manifest.set_resource(id),
744 )?;
745 set_manifest_ids(
746 capabilities.get("prompts"),
747 "prompt",
748 CapabilityManifest::MAX_PROMPTS,
749 |id| manifest.set_prompt(id),
750 )?;
751 set_extension_ids(capabilities.get("extensions"), &mut manifest)?;
752 }
753
754 Ok(manifest)
755 }
756
757 fn is_notification_only_method(method: &str) -> bool {
758 matches!(
759 method,
760 "notifications/initialized" | "notifications/cancelled"
761 )
762 }
763
764 fn is_server_to_client_request_method(method: &str) -> bool {
765 matches!(
766 method,
767 "roots/list" | "elicitation/create" | "sampling/createMessage"
768 )
769 }
770
771 fn is_remote_shutdown_method(method: &str) -> bool {
772 matches!(
773 method,
774 "shutdown" | "exit" | "terminate" | "server/shutdown" | "notifications/shutdown"
775 )
776 }
777
778 fn requires_initialized(method: &str) -> bool {
779 matches!(
780 method,
781 "roots/list"
782 | "elicitation/create"
783 | "tools/list"
784 | "tools/call"
785 | "resources/list"
786 | "resources/read"
787 | "resources/subscribe"
788 | "resources/unsubscribe"
789 | "prompts/list"
790 | "prompts/get"
791 | "logging/setLevel"
792 | "sampling/createMessage"
793 | "completion/complete"
794 )
795 }
796
797 async fn require_initialized_for_method(
798 &self,
799 request: &JsonRpcRequest,
800 ) -> Result<(), CompleteAdapterError> {
801 if !Self::requires_initialized(&request.method) || *self.initialized_seen.read().await {
802 return Ok(());
803 }
804
805 Err(CompleteAdapterError::LifecycleNotInitialized)
806 }
807
808 fn require_request_method(
809 &self,
810 request: &JsonRpcRequest,
811 expected_method: &'static str,
812 ) -> Result<(), CompleteAdapterError> {
813 if request.method == expected_method && !request.is_notification() {
814 return Ok(());
815 }
816
817 self.audit_request(
818 SecurityAuditAction::ValidationRejected,
819 "invalid_request",
820 request,
821 );
822 Err(CompleteAdapterError::InvalidRequest(format!(
823 "{expected_method} requires a JSON-RPC request id"
824 )))
825 }
826
827 fn require_notification_method(
828 &self,
829 request: &JsonRpcRequest,
830 expected_method: &'static str,
831 ) -> Result<(), CompleteAdapterError> {
832 if request.method == expected_method && request.is_notification() {
833 return Ok(());
834 }
835
836 self.audit_request(
837 SecurityAuditAction::ValidationRejected,
838 "invalid_request",
839 request,
840 );
841 Err(CompleteAdapterError::InvalidRequest(format!(
842 "{expected_method} requires a JSON-RPC notification"
843 )))
844 }
845
846 fn list_cursor<'a>(
847 &self,
848 request: &'a JsonRpcRequest,
849 ) -> Result<Option<&'a str>, CompleteAdapterError> {
850 let Some(params) = optional_object(request.params.as_ref(), "list params")? else {
851 return Ok(None);
852 };
853
854 for key in params.keys() {
855 if key != "cursor" && key != "_meta" {
856 return Err(CompleteAdapterError::InvalidParams(
857 "unsupported list param".into(),
858 ));
859 }
860 }
861
862 if let Some(meta) = params.get("_meta") {
863 if !meta.is_object() {
864 return Err(CompleteAdapterError::InvalidParams(
865 "list _meta must be object".into(),
866 ));
867 }
868 }
869
870 match params.get("cursor") {
871 None => Ok(None),
872 Some(Value::String(cursor)) => Ok(Some(cursor.as_str())),
873 Some(_) => Err(CompleteAdapterError::InvalidParams(
874 "list cursor must be string".into(),
875 )),
876 }
877 }
878
879 fn request_id_for_audit(id: &RequestId) -> Option<String> {
880 match id {
881 RequestId::String(s) => Some(s.clone()),
882 RequestId::Number(n) => Some(n.to_string()),
883 RequestId::Null => Some("null".to_string()),
884 RequestId::Missing => None,
885 }
886 }
887
888 fn audit_request(
889 &self,
890 action: SecurityAuditAction,
891 reason: &'static str,
892 request: &JsonRpcRequest,
893 ) {
894 let mut event = SecurityAuditEvent::new(action, reason).with_method(request.method.clone());
895 if let Some(request_id) = Self::request_id_for_audit(&request.id) {
896 event = event.with_request_id(request_id);
897 }
898 self.security_audit.record(event);
899 }
900
901 async fn record_tool_call_request_id(
902 &self,
903 request: &JsonRpcRequest,
904 ) -> Result<bool, CompleteAdapterError> {
905 let Some(request_id) = Self::request_id_for_audit(&request.id) else {
906 return Ok(true);
907 };
908
909 let mut guard = self.tool_call_replay_guard.write().await;
910 Ok(guard.check_and_record(replay_key("tools/call", &request_id)))
911 }
912
913 fn server_to_client_method_response(
914 &self,
915 request: &JsonRpcRequest,
916 ) -> Result<String, CompleteAdapterError> {
917 self.audit_request(
918 SecurityAuditAction::RequestRejected,
919 "server_to_client_method",
920 request,
921 );
922 self.format_error(request.id.clone(), JsonRpcError::method_not_found())
923 }
924
925 fn adapter_error_response(
926 &self,
927 request: &JsonRpcRequest,
928 error: CompleteAdapterError,
929 ) -> Result<String, CompleteAdapterError> {
930 let (action, reason, json_error) = match error {
931 CompleteAdapterError::InvalidParams(_) => (
932 SecurityAuditAction::ValidationRejected,
933 "invalid_params",
934 JsonRpcError::invalid_params(),
935 ),
936 CompleteAdapterError::ParseError(JsonRpcParseError::InvalidJson(_)) => (
937 SecurityAuditAction::ValidationRejected,
938 "parse_error",
939 JsonRpcError::parse_error(),
940 ),
941 CompleteAdapterError::ParseError(JsonRpcParseError::RequestIdTooLarge) => (
942 SecurityAuditAction::ValidationRejected,
943 "request_id_too_large",
944 JsonRpcError::invalid_request(),
945 ),
946 CompleteAdapterError::ParseError(JsonRpcParseError::RequestIdSensitive) => (
947 SecurityAuditAction::ValidationRejected,
948 "request_id_sensitive",
949 JsonRpcError::invalid_request(),
950 ),
951 CompleteAdapterError::ParseError(JsonRpcParseError::BatchUnsupported) => (
952 SecurityAuditAction::ValidationRejected,
953 "batch_unsupported",
954 JsonRpcError::invalid_request(),
955 ),
956 CompleteAdapterError::ParseError(_) => (
957 SecurityAuditAction::ValidationRejected,
958 "invalid_request",
959 JsonRpcError::invalid_request(),
960 ),
961 CompleteAdapterError::InvalidRequest(_) => (
962 SecurityAuditAction::ValidationRejected,
963 "invalid_request",
964 JsonRpcError::invalid_request(),
965 ),
966 CompleteAdapterError::LifecycleNotInitialized => (
967 SecurityAuditAction::RequestRejected,
968 "lifecycle_not_initialized",
969 JsonRpcError::invalid_request(),
970 ),
971 CompleteAdapterError::LifecycleAlreadyInitialized => (
972 SecurityAuditAction::RequestRejected,
973 "lifecycle_already_initialized",
974 JsonRpcError::invalid_request(),
975 ),
976 CompleteAdapterError::UnknownTool(_) => (
977 SecurityAuditAction::CapabilityDenied,
978 "unknown_tool",
979 JsonRpcError::invalid_params(),
980 ),
981 CompleteAdapterError::CapabilityDenied => (
982 SecurityAuditAction::CapabilityDenied,
983 "capability_denied",
984 JsonRpcError::new(-32001, "Capability denied"),
985 ),
986 CompleteAdapterError::CapacityExceeded { .. } => (
987 SecurityAuditAction::ValidationRejected,
988 "capacity_exceeded",
989 JsonRpcError::invalid_request(),
990 ),
991 CompleteAdapterError::DcpError(_) => (
992 SecurityAuditAction::RequestRejected,
993 "dcp_error",
994 JsonRpcError::new(-32000, "Request failed"),
995 ),
996 CompleteAdapterError::ResourceError(_) | CompleteAdapterError::PromptError(_) => (
997 SecurityAuditAction::RequestRejected,
998 "not_found",
999 JsonRpcError::invalid_params(),
1000 ),
1001 CompleteAdapterError::SerializationError(_) => (
1002 SecurityAuditAction::RequestRejected,
1003 "serialization_error",
1004 JsonRpcError::internal_error(),
1005 ),
1006 };
1007
1008 self.audit_request(action, reason, request);
1009 self.format_error(request.id.clone(), json_error)
1010 }
1011
1012 async fn require_any_tool_capability(
1013 &self,
1014 _request: &JsonRpcRequest,
1015 ) -> Result<(), CompleteAdapterError> {
1016 if self.negotiated_manifest.read().await.tool_count() > 0 {
1017 Ok(())
1018 } else {
1019 Err(CompleteAdapterError::CapabilityDenied)
1020 }
1021 }
1022
1023 async fn require_any_resource_capability(
1024 &self,
1025 _request: &JsonRpcRequest,
1026 ) -> Result<(), CompleteAdapterError> {
1027 if self.negotiated_manifest.read().await.resource_count() > 0 {
1028 Ok(())
1029 } else {
1030 Err(CompleteAdapterError::CapabilityDenied)
1031 }
1032 }
1033
1034 async fn require_any_prompt_capability(
1035 &self,
1036 _request: &JsonRpcRequest,
1037 ) -> Result<(), CompleteAdapterError> {
1038 if self.negotiated_manifest.read().await.prompt_count() > 0 {
1039 Ok(())
1040 } else {
1041 Err(CompleteAdapterError::CapabilityDenied)
1042 }
1043 }
1044
1045 async fn require_resource_uri_capability(&self, uri: &str) -> Result<(), CompleteAdapterError> {
1046 let resource_id = self
1047 .resources
1048 .read()
1049 .await
1050 .handler_id_for_uri(uri)
1051 .ok_or(CompleteAdapterError::CapabilityDenied)?;
1052
1053 if self
1054 .negotiated_manifest
1055 .read()
1056 .await
1057 .has_resource(resource_id)
1058 {
1059 Ok(())
1060 } else {
1061 Err(CompleteAdapterError::CapabilityDenied)
1062 }
1063 }
1064
1065 async fn require_tool_capability(&self, tool_id: u16) -> Result<(), CompleteAdapterError> {
1066 match self.negotiated_manifest.read().await.require_tool(tool_id) {
1067 Ok(()) => Ok(()),
1068 Err(SecurityError::InsufficientCapabilities) => {
1069 Err(CompleteAdapterError::CapabilityDenied)
1070 }
1071 Err(_) => Err(CompleteAdapterError::CapabilityDenied),
1072 }
1073 }
1074
1075 pub fn parse_request(&self, json: &str) -> Result<JsonRpcRequest, CompleteAdapterError> {
1077 Ok(JsonRpcParser::parse_request_with_limit(
1078 json,
1079 self.max_request_size,
1080 )?)
1081 }
1082
1083 pub fn format_success(
1085 &self,
1086 id: RequestId,
1087 result: Value,
1088 ) -> Result<String, CompleteAdapterError> {
1089 let response = JsonRpcResponse::success(id, result);
1090 JsonRpcParser::format_response(&response)
1091 .map_err(|e| CompleteAdapterError::SerializationError(e.to_string()))
1092 }
1093
1094 pub fn format_error(
1096 &self,
1097 id: RequestId,
1098 error: JsonRpcError,
1099 ) -> Result<String, CompleteAdapterError> {
1100 let response = JsonRpcResponse::error(id, error);
1101 JsonRpcParser::format_response(&response)
1102 .map_err(|e| CompleteAdapterError::SerializationError(e.to_string()))
1103 }
1104
1105 pub async fn handle_initialize(
1111 &self,
1112 request: &JsonRpcRequest,
1113 ) -> Result<String, CompleteAdapterError> {
1114 self.require_request_method(request, "initialize")?;
1115
1116 let client_manifest = Self::client_manifest_from_initialize(request.params.as_ref())?;
1117 let requested_version = match request
1118 .params
1119 .as_ref()
1120 .and_then(|p| p.get("protocolVersion"))
1121 {
1122 Some(Value::String(version)) => version.as_str(),
1123 Some(_) => {
1124 return Err(CompleteAdapterError::InvalidParams(
1125 "protocolVersion must be string".into(),
1126 ));
1127 }
1128 None => "2024-11-05",
1129 };
1130
1131 let negotiated = self
1132 .version_negotiator
1133 .try_negotiate(requested_version)
1134 .ok_or_else(|| {
1135 CompleteAdapterError::InvalidParams("unsupported protocolVersion".into())
1136 })?;
1137
1138 {
1139 let mut initialize_seen = self.initialize_seen.write().await;
1140 if *initialize_seen {
1141 return Err(CompleteAdapterError::InvalidRequest(
1142 "already initialized".to_string(),
1143 ));
1144 }
1145 *initialize_seen = true;
1146 }
1147
1148 *self.protocol_version.write().await = negotiated;
1150
1151 let mut server_manifest = self.server_manifest.clone();
1152 for resource_id in self.resources.read().await.handler_ids() {
1153 server_manifest.set_resource(resource_id);
1154 }
1155 if let Some(authorization_policy) = self.authorization_policy.as_ref() {
1156 server_manifest = server_manifest.intersect(authorization_policy);
1157 }
1158 *self.negotiated_manifest.write().await =
1159 CapabilityManifest::negotiate(&client_manifest, &server_manifest);
1160
1161 let negotiated_manifest = self.negotiated_manifest.read().await;
1163 let mut capabilities = serde_json::json!({
1164 "logging": {}
1165 });
1166
1167 if negotiated_manifest.tool_count() > 0 {
1168 capabilities["tools"] = serde_json::json!({ "listChanged": true });
1169 }
1170 if negotiated_manifest.resource_count() > 0 {
1171 let supports_subscribe = self
1172 .resources
1173 .read()
1174 .await
1175 .any_allowed_supports_subscribe(|id| negotiated_manifest.has_resource(id));
1176 let mut resource_capabilities = Map::new();
1177 resource_capabilities.insert("listChanged".into(), Value::Bool(true));
1178 if supports_subscribe {
1179 resource_capabilities.insert("subscribe".into(), Value::Bool(true));
1180 }
1181 capabilities["resources"] = Value::Object(resource_capabilities);
1182 }
1183 if negotiated_manifest.prompt_count() > 0 {
1184 capabilities["prompts"] = serde_json::json!({ "listChanged": true });
1185 }
1186 drop(negotiated_manifest);
1187
1188 let result = serde_json::json!({
1189 "protocolVersion": negotiated.as_str(),
1190 "capabilities": capabilities,
1191 "serverInfo": {
1192 "name": self.server_name,
1193 "version": self.server_version
1194 }
1195 });
1196 self.format_success(request.id.clone(), result)
1197 }
1198
1199 pub async fn protocol_version(&self) -> super::mcp2025::ProtocolVersion {
1201 *self.protocol_version.read().await
1202 }
1203
1204 pub async fn handle_initialized(
1206 &self,
1207 request: &JsonRpcRequest,
1208 ) -> Result<Option<String>, CompleteAdapterError> {
1209 if request.method != "notifications/initialized" || !request.is_notification() {
1210 self.audit_request(
1211 SecurityAuditAction::ValidationRejected,
1212 "invalid_request",
1213 request,
1214 );
1215 return Err(CompleteAdapterError::InvalidRequest(
1216 "initialized must be a notifications/initialized notification".to_string(),
1217 ));
1218 }
1219
1220 if request.params.is_some() {
1221 return Err(CompleteAdapterError::InvalidParams(
1222 "initialized notification must not include params".to_string(),
1223 ));
1224 }
1225
1226 if !*self.initialize_seen.read().await {
1227 return Err(CompleteAdapterError::LifecycleNotInitialized);
1228 }
1229
1230 let mut initialized_seen = self.initialized_seen.write().await;
1231 if *initialized_seen {
1232 return Err(CompleteAdapterError::LifecycleAlreadyInitialized);
1233 }
1234 *initialized_seen = true;
1235
1236 Ok(None)
1238 }
1239
1240 pub fn roots(&self) -> Arc<super::mcp2025::RootsRegistry> {
1246 Arc::clone(&self.roots)
1247 }
1248
1249 pub async fn handle_roots_list(
1251 &self,
1252 request: &JsonRpcRequest,
1253 ) -> Result<String, CompleteAdapterError> {
1254 self.require_request_method(request, "roots/list")?;
1255 self.require_initialized_for_method(request).await?;
1256 self.server_to_client_method_response(request)
1257 }
1258
1259 pub fn elicitation(&self) -> Arc<super::mcp2025::ElicitationHandler> {
1265 Arc::clone(&self.elicitation)
1266 }
1267
1268 pub async fn handle_elicitation_create(
1270 &self,
1271 request: &JsonRpcRequest,
1272 ) -> Result<String, CompleteAdapterError> {
1273 self.require_request_method(request, "elicitation/create")?;
1274 self.require_initialized_for_method(request).await?;
1275 self.server_to_client_method_response(request)
1276 }
1277
1278 pub fn resource_templates(&self) -> Arc<super::mcp2025::ResourceTemplateRegistry> {
1284 Arc::clone(&self.resource_templates)
1285 }
1286
1287 pub fn notifications(&self) -> Arc<super::mcp2025::NotificationManager> {
1289 Arc::clone(&self.notifications)
1290 }
1291
1292 pub fn cancellation(&self) -> Arc<super::mcp2025::CancellationManager> {
1294 Arc::clone(&self.cancellation)
1295 }
1296
1297 pub fn progress(&self) -> Arc<super::mcp2025::ProgressTracker> {
1299 Arc::clone(&self.progress)
1300 }
1301
1302 pub async fn handle_cancelled(
1308 &self,
1309 request: &JsonRpcRequest,
1310 ) -> Result<Option<String>, CompleteAdapterError> {
1311 self.require_notification_method(request, "notifications/cancelled")?;
1312
1313 let params = request
1314 .params
1315 .as_ref()
1316 .and_then(Value::as_object)
1317 .ok_or(CompleteAdapterError::InvalidParams("missing params".into()))?;
1318
1319 if params
1320 .keys()
1321 .any(|key| key != "requestId" && key != "reason")
1322 {
1323 return Err(CompleteAdapterError::InvalidParams(
1324 "unsupported cancelled param".into(),
1325 ));
1326 }
1327
1328 let request_id = params
1329 .get("requestId")
1330 .ok_or(CompleteAdapterError::InvalidParams(
1331 "missing requestId".into(),
1332 ))?;
1333
1334 let request_id = match RequestId::try_from_json_value(request_id) {
1335 Ok(id @ (RequestId::Number(_) | RequestId::String(_))) => id,
1336 Ok(_) | Err(JsonRpcParseError::InvalidStructure) => {
1337 return Err(CompleteAdapterError::InvalidParams(
1338 "invalid requestId".into(),
1339 ));
1340 }
1341 Err(JsonRpcParseError::RequestIdTooLarge) => {
1342 return Err(CompleteAdapterError::ParseError(
1343 JsonRpcParseError::RequestIdTooLarge,
1344 ));
1345 }
1346 Err(JsonRpcParseError::RequestIdSensitive) => {
1347 return Err(CompleteAdapterError::ParseError(
1348 JsonRpcParseError::RequestIdSensitive,
1349 ));
1350 }
1351 Err(_) => {
1352 return Err(CompleteAdapterError::InvalidParams(
1353 "invalid requestId".into(),
1354 ));
1355 }
1356 };
1357
1358 let reason = match params.get("reason") {
1359 None => None,
1360 Some(Value::String(reason)) => Some(sanitize_text(reason)),
1361 Some(_) => {
1362 return Err(CompleteAdapterError::InvalidParams(
1363 "invalid cancellation reason".into(),
1364 ));
1365 }
1366 };
1367
1368 self.cancellation.cancel(&request_id, reason).await;
1370
1371 Ok(None)
1373 }
1374
1375 pub fn extract_progress_token(request: &JsonRpcRequest) -> Option<String> {
1381 request
1382 .params
1383 .as_ref()
1384 .and_then(|p| p.get("_meta"))
1385 .and_then(|m| m.get("progressToken"))
1386 .and_then(|t| t.as_str())
1387 .map(|s| s.to_string())
1388 }
1389
1390 pub fn handle_ping(&self, request: &JsonRpcRequest) -> Result<String, CompleteAdapterError> {
1396 self.require_request_method(request, "ping")?;
1397 self.format_success(request.id.clone(), serde_json::json!({}))
1399 }
1400
1401 pub async fn handle_tools_list(
1407 &self,
1408 request: &JsonRpcRequest,
1409 router: &BinaryTrieRouter,
1410 ) -> Result<String, CompleteAdapterError> {
1411 self.require_request_method(request, "tools/list")?;
1412 self.require_initialized_for_method(request).await?;
1413 let _ = self.list_cursor(request)?;
1414 self.require_any_tool_capability(request).await?;
1415
1416 let negotiated = self.negotiated_manifest.read().await;
1417 let tools: Vec<Value> = self
1418 .tool_cache
1419 .iter()
1420 .filter_map(|(name, tool_id)| {
1421 if !negotiated.has_tool(*tool_id) {
1422 return None;
1423 }
1424
1425 let (description, input_schema) = router
1426 .tool_schema(*tool_id)
1427 .map(|schema| {
1428 (
1429 schema.description.to_string(),
1430 tool_input_schema_json(&schema.input),
1431 )
1432 })
1433 .unwrap_or_else(|| {
1434 (
1435 format!("Tool: {}", name),
1436 serde_json::json!({
1437 "type": "object",
1438 "properties": {},
1439 "additionalProperties": false
1440 }),
1441 )
1442 });
1443
1444 Some(serde_json::json!({
1445 "name": name,
1446 "description": description,
1447 "inputSchema": input_schema
1448 }))
1449 })
1450 .collect();
1451
1452 self.format_success(request.id.clone(), serde_json::json!({ "tools": tools }))
1453 }
1454
1455 pub async fn handle_tools_call(
1457 &self,
1458 request: &JsonRpcRequest,
1459 router: &BinaryTrieRouter,
1460 ) -> Result<String, CompleteAdapterError> {
1461 self.require_request_method(request, "tools/call")?;
1462 self.require_initialized_for_method(request).await?;
1463
1464 let params = tools_call_params_object(request.params.as_ref())?;
1465
1466 let tool_name = params.get("name").and_then(|v| v.as_str()).ok_or(
1467 CompleteAdapterError::InvalidParams("missing tool name".into()),
1468 )?;
1469
1470 let tool_id = self
1471 .tool_cache
1472 .get(tool_name)
1473 .ok_or(CompleteAdapterError::CapabilityDenied)?;
1474
1475 self.require_tool_capability(*tool_id).await?;
1476
1477 let schema = router
1478 .tool_schema(*tool_id)
1479 .ok_or(CompleteAdapterError::DcpError(DCPError::ToolNotFound))?;
1480 let (args_bytes, arg_layout) =
1481 encode_json_arguments(&schema.input, params.get("arguments"))?;
1482 if !self.record_tool_call_request_id(request).await? {
1483 self.audit_request(
1484 SecurityAuditAction::ReplayRejected,
1485 "request_replay",
1486 request,
1487 );
1488 return self.format_error(
1489 request.id.clone(),
1490 JsonRpcError::new(-32002, "Request replay rejected"),
1491 );
1492 }
1493
1494 let shared_args = SharedArgs::new(&args_bytes, arg_layout);
1495 let capabilities = self.negotiated_manifest.read().await.clone();
1496 let result = router
1497 .execute_authorized(&capabilities, *tool_id, &shared_args)
1498 .map_err(|error| match error {
1499 SecurityError::ValidationFailed => {
1500 CompleteAdapterError::InvalidParams("argument validation failed".into())
1501 }
1502 SecurityError::InsufficientCapabilities => CompleteAdapterError::CapabilityDenied,
1503 _ => CompleteAdapterError::CapabilityDenied,
1504 })?;
1505 let result_value = match result {
1506 ToolResult::Success(data) => serde_json::from_slice(&data)
1507 .unwrap_or(Value::String(String::from_utf8_lossy(&data).into())),
1508 ToolResult::Empty => Value::Null,
1509 ToolResult::Error(e) => serde_json::json!({"error": e.to_string()}),
1510 };
1511
1512 self.format_success(request.id.clone(), serde_json::json!({
1513 "content": [{ "type": "text", "text": serde_json::to_string(&result_value).unwrap_or_default() }]
1514 }))
1515 }
1516
1517 pub async fn handle_resources_list(
1523 &self,
1524 request: &JsonRpcRequest,
1525 ) -> Result<String, CompleteAdapterError> {
1526 self.require_request_method(request, "resources/list")?;
1527 self.require_initialized_for_method(request).await?;
1528 let cursor = self.list_cursor(request)?;
1529 self.require_any_resource_capability(request).await?;
1530
1531 let negotiated = self.negotiated_manifest.read().await;
1532 let registry = self.resources.read().await;
1533 let list =
1534 registry.list_allowed(cursor, |resource_id| negotiated.has_resource(resource_id))?;
1535 let allowed_templates: HashSet<String> = registry
1536 .allowed_uri_templates(|resource_id| negotiated.has_resource(resource_id))
1537 .into_iter()
1538 .collect();
1539 drop(registry);
1540 drop(negotiated);
1541
1542 let mut resources = Vec::new();
1543 for resource in list.resources {
1544 if !self.roots.allows_uri(&resource.uri).await {
1545 continue;
1546 }
1547
1548 resources.push(serde_json::json!({
1549 "uri": resource.uri,
1550 "name": resource.name,
1551 "description": resource.description,
1552 "mimeType": resource.mime_type
1553 }));
1554 }
1555
1556 let mut result = serde_json::json!({ "resources": resources });
1557 if let Some(cursor) = list.next_cursor {
1558 result["nextCursor"] = Value::String(cursor);
1559 }
1560
1561 let version = *self.protocol_version.read().await;
1563 if version.supports_roots() {
1564 let templates: Vec<_> = self
1565 .resource_templates
1566 .list()
1567 .await
1568 .into_iter()
1569 .filter(|template| allowed_templates.contains(&template.uri_template))
1570 .collect();
1571 if !templates.is_empty() {
1572 result["resourceTemplates"] =
1573 serde_json::to_value(&templates).unwrap_or(Value::Array(vec![]));
1574 }
1575 }
1576
1577 self.format_success(request.id.clone(), result)
1578 }
1579
1580 pub async fn handle_resources_read(
1582 &self,
1583 request: &JsonRpcRequest,
1584 ) -> Result<String, CompleteAdapterError> {
1585 self.require_request_method(request, "resources/read")?;
1586 self.require_initialized_for_method(request).await?;
1587
1588 let uri = request
1589 .params
1590 .as_ref()
1591 .and_then(|p| p.get("uri"))
1592 .and_then(|v| v.as_str())
1593 .ok_or(CompleteAdapterError::InvalidParams("missing uri".into()))?;
1594
1595 self.require_resource_uri_capability(uri).await?;
1596
1597 if !self.roots.allows_uri(uri).await {
1598 return Err(CompleteAdapterError::InvalidParams(
1599 "resource outside configured roots".into(),
1600 ));
1601 }
1602
1603 let registry = self.resources.read().await;
1604 let content = registry.read(uri)?;
1605
1606 let content_value = match content {
1607 ResourceContent::Text {
1608 uri,
1609 mime_type,
1610 text,
1611 } => serde_json::json!({
1612 "uri": uri,
1613 "mimeType": mime_type,
1614 "text": text
1615 }),
1616 ResourceContent::Blob {
1617 uri,
1618 mime_type,
1619 blob,
1620 } => serde_json::json!({
1621 "uri": uri,
1622 "mimeType": mime_type,
1623 "blob": blob
1624 }),
1625 };
1626
1627 self.format_success(
1628 request.id.clone(),
1629 serde_json::json!({
1630 "contents": [content_value]
1631 }),
1632 )
1633 }
1634
1635 pub async fn handle_resources_subscribe(
1637 &self,
1638 request: &JsonRpcRequest,
1639 ) -> Result<String, CompleteAdapterError> {
1640 self.require_request_method(request, "resources/subscribe")?;
1641 self.require_initialized_for_method(request).await?;
1642
1643 let uri = request
1644 .params
1645 .as_ref()
1646 .and_then(|p| p.get("uri"))
1647 .and_then(|v| v.as_str())
1648 .ok_or(CompleteAdapterError::InvalidParams("missing uri".into()))?;
1649
1650 self.require_resource_uri_capability(uri).await?;
1651
1652 if !self.roots.allows_uri(uri).await {
1653 return Err(CompleteAdapterError::InvalidParams(
1654 "resource outside configured roots".into(),
1655 ));
1656 }
1657
1658 self.resources
1659 .read()
1660 .await
1661 .ensure_subscribable(uri)
1662 .map_err(CompleteAdapterError::from)?;
1663
1664 if !self.subscriptions.subscribe(uri, "default").await {
1666 return Err(CompleteAdapterError::CapacityExceeded {
1667 kind: "subscription",
1668 max: super::mcp2025::SubscriptionTracker::DEFAULT_MAX_SUBSCRIBERS_PER_RESOURCE,
1669 });
1670 }
1671
1672 self.format_success(request.id.clone(), serde_json::json!({}))
1673 }
1674
1675 pub async fn handle_resources_unsubscribe(
1677 &self,
1678 request: &JsonRpcRequest,
1679 ) -> Result<String, CompleteAdapterError> {
1680 self.require_request_method(request, "resources/unsubscribe")?;
1681 self.require_initialized_for_method(request).await?;
1682
1683 let uri = request
1684 .params
1685 .as_ref()
1686 .and_then(|p| p.get("uri"))
1687 .and_then(|v| v.as_str())
1688 .ok_or(CompleteAdapterError::InvalidParams("missing uri".into()))?;
1689
1690 self.require_resource_uri_capability(uri).await?;
1691
1692 if !self.roots.allows_uri(uri).await {
1693 return Err(CompleteAdapterError::InvalidParams(
1694 "resource outside configured roots".into(),
1695 ));
1696 }
1697
1698 self.subscriptions.unsubscribe(uri, "default").await;
1700
1701 self.format_success(request.id.clone(), serde_json::json!({}))
1702 }
1703
1704 pub async fn handle_prompts_list(
1710 &self,
1711 request: &JsonRpcRequest,
1712 ) -> Result<String, CompleteAdapterError> {
1713 self.require_request_method(request, "prompts/list")?;
1714 self.require_initialized_for_method(request).await?;
1715 let _ = self.list_cursor(request)?;
1716 self.require_any_prompt_capability(request).await?;
1717
1718 let negotiated = self.negotiated_manifest.read().await;
1719 let prompts: Vec<Value> = self
1720 .prompts
1721 .iter()
1722 .filter(|(name, _)| {
1723 self.prompt_ids
1724 .get(*name)
1725 .map(|prompt_id| negotiated.has_prompt(*prompt_id))
1726 .unwrap_or(false)
1727 })
1728 .map(|(_, p)| {
1729 let args: Vec<Value> = p
1730 .arguments
1731 .iter()
1732 .map(|a| {
1733 serde_json::json!({
1734 "name": a.name,
1735 "description": a.description,
1736 "required": a.required
1737 })
1738 })
1739 .collect();
1740 serde_json::json!({
1741 "name": p.name,
1742 "description": p.description,
1743 "arguments": args
1744 })
1745 })
1746 .collect();
1747
1748 self.format_success(
1749 request.id.clone(),
1750 serde_json::json!({ "prompts": prompts }),
1751 )
1752 }
1753
1754 pub async fn handle_prompts_get(
1756 &self,
1757 request: &JsonRpcRequest,
1758 ) -> Result<String, CompleteAdapterError> {
1759 self.require_request_method(request, "prompts/get")?;
1760 self.require_initialized_for_method(request).await?;
1761
1762 let params = request
1763 .params
1764 .as_ref()
1765 .ok_or(CompleteAdapterError::InvalidParams("missing params".into()))?;
1766
1767 let name = params.get("name").and_then(|v| v.as_str()).ok_or(
1768 CompleteAdapterError::InvalidParams("missing prompt name".into()),
1769 )?;
1770
1771 let prompt_id = self
1772 .prompt_ids
1773 .get(name)
1774 .copied()
1775 .ok_or(CompleteAdapterError::CapabilityDenied)?;
1776 if !self.negotiated_manifest.read().await.has_prompt(prompt_id) {
1777 return Err(CompleteAdapterError::CapabilityDenied);
1778 }
1779
1780 let template = self
1781 .prompts
1782 .get(name)
1783 .ok_or_else(|| CompleteAdapterError::PromptError("prompt not found".to_string()))?;
1784
1785 let args = if let Some(arguments) = params.get("arguments") {
1786 let arguments = arguments.as_object().ok_or_else(|| {
1787 CompleteAdapterError::InvalidParams("invalid prompt arguments".into())
1788 })?;
1789 let mut args = HashMap::new();
1790 for (key, value) in arguments {
1791 if !template.has_argument(key) {
1792 return Err(CompleteAdapterError::InvalidParams(
1793 "unknown prompt argument".into(),
1794 ));
1795 }
1796 let value = value.as_str().ok_or_else(|| {
1797 CompleteAdapterError::InvalidParams("invalid prompt argument".into())
1798 })?;
1799 args.insert(key.clone(), value.to_string());
1800 }
1801 args
1802 } else {
1803 HashMap::new()
1804 };
1805
1806 let rendered = template.render(&args)?;
1807
1808 self.format_success(
1809 request.id.clone(),
1810 serde_json::json!({
1811 "description": template.description,
1812 "messages": [{
1813 "role": "user",
1814 "content": { "type": "text", "text": rendered }
1815 }]
1816 }),
1817 )
1818 }
1819
1820 pub async fn handle_logging_set_level(
1826 &self,
1827 request: &JsonRpcRequest,
1828 ) -> Result<String, CompleteAdapterError> {
1829 self.require_request_method(request, "logging/setLevel")?;
1830 self.require_initialized_for_method(request).await?;
1831
1832 let level_str = request
1833 .params
1834 .as_ref()
1835 .and_then(|p| p.get("level"))
1836 .and_then(|v| v.as_str())
1837 .ok_or(CompleteAdapterError::InvalidParams("missing level".into()))?;
1838
1839 let level = LogLevel::from_str(level_str)
1840 .ok_or_else(|| CompleteAdapterError::InvalidParams("invalid log level".to_string()))?;
1841
1842 *self.log_level.write().await = level;
1843
1844 self.format_success(request.id.clone(), serde_json::json!({}))
1845 }
1846
1847 pub async fn handle_sampling_create_message(
1853 &self,
1854 request: &JsonRpcRequest,
1855 ) -> Result<String, CompleteAdapterError> {
1856 self.require_request_method(request, "sampling/createMessage")?;
1857 self.require_initialized_for_method(request).await?;
1858 self.server_to_client_method_response(request)
1859 }
1860
1861 pub async fn handle_completion_complete(
1867 &self,
1868 request: &JsonRpcRequest,
1869 ) -> Result<String, CompleteAdapterError> {
1870 self.require_request_method(request, "completion/complete")?;
1871 self.require_initialized_for_method(request).await?;
1872
1873 let params = request
1874 .params
1875 .as_ref()
1876 .and_then(Value::as_object)
1877 .ok_or_else(|| CompleteAdapterError::InvalidParams("missing params".into()))?;
1878
1879 let reference = params
1880 .get("ref")
1881 .and_then(Value::as_object)
1882 .ok_or_else(|| CompleteAdapterError::InvalidParams("invalid ref".into()))?;
1883 let ref_type = reference
1884 .get("type")
1885 .and_then(Value::as_str)
1886 .ok_or_else(|| CompleteAdapterError::InvalidParams("invalid ref".into()))?;
1887
1888 let argument = params
1889 .get("argument")
1890 .and_then(Value::as_object)
1891 .ok_or_else(|| CompleteAdapterError::InvalidParams("invalid argument".into()))?;
1892 let argument_name = argument
1893 .get("name")
1894 .and_then(Value::as_str)
1895 .ok_or_else(|| CompleteAdapterError::InvalidParams("invalid argument".into()))?;
1896 let argument_value = argument
1897 .get("value")
1898 .and_then(Value::as_str)
1899 .ok_or_else(|| CompleteAdapterError::InvalidParams("invalid argument".into()))?;
1900 if argument_name.is_empty() || argument_value.is_empty() {
1901 return Err(CompleteAdapterError::InvalidParams(
1902 "invalid argument".into(),
1903 ));
1904 }
1905
1906 match ref_type {
1907 "ref/prompt" => {
1908 let name = reference
1909 .get("name")
1910 .and_then(Value::as_str)
1911 .ok_or_else(|| CompleteAdapterError::InvalidParams("invalid ref".into()))?;
1912 if name.is_empty() {
1913 return Err(CompleteAdapterError::InvalidParams("invalid ref".into()));
1914 }
1915 let prompt_id = self
1916 .prompt_ids
1917 .get(name)
1918 .copied()
1919 .ok_or(CompleteAdapterError::CapabilityDenied)?;
1920 if !self.negotiated_manifest.read().await.has_prompt(prompt_id) {
1921 return Err(CompleteAdapterError::CapabilityDenied);
1922 }
1923 let template = self.prompts.get(name).ok_or_else(|| {
1924 CompleteAdapterError::PromptError("prompt not found".to_string())
1925 })?;
1926 let argument_name = params
1927 .get("argument")
1928 .and_then(Value::as_object)
1929 .and_then(|argument| argument.get("name"))
1930 .and_then(Value::as_str)
1931 .ok_or_else(|| {
1932 CompleteAdapterError::InvalidParams("invalid argument".into())
1933 })?;
1934 if !template.has_argument(argument_name) {
1935 return Err(CompleteAdapterError::InvalidParams(
1936 "unknown prompt argument".into(),
1937 ));
1938 }
1939 }
1940 "ref/resource" => {
1941 let uri = reference
1942 .get("uri")
1943 .and_then(Value::as_str)
1944 .ok_or_else(|| CompleteAdapterError::InvalidParams("invalid ref".into()))?;
1945 if uri.is_empty() {
1946 return Err(CompleteAdapterError::InvalidParams("invalid ref".into()));
1947 }
1948 self.require_resource_uri_capability(uri).await?;
1949 if !self.roots.allows_uri(uri).await {
1950 return Err(CompleteAdapterError::InvalidParams(
1951 "resource outside configured roots".into(),
1952 ));
1953 }
1954 }
1955 _ => return Err(CompleteAdapterError::InvalidParams("invalid ref".into())),
1956 }
1957
1958 self.format_success(
1960 request.id.clone(),
1961 serde_json::json!({
1962 "completion": { "values": [], "hasMore": false }
1963 }),
1964 )
1965 }
1966
1967 pub async fn handle_request(
1973 &self,
1974 json: &str,
1975 router: &BinaryTrieRouter,
1976 ) -> Result<Option<String>, CompleteAdapterError> {
1977 let request = match self.parse_request(json) {
1978 Ok(request) => request,
1979 Err(error) => {
1980 let (reason, json_error) = match error {
1981 CompleteAdapterError::ParseError(JsonRpcParseError::InvalidJson(_)) => {
1982 ("parse_error", JsonRpcError::parse_error())
1983 }
1984 CompleteAdapterError::ParseError(JsonRpcParseError::RequestTooLarge) => {
1985 ("request_too_large", JsonRpcError::invalid_request())
1986 }
1987 CompleteAdapterError::ParseError(JsonRpcParseError::RequestIdTooLarge) => {
1988 ("request_id_too_large", JsonRpcError::invalid_request())
1989 }
1990 CompleteAdapterError::ParseError(JsonRpcParseError::RequestIdSensitive) => {
1991 ("request_id_sensitive", JsonRpcError::invalid_request())
1992 }
1993 CompleteAdapterError::ParseError(JsonRpcParseError::BatchUnsupported) => {
1994 ("batch_unsupported", JsonRpcError::invalid_request())
1995 }
1996 _ => ("invalid_request", JsonRpcError::invalid_request()),
1997 };
1998 self.security_audit.record(SecurityAuditEvent::new(
1999 SecurityAuditAction::ValidationRejected,
2000 reason,
2001 ));
2002 return self
2003 .format_error(RequestId::Null, json_error)
2004 .map(|response| Some(response));
2005 }
2006 };
2007
2008 let is_notification = request.is_notification();
2010
2011 if Self::is_remote_shutdown_method(&request.method) {
2012 self.audit_request(
2013 SecurityAuditAction::ShutdownRejected,
2014 "remote_shutdown_rejected",
2015 &request,
2016 );
2017
2018 if is_notification {
2019 return Ok(None);
2020 }
2021
2022 return self
2023 .format_error(request.id.clone(), JsonRpcError::method_not_found())
2024 .map(Some);
2025 }
2026
2027 if is_notification && !Self::is_notification_only_method(&request.method) {
2028 self.audit_request(
2029 SecurityAuditAction::RequestRejected,
2030 "notification_not_allowed",
2031 &request,
2032 );
2033 return Ok(None);
2034 }
2035
2036 if !is_notification && Self::is_notification_only_method(&request.method) {
2037 self.audit_request(
2038 SecurityAuditAction::ValidationRejected,
2039 "notification_method_with_id",
2040 &request,
2041 );
2042 return self
2043 .format_error(request.id.clone(), JsonRpcError::invalid_request())
2044 .map(Some);
2045 }
2046
2047 if !is_notification {
2048 if let Err(error) = self.require_initialized_for_method(&request).await {
2049 return self.adapter_error_response(&request, error).map(Some);
2050 }
2051 }
2052
2053 let result = match request.method.as_str() {
2054 "initialize" => self.handle_initialize(&request).await.map(Some),
2056 "notifications/initialized" => self.handle_initialized(&request).await,
2057
2058 "ping" => self.handle_ping(&request).map(Some),
2060
2061 method if Self::is_server_to_client_request_method(method) => {
2063 self.server_to_client_method_response(&request).map(Some)
2064 }
2065
2066 "tools/list" => match self.require_any_tool_capability(&request).await {
2068 Ok(()) => self.handle_tools_list(&request, router).await.map(Some),
2069 Err(error) => Err(error),
2070 },
2071 "tools/call" => self.handle_tools_call(&request, router).await.map(Some),
2072
2073 "resources/list" => match self.require_any_resource_capability(&request).await {
2075 Ok(()) => self.handle_resources_list(&request).await.map(Some),
2076 Err(error) => Err(error),
2077 },
2078 "resources/read" => match self.require_any_resource_capability(&request).await {
2079 Ok(()) => self.handle_resources_read(&request).await.map(Some),
2080 Err(error) => Err(error),
2081 },
2082 "resources/subscribe" => match self.require_any_resource_capability(&request).await {
2083 Ok(()) => self.handle_resources_subscribe(&request).await.map(Some),
2084 Err(error) => Err(error),
2085 },
2086 "resources/unsubscribe" => match self.require_any_resource_capability(&request).await {
2087 Ok(()) => self.handle_resources_unsubscribe(&request).await.map(Some),
2088 Err(error) => Err(error),
2089 },
2090
2091 "prompts/list" => match self.require_any_prompt_capability(&request).await {
2093 Ok(()) => self.handle_prompts_list(&request).await.map(Some),
2094 Err(error) => Err(error),
2095 },
2096 "prompts/get" => match self.require_any_prompt_capability(&request).await {
2097 Ok(()) => self.handle_prompts_get(&request).await.map(Some),
2098 Err(error) => Err(error),
2099 },
2100
2101 "logging/setLevel" => self.handle_logging_set_level(&request).await.map(Some),
2103
2104 "completion/complete" => self.handle_completion_complete(&request).await.map(Some),
2106
2107 "notifications/cancelled" => self.handle_cancelled(&request).await,
2109
2110 _ => {
2112 self.audit_request(
2113 SecurityAuditAction::RequestRejected,
2114 "method_not_found",
2115 &request,
2116 );
2117 self.format_error(request.id.clone(), JsonRpcError::method_not_found())
2118 .map(Some)
2119 }
2120 };
2121
2122 if is_notification {
2124 if let Err(error) = result {
2125 let _ = self.adapter_error_response(&request, error);
2126 }
2127 Ok(None)
2128 } else if let Err(error) = result {
2129 self.adapter_error_response(&request, error)
2130 .map(|response| Some(response))
2131 } else {
2132 result
2133 }
2134 }
2135}
2136
2137impl Default for CompleteMcpAdapter {
2138 fn default() -> Self {
2139 Self::new()
2140 }
2141}
2142
2143#[cfg(test)]
2144mod tests {
2145 use super::*;
2146
2147 #[test]
2148 fn test_prompt_template_render() {
2149 let template = PromptTemplate::new("test", "Test prompt", "Hello {{name}}!")
2150 .with_argument("name", "The name", true);
2151
2152 let mut args = HashMap::new();
2153 args.insert("name".to_string(), "World".to_string());
2154
2155 let rendered = template.render(&args).unwrap();
2156 assert_eq!(rendered, "Hello World!");
2157 }
2158
2159 #[test]
2160 fn test_prompt_template_missing_required() {
2161 let template = PromptTemplate::new("test", "Test prompt", "Hello {{name}}!")
2162 .with_argument("name", "The name", true);
2163
2164 let args = HashMap::new();
2165 let result = template.render(&args);
2166 assert!(result.is_err());
2167 }
2168
2169 #[tokio::test]
2170 async fn test_complete_adapter_initialize() {
2171 let adapter = CompleteMcpAdapter::new();
2172
2173 let request = JsonRpcRequest::new("initialize", None, RequestId::Number(1));
2175 let response = adapter.handle_initialize(&request).await.unwrap();
2176 let parsed = JsonRpcParser::parse_response(&response).unwrap();
2177
2178 assert!(parsed.is_success());
2179 let result = parsed.result.unwrap();
2180 assert_eq!(result["protocolVersion"], "2024-11-05");
2182 assert!(result["capabilities"]["roots"].is_null());
2184 assert!(result["capabilities"]["elicitation"].is_null());
2185 }
2186
2187 #[tokio::test]
2188 async fn test_complete_adapter_initialize_version_negotiation() {
2189 let adapter = CompleteMcpAdapter::new();
2191 let request = JsonRpcRequest::new(
2192 "initialize",
2193 Some(serde_json::json!({"protocolVersion": "2024-11-05"})),
2194 RequestId::Number(1),
2195 );
2196 let response = adapter.handle_initialize(&request).await.unwrap();
2197 let parsed = JsonRpcParser::parse_response(&response).unwrap();
2198 let result = parsed.result.unwrap();
2199 assert_eq!(result["protocolVersion"], "2024-11-05");
2200 assert!(result["capabilities"]["roots"].is_null());
2202
2203 let adapter = CompleteMcpAdapter::new();
2205 let request = JsonRpcRequest::new(
2206 "initialize",
2207 Some(serde_json::json!({"protocolVersion": "2025-03-26"})),
2208 RequestId::Number(2),
2209 );
2210 let response = adapter.handle_initialize(&request).await.unwrap();
2211 let parsed = JsonRpcParser::parse_response(&response).unwrap();
2212 let result = parsed.result.unwrap();
2213 assert_eq!(result["protocolVersion"], "2025-03-26");
2214 assert!(result["capabilities"]["roots"].is_null());
2216 assert!(result["capabilities"]["elicitation"].is_null());
2217
2218 let adapter = CompleteMcpAdapter::new();
2220 let request = JsonRpcRequest::new(
2221 "initialize",
2222 Some(serde_json::json!({"protocolVersion": "2025-06-18"})),
2223 RequestId::Number(3),
2224 );
2225 let response = adapter.handle_initialize(&request).await.unwrap();
2226 let parsed = JsonRpcParser::parse_response(&response).unwrap();
2227 let result = parsed.result.unwrap();
2228 assert_eq!(result["protocolVersion"], "2025-06-18");
2229 assert!(result["capabilities"]["roots"].is_null());
2231 assert!(result["capabilities"]["elicitation"].is_null());
2232 }
2233
2234 #[tokio::test]
2235 async fn test_complete_adapter_prompts() {
2236 let mut adapter = CompleteMcpAdapter::new();
2237 adapter
2238 .register_prompt(
2239 PromptTemplate::new("greet", "Greeting prompt", "Hello {{name}}!").with_argument(
2240 "name",
2241 "Name to greet",
2242 true,
2243 ),
2244 )
2245 .unwrap();
2246
2247 let init = JsonRpcRequest::new(
2248 "initialize",
2249 Some(serde_json::json!({"capabilities": {"dcp": {"prompts": [0]}}})),
2250 RequestId::Number(0),
2251 );
2252 adapter.handle_initialize(&init).await.unwrap();
2253 adapter
2254 .handle_initialized(&JsonRpcRequest::notification(
2255 "notifications/initialized",
2256 None,
2257 ))
2258 .await
2259 .unwrap();
2260
2261 let request = JsonRpcRequest::new("prompts/list", None, RequestId::Number(1));
2263 let response = adapter.handle_prompts_list(&request).await.unwrap();
2264 let parsed = JsonRpcParser::parse_response(&response).unwrap();
2265 assert!(parsed.is_success());
2266
2267 let request = JsonRpcRequest::new(
2269 "prompts/get",
2270 Some(serde_json::json!({"name": "greet", "arguments": {"name": "World"}})),
2271 RequestId::Number(2),
2272 );
2273 let response = adapter.handle_prompts_get(&request).await.unwrap();
2274 let parsed = JsonRpcParser::parse_response(&response).unwrap();
2275 assert!(parsed.is_success());
2276 }
2277
2278 #[tokio::test]
2279 async fn test_notification_no_response() {
2280 let adapter = CompleteMcpAdapter::new();
2281 let router = BinaryTrieRouter::new();
2282
2283 let json = r#"{"jsonrpc":"2.0","method":"notifications/initialized"}"#;
2285 let result = adapter.handle_request(json, &router).await.unwrap();
2286 assert!(result.is_none());
2287 }
2288
2289 #[tokio::test]
2290 async fn test_unknown_method() {
2291 let adapter = CompleteMcpAdapter::new();
2292 let router = BinaryTrieRouter::new();
2293
2294 let json = r#"{"jsonrpc":"2.0","method":"unknown/method","id":1}"#;
2295 let result = adapter.handle_request(json, &router).await.unwrap();
2296
2297 let response = JsonRpcParser::parse_response(&result.unwrap()).unwrap();
2298 assert!(response.is_error());
2299 assert_eq!(response.error.unwrap().code, -32601);
2300 }
2301}