1use rmcp::model::{ElicitResult, ElicitationAction};
5
6use super::{Agent, Channel, LlmProvider};
7
8impl<C: Channel> Agent<C> {
9 #[tracing::instrument(skip_all, name = "core.agent.handle_mcp_command")]
15 pub(super) async fn handle_mcp_command(
16 &mut self,
17 args: &str,
18 ) -> Result<String, super::error::AgentError> {
19 let parts: Vec<&str> = args.split_whitespace().collect();
20 match parts.first().copied() {
21 Some("add") => self.handle_mcp_add(&parts[1..]).await,
22 Some("list") => self.handle_mcp_list().await,
23 Some("tools") => Ok(self.handle_mcp_tools(parts.get(1).copied())),
24 Some("remove") => self.handle_mcp_remove(parts.get(1).copied()).await,
25 _ => Ok("Usage: /mcp add|list|tools|remove".to_owned()),
26 }
27 }
28
29 async fn handle_mcp_add(&mut self, args: &[&str]) -> Result<String, super::error::AgentError> {
30 if args.len() < 2 {
31 return Ok("Usage: /mcp add <id> <command> [args...] | /mcp add <id> <url>".to_owned());
32 }
33
34 let Some(manager) = self.services.mcp.manager.clone() else {
36 return Ok("MCP is not enabled.".to_owned());
37 };
38
39 let target = args[1];
40 if let Some(err) = validate_mcp_command(target, &self.services.mcp.allowed_commands) {
41 return Ok(err);
42 }
43
44 let current_count = manager.list_servers().await.len();
46 if current_count >= self.services.mcp.max_dynamic {
47 return Ok(format!(
48 "Server limit reached ({}/{}).",
49 current_count, self.services.mcp.max_dynamic
50 ));
51 }
52
53 let entry = build_server_entry(args[0], target, &args[2..]);
54
55 match manager.add_server(&entry).await {
56 Ok(tools) => {
57 let count = tools.len();
58 self.services
59 .mcp
60 .server_outcomes
61 .push(zeph_mcp::ServerConnectOutcome {
62 id: entry.id.clone(),
63 connected: true,
64 tool_count: count,
65 error: String::new(),
66 input_schemas_dropped: 0,
70 output_schemas_dropped: 0,
71 });
72 self.services.mcp.tools.extend(tools);
73 self.services.mcp.sync_executor_tools();
74 self.services.mcp.pruning_cache.reset();
75 self.services.mcp.pending_semantic_rebuild = true;
78 self.update_mcp_metrics();
79 Ok(format!(
80 "Connected MCP server '{}' ({count} tool(s))",
81 entry.id
82 ))
83 }
84 Err(e) => {
85 tracing::warn!(server_id = entry.id, "MCP add failed: {e:#}");
86 Ok(format!("Failed to connect server '{}': {e}", entry.id))
87 }
88 }
89 }
90
91 async fn handle_mcp_list(&mut self) -> Result<String, super::error::AgentError> {
92 use std::fmt::Write;
93
94 let Some(manager) = self.services.mcp.manager.clone() else {
95 return Ok("MCP is not enabled.".to_owned());
96 };
97
98 let server_ids = manager.list_servers().await;
99 if server_ids.is_empty() {
100 return Ok("No MCP servers connected.".to_owned());
101 }
102
103 let mut output = String::from("Connected MCP servers:\n");
104 let mut total = 0usize;
105 for id in &server_ids {
106 let count = self
107 .services
108 .mcp
109 .tools
110 .iter()
111 .filter(|t| t.server_id == *id)
112 .count();
113 total += count;
114 let _ = writeln!(output, "- {id} ({count} tools)");
115 }
116 let _ = write!(output, "Total: {total} tool(s)");
117
118 Ok(output)
119 }
120
121 fn handle_mcp_tools(&mut self, server_id: Option<&str>) -> String {
122 use std::fmt::Write;
123
124 let Some(server_id) = server_id else {
125 return "Usage: /mcp tools <server_id>".to_owned();
126 };
127
128 let tools: Vec<_> = self
129 .services
130 .mcp
131 .tools
132 .iter()
133 .filter(|t| t.server_id == server_id)
134 .collect();
135
136 if tools.is_empty() {
137 return format!("No tools found for server '{server_id}'.");
138 }
139
140 let mut output = format!("Tools for '{server_id}' ({} total):\n", tools.len());
141 for t in &tools {
142 if t.description.is_empty() {
143 let _ = writeln!(output, "- {}", t.name);
144 } else {
145 let _ = writeln!(output, "- {} — {}", t.name, t.description);
146 }
147 }
148 output
149 }
150
151 async fn handle_mcp_remove(
152 &mut self,
153 server_id: Option<&str>,
154 ) -> Result<String, super::error::AgentError> {
155 let Some(server_id) = server_id else {
156 return Ok("Usage: /mcp remove <id>".to_owned());
157 };
158
159 let Some(manager) = self.services.mcp.manager.clone() else {
161 return Ok("MCP is not enabled.".to_owned());
162 };
163
164 match manager.remove_server(server_id).await {
165 Ok(()) => {
166 let before = self.services.mcp.tools.len();
167 self.services.mcp.tools.retain(|t| t.server_id != server_id);
168 let removed = before - self.services.mcp.tools.len();
169 self.services
170 .mcp
171 .server_outcomes
172 .retain(|o| o.id != server_id);
173 self.services.mcp.sync_executor_tools();
174 self.services.mcp.pruning_cache.reset();
175 self.services.mcp.pending_semantic_rebuild = true;
178 self.update_mcp_metrics();
179 let sid = server_id.to_owned();
180 self.update_metrics(|m| {
181 m.active_mcp_tools
182 .retain(|name| !name.starts_with(&format!("{sid}:")));
183 });
184 Ok(format!(
185 "Disconnected MCP server '{server_id}' (removed {removed} tools)"
186 ))
187 }
188 Err(e) => {
189 tracing::warn!(server_id, "MCP remove failed: {e:#}");
190 Ok(format!("Failed to remove server '{server_id}': {e}"))
191 }
192 }
193 }
194
195 pub(super) async fn append_mcp_prompt(&mut self, query: &str, system_prompt: &mut String) {
196 let matched_tools = self.match_mcp_tools(query).await;
197 let active_mcp: Vec<String> = matched_tools
198 .iter()
199 .map(zeph_mcp::McpTool::qualified_name)
200 .collect();
201 let mcp_total = self.services.mcp.tools.len();
202 let (mcp_server_count, mcp_connected_count) =
203 if self.services.mcp.server_outcomes.is_empty() {
204 let connected = self
205 .services
206 .mcp
207 .tools
208 .iter()
209 .map(|t| &t.server_id)
210 .collect::<std::collections::HashSet<_>>()
211 .len();
212 (connected, connected)
213 } else {
214 let total = self.services.mcp.server_outcomes.len();
215 let connected = self
216 .services
217 .mcp
218 .server_outcomes
219 .iter()
220 .filter(|o| o.connected)
221 .count();
222 (total, connected)
223 };
224 self.update_metrics(|m| {
225 m.active_mcp_tools = active_mcp;
226 m.mcp_tool_count = mcp_total;
227 m.mcp_server_count = mcp_server_count;
228 m.mcp_connected_count = mcp_connected_count;
229 });
230 if let Some(ref manager) = self.services.mcp.manager {
231 let instructions = manager.all_server_instructions().await;
232 if !instructions.is_empty() {
233 system_prompt.push_str("\n\n");
234 system_prompt.push_str(&instructions);
235 }
236 }
237 if !matched_tools.is_empty() {
238 let tool_names: Vec<&str> = matched_tools.iter().map(|t| t.name.as_str()).collect();
239 tracing::debug!(
240 skills = ?self.services.skill.active_skill_names,
241 mcp_tools = ?tool_names,
242 "matched items"
243 );
244 let tools_prompt = zeph_mcp::format_mcp_tools_prompt(&matched_tools);
245 if !tools_prompt.is_empty() {
246 system_prompt.push_str("\n\n");
247 system_prompt.push_str(&tools_prompt);
248 }
249 }
250 }
251
252 async fn match_mcp_tools(&self, query: &str) -> Vec<zeph_mcp::McpTool> {
253 let Some(ref registry) = self.services.mcp.registry else {
254 return self.services.mcp.tools.clone();
255 };
256 let provider = self.embedding_provider.clone();
257 let hits = registry
258 .search(query, self.services.skill.max_active_skills, |text| {
259 let owned = text.to_owned();
260 let p = provider.clone();
261 Box::pin(async move { p.embed(&owned).await })
262 })
263 .await;
264 self.rehydrate_mcp_tools(hits)
265 }
266
267 fn rehydrate_mcp_tools(&self, hits: Vec<zeph_mcp::McpTool>) -> Vec<zeph_mcp::McpTool> {
279 hits.into_iter()
280 .filter_map(|hit| {
281 let live = self
282 .services
283 .mcp
284 .tools
285 .iter()
286 .find(|t| t.server_id == hit.server_id && t.name == hit.name)
287 .cloned();
288 if live.is_none() {
289 tracing::warn!(
290 server_id = hit.server_id,
291 tool = hit.name,
292 "MCP tool from semantic search has no live match; dropping stale Qdrant hit"
293 );
294 }
295 live
296 })
297 .collect()
298 }
299
300 pub(super) async fn check_tool_refresh(&mut self) {
315 if self.services.mcp.pending_semantic_rebuild {
317 self.services.mcp.pending_semantic_rebuild = false;
318 self.refresh_mcp_tool_ids();
319 self.rebuild_semantic_index().await;
320 self.sync_mcp_registry().await;
321 self.refresh_shadow_sentinel_mcp_tool_ids();
322 let mcp_total = self.services.mcp.tools.len();
323 let mcp_servers = self
324 .services
325 .mcp
326 .tools
327 .iter()
328 .map(|t| &t.server_id)
329 .collect::<std::collections::HashSet<_>>()
330 .len();
331 self.update_metrics(|m| {
332 m.mcp_tool_count = mcp_total;
333 m.mcp_server_count = mcp_servers;
334 });
335 }
336
337 let Some(ref mut rx) = self.services.mcp.tool_rx else {
338 return;
339 };
340 if !rx.has_changed().unwrap_or(false) {
341 return;
342 }
343 let new_tools = rx.borrow_and_update().clone();
344 if new_tools.is_empty() {
345 return;
355 }
356 tracing::info!(
357 tools = new_tools.len(),
358 "tools/list_changed: agent tool list refreshed"
359 );
360 self.services.mcp.tools = new_tools;
361 self.services.mcp.sync_executor_tools();
362 self.services.mcp.pruning_cache.reset();
363 self.refresh_mcp_tool_ids();
364 self.rebuild_semantic_index().await;
365 self.sync_mcp_registry().await;
366 self.refresh_shadow_sentinel_mcp_tool_ids();
367 let mcp_total = self.services.mcp.tools.len();
368 let mcp_servers = self
369 .services
370 .mcp
371 .tools
372 .iter()
373 .map(|t| &t.server_id)
374 .collect::<std::collections::HashSet<_>>()
375 .len();
376 self.update_metrics(|m| {
377 m.mcp_tool_count = mcp_total;
378 m.mcp_server_count = mcp_servers;
379 });
380 }
381
382 fn refresh_shadow_sentinel_mcp_tool_ids(&self) {
390 let Some(ref sentinel) = self.services.security.shadow_sentinel else {
391 return;
392 };
393 let ids: std::collections::HashSet<String> = self
394 .services
395 .mcp
396 .tools
397 .iter()
398 .map(zeph_mcp::McpTool::sanitized_id)
399 .collect();
400 *sentinel.mcp_tool_ids_handle().write() = ids;
401 }
402
403 fn refresh_mcp_tool_ids(&self) {
410 let Some(ref handle) = self.services.security.mcp_tool_ids else {
411 return;
412 };
413 let ids: std::collections::HashSet<String> = self
414 .services
415 .mcp
416 .tools
417 .iter()
418 .map(zeph_mcp::McpTool::sanitized_id)
419 .collect();
420 *handle.write() = ids;
421 }
422
423 pub(super) async fn sync_mcp_registry(&mut self) {
424 if self.services.mcp.registry.is_none() {
425 return;
426 }
427 if !self.embedding_provider.supports_embeddings() {
428 return;
429 }
430 let tools = self.services.mcp.tools.clone();
432 let provider = self.embedding_provider.clone();
433 let embedding_model = self.services.skill.embedding_model.clone();
434 let embed_timeout =
435 std::time::Duration::from_secs(self.runtime.config.timeouts.embedding_seconds);
436 let embed_fn = move |text: &str| -> zeph_mcp::registry::EmbedFuture {
437 let owned = text.to_owned();
438 let p = provider.clone();
439 Box::pin(async move {
440 if let Ok(result) = tokio::time::timeout(embed_timeout, p.embed(&owned)).await {
441 result
442 } else {
443 tracing::warn!(
444 timeout_secs = embed_timeout.as_secs(),
445 "MCP registry: embedding timed out"
446 );
447 Err(zeph_llm::LlmError::Timeout)
448 }
449 })
450 };
451 let Some(mut registry) = self.services.mcp.registry.take() else {
454 return;
455 };
456 if let Err(e) = registry.sync(&tools, &embedding_model, embed_fn).await {
457 tracing::warn!("failed to sync MCP tool registry: {e:#}");
458 }
459 self.services.mcp.registry = Some(registry);
460 }
461
462 pub async fn init_semantic_index(&mut self) {
469 self.rebuild_semantic_index().await;
470 }
471
472 pub(super) async fn process_pending_elicitations(&mut self) {
477 loop {
478 let Some(ref mut rx) = self.services.mcp.elicitation_rx else {
479 return;
480 };
481 match rx.try_recv() {
482 Ok(event) => {
483 self.handle_elicitation_event(event).await;
484 }
485 Err(tokio::sync::mpsc::error::TryRecvError::Empty) => return,
486 Err(tokio::sync::mpsc::error::TryRecvError::Disconnected) => {
487 self.services.mcp.elicitation_rx = None;
488 return;
489 }
490 }
491 }
492 }
493
494 pub(super) async fn handle_elicitation_event(&mut self, event: zeph_mcp::ElicitationEvent) {
496 use crate::channel::{ElicitationRequest, ElicitationResponse};
497
498 let decline = ElicitResult::new(ElicitationAction::Decline);
499
500 let channel_request = match &event.request {
501 rmcp::model::ElicitRequestParams::FormElicitationParams {
502 message,
503 requested_schema,
504 ..
505 } => {
506 let fields = build_elicitation_fields(requested_schema);
507 ElicitationRequest {
508 server_name: event.server_id.clone(),
509 message: sanitize_elicitation_message(message),
510 fields,
511 }
512 }
513 rmcp::model::ElicitRequestParams::UrlElicitationParams { .. } => {
514 tracing::debug!(
516 server_id = event.server_id,
517 "URL elicitation not supported, declining"
518 );
519 let _ = event.response_tx.send(decline);
520 return;
521 }
522 _ => {
524 tracing::debug!(
525 server_id = event.server_id,
526 "unknown elicitation request variant, declining"
527 );
528 let _ = event.response_tx.send(decline);
529 return;
530 }
531 };
532
533 if self.services.mcp.elicitation_warn_sensitive_fields {
534 let sensitive: Vec<&str> = channel_request
535 .fields
536 .iter()
537 .filter(|f| is_sensitive_field(&f.name))
538 .map(|f| f.name.as_str())
539 .collect();
540 if !sensitive.is_empty() {
541 let fields_list = sensitive.join(", ");
542 let warning = format!(
543 "Warning: [{}] is requesting sensitive information (field: {}). \
544 Only proceed if you trust this server.",
545 channel_request.server_name, fields_list,
546 );
547 tracing::warn!(
548 server_id = event.server_id,
549 fields = %fields_list,
550 "elicitation requests sensitive fields"
551 );
552 let _ = self.channel.send(&warning).await;
553 }
554 }
555
556 self.channel
557 .send_status_best_effort("MCP server requesting input…")
558 .await;
559 let response = match self.channel.elicit(channel_request).await {
560 Ok(r) => r,
561 Err(e) => {
562 tracing::warn!(
563 server_id = event.server_id,
564 "elicitation channel error: {e:#}"
565 );
566 self.channel.send_status_best_effort("").await;
567 let _ = event.response_tx.send(decline);
568 return;
569 }
570 };
571 self.channel.send_status_best_effort("").await;
572
573 let result = match response {
574 ElicitationResponse::Accepted(value) => {
575 ElicitResult::new(ElicitationAction::Accept).with_content(value)
576 }
577 ElicitationResponse::Declined => ElicitResult::new(ElicitationAction::Decline),
578 ElicitationResponse::Cancelled => ElicitResult::new(ElicitationAction::Cancel),
579 };
580
581 if event.response_tx.send(result).is_err() {
582 tracing::warn!(
583 server_id = event.server_id,
584 "elicitation response dropped — handler disconnected"
585 );
586 }
587 }
588
589 fn update_mcp_metrics(&mut self) {
590 let mcp_total = self.services.mcp.tools.len();
591 let mcp_server_count = self.services.mcp.server_outcomes.len();
592 let mcp_connected_count = self
593 .services
594 .mcp
595 .server_outcomes
596 .iter()
597 .filter(|o| o.connected)
598 .count();
599 let mcp_servers: Vec<crate::metrics::McpServerStatus> = self
600 .services
601 .mcp
602 .server_outcomes
603 .iter()
604 .map(|o| crate::metrics::McpServerStatus {
605 id: o.id.clone(),
606 status: if o.connected {
607 crate::metrics::McpServerConnectionStatus::Connected
608 } else {
609 crate::metrics::McpServerConnectionStatus::Failed
610 },
611 tool_count: o.tool_count,
612 error: o.error.clone(),
613 input_schemas_dropped: o.input_schemas_dropped,
614 output_schemas_dropped: o.output_schemas_dropped,
615 })
616 .collect();
617 self.update_metrics(|m| {
618 m.mcp_tool_count = mcp_total;
619 m.mcp_server_count = mcp_server_count;
620 m.mcp_connected_count = mcp_connected_count;
621 m.mcp_servers = mcp_servers;
622 });
623 }
624
625 pub(in crate::agent) async fn rebuild_semantic_index(&mut self) {
635 if self.services.mcp.discovery_strategy != zeph_mcp::ToolDiscoveryStrategy::Embedding {
636 return;
637 }
638
639 if self.services.mcp.tools.is_empty() {
640 self.services.mcp.semantic_index = None;
641 return;
642 }
643
644 let provider = self
646 .services
647 .mcp
648 .discovery_provider
649 .clone()
650 .unwrap_or_else(|| self.embedding_provider.clone());
651
652 let inner_embed = provider.embed_fn();
653 let embed_timeout =
654 std::time::Duration::from_secs(self.runtime.config.timeouts.embedding_seconds);
655 let embed_fn = move |text: &str| -> zeph_llm::provider::EmbedFuture {
656 let fut = inner_embed(text);
657 Box::pin(async move {
658 if let Ok(result) = tokio::time::timeout(embed_timeout, fut).await {
659 result
660 } else {
661 tracing::warn!(
662 timeout_secs = embed_timeout.as_secs(),
663 "semantic index: embedding probe timed out"
664 );
665 Err(zeph_llm::LlmError::Timeout)
666 }
667 })
668 };
669
670 let tools = self.services.mcp.tools.clone();
672 match zeph_mcp::SemanticToolIndex::build(&tools, &embed_fn).await {
673 Ok(idx) => {
674 tracing::info!(
675 indexed = idx.len(),
676 total = self.services.mcp.tools.len(),
677 "semantic tool index built"
678 );
679 self.services.mcp.semantic_index = Some(idx);
680 }
681 Err(e) => {
682 tracing::warn!(
683 "semantic tool index build failed, falling back to all tools: {e:#}"
684 );
685 self.services.mcp.semantic_index = None;
686 }
687 }
688 }
689}
690
691fn validate_mcp_command(target: &str, allowed_commands: &[String]) -> Option<String> {
695 let is_url = target.starts_with("http://") || target.starts_with("https://");
696 if !is_url && !allowed_commands.is_empty() && !allowed_commands.iter().any(|c| c == target) {
697 Some(format!(
698 "Command '{target}' is not allowed. Permitted: {}",
699 allowed_commands.join(", ")
700 ))
701 } else {
702 None
703 }
704}
705
706fn build_server_entry(id: &str, target: &str, extra_args: &[&str]) -> zeph_mcp::ServerEntry {
708 let is_url = target.starts_with("http://") || target.starts_with("https://");
709 let transport = if is_url {
710 zeph_mcp::McpTransport::Http {
711 url: target.to_owned(),
712 headers: std::collections::HashMap::new(),
713 }
714 } else {
715 zeph_mcp::McpTransport::Stdio {
716 command: target.to_owned(),
717 args: extra_args.iter().map(|&s| s.to_owned()).collect(),
718 env: std::collections::HashMap::new(),
719 }
720 };
721 zeph_mcp::ServerEntry {
722 id: id.to_owned(),
723 transport,
724 timeout: std::time::Duration::from_secs(30),
725 trust_level: zeph_config::McpTrustLevel::Untrusted,
726 tool_allowlist: None,
727 expected_tools: Vec::new(),
728 roots: Vec::new(),
729 tool_metadata: std::collections::HashMap::new(),
730 elicitation_enabled: false,
731 elicitation_timeout_secs: 120,
732 env_isolation: false,
733 media_passthrough: false,
734 }
735}
736
737fn build_elicitation_fields(
739 schema: &rmcp::model::ElicitationSchema,
740) -> Vec<crate::channel::ElicitationField> {
741 use crate::channel::{ElicitationField, ElicitationFieldType};
742 use rmcp::model::PrimitiveSchemaDefinition;
743
744 schema
745 .properties
746 .iter()
747 .map(|(name, prop)| {
748 let json = serde_json::to_value(prop).unwrap_or_default();
753 let description = json
754 .get("description")
755 .and_then(|v| v.as_str())
756 .map(sanitize_elicitation_message);
757
758 let field_type = match prop {
759 PrimitiveSchemaDefinition::Boolean(_) => ElicitationFieldType::Boolean,
760 PrimitiveSchemaDefinition::Integer(_) => ElicitationFieldType::Integer,
761 PrimitiveSchemaDefinition::Number(_) => ElicitationFieldType::Number,
762 PrimitiveSchemaDefinition::Enum(_) => {
763 let vals = json
766 .get("enum")
767 .and_then(|v| v.as_array())
768 .map(|arr| {
769 arr.iter()
770 .filter_map(|v| v.as_str())
771 .map(sanitize_elicitation_message)
772 .collect::<Vec<_>>()
773 })
774 .unwrap_or_default();
775 ElicitationFieldType::Enum(vals)
776 }
777 PrimitiveSchemaDefinition::String(_) => ElicitationFieldType::String,
778 _ => {
781 tracing::debug!(
782 "unknown PrimitiveSchemaDefinition variant, defaulting to String"
783 );
784 ElicitationFieldType::String
785 }
786 };
787 let required = schema.required.as_deref().is_some_and(|r| r.contains(name));
788 ElicitationField {
789 name: name.clone(),
792 description,
793 field_type,
794 required,
795 }
796 })
797 .collect()
798}
799
800const SENSITIVE_FIELD_PATTERNS: &[&str] = &[
802 "password",
803 "passwd",
804 "token",
805 "secret",
806 "key",
807 "credential",
808 "apikey",
809 "api_key",
810 "auth",
811 "authorization",
812 "private",
813 "passphrase",
814 "pin",
815];
816
817fn is_sensitive_field(field_name: &str) -> bool {
819 let lower = field_name.to_lowercase();
820 SENSITIVE_FIELD_PATTERNS
821 .iter()
822 .any(|pattern| lower.contains(pattern))
823}
824
825fn sanitize_elicitation_message(message: &str) -> String {
827 const MAX_CHARS: usize = 500;
828 message
830 .chars()
831 .filter(|c| !c.is_control() || *c == '\n' || *c == '\t')
832 .take(MAX_CHARS)
833 .collect()
834}
835
836#[cfg(test)]
837mod tests {
838 use super::super::agent_tests::{
839 MockChannel, MockToolExecutor, create_test_registry, mock_provider,
840 };
841 use super::*;
842 use std::assert_matches;
843
844 #[tokio::test]
845 async fn handle_mcp_command_unknown_subcommand_shows_usage() {
846 let provider = mock_provider(vec![]);
847 let channel = MockChannel::new(vec![]);
848 let registry = create_test_registry();
849 let executor = MockToolExecutor::no_tools();
850 let mut agent = Agent::new(provider, channel, registry, None, 5, executor);
851
852 let result = agent.handle_mcp_command("unknown").await.unwrap();
853 assert!(
854 result.contains("Usage: /mcp"),
855 "expected usage message, got: {result:?}"
856 );
857 }
858
859 #[tokio::test]
860 async fn handle_mcp_list_no_manager_shows_disabled() {
861 let provider = mock_provider(vec![]);
862 let channel = MockChannel::new(vec![]);
863 let registry = create_test_registry();
864 let executor = MockToolExecutor::no_tools();
865 let mut agent = Agent::new(provider, channel, registry, None, 5, executor);
866
867 let result = agent.handle_mcp_command("list").await.unwrap();
868 assert!(
869 result.contains("MCP is not enabled"),
870 "expected not-enabled message, got: {result:?}"
871 );
872 }
873
874 #[tokio::test]
875 async fn handle_mcp_tools_no_server_id_shows_usage() {
876 let provider = mock_provider(vec![]);
877 let channel = MockChannel::new(vec![]);
878 let registry = create_test_registry();
879 let executor = MockToolExecutor::no_tools();
880 let mut agent = Agent::new(provider, channel, registry, None, 5, executor);
881
882 let result = agent.handle_mcp_command("tools").await.unwrap();
883 assert!(
884 result.contains("Usage: /mcp tools"),
885 "expected tools usage message, got: {result:?}"
886 );
887 }
888
889 #[tokio::test]
890 async fn handle_mcp_remove_no_server_id_shows_usage() {
891 let provider = mock_provider(vec![]);
892 let channel = MockChannel::new(vec![]);
893 let registry = create_test_registry();
894 let executor = MockToolExecutor::no_tools();
895 let mut agent = Agent::new(provider, channel, registry, None, 5, executor);
896
897 let result = agent.handle_mcp_command("remove").await.unwrap();
898 assert!(
899 result.contains("Usage: /mcp remove"),
900 "expected remove usage message, got: {result:?}"
901 );
902 }
903
904 #[tokio::test]
905 async fn handle_mcp_remove_no_manager_shows_disabled() {
906 let provider = mock_provider(vec![]);
907 let channel = MockChannel::new(vec![]);
908 let registry = create_test_registry();
909 let executor = MockToolExecutor::no_tools();
910 let mut agent = Agent::new(provider, channel, registry, None, 5, executor);
911
912 let result = agent.handle_mcp_command("remove my-server").await.unwrap();
913 assert!(
914 result.contains("MCP is not enabled"),
915 "expected not-enabled message, got: {result:?}"
916 );
917 }
918
919 #[tokio::test]
920 async fn handle_mcp_add_insufficient_args_shows_usage() {
921 let provider = mock_provider(vec![]);
922 let channel = MockChannel::new(vec![]);
923 let registry = create_test_registry();
924 let executor = MockToolExecutor::no_tools();
925 let mut agent = Agent::new(provider, channel, registry, None, 5, executor);
926
927 let result = agent.handle_mcp_command("add server-id").await.unwrap();
929 assert!(
930 result.contains("Usage: /mcp add"),
931 "expected add usage message, got: {result:?}"
932 );
933 }
934
935 #[tokio::test]
936 async fn handle_mcp_tools_with_unknown_server_shows_no_tools() {
937 let provider = mock_provider(vec![]);
938 let channel = MockChannel::new(vec![]);
939 let registry = create_test_registry();
940 let executor = MockToolExecutor::no_tools();
941 let mut agent = Agent::new(provider, channel, registry, None, 5, executor);
942
943 let result = agent
945 .handle_mcp_command("tools nonexistent-server")
946 .await
947 .unwrap();
948 assert!(
949 result.contains("No tools found"),
950 "expected no-tools message, got: {result:?}"
951 );
952 }
953
954 #[tokio::test]
955 async fn mcp_tool_count_starts_at_zero() {
956 let provider = mock_provider(vec![]);
957 let channel = MockChannel::new(vec![]);
958 let registry = create_test_registry();
959 let executor = MockToolExecutor::no_tools();
960 let agent = Agent::new(provider, channel, registry, None, 5, executor);
961
962 assert_eq!(agent.services.mcp.tool_count(), 0);
963 }
964
965 fn test_mcp_tool(
966 server_id: &str,
967 name: &str,
968 input_schema: serde_json::Value,
969 ) -> zeph_mcp::McpTool {
970 zeph_mcp::McpTool {
971 server_id: server_id.to_owned(),
972 name: name.to_owned(),
973 description: format!("{name} description"),
974 input_schema,
975 output_schema: None,
976 security_meta: zeph_config::mcp_security::ToolSecurityMeta::default(),
977 }
978 }
979
980 #[tokio::test]
984 async fn rehydrate_mcp_tools_replaces_stub_with_live_schema() {
985 let provider = mock_provider(vec![]);
986 let channel = MockChannel::new(vec![]);
987 let registry = create_test_registry();
988 let executor = MockToolExecutor::no_tools();
989 let mut agent = Agent::new(provider, channel, registry, None, 5, executor);
990
991 let real_schema =
992 serde_json::json!({"type": "object", "properties": {"path": {"type": "string"}}});
993 agent.services.mcp.tools = vec![test_mcp_tool("fs", "read_file", real_schema.clone())];
994
995 let stub = test_mcp_tool("fs", "read_file", serde_json::json!({}));
998
999 let rehydrated = agent.rehydrate_mcp_tools(vec![stub]);
1000
1001 assert_eq!(rehydrated.len(), 1);
1002 assert_eq!(rehydrated[0].input_schema, real_schema);
1003 }
1004
1005 #[tokio::test]
1009 async fn rehydrate_mcp_tools_drops_hit_with_no_live_match() {
1010 let provider = mock_provider(vec![]);
1011 let channel = MockChannel::new(vec![]);
1012 let registry = create_test_registry();
1013 let executor = MockToolExecutor::no_tools();
1014 let mut agent = Agent::new(provider, channel, registry, None, 5, executor);
1015
1016 agent.services.mcp.tools = vec![test_mcp_tool("fs", "other_tool", serde_json::json!({}))];
1017
1018 let stub = test_mcp_tool("fs", "read_file", serde_json::json!({}));
1019
1020 let rehydrated = agent.rehydrate_mcp_tools(vec![stub]);
1021
1022 assert!(
1023 rehydrated.is_empty(),
1024 "stale hit with no live match must be dropped, not passed through with an empty schema"
1025 );
1026 }
1027
1028 #[tokio::test]
1029 async fn rehydrate_mcp_tools_mixed_batch_keeps_only_matches() {
1030 let provider = mock_provider(vec![]);
1031 let channel = MockChannel::new(vec![]);
1032 let registry = create_test_registry();
1033 let executor = MockToolExecutor::no_tools();
1034 let mut agent = Agent::new(provider, channel, registry, None, 5, executor);
1035
1036 let schema_a =
1037 serde_json::json!({"type": "object", "properties": {"a": {"type": "string"}}});
1038 let schema_c =
1039 serde_json::json!({"type": "object", "properties": {"c": {"type": "number"}}});
1040 agent.services.mcp.tools = vec![
1041 test_mcp_tool("srv1", "tool_a", schema_a.clone()),
1042 test_mcp_tool("srv2", "tool_c", schema_c.clone()),
1043 ];
1044
1045 let hits = vec![
1046 test_mcp_tool("srv1", "tool_a", serde_json::json!({})),
1047 test_mcp_tool("srv1", "tool_b", serde_json::json!({})), test_mcp_tool("srv2", "tool_c", serde_json::json!({})),
1049 ];
1050
1051 let rehydrated = agent.rehydrate_mcp_tools(hits);
1052
1053 assert_eq!(rehydrated.len(), 2);
1054 assert_eq!(rehydrated[0].name, "tool_a");
1055 assert_eq!(rehydrated[0].input_schema, schema_a);
1056 assert_eq!(rehydrated[1].name, "tool_c");
1057 assert_eq!(rehydrated[1].input_schema, schema_c);
1058 }
1059
1060 #[tokio::test]
1064 async fn rehydrated_tool_schema_reaches_llm_prompt() {
1065 let provider = mock_provider(vec![]);
1066 let channel = MockChannel::new(vec![]);
1067 let registry = create_test_registry();
1068 let executor = MockToolExecutor::no_tools();
1069 let mut agent = Agent::new(provider, channel, registry, None, 5, executor);
1070
1071 let real_schema = serde_json::json!({
1072 "type": "object",
1073 "properties": {"query": {"type": "string"}},
1074 "required": ["query"]
1075 });
1076 agent.services.mcp.tools = vec![test_mcp_tool("search", "web_search", real_schema)];
1077
1078 let stub = test_mcp_tool("search", "web_search", serde_json::json!({}));
1079 let rehydrated = agent.rehydrate_mcp_tools(vec![stub]);
1080
1081 let prompt = zeph_mcp::format_mcp_tools_prompt(&rehydrated);
1082
1083 assert!(
1084 !prompt.contains("<parameters>{}</parameters>"),
1085 "expected real schema in prompt, got empty parameters block: {prompt}"
1086 );
1087 assert!(
1088 prompt.contains("\"query\""),
1089 "expected real schema fields in prompt: {prompt}"
1090 );
1091 }
1092
1093 #[tokio::test]
1094 async fn check_tool_refresh_no_rx_is_noop() {
1095 let provider = mock_provider(vec![]);
1096 let channel = MockChannel::new(vec![]);
1097 let registry = create_test_registry();
1098 let executor = MockToolExecutor::no_tools();
1099 let mut agent = Agent::new(provider, channel, registry, None, 5, executor);
1100 agent.check_tool_refresh().await;
1102 assert_eq!(agent.services.mcp.tool_count(), 0);
1103 }
1104
1105 #[tokio::test]
1106 async fn check_tool_refresh_no_change_is_noop() {
1107 let provider = mock_provider(vec![]);
1108 let channel = MockChannel::new(vec![]);
1109 let registry = create_test_registry();
1110 let executor = MockToolExecutor::no_tools();
1111 let mut agent = Agent::new(provider, channel, registry, None, 5, executor);
1112
1113 let (tx, rx) = tokio::sync::watch::channel(Vec::new());
1114 agent.services.mcp.tool_rx = Some(rx);
1115 agent.check_tool_refresh().await;
1117 assert_eq!(agent.services.mcp.tool_count(), 0);
1118 drop(tx);
1119 }
1120
1121 #[tokio::test]
1122 async fn check_tool_refresh_with_empty_initial_value_does_not_replace_tools() {
1123 let provider = mock_provider(vec![]);
1124 let channel = MockChannel::new(vec![]);
1125 let registry = create_test_registry();
1126 let executor = MockToolExecutor::no_tools();
1127 let mut agent = Agent::new(provider, channel, registry, None, 5, executor);
1128 agent.services.mcp.tools = vec![zeph_mcp::McpTool {
1129 server_id: "srv".into(),
1130 name: "existing_tool".into(),
1131 description: String::new(),
1132 input_schema: serde_json::json!({}),
1133 output_schema: None,
1134 security_meta: zeph_config::mcp_security::ToolSecurityMeta::default(),
1135 }];
1136
1137 let (_tx, rx) = tokio::sync::watch::channel(Vec::<zeph_mcp::McpTool>::new());
1138 agent.services.mcp.tool_rx = Some(rx);
1139 agent.check_tool_refresh().await;
1141 assert_eq!(agent.services.mcp.tool_count(), 1);
1142 }
1143
1144 #[tokio::test]
1145 async fn check_tool_refresh_applies_update() {
1146 let provider = mock_provider(vec![]);
1147 let channel = MockChannel::new(vec![]);
1148 let registry = create_test_registry();
1149 let executor = MockToolExecutor::no_tools();
1150 let mut agent = Agent::new(provider, channel, registry, None, 5, executor);
1151
1152 let (tx, rx) = tokio::sync::watch::channel(Vec::<zeph_mcp::McpTool>::new());
1153 agent.services.mcp.tool_rx = Some(rx);
1154
1155 let new_tools = vec![zeph_mcp::McpTool {
1156 server_id: "srv".into(),
1157 name: "refreshed_tool".into(),
1158 description: String::new(),
1159 input_schema: serde_json::json!({}),
1160 output_schema: None,
1161 security_meta: zeph_config::mcp_security::ToolSecurityMeta::default(),
1162 }];
1163 tx.send(new_tools).unwrap();
1164
1165 agent.check_tool_refresh().await;
1166 assert_eq!(agent.services.mcp.tool_count(), 1);
1167 assert_eq!(agent.services.mcp.tools[0].name, "refreshed_tool");
1168 }
1169
1170 #[tokio::test]
1175 async fn check_tool_refresh_updates_shadow_sentinel_mcp_tool_ids() {
1176 use crate::agent::shadow_sentinel::{
1177 ProbeVerdict, SafetyProbe, ShadowEventStore, ShadowSentinel,
1178 };
1179
1180 struct NoopProbe;
1181 impl SafetyProbe for NoopProbe {
1182 fn evaluate<'a>(
1183 &'a self,
1184 _: &'a str,
1185 _: &'a serde_json::Value,
1186 _: &'a [crate::agent::shadow_sentinel::SentinelEvent],
1187 ) -> std::pin::Pin<Box<dyn std::future::Future<Output = ProbeVerdict> + Send + 'a>>
1188 {
1189 Box::pin(async { ProbeVerdict::Allow })
1190 }
1191 }
1192
1193 let provider = mock_provider(vec![]);
1194 let channel = MockChannel::new(vec![]);
1195 let registry = create_test_registry();
1196 let executor = MockToolExecutor::no_tools();
1197 let mut agent = Agent::new(provider, channel, registry, None, 5, executor);
1198
1199 let pool = zeph_db::DbConfig {
1200 url: ":memory:".to_owned(),
1201 ..Default::default()
1202 }
1203 .connect()
1204 .await
1205 .expect("connect + migrate in-memory sqlite pool");
1206 let store = ShadowEventStore::new(pool);
1207 let sentinel = std::sync::Arc::new(ShadowSentinel::new(
1208 store,
1209 Box::new(NoopProbe),
1210 zeph_config::ShadowSentinelConfig::default(),
1211 "test-session",
1212 ));
1213 agent.services.security.shadow_sentinel = Some(std::sync::Arc::clone(&sentinel));
1214
1215 let (tx, rx) = tokio::sync::watch::channel(Vec::<zeph_mcp::McpTool>::new());
1216 agent.services.mcp.tool_rx = Some(rx);
1217
1218 let new_tool = zeph_mcp::McpTool {
1219 server_id: "srv".into(),
1220 name: "refreshed_tool".into(),
1221 description: String::new(),
1222 input_schema: serde_json::json!({}),
1223 output_schema: None,
1224 security_meta: zeph_config::mcp_security::ToolSecurityMeta::default(),
1225 };
1226 let expected_id = new_tool.sanitized_id();
1227 tx.send(vec![new_tool]).unwrap();
1228
1229 agent.check_tool_refresh().await;
1230
1231 assert!(
1232 sentinel.mcp_tool_ids_handle().read().contains(&expected_id),
1233 "ShadowSentinel's mcp_tool_ids must be refreshed after a tools/list_changed event"
1234 );
1235 }
1236
1237 #[tokio::test]
1238 async fn check_tool_refresh_without_mcp_tool_ids_handle_does_not_panic() {
1239 let provider = mock_provider(vec![]);
1240 let channel = MockChannel::new(vec![]);
1241 let registry = create_test_registry();
1242 let executor = MockToolExecutor::no_tools();
1243 let mut agent = Agent::new(provider, channel, registry, None, 5, executor);
1244 assert!(agent.services.security.mcp_tool_ids.is_none());
1247
1248 let (tx, rx) = tokio::sync::watch::channel(Vec::<zeph_mcp::McpTool>::new());
1249 agent.services.mcp.tool_rx = Some(rx);
1250 let new_tools = vec![zeph_mcp::McpTool {
1251 server_id: "srv".into(),
1252 name: "refreshed_tool".into(),
1253 description: String::new(),
1254 input_schema: serde_json::json!({}),
1255 output_schema: None,
1256 security_meta: zeph_config::mcp_security::ToolSecurityMeta::default(),
1257 }];
1258 tx.send(new_tools).unwrap();
1259
1260 agent.check_tool_refresh().await;
1261 assert_eq!(agent.services.mcp.tool_count(), 1);
1262 }
1263
1264 #[tokio::test]
1265 async fn check_tool_refresh_updates_attached_mcp_tool_ids_handle() {
1266 let provider = mock_provider(vec![]);
1267 let channel = MockChannel::new(vec![]);
1268 let registry = create_test_registry();
1269 let executor = MockToolExecutor::no_tools();
1270 let mut agent = Agent::new(provider, channel, registry, None, 5, executor);
1271
1272 let handle =
1273 std::sync::Arc::new(parking_lot::RwLock::new(std::collections::HashSet::new()));
1274 agent.services.security.mcp_tool_ids = Some(std::sync::Arc::clone(&handle));
1275
1276 let (tx, rx) = tokio::sync::watch::channel(Vec::<zeph_mcp::McpTool>::new());
1277 agent.services.mcp.tool_rx = Some(rx);
1278 let new_tools = vec![zeph_mcp::McpTool {
1279 server_id: "srv".into(),
1280 name: "refreshed_tool".into(),
1281 description: String::new(),
1282 input_schema: serde_json::json!({}),
1283 output_schema: None,
1284 security_meta: zeph_config::mcp_security::ToolSecurityMeta::default(),
1285 }];
1286 tx.send(new_tools).unwrap();
1287
1288 agent.check_tool_refresh().await;
1289
1290 assert!(
1291 handle.read().contains("srv_refreshed_tool"),
1292 "expected the sanitized id of the newly-connected tool in the handle, got: {:?}",
1293 *handle.read()
1294 );
1295 }
1296
1297 #[tokio::test]
1298 async fn check_tool_refresh_drops_disconnected_tool_from_mcp_tool_ids_handle() {
1299 let provider = mock_provider(vec![]);
1300 let channel = MockChannel::new(vec![]);
1301 let registry = create_test_registry();
1302 let executor = MockToolExecutor::no_tools();
1303 let mut agent = Agent::new(provider, channel, registry, None, 5, executor);
1304
1305 let handle =
1308 std::sync::Arc::new(parking_lot::RwLock::new(std::collections::HashSet::from([
1309 "stale_server_old_tool".to_owned(),
1310 ])));
1311 agent.services.security.mcp_tool_ids = Some(std::sync::Arc::clone(&handle));
1312
1313 let (tx, rx) = tokio::sync::watch::channel(Vec::<zeph_mcp::McpTool>::new());
1314 agent.services.mcp.tool_rx = Some(rx);
1315 let new_tools = vec![zeph_mcp::McpTool {
1316 server_id: "srv".into(),
1317 name: "refreshed_tool".into(),
1318 description: String::new(),
1319 input_schema: serde_json::json!({}),
1320 output_schema: None,
1321 security_meta: zeph_config::mcp_security::ToolSecurityMeta::default(),
1322 }];
1323 tx.send(new_tools).unwrap();
1324
1325 agent.check_tool_refresh().await;
1326
1327 let ids = handle.read();
1328 assert!(
1329 !ids.contains("stale_server_old_tool"),
1330 "disconnected server's tool id must be dropped (replace, not union), got: {ids:?}"
1331 );
1332 assert!(ids.contains("srv_refreshed_tool"));
1333 }
1334
1335 #[tokio::test]
1339 async fn check_tool_refresh_updates_mcp_tool_ids_handle_via_pending_semantic_rebuild() {
1340 let provider = mock_provider(vec![]);
1341 let channel = MockChannel::new(vec![]);
1342 let registry = create_test_registry();
1343 let executor = MockToolExecutor::no_tools();
1344 let mut agent = Agent::new(provider, channel, registry, None, 5, executor);
1345
1346 let handle =
1347 std::sync::Arc::new(parking_lot::RwLock::new(std::collections::HashSet::new()));
1348 agent.services.security.mcp_tool_ids = Some(std::sync::Arc::clone(&handle));
1349
1350 agent.services.mcp.tools = vec![zeph_mcp::McpTool {
1353 server_id: "srv".into(),
1354 name: "added_tool".into(),
1355 description: String::new(),
1356 input_schema: serde_json::json!({}),
1357 output_schema: None,
1358 security_meta: zeph_config::mcp_security::ToolSecurityMeta::default(),
1359 }];
1360 agent.services.mcp.pending_semantic_rebuild = true;
1361
1362 agent.check_tool_refresh().await;
1363
1364 assert!(
1365 handle.read().contains("srv_added_tool"),
1366 "expected the sanitized id of the /mcp add-connected tool in the handle, got: {:?}",
1367 *handle.read()
1368 );
1369 assert!(!agent.services.mcp.pending_semantic_rebuild);
1370 }
1371
1372 #[test]
1373 fn sanitize_elicitation_message_strips_control_chars() {
1374 let input = "hello\x01world\x1b[31mred\x1b[0m";
1375 let output = sanitize_elicitation_message(input);
1376 assert!(!output.contains('\x01'));
1377 assert!(!output.contains('\x1b'));
1378 assert!(output.contains("hello"));
1379 assert!(output.contains("world"));
1380 }
1381
1382 #[test]
1383 fn sanitize_elicitation_message_preserves_newline_and_tab() {
1384 let input = "line1\nline2\ttabbed";
1385 let output = sanitize_elicitation_message(input);
1386 assert_eq!(output, "line1\nline2\ttabbed");
1387 }
1388
1389 #[test]
1390 fn sanitize_elicitation_message_caps_at_500_chars() {
1391 let input: String = "a".repeat(600);
1393 let output = sanitize_elicitation_message(&input);
1394 assert_eq!(output.chars().count(), 500);
1395 }
1396
1397 #[test]
1398 fn sanitize_elicitation_message_handles_multibyte_boundary() {
1399 let input: String = "é".repeat(300); let output = sanitize_elicitation_message(&input);
1402 assert_eq!(output.chars().count(), 300);
1404 }
1405
1406 #[test]
1407 fn build_elicitation_fields_maps_primitive_types() {
1408 use crate::channel::ElicitationFieldType;
1409 use rmcp::model::{
1410 BooleanSchema, ElicitationSchema, IntegerSchema, NumberSchema,
1411 PrimitiveSchemaDefinition, StringSchema,
1412 };
1413 use std::collections::BTreeMap;
1414
1415 let mut props = BTreeMap::new();
1416 props.insert(
1417 "flag".to_owned(),
1418 PrimitiveSchemaDefinition::Boolean(BooleanSchema::new()),
1419 );
1420 props.insert(
1421 "count".to_owned(),
1422 PrimitiveSchemaDefinition::Integer(IntegerSchema::new()),
1423 );
1424 props.insert(
1425 "ratio".to_owned(),
1426 PrimitiveSchemaDefinition::Number(NumberSchema::new()),
1427 );
1428 props.insert(
1429 "name".to_owned(),
1430 PrimitiveSchemaDefinition::String(StringSchema::new()),
1431 );
1432
1433 let schema = ElicitationSchema::new(props);
1434 let fields = build_elicitation_fields(&schema);
1435
1436 let get = |n: &str| fields.iter().find(|f| f.name == n).unwrap();
1437 assert_matches!(get("flag").field_type, ElicitationFieldType::Boolean);
1438 assert_matches!(get("count").field_type, ElicitationFieldType::Integer);
1439 assert_matches!(get("ratio").field_type, ElicitationFieldType::Number);
1440 assert_matches!(get("name").field_type, ElicitationFieldType::String);
1441 }
1442
1443 #[test]
1444 fn build_elicitation_fields_required_flag() {
1445 use rmcp::model::{ElicitationSchema, PrimitiveSchemaDefinition, StringSchema};
1446 use std::collections::BTreeMap;
1447
1448 let mut props = BTreeMap::new();
1449 props.insert(
1450 "req".to_owned(),
1451 PrimitiveSchemaDefinition::String(StringSchema::new()),
1452 );
1453 props.insert(
1454 "opt".to_owned(),
1455 PrimitiveSchemaDefinition::String(StringSchema::new()),
1456 );
1457
1458 let mut schema = ElicitationSchema::new(props);
1459 schema.required = Some(vec!["req".to_owned()]);
1460
1461 let fields = build_elicitation_fields(&schema);
1462 let req = fields.iter().find(|f| f.name == "req").unwrap();
1463 let opt = fields.iter().find(|f| f.name == "opt").unwrap();
1464 assert!(req.required);
1465 assert!(!opt.required);
1466 }
1467
1468 #[test]
1469 fn is_sensitive_field_detects_common_patterns() {
1470 assert!(is_sensitive_field("password"));
1471 assert!(is_sensitive_field("PASSWORD"));
1472 assert!(is_sensitive_field("user_password"));
1473 assert!(is_sensitive_field("api_token"));
1474 assert!(is_sensitive_field("SECRET_KEY"));
1475 assert!(is_sensitive_field("auth_header"));
1476 assert!(is_sensitive_field("private_key"));
1477 }
1478
1479 #[test]
1480 fn is_sensitive_field_allows_non_sensitive_names() {
1481 assert!(!is_sensitive_field("username"));
1482 assert!(!is_sensitive_field("email"));
1483 assert!(!is_sensitive_field("message"));
1484 assert!(!is_sensitive_field("description"));
1485 assert!(!is_sensitive_field("subject"));
1486 }
1487}