1use crate::bridge::MemoryBridge;
8use crate::tools::*;
9use rmcp::{
10 handler::server::{router::tool::ToolRouter, wrapper::Parameters},
11 tool, tool_handler, tool_router, ErrorData, ServerHandler,
12};
13use std::collections::HashSet;
14use std::sync::Arc;
15use tokio::runtime::Handle;
16
17use crate::tools::{
19 AddGraphEdgeParams, CommunityParams, FactorGraphParams, InvalidateGraphEdgeParams,
20 ListGraphEdgesParams, RecordOutcomeParams, TopologyParams,
21};
22use crate::tools::{
23 AddSupportAdmissionParams, ClassifyQueryParams, EntityLookupParams,
24 EvaluateProofDebtGateParams, ExportClaimBundleParams, PlanQueryParams, ProjectionHealthParams,
25 ProofDebtStatusParams, QueryOrchestratedParams, QueryTemporalKParams,
26 RecordContradictionParams, ResolveContradictionParams, SupersedeClaimParams,
27 VerifyLedgerParams,
28};
29
30pub struct SemanticMemoryServer {
31 bridge: Arc<MemoryBridge>,
32 tool_router: ToolRouter<Self>,
33 #[cfg(feature = "orchestration")]
34 runtime: Option<knowledge_runtime::KnowledgeRuntime>,
35}
36
37impl SemanticMemoryServer {
38 pub fn new(bridge: MemoryBridge, tool_profile: &str) -> Self {
39 let mut router = Self::tool_router();
40
41 let admin_tools = [
44 "sm_reconcile",
46 "sm_vacuum",
47 "sm_reembed_all",
48 "sm_embeddings_are_dirty",
49 "sm_get_search_receipt",
51 "sm_replay_search_receipt",
52 "sm_query_claim_versions",
54 "sm_query_relation_versions",
55 "sm_query_episodes",
56 "sm_query_entity_aliases",
57 "sm_query_evidence_refs",
58 "sm_import_envelope",
60 "sm_import_status",
61 "sm_list_imports",
62 "sm_query_temporal",
64 "sm_projection_health",
65 "sm_proof_debt_status",
67 "sm_evaluate_proof_debt_gate",
68 "sm_resolve_contradiction",
69 "sm_verify_ledger",
70 "sm_export_claim_bundle",
71 ];
72
73 match tool_profile {
74 "full" => { }
75 "standard" => {
76 for t in &admin_tools {
78 if t.starts_with("sm_import_")
79 || *t == "sm_list_imports"
80 || *t == "sm_query_temporal"
81 || *t == "sm_projection_health"
82 || *t == "sm_proof_debt_status"
83 || *t == "sm_evaluate_proof_debt_gate"
84 || *t == "sm_resolve_contradiction"
85 || *t == "sm_verify_ledger"
86 || *t == "sm_export_claim_bundle"
87 {
88 router.disable_route(*t);
89 }
90 }
91 }
92 _ => {
93 for t in &admin_tools {
95 router.disable_route(*t);
96 }
97 }
98 }
99
100 eprintln!(
101 "Tool profile: {} ({} tools visible)",
102 tool_profile,
103 router.list_all().len()
104 );
105
106 #[cfg(feature = "orchestration")]
108 let runtime = {
109 let adapter = knowledge_runtime::adapters::semantic_memory::SemanticMemoryAdapter::new(
110 bridge.store.clone(),
111 );
112 let config = knowledge_runtime::RuntimeConfig {
113 default_scope: knowledge_runtime::Scope::new("general"),
114 query: knowledge_runtime::config::QueryConfig::default(),
115 entity: knowledge_runtime::config::EntityConfig::default(),
116 projection: knowledge_runtime::config::ProjectionConfig::default(),
117 strict_temporal: false,
118 strict_scope: false,
119 };
120 knowledge_runtime::KnowledgeRuntime::new(config, adapter).ok()
121 };
122
123 Self {
124 bridge: Arc::new(bridge),
125 tool_router: router,
126 #[cfg(feature = "orchestration")]
127 runtime,
128 }
129 }
130}
131
132fn load_stored_edge_refs(
135 store: &semantic_memory::MemoryStore,
136) -> Result<Vec<semantic_memory::discord::GraphEdgeRef>, ErrorData> {
137 let edges =
138 tokio::task::block_in_place(|| Handle::current().block_on(store.list_all_graph_edges()))
139 .map_err(|e| {
140 ErrorData::internal_error(format!("Failed to load graph edges: {e}"), None)
141 })?;
142 let refs = edges
143 .iter()
144 .map(|edge| {
145 let parsed_type = edge
146 .edge_type_parsed
147 .clone()
148 .or_else(|| serde_json::from_str(&edge.edge_type).ok())
149 .unwrap_or(semantic_memory::GraphEdgeType::Entity {
150 relation: "unknown".to_string(),
151 });
152 let type_str = match parsed_type {
153 semantic_memory::GraphEdgeType::Semantic { .. } => "semantic",
154 semantic_memory::GraphEdgeType::Temporal { .. } => "temporal",
155 semantic_memory::GraphEdgeType::Causal { .. } => "causal",
156 semantic_memory::GraphEdgeType::Entity { .. } => "entity",
157 };
158 semantic_memory::discord::GraphEdgeRef {
159 source: edge.source.clone(),
160 target: edge.target.clone(),
161 edge_type: type_str.to_string(),
162 weight: edge.weight,
163 }
164 })
165 .collect();
166 Ok(refs)
167}
168
169fn load_stored_factor_edges(
172 store: &semantic_memory::MemoryStore,
173) -> Result<
174 Vec<(
175 String,
176 String,
177 semantic_memory::GraphEdgeType,
178 f64,
179 Option<String>,
180 )>,
181 ErrorData,
182> {
183 let edges =
184 tokio::task::block_in_place(|| Handle::current().block_on(store.list_all_graph_edges()))
185 .map_err(|e| {
186 ErrorData::internal_error(format!("Failed to load graph edges: {e}"), None)
187 })?;
188 let raw = edges
189 .iter()
190 .map(|edge| {
191 let parsed_type = edge
192 .edge_type_parsed
193 .clone()
194 .or_else(|| serde_json::from_str(&edge.edge_type).ok())
195 .unwrap_or(semantic_memory::GraphEdgeType::Entity {
196 relation: "unknown".to_string(),
197 });
198 (
199 edge.source.clone(),
200 edge.target.clone(),
201 parsed_type,
202 edge.weight,
203 edge.metadata.clone(),
204 )
205 })
206 .collect();
207 Ok(raw)
208}
209
210fn load_stored_edge_pairs(
212 store: &semantic_memory::MemoryStore,
213) -> Result<Vec<(String, String)>, ErrorData> {
214 let edges =
215 tokio::task::block_in_place(|| Handle::current().block_on(store.list_all_graph_edges()))
216 .map_err(|e| {
217 ErrorData::internal_error(format!("Failed to load graph edges: {e}"), None)
218 })?;
219 let pairs = edges
220 .iter()
221 .map(|edge| (edge.source.clone(), edge.target.clone()))
222 .collect();
223 Ok(pairs)
224}
225
226fn load_neighborhood_edge_pairs(
230 store: &semantic_memory::MemoryStore,
231 seed_ids: &[String],
232) -> Result<Vec<(String, String)>, ErrorData> {
233 if seed_ids.is_empty() {
234 return load_stored_edge_pairs(store);
235 }
236 let edges = tokio::task::block_in_place(|| {
237 Handle::current().block_on(store.list_graph_edges_for_neighborhood(
238 seed_ids.to_vec(),
239 2,
240 200,
241 ))
242 })
243 .map_err(|e| {
244 ErrorData::internal_error(format!("Failed to load neighborhood edges: {e}"), None)
245 })?;
246 let pairs = edges
247 .iter()
248 .map(|edge| (edge.source.clone(), edge.target.clone()))
249 .collect();
250 Ok(pairs)
251}
252
253fn load_neighborhood_edge_refs(
255 store: &semantic_memory::MemoryStore,
256 seed_ids: &[String],
257) -> Result<Vec<semantic_memory::discord::GraphEdgeRef>, ErrorData> {
258 if seed_ids.is_empty() {
259 return load_stored_edge_refs(store);
260 }
261 let edges = tokio::task::block_in_place(|| {
262 Handle::current().block_on(store.list_graph_edges_for_neighborhood(
263 seed_ids.to_vec(),
264 2,
265 200,
266 ))
267 })
268 .map_err(|e| {
269 ErrorData::internal_error(format!("Failed to load neighborhood edges: {e}"), None)
270 })?;
271 let refs = edges
272 .iter()
273 .map(|edge| {
274 let parsed_type = edge
275 .edge_type_parsed
276 .clone()
277 .or_else(|| serde_json::from_str(&edge.edge_type).ok())
278 .unwrap_or(semantic_memory::GraphEdgeType::Entity {
279 relation: "unknown".to_string(),
280 });
281 let type_str = match parsed_type {
282 semantic_memory::GraphEdgeType::Semantic { .. } => "semantic",
283 semantic_memory::GraphEdgeType::Temporal { .. } => "temporal",
284 semantic_memory::GraphEdgeType::Causal { .. } => "causal",
285 semantic_memory::GraphEdgeType::Entity { .. } => "entity",
286 };
287 semantic_memory::discord::GraphEdgeRef {
288 source: edge.source.clone(),
289 target: edge.target.clone(),
290 edge_type: type_str.to_string(),
291 weight: edge.weight,
292 }
293 })
294 .collect();
295 Ok(refs)
296}
297
298fn load_neighborhood_factor_edges(
300 store: &semantic_memory::MemoryStore,
301 seed_ids: &[String],
302) -> Result<
303 Vec<(
304 String,
305 String,
306 semantic_memory::GraphEdgeType,
307 f64,
308 Option<String>,
309 )>,
310 ErrorData,
311> {
312 if seed_ids.is_empty() {
313 return load_stored_factor_edges(store);
314 }
315 let edges = tokio::task::block_in_place(|| {
316 Handle::current().block_on(store.list_graph_edges_for_neighborhood(
317 seed_ids.to_vec(),
318 2,
319 200,
320 ))
321 })
322 .map_err(|e| {
323 ErrorData::internal_error(format!("Failed to load neighborhood edges: {e}"), None)
324 })?;
325 let raw = edges
326 .iter()
327 .map(|edge| {
328 let parsed_type = edge
329 .edge_type_parsed
330 .clone()
331 .or_else(|| serde_json::from_str(&edge.edge_type).ok())
332 .unwrap_or(semantic_memory::GraphEdgeType::Entity {
333 relation: "unknown".to_string(),
334 });
335 (
336 edge.source.clone(),
337 edge.target.clone(),
338 parsed_type,
339 edge.weight,
340 edge.metadata.clone(),
341 )
342 })
343 .collect();
344 Ok(raw)
345}
346
347fn load_superseded_targets(
349 store: &semantic_memory::MemoryStore,
350) -> Result<HashSet<String>, ErrorData> {
351 let edges =
352 tokio::task::block_in_place(|| Handle::current().block_on(store.list_all_graph_edges()))
353 .map_err(|e| {
354 ErrorData::internal_error(format!("Failed to load graph edges: {e}"), None)
355 })?;
356 let mut targets = HashSet::new();
357 for edge in edges {
358 let parsed_type = edge
359 .edge_type_parsed
360 .clone()
361 .or_else(|| serde_json::from_str(&edge.edge_type).ok());
362 if let Some(semantic_memory::GraphEdgeType::Entity { relation }) = parsed_type {
363 if relation == "supersedes" {
364 targets.insert(edge.target);
365 }
366 }
367 }
368 Ok(targets)
369}
370
371fn query_allows_superseded(query: &str) -> bool {
372 let q = query.to_lowercase();
373 q.contains("supersed")
374 || q.contains("stale")
375 || q.contains("obsolete")
376 || q.contains("histor")
377 || q.contains("old fact")
378 || q.contains("previous fact")
379}
380
381fn build_projection_query(params: ProjectionQueryParams) -> semantic_memory::ProjectionQuery {
388 use stack_ids::{ClaimId, ClaimVersionId, EntityId, ScopeKey};
389
390 let scope = ScopeKey {
391 namespace: params.namespace,
392 domain: params.domain,
393 workspace_id: params.workspace_id,
394 repo_id: params.repo_id,
395 };
396
397 let limit = params.limit.unwrap_or(10) as usize;
398
399 semantic_memory::ProjectionQuery {
400 scope,
401 text_query: params.text_query,
402 valid_at: params.valid_at,
403 recorded_at_or_before: params.recorded_at_or_before,
404 subject_entity_id: params.subject_entity_id.map(EntityId::new),
405 canonical_entity_id: params.canonical_entity_id.map(EntityId::new),
406 claim_state: params.claim_state,
407 claim_id: params.claim_id.map(ClaimId::new),
408 claim_version_id: params.claim_version_id.map(ClaimVersionId::new),
409 limit,
410 }
411}
412
413fn json_to_string(value: &serde_json::Value) -> Result<String, ErrorData> {
414 serde_json::to_string_pretty(value)
415 .map_err(|e| ErrorData::internal_error(format!("Serialization error: {e}"), None))
416}
417
418#[tool_router]
419impl SemanticMemoryServer {
420 #[tool(
423 description = "Semantic hybrid search (BM25 + vector + RRF). Returns ranked results with content, scores, and stable result IDs.",
424 annotations(read_only_hint = true)
425 )]
426 fn sm_search(
427 &self,
428 Parameters(SearchParams {
429 query,
430 top_k,
431 namespaces,
432 }): Parameters<SearchParams>,
433 ) -> Result<String, ErrorData> {
434 let requested_k = top_k.map(|v| v as usize).unwrap_or(5);
435 let allow_superseded = query_allows_superseded(&query);
436 let search_k = if allow_superseded {
437 requested_k
438 } else {
439 (requested_k * 4).max(20)
440 };
441 let ns: Option<Vec<&str>> = namespaces
442 .as_ref()
443 .map(|v| v.iter().map(|s| s.as_str()).collect());
444
445 let store = &self.bridge.store;
446 let result = tokio::task::block_in_place(|| {
447 Handle::current().block_on(store.search(&query, Some(search_k), ns.as_deref(), None))
448 });
449
450 match result {
451 Ok(results) => {
452 let superseded_targets = if allow_superseded {
453 HashSet::new()
454 } else {
455 load_superseded_targets(store)?
456 };
457 let fresh_results: Vec<_> = results
458 .iter()
459 .filter(|r| !superseded_targets.contains(&r.source.result_id()))
460 .collect();
461 let result_refs: Vec<_> =
462 if superseded_targets.is_empty() || fresh_results.is_empty() {
463 results.iter().collect()
464 } else {
465 fresh_results
466 };
467 let superseded_filtered_count = results.len().saturating_sub(result_refs.len());
468 let json_results: Vec<serde_json::Value> = result_refs
469 .iter()
470 .take(requested_k)
471 .map(|r| {
472 serde_json::json!({
473 "result_id": r.source.result_id(),
474 "content": r.content,
475 "source": format!("{:?}", r.source),
476 "score": r.score,
477 "bm25_rank": r.bm25_rank,
478 "vector_rank": r.vector_rank,
479 "cosine_similarity": r.cosine_similarity,
480 })
481 })
482 .collect();
483 json_to_string(&serde_json::json!({
484 "ok": true,
485 "results": json_results,
486 "count": json_results.len(),
487 "superseded_filtered_count": superseded_filtered_count,
488 }))
489 }
490 Err(e) => Err(ErrorData::internal_error(
491 format!("Search error: {e}"),
492 None,
493 )),
494 }
495 }
496
497 #[tool(
498 description = "Search with full score breakdown showing how BM25 and vector scores combine. Useful for debugging retrieval quality.",
499 annotations(read_only_hint = true)
500 )]
501 #[allow(dead_code)]
502 fn sm_search_explained(
503 &self,
504 Parameters(SearchExplainedParams { query, top_k }): Parameters<SearchExplainedParams>,
505 ) -> Result<String, ErrorData> {
506 let requested_k = top_k.map(|v| v as usize).unwrap_or(5);
507 let allow_superseded = query_allows_superseded(&query);
508 let search_k = if allow_superseded {
509 requested_k
510 } else {
511 (requested_k * 4).max(20)
512 };
513 let store = &self.bridge.store;
514 let result = tokio::task::block_in_place(|| {
515 Handle::current().block_on(store.search_explained(&query, Some(search_k), None, None))
516 });
517
518 match result {
519 Ok(results) => {
520 let superseded_targets = if allow_superseded {
521 HashSet::new()
522 } else {
523 load_superseded_targets(store)?
524 };
525 let fresh_results: Vec<_> = results
526 .iter()
527 .filter(|r| !superseded_targets.contains(&r.result.source.result_id()))
528 .collect();
529 let result_refs: Vec<_> =
530 if superseded_targets.is_empty() || fresh_results.is_empty() {
531 results.iter().collect()
532 } else {
533 fresh_results
534 };
535 let superseded_filtered_count = results.len().saturating_sub(result_refs.len());
536 let json_results: Vec<serde_json::Value> = result_refs
537 .iter()
538 .take(requested_k)
539 .map(|r| {
540 serde_json::json!({
541 "result_id": r.result.source.result_id(),
542 "content": r.result.content,
543 "source": format!("{:?}", r.result.source),
544 "score": r.result.score,
545 "bm25_rank": r.result.bm25_rank,
546 "vector_rank": r.result.vector_rank,
547 "cosine_similarity": r.result.cosine_similarity,
548 "breakdown": {
549 "rrf_score": r.breakdown.rrf_score,
550 "bm25_score": r.breakdown.bm25_score,
551 "vector_score": r.breakdown.vector_score,
552 "recency_score": r.breakdown.recency_score,
553 "bm25_rank": r.breakdown.bm25_rank,
554 "vector_rank": r.breakdown.vector_rank,
555 "vector_source_rank": r.breakdown.vector_source_rank,
556 "vector_source_score": r.breakdown.vector_source_score,
557 "bm25_contribution": r.breakdown.bm25_contribution,
558 "vector_contribution": r.breakdown.vector_contribution,
559 "vector_reranked_from_f32": r.breakdown.vector_reranked_from_f32,
560 "bm25_weight": r.breakdown.bm25_weight,
561 "vector_weight": r.breakdown.vector_weight,
562 "recency_weight": r.breakdown.recency_weight,
563 "rrf_k": r.breakdown.rrf_k,
564 },
565 })
566 })
567 .collect();
568 json_to_string(&serde_json::json!({
569 "ok": true,
570 "results": json_results,
571 "count": json_results.len(),
572 "superseded_filtered_count": superseded_filtered_count,
573 }))
574 }
575 Err(e) => Err(ErrorData::internal_error(
576 format!("Search error: {e}"),
577 None,
578 )),
579 }
580 }
581
582 #[tool(
583 description = "Add a fact to the knowledge base. Embedded and indexed for semantic search. Returns fact ID and content digest.",
584 annotations(idempotent_hint = true)
585 )]
586 fn sm_add_fact(
587 &self,
588 Parameters(AddFactParams {
589 content,
590 namespace,
591 source,
592 extract_entities,
593 memory_kind,
594 sensitivity,
595 evidence_refs,
596 }): Parameters<AddFactParams>,
597 ) -> Result<String, ErrorData> {
598 let store = &self.bridge.store;
599 let src = source.as_deref();
600
601 let sens = sensitivity.unwrap_or_else(|| "internal".to_string());
603 let kind = memory_kind.unwrap_or_else(|| "durable_fact".to_string());
604
605 if sens == "confidential" || sens == "restricted" {
607 return Err(ErrorData::invalid_params(
608 format!("Admission gate BLOCKED: sensitivity='{sens}' content cannot be stored without explicit user request"),
609 None,
610 ));
611 }
612
613 if kind == "ephemeral_inference" {
615 let refs = evidence_refs.as_ref().map(|v| v.len()).unwrap_or(0);
616 if refs == 0 {
617 return Err(ErrorData::invalid_params(
618 "Admission gate BLOCKED: ephemeral_inference requires evidence_refs to promote to durable".to_string(),
619 None,
620 ));
621 }
622 }
623
624 let mut meta = serde_json::Map::new();
626 meta.insert("memory_kind".to_string(), serde_json::json!(kind));
627 meta.insert("sensitivity".to_string(), serde_json::json!(sens));
628 if let Some(refs) = evidence_refs {
629 meta.insert("evidence_refs".to_string(), serde_json::json!(refs));
630 }
631 let _metadata_str = serde_json::to_string(&serde_json::Value::Object(meta)).ok();
632
633 let result = tokio::task::block_in_place(|| {
634 Handle::current().block_on(store.add_fact(&namespace, &content, src, None))
635 });
636
637 match result {
638 Ok(id) => {
639 if extract_entities == Some(true) {
641 let prompt = format!(
642 "Extract entities from this text as JSON. Format: {{\"entities\": [{{\"name\": \"...\", \"type\": \"person|project|concept|tool|version|path\"}}]}}\nText: {content}\nJSON:"
643 );
644 let body = serde_json::json!({
645 "model": "granite4.1:3b",
646 "prompt": prompt,
647 "stream": false,
648 "options": {"temperature": 0, "num_predict": 200}
649 });
650 if let Ok(resp) = reqwest::blocking::Client::new()
651 .post("http://127.0.0.1:11434/api/generate")
652 .json(&body)
653 .send()
654 {
655 if let Ok(v) = resp.json::<serde_json::Value>() {
656 if let Some(response_str) = v.get("response").and_then(|r| r.as_str()) {
657 let parsed_result =
659 boundary_compiler::parse_with_dup_check(response_str.trim());
660 if let Ok(parsed) = parsed_result {
661 if let Some(entities) =
662 parsed.get("entities").and_then(|e| e.as_array())
663 {
664 let fact_node = format!("fact:{id}");
665 for entity in entities {
666 if let Some(name) =
667 entity.get("name").and_then(|n| n.as_str())
668 {
669 let entity_node = format!("entity:{name}");
670 let _ = tokio::task::block_in_place(|| {
671 Handle::current()
672 .block_on(store.add_graph_edge(
673 &fact_node,
674 &entity_node,
675 semantic_memory::GraphEdgeType::Entity {
676 relation: "mentions".to_string(),
677 },
678 1.0,
679 None,
680 ))
681 });
682 }
683 }
684 }
685 }
686 }
687 }
688 }
689 }
690
691 json_to_string(&serde_json::json!({
692 "ok": true,
693 "fact_id": id,
694 "namespace": namespace,
695 "message": "Fact added successfully",
696 }))
697 }
698 Err(e) => Err(ErrorData::internal_error(
699 format!("Error adding fact: {e}"),
700 None,
701 )),
702 }
703 }
704
705 #[tool(
706 description = "Ingest a document with automatic chunking. Splits into chunks, each embedded and indexed. Returns document ID and chunk count.",
707 annotations(idempotent_hint = true)
708 )]
709 fn sm_ingest_document(
710 &self,
711 Parameters(IngestDocumentParams {
712 content,
713 title,
714 namespace,
715 }): Parameters<IngestDocumentParams>,
716 ) -> Result<String, ErrorData> {
717 let store = &self.bridge.store;
718 let result = tokio::task::block_in_place(|| {
719 Handle::current()
720 .block_on(store.ingest_document(&title, &content, &namespace, None, None))
721 });
722
723 match result {
724 Ok(doc_id) => {
725 let chunk_count = tokio::task::block_in_place(|| {
726 Handle::current().block_on(store.count_chunks_for_document(&doc_id))
727 })
728 .unwrap_or(0);
729 json_to_string(&serde_json::json!({
730 "ok": true,
731 "document_id": doc_id,
732 "title": title,
733 "chunk_count": chunk_count,
734 "message": "Document ingested successfully",
735 }))
736 }
737 Err(e) => Err(ErrorData::internal_error(
738 format!("Error ingesting document: {e}"),
739 None,
740 )),
741 }
742 }
743
744 #[tool(
745 description = "Get knowledge base statistics: fact/chunk/document/session counts, DB size, embedding model, and graph edge count.",
746 annotations(read_only_hint = true)
747 )]
748 fn sm_stats(&self) -> Result<String, ErrorData> {
749 let store = &self.bridge.store;
750 let result = tokio::task::block_in_place(|| Handle::current().block_on(store.stats()));
751
752 match result {
753 Ok(stats) => {
754 let graph_edge_count = tokio::task::block_in_place(|| {
757 Handle::current().block_on(store.list_all_graph_edges())
758 })
759 .map(|edges| edges.len())
760 .unwrap_or_else(|e| {
761 tracing::warn!("graph_edges table unavailable: {e}");
762 0
763 });
764 json_to_string(&serde_json::json!({
765 "ok": true,
766 "facts": stats.total_facts,
767 "chunks": stats.total_chunks,
768 "documents": stats.total_documents,
769 "sessions": stats.total_sessions,
770 "messages": stats.total_messages,
771 "graph_edges": graph_edge_count,
772 "db_size_bytes": stats.database_size_bytes,
773 "db_size_mb": (stats.database_size_bytes as f64 / 1_048_576.0 * 100.0).round() / 100.0,
774 "embedding_model": stats.embedding_model,
775 "embedding_dimensions": stats.embedding_dimensions,
776 }))
777 }
778 Err(e) => Err(ErrorData::internal_error(format!("Stats error: {e}"), None)),
779 }
780 }
781
782 #[tool(
783 description = "Find shortest path between two items in the knowledge graph. Traverses all edge types. Returns node IDs with edge evidence per hop.",
784 annotations(read_only_hint = true)
785 )]
786 fn sm_graph_path(
787 &self,
788 Parameters(GraphPathParams {
789 from_id,
790 to_id,
791 max_depth,
792 }): Parameters<GraphPathParams>,
793 ) -> Result<String, ErrorData> {
794 let depth = max_depth.map(|v| v as usize).unwrap_or(5);
795 let store = &self.bridge.store;
796 let g = store.graph_view();
797
798 match g.path(&from_id, &to_id, depth) {
799 Ok(Some(path)) => {
800 let path_segments = build_path_segments(store, &path);
802 json_to_string(&serde_json::json!({
803 "ok": true,
804 "from": from_id,
805 "to": to_id,
806 "path": path,
807 "path_length": path.len(),
808 "segments": path_segments,
809 }))
810 }
811 Ok(None) => json_to_string(&serde_json::json!({
812 "ok": true,
813 "from": from_id,
814 "to": to_id,
815 "path": null,
816 "message": format!("No path found from {from_id} to {to_id} within depth {depth}"),
817 })),
818 Err(e) => Err(ErrorData::internal_error(
819 format!("Graph view error: {e}"),
820 None,
821 )),
822 }
823 }
824
825 #[tool(
828 description = "Fetch one fact by id (bare UUID or prefixed 'fact:<uuid>'). Returns full content, namespace, source, timestamps, and metadata.",
829 annotations(read_only_hint = true)
830 )]
831 fn sm_get_fact(
832 &self,
833 Parameters(GetFactParams { fact_id }): Parameters<GetFactParams>,
834 ) -> Result<String, ErrorData> {
835 let bare = fact_id
836 .strip_prefix("fact:")
837 .unwrap_or(&fact_id)
838 .to_string();
839 let store = &self.bridge.store;
840 let result =
841 tokio::task::block_in_place(|| Handle::current().block_on(store.get_fact(&bare)));
842 match result {
843 Ok(Some(f)) => json_to_string(&serde_json::json!({
844 "ok": true,
845 "found": true,
846 "fact": {
847 "result_id": format!("fact:{}", f.id),
848 "id": f.id,
849 "namespace": f.namespace,
850 "content": f.content,
851 "source": f.source,
852 "created_at": f.created_at,
853 "updated_at": f.updated_at,
854 "metadata": f.metadata,
855 },
856 })),
857 Ok(None) => json_to_string(&serde_json::json!({
858 "ok": true,
859 "found": false,
860 "message": format!("No fact with id '{fact_id}'"),
861 })),
862 Err(e) => Err(ErrorData::internal_error(
863 format!("get_fact error: {e}"),
864 None,
865 )),
866 }
867 }
868
869 #[tool(
870 description = "Enumerate facts in a namespace (newest first) with pagination. Exhaustive, not similarity-ranked — for browsing, auditing, or deduping.",
871 annotations(read_only_hint = true)
872 )]
873 fn sm_list_facts(
874 &self,
875 Parameters(ListFactsParams {
876 namespace,
877 limit,
878 offset,
879 }): Parameters<ListFactsParams>,
880 ) -> Result<String, ErrorData> {
881 let lim = limit.map(|v| v as usize).unwrap_or(50);
882 let off = offset.map(|v| v as usize).unwrap_or(0);
883 let store = &self.bridge.store;
884 let result = tokio::task::block_in_place(|| {
885 Handle::current().block_on(store.list_facts(&namespace, lim, off))
886 });
887 match result {
888 Ok(facts) => {
889 let arr: Vec<serde_json::Value> = facts
890 .iter()
891 .map(|f| {
892 serde_json::json!({
893 "result_id": format!("fact:{}", f.id),
894 "id": f.id,
895 "namespace": f.namespace,
896 "content": f.content,
897 "source": f.source,
898 "updated_at": f.updated_at,
899 })
900 })
901 .collect();
902 json_to_string(&serde_json::json!({
903 "ok": true,
904 "namespace": namespace,
905 "count": arr.len(),
906 "limit": lim,
907 "offset": off,
908 "facts": arr,
909 }))
910 }
911 Err(e) => Err(ErrorData::internal_error(
912 format!("list_facts error: {e}"),
913 None,
914 )),
915 }
916 }
917
918 #[tool(
919 description = "List namespaces that currently contain facts. Use before sm_list_facts to discover what is stored.",
920 annotations(read_only_hint = true)
921 )]
922 fn sm_list_namespaces(&self) -> Result<String, ErrorData> {
923 let store = &self.bridge.store;
924 let result = tokio::task::block_in_place(|| {
925 Handle::current().block_on(store.list_fact_namespaces())
926 });
927 match result {
928 Ok(ns) => json_to_string(&serde_json::json!({
929 "ok": true,
930 "count": ns.len(),
931 "namespaces": ns,
932 })),
933 Err(e) => Err(ErrorData::internal_error(
934 format!("list_namespaces error: {e}"),
935 None,
936 )),
937 }
938 }
939
940 #[tool(
941 description = "Fetch a fact plus its graph neighbors WITH their content in one call. Hydrates neighbor facts for ids returned by graph tools.",
942 annotations(read_only_hint = true)
943 )]
944 fn sm_get_fact_neighbors(
945 &self,
946 Parameters(GetFactNeighborsParams { item_id }): Parameters<GetFactNeighborsParams>,
947 ) -> Result<String, ErrorData> {
948 let node_id = if item_id.contains(':') {
949 item_id.clone()
950 } else {
951 format!("fact:{item_id}")
952 };
953 let bare = node_id
954 .strip_prefix("fact:")
955 .unwrap_or(&node_id)
956 .to_string();
957 let store = &self.bridge.store;
958
959 let center =
960 tokio::task::block_in_place(|| Handle::current().block_on(store.get_fact(&bare)))
961 .map_err(|e| ErrorData::internal_error(format!("get_fact error: {e}"), None))?;
962 let edges = tokio::task::block_in_place(|| {
963 Handle::current().block_on(store.list_graph_edges_for_node(&node_id))
964 })
965 .map_err(|e| ErrorData::internal_error(format!("list edges error: {e}"), None))?;
966
967 let mut neighbors: Vec<serde_json::Value> = Vec::new();
968 for e in &edges {
969 let outgoing = e.source == node_id;
970 let other = if outgoing { &e.target } else { &e.source };
971 let other_bare = other.strip_prefix("fact:").unwrap_or(other).to_string();
972 let content = tokio::task::block_in_place(|| {
973 Handle::current().block_on(store.get_fact(&other_bare))
974 })
975 .ok()
976 .flatten()
977 .map(|f| f.content);
978 neighbors.push(serde_json::json!({
979 "neighbor_id": other,
980 "direction": if outgoing { "out" } else { "in" },
981 "edge_type": e.edge_type,
982 "weight": e.weight,
983 "content": content,
984 }));
985 }
986 json_to_string(&serde_json::json!({
987 "ok": true,
988 "item_id": node_id,
989 "center_content": center.map(|f| f.content),
990 "neighbor_count": neighbors.len(),
991 "neighbors": neighbors,
992 }))
993 }
994
995 #[tool(
996 description = "Create a replacement fact and link it to a stale fact via 'supersedes' edge. Use instead of deleting outdated facts. Returns new fact id and edge id.",
997 annotations(idempotent_hint = true)
998 )]
999 fn sm_supersede_fact(
1000 &self,
1001 Parameters(SupersedeFactParams {
1002 old_fact_id,
1003 content,
1004 namespace,
1005 source,
1006 reason,
1007 }): Parameters<SupersedeFactParams>,
1008 ) -> Result<String, ErrorData> {
1009 use semantic_memory::GraphEdgeType;
1010
1011 let old_bare = old_fact_id
1012 .strip_prefix("fact:")
1013 .unwrap_or(&old_fact_id)
1014 .to_string();
1015 let old_node = format!("fact:{old_bare}");
1016 let store = &self.bridge.store;
1017 let old =
1018 tokio::task::block_in_place(|| Handle::current().block_on(store.get_fact(&old_bare)))
1019 .map_err(|e| ErrorData::internal_error(format!("get old fact error: {e}"), None))?;
1020 let Some(old_fact) = old else {
1021 return Err(ErrorData::invalid_params(
1022 format!("No fact with id '{old_fact_id}'"),
1023 None,
1024 ));
1025 };
1026
1027 let ns = namespace.unwrap_or_else(|| old_fact.namespace.clone());
1028 let new_id = tokio::task::block_in_place(|| {
1029 Handle::current().block_on(store.add_fact(&ns, &content, source.as_deref(), None))
1030 })
1031 .map_err(|e| ErrorData::internal_error(format!("add replacement fact error: {e}"), None))?;
1032 let new_node = format!("fact:{new_id}");
1033 let metadata = serde_json::json!({
1034 "reason": reason.unwrap_or_else(|| "replacement fact supersedes stale fact".to_string()),
1035 "old_fact_id": old_bare,
1036 });
1037 let edge = tokio::task::block_in_place(|| {
1038 Handle::current().block_on(store.add_graph_edge(
1039 &new_node,
1040 &old_node,
1041 GraphEdgeType::Entity {
1042 relation: "supersedes".to_string(),
1043 },
1044 1.0,
1045 Some(metadata),
1046 ))
1047 })
1048 .map_err(|e| ErrorData::internal_error(format!("add supersedes edge error: {e}"), None))?;
1049
1050 json_to_string(&serde_json::json!({
1051 "ok": true,
1052 "new_fact_id": new_id,
1053 "new_result_id": new_node,
1054 "old_fact_id": old_bare,
1055 "old_result_id": old_node,
1056 "namespace": ns,
1057 "edge_id": edge.id,
1058 "relation": "supersedes",
1059 }))
1060 }
1061
1062 #[tool(
1065 description = "Create a conversation session (container for messages). Returns session id. Use to persist history recallable via sm_search_conversations.",
1066 annotations(idempotent_hint = true)
1067 )]
1068 fn sm_create_session(
1069 &self,
1070 Parameters(CreateSessionParams { channel, metadata }): Parameters<CreateSessionParams>,
1071 ) -> Result<String, ErrorData> {
1072 let meta: Option<serde_json::Value> = metadata
1073 .as_deref()
1074 .and_then(|s| serde_json::from_str(s).ok());
1075 let store = &self.bridge.store;
1076 let result = tokio::task::block_in_place(|| {
1077 Handle::current().block_on(store.create_session_with_metadata(&channel, meta))
1078 });
1079 match result {
1080 Ok(id) => json_to_string(
1081 &serde_json::json!({"ok": true, "session_id": id, "channel": channel}),
1082 ),
1083 Err(e) => Err(ErrorData::internal_error(
1084 format!("create_session error: {e}"),
1085 None,
1086 )),
1087 }
1088 }
1089
1090 #[tool(
1091 description = "Append a message to a session. role: user|assistant|system|tool. Message is embedded and FTS-indexed. Returns message id."
1092 )]
1093 fn sm_add_message(
1094 &self,
1095 Parameters(AddMessageParams {
1096 session_id,
1097 role,
1098 content,
1099 }): Parameters<AddMessageParams>,
1100 ) -> Result<String, ErrorData> {
1101 let parsed_role = match role.to_lowercase().as_str() {
1102 "user" => semantic_memory::types::Role::User,
1103 "assistant" => semantic_memory::types::Role::Assistant,
1104 "system" => semantic_memory::types::Role::System,
1105 "tool" => semantic_memory::types::Role::Tool,
1106 other => {
1107 return Err(ErrorData::invalid_params(
1108 format!("invalid role '{other}' (use user|assistant|system|tool)"),
1109 None,
1110 ))
1111 }
1112 };
1113 let store = &self.bridge.store;
1114 let result = tokio::task::block_in_place(|| {
1115 Handle::current().block_on(store.add_message_embedded(
1116 &session_id,
1117 parsed_role,
1118 &content,
1119 None,
1120 None,
1121 ))
1122 });
1123 match result {
1124 Ok(id) => json_to_string(
1125 &serde_json::json!({"ok": true, "message_id": id, "session_id": session_id}),
1126 ),
1127 Err(e) => Err(ErrorData::internal_error(
1128 format!("add_message error: {e}"),
1129 None,
1130 )),
1131 }
1132 }
1133
1134 #[tool(description = "List recent conversation sessions (newest first) with message counts.", annotations(read_only_hint = true))]
1135 fn sm_list_sessions(
1136 &self,
1137 Parameters(ListSessionsParams { limit, offset }): Parameters<ListSessionsParams>,
1138 ) -> Result<String, ErrorData> {
1139 let lim = limit.map(|v| v as usize).unwrap_or(20);
1140 let off = offset.map(|v| v as usize).unwrap_or(0);
1141 let store = &self.bridge.store;
1142 let result = tokio::task::block_in_place(|| {
1143 Handle::current().block_on(store.list_sessions(lim, off))
1144 });
1145 match result {
1146 Ok(sessions) => json_to_string(&serde_json::json!({
1147 "ok": true,
1148 "count": sessions.len(),
1149 "sessions": sessions.iter().map(|s| serde_json::json!({
1150 "session_id": s.id,
1151 "channel": s.channel,
1152 "message_count": s.message_count,
1153 "created_at": s.created_at,
1154 "updated_at": s.updated_at,
1155 })).collect::<Vec<_>>(),
1156 })),
1157 Err(e) => Err(ErrorData::internal_error(
1158 format!("list_sessions error: {e}"),
1159 None,
1160 )),
1161 }
1162 }
1163
1164 #[tool(
1165 description = "Get most recent messages from a session within a token budget (default 4000), chronological order. Returns role, content, timestamps.",
1166 annotations(read_only_hint = true)
1167 )]
1168 fn sm_get_messages(
1169 &self,
1170 Parameters(GetMessagesParams {
1171 session_id,
1172 max_tokens,
1173 }): Parameters<GetMessagesParams>,
1174 ) -> Result<String, ErrorData> {
1175 let budget = max_tokens.unwrap_or(4000);
1176 let store = &self.bridge.store;
1177 let result = tokio::task::block_in_place(|| {
1178 Handle::current().block_on(store.get_messages_within_budget(&session_id, budget))
1179 });
1180 match result {
1181 Ok(msgs) => json_to_string(&serde_json::json!({
1182 "ok": true,
1183 "session_id": session_id,
1184 "count": msgs.len(),
1185 "messages": msgs.iter().map(|m| serde_json::json!({
1186 "id": m.id,
1187 "role": m.role,
1188 "content": m.content,
1189 "token_count": m.token_count,
1190 "created_at": m.created_at,
1191 })).collect::<Vec<_>>(),
1192 })),
1193 Err(e) => Err(ErrorData::internal_error(
1194 format!("get_messages error: {e}"),
1195 None,
1196 )),
1197 }
1198 }
1199
1200 #[tool(
1201 description = "Hybrid semantic search over stored conversation MESSAGES (not facts). Recall what was discussed in past sessions. Returns ranked messages.",
1202 annotations(read_only_hint = true)
1203 )]
1204 fn sm_search_conversations(
1205 &self,
1206 Parameters(SearchConversationsParams { query, top_k }): Parameters<
1207 SearchConversationsParams,
1208 >,
1209 ) -> Result<String, ErrorData> {
1210 let k = top_k.map(|v| v as usize);
1211 let store = &self.bridge.store;
1212 let result = tokio::task::block_in_place(|| {
1213 Handle::current().block_on(store.search_conversations(&query, k, None))
1214 });
1215 match result {
1216 Ok(results) => json_to_string(&serde_json::json!({
1217 "ok": true,
1218 "count": results.len(),
1219 "results": results.iter().map(|r| serde_json::json!({
1220 "result_id": r.source.result_id(),
1221 "content": r.content,
1222 "score": r.score,
1223 "cosine_similarity": r.cosine_similarity,
1224 })).collect::<Vec<_>>(),
1225 })),
1226 Err(e) => Err(ErrorData::internal_error(
1227 format!("search_conversations error: {e}"),
1228 None,
1229 )),
1230 }
1231 }
1232
1233 #[tool(
1240 description = "Profile a query and get an adaptive routing decision. Determines which retrieval stages (BM25, vector, rerank, graph, decoder, discord) to activate.",
1241 annotations(read_only_hint = true)
1242 )]
1243 fn sm_route_query(
1244 &self,
1245 Parameters(RouteQueryParams { query }): Parameters<RouteQueryParams>,
1246 ) -> Result<String, ErrorData> {
1247 use semantic_memory::routing::RetrievalRouter;
1248
1249 let router = RetrievalRouter {
1250 decoder_enabled: true,
1251 discord_enabled: true,
1252 corpus_density: 0.5,
1253 ..Default::default()
1254 };
1255
1256 let decision = router.route_query(&query);
1257 json_to_string(&serde_json::json!({
1258 "ok": true,
1259 "bm25_coarse": decision.bm25_coarse,
1260 "vector_medium": decision.vector_medium,
1261 "rerank_fine": decision.rerank_fine,
1262 "graph_expansion": decision.graph_expansion,
1263 "decoder": decision.decoder,
1264 "discord": decision.discord,
1265 "no_retrieval": decision.no_retrieval,
1266 "reasoning": decision.reasoning,
1267 }))
1268 }
1269
1270 #[tool(
1271 description = "Adaptive search: profiles query, routes to appropriate stages, applies factor graph belief propagation if decoder is activated. Returns results with stable IDs.",
1272 annotations(read_only_hint = true)
1273 )]
1274 fn sm_search_with_routing(
1275 &self,
1276 Parameters(SearchWithRoutingParams {
1277 query,
1278 top_k,
1279 contradictions,
1280 group_by_community,
1281 }): Parameters<SearchWithRoutingParams>,
1282 ) -> Result<String, ErrorData> {
1283 use semantic_memory::integration::plan_execution;
1284 use semantic_memory::rl_routing::route_with_rl;
1285 use semantic_memory::routing::QueryProfile;
1286
1287 let k = top_k.map(|v| v as usize).unwrap_or(5);
1288 let allow_superseded = query_allows_superseded(&query);
1289 let search_k = if allow_superseded { k } else { (k * 4).max(20) };
1290
1291 let store = &self.bridge.store;
1293 let policy =
1294 tokio::task::block_in_place(|| Handle::current().block_on(store.load_routing_policy()))
1295 .ok()
1296 .flatten()
1297 .unwrap_or_default();
1298 let profile = QueryProfile::from_query(&query);
1299 let decision = route_with_rl(&policy, &profile);
1300 let contras = contradictions.unwrap_or_default();
1301 let plan = plan_execution(&decision, contras.clone());
1302
1303 let store = &self.bridge.store;
1304 let search_result = tokio::task::block_in_place(|| {
1305 Handle::current().block_on(store.search(&query, Some(search_k), None, None))
1306 });
1307
1308 match search_result {
1309 Ok(results) => {
1310 let superseded_targets = if allow_superseded {
1311 HashSet::new()
1312 } else {
1313 load_superseded_targets(store)?
1314 };
1315 let fresh_results: Vec<_> = results
1316 .iter()
1317 .filter(|r| !superseded_targets.contains(&r.source.result_id()))
1318 .collect();
1319 let result_refs: Vec<_> =
1320 if superseded_targets.is_empty() || fresh_results.is_empty() {
1321 results.iter().collect()
1322 } else {
1323 fresh_results
1324 };
1325 let superseded_filtered_count = results.len().saturating_sub(result_refs.len());
1326 let json_results: Vec<serde_json::Value> = result_refs
1327 .iter()
1328 .take(k)
1329 .map(|r| {
1330 serde_json::json!({
1331 "result_id": r.source.result_id(),
1332 "content": r.content,
1333 "score": r.score,
1334 })
1335 })
1336 .collect();
1337
1338 let mut factor_graph_payload = serde_json::json!({
1339 "enabled": false,
1340 });
1341
1342 let mut decoder_executed = false;
1343 let mut discord_executed = false;
1344 let mut discord_results_payload: Vec<serde_json::Value> = Vec::new();
1345
1346 if decision.decoder {
1347 #[cfg(feature = "full")]
1348 {
1349 use semantic_memory::factor_graph::{
1350 factors_from_edges, FactorGraph, FactorGraphConfig,
1351 };
1352
1353 let graph_edges = tokio::task::block_in_place(|| {
1354 Handle::current().block_on(store.list_all_graph_edges())
1355 });
1356
1357 match graph_edges {
1358 Ok(edges) => {
1359 let raw_edges: Vec<(
1360 String,
1361 String,
1362 semantic_memory::GraphEdgeType,
1363 f64,
1364 Option<String>,
1365 )> = edges
1366 .iter()
1367 .map(|edge| {
1368 let parsed_type = edge
1369 .edge_type_parsed
1370 .clone()
1371 .or_else(|| serde_json::from_str(&edge.edge_type).ok())
1372 .unwrap_or(semantic_memory::GraphEdgeType::Entity {
1373 relation: "unknown".to_string(),
1374 });
1375 (
1376 edge.source.clone(),
1377 edge.target.clone(),
1378 parsed_type,
1379 edge.weight,
1380 edge.metadata.clone(),
1381 )
1382 })
1383 .collect();
1384
1385 let nodes: Vec<(String, f64)> = result_refs
1386 .iter()
1387 .map(|r| (r.source.result_id(), r.score))
1388 .collect();
1389 let factors = factors_from_edges(&raw_edges);
1390 let graph =
1391 FactorGraph::new(&nodes, factors, FactorGraphConfig::default());
1392 let propagated = graph.propagate();
1393 let top_beliefs = propagated.top_k(k);
1394
1395 factor_graph_payload = serde_json::json!({
1396 "enabled": true,
1397 "top_k_beliefs": top_beliefs
1398 .into_iter()
1399 .map(|(item_id, belief)| serde_json::json!({
1400 "item_id": item_id,
1401 "belief": belief,
1402 }))
1403 .collect::<Vec<_>>(),
1404 "iterations": propagated.iterations,
1405 "converged": propagated.converged,
1406 "elapsed_ms": propagated.elapsed_ms,
1407 "factor_counts": {
1408 "semantic": propagated.factor_counts.semantic,
1409 "temporal": propagated.factor_counts.temporal,
1410 "causal": propagated.factor_counts.causal,
1411 "entity": propagated.factor_counts.entity,
1412 "total": propagated.factor_counts.total(),
1413 },
1414 });
1415 decoder_executed = true;
1416 }
1417 Err(e) => {
1418 factor_graph_payload = serde_json::json!({
1419 "enabled": false,
1420 "error": format!("factor graph analysis failed: {e}"),
1421 });
1422 }
1423 }
1424 }
1425
1426 #[cfg(not(feature = "full"))]
1427 {
1428 factor_graph_payload = serde_json::json!({
1429 "enabled": false,
1430 "reason": "factor graph analysis requires the `full` feature",
1431 });
1432 }
1433
1434 if !plan.contradictions.is_empty() {
1435 use semantic_memory::decoder::{compute_correction, detect_syndromes};
1436 let result_scores: Vec<(String, f64)> = result_refs
1437 .iter()
1438 .map(|r| (r.source.result_id(), r.score))
1439 .collect();
1440 let syndromes = detect_syndromes(&result_scores, &plan.contradictions);
1441 let _ = compute_correction(&syndromes, 10.0);
1442 decoder_executed = true;
1443 }
1444 }
1445
1446 if plan.use_discord {
1447 use semantic_memory::discord::DiscordScorer;
1448 let direct_ids: Vec<String> =
1449 result_refs.iter().map(|r| r.source.result_id()).collect();
1450 let existing_ids: std::collections::HashSet<String> =
1451 direct_ids.iter().cloned().collect();
1452 if let Ok(edges) = load_neighborhood_edge_refs(&self.bridge.store, &direct_ids)
1453 {
1454 let scorer = DiscordScorer::with_defaults();
1455 let discord_hits = scorer.score(&direct_ids, &edges);
1456 for hit in &discord_hits {
1457 if !existing_ids.contains(&hit.item_id) {
1458 discord_results_payload.push(serde_json::json!({
1459 "result_id": hit.item_id,
1460 "discord_score": hit.discord_score,
1461 "anchor_ids": hit.anchor_ids,
1462 "relationship_types": hit.relationship_types,
1463 }));
1464 }
1465 }
1466 discord_executed = true;
1467 }
1468 }
1469
1470 let mut matryoshka_payload = serde_json::json!({
1471 "enabled": false,
1472 });
1473 if decision.vector_medium {
1474 #[cfg(feature = "full")]
1475 {
1476 use semantic_memory::integration::multi_resolution_route;
1477 use semantic_memory::matryoshka::MatryoshkaConfig;
1478 use semantic_memory::routing::QueryProfile;
1479
1480 let route_profile = QueryProfile::from_query(&query);
1481 let route_decision =
1482 multi_resolution_route(&route_profile, &MatryoshkaConfig::default());
1483 matryoshka_payload = serde_json::json!({
1484 "enabled": true,
1485 "candidate_dim": route_decision.candidate_dim,
1486 "heuristic_recall_estimate": route_decision.estimated_recall,
1487 "recall_basis": "heuristic_dimensional_model_not_corpus_measured",
1488 "embedding_dim": route_decision.embedding_dim,
1489 "reasoning": route_decision.reasoning,
1490 });
1491 }
1492
1493 #[cfg(not(feature = "full"))]
1494 {
1495 matryoshka_payload = serde_json::json!({
1496 "enabled": false,
1497 "reason": "matryoshka routing requires the `full` feature",
1498 });
1499 }
1500 }
1501
1502 let grouped_results_payload: serde_json::Value = if group_by_community == Some(true)
1504 {
1505 let seed_ids: Vec<String> = result_refs
1506 .iter()
1507 .take(k)
1508 .map(|r| r.source.result_id())
1509 .collect();
1510 let edges = load_neighborhood_edge_pairs(store, &seed_ids).unwrap_or_default();
1511 if !edges.is_empty() {
1512 use semantic_memory::community::detect_communities;
1513 let communities = detect_communities(&edges, 1.0, 42);
1514 let mut member_to_comm: std::collections::HashMap<String, String> =
1515 std::collections::HashMap::new();
1516 for c in &communities {
1517 for m in &c.members {
1518 member_to_comm.insert(m.clone(), c.id.clone());
1519 }
1520 }
1521 let mut groups: std::collections::HashMap<String, Vec<serde_json::Value>> =
1522 std::collections::HashMap::new();
1523 let mut ungrouped: Vec<serde_json::Value> = Vec::new();
1524 for r in &json_results {
1525 if let Some(rid) = r.get("result_id").and_then(|v| v.as_str()) {
1526 match member_to_comm.get(rid).cloned() {
1527 Some(cid) => groups.entry(cid).or_default().push(r.clone()),
1528 None => ungrouped.push(r.clone()),
1529 }
1530 }
1531 }
1532 let mut map = serde_json::Map::new();
1533 for (cid, items) in groups {
1534 map.insert(format!("community_{cid}"), serde_json::json!(items));
1535 }
1536 if !ungrouped.is_empty() {
1537 map.insert("ungrouped".to_string(), serde_json::json!(ungrouped));
1538 }
1539 serde_json::Value::Object(map)
1540 } else {
1541 serde_json::Value::Null
1542 }
1543 } else {
1544 serde_json::Value::Null
1545 };
1546
1547 let mut topology_payload = serde_json::json!({ "auto_called": false });
1549 {
1550 use semantic_memory::routing::{QueryComplexityClass, QueryProfile};
1551 let route_profile = QueryProfile::from_query(&query);
1552 if route_profile.complexity_class == QueryComplexityClass::Synthesis
1553 && result_refs.len() > 10
1554 {
1555 #[cfg(feature = "full")]
1556 {
1557 use semantic_memory::topology::{compute_betti_numbers, find_voids};
1558 let edges = load_stored_edge_pairs(store).unwrap_or_default();
1559 if !edges.is_empty() {
1560 let mut adjacency: std::collections::HashMap<String, Vec<String>> =
1561 std::collections::HashMap::new();
1562 for (src, tgt) in &edges {
1563 adjacency.entry(src.clone()).or_default().push(tgt.clone());
1564 adjacency.entry(tgt.clone()).or_default().push(src.clone());
1565 }
1566 let betti = compute_betti_numbers(&adjacency);
1567 let voids = find_voids(&edges);
1568 topology_payload = serde_json::json!({
1569 "auto_called": true,
1570 "trigger": "synthesis_class_with_10_plus_results",
1571 "betti_numbers": {
1572 "betti_0": betti.betti_0,
1573 "betti_1": betti.betti_1,
1574 },
1575 "void_count": voids.len(),
1576 "voids": voids.iter().map(|v| serde_json::json!({
1577 "description": v.description,
1578 "void_type": format!("{:?}", v.void_type),
1579 "nearby_items": v.nearby_items,
1580 "suggested_connections": v.suggested_connections,
1581 })).collect::<Vec<_>>(),
1582 });
1583 } else {
1584 topology_payload = serde_json::json!({
1585 "auto_called": true,
1586 "trigger": "synthesis_class_with_10_plus_results",
1587 "note": "no graph edges in store",
1588 });
1589 }
1590 }
1591 #[cfg(not(feature = "full"))]
1592 {
1593 topology_payload = serde_json::json!({
1594 "auto_called": true,
1595 "trigger": "synthesis_class_with_10_plus_results",
1596 "error": "topology requires the full feature",
1597 });
1598 }
1599 }
1600 }
1601
1602 json_to_string(&serde_json::json!({
1603 "ok": true,
1604 "routing_decision": {
1605 "bm25_coarse": decision.bm25_coarse,
1606 "vector_medium": decision.vector_medium,
1607 "rerank_fine": decision.rerank_fine,
1608 "graph_expansion": decision.graph_expansion,
1609 "decoder": decision.decoder,
1610 "discord": decision.discord,
1611 "no_retrieval": decision.no_retrieval,
1612 "reasoning": decision.reasoning,
1613 },
1614 "results": json_results,
1615 "count": json_results.len(),
1616 "superseded_filtered_count": superseded_filtered_count,
1617 "decoder_planned": plan.use_decoder,
1618 "decoder_executed": decoder_executed,
1619 "discord_planned": plan.use_discord,
1620 "discord_executed": discord_executed,
1621 "discord_results": discord_results_payload,
1622 "factor_graph": factor_graph_payload,
1623 "matryoshka": matryoshka_payload,
1624 "grouped_results": grouped_results_payload,
1625 "topology": topology_payload,
1626 }))
1627 }
1628 Err(e) => Err(ErrorData::internal_error(
1629 format!("Search error: {e}"),
1630 None,
1631 )),
1632 }
1633 }
1634
1635 #[tool(
1636 description = "Detect contradictions in search results. Runs syndrome detection, computes corrections, and applies belief propagation to refine confidence scores.",
1637 annotations(read_only_hint = true)
1638 )]
1639 fn sm_decoder_analyze(
1640 &self,
1641 Parameters(DecoderAnalyzeParams {
1642 results,
1643 contradictions,
1644 }): Parameters<DecoderAnalyzeParams>,
1645 ) -> Result<String, ErrorData> {
1646 use semantic_memory::decoder::{
1647 compute_correction, detect_syndromes, pass_messages, ConflictGraph,
1648 };
1649
1650 let contras = contradictions.unwrap_or_default();
1651 let syndromes = detect_syndromes(&results, &contras);
1652 let corrections = compute_correction(&syndromes, 10.0);
1653 let graph = ConflictGraph::from_syndromes(&results, &syndromes);
1654 let mp = pass_messages(&graph, 50, 0.001);
1655
1656 json_to_string(&serde_json::json!({
1657 "ok": true,
1658 "syndromes": syndromes.iter().map(|s| serde_json::json!({
1659 "id": s.id,
1660 "severity": format!("{:?}", s.severity),
1661 "items": s.items,
1662 "description": s.description,
1663 "type": format!("{:?}", s.syndrome_type),
1664 })).collect::<Vec<_>>(),
1665 "syndrome_count": syndromes.len(),
1666 "corrections": corrections.iter().map(|c| serde_json::json!({
1667 "id": c.id,
1668 "confidence": c.confidence,
1669 "cost": c.cost,
1670 "operations": c.operations.len(),
1671 })).collect::<Vec<_>>(),
1672 "correction_count": corrections.len(),
1673 "message_passing": {
1674 "iterations": mp.iterations,
1675 "converged": mp.converged,
1676 "elapsed_ms": mp.elapsed_ms,
1677 },
1678 }))
1679 }
1680
1681 #[tool(
1682 description = "Detect contradictions among the top results for a query from their CONTENT (numeric, value, negation, or antonym disagreement) — no pre-asserted edges required. Returns candidate conflicting pairs, each with the signals that fired and a human-readable reason. Persist a confirmed pair with sm_add_graph_edge(edge_type=\"contradicts\") so the decoder/community/factor-graph tools pick it up.",
1683 annotations(read_only_hint = true)
1684 )]
1685 fn sm_detect_contradictions(
1686 &self,
1687 Parameters(DetectContradictionsParams { query, top_k }): Parameters<
1688 DetectContradictionsParams,
1689 >,
1690 ) -> Result<String, ErrorData> {
1691 use semantic_memory::contradiction_detect::{detect_contradictions, DetectorConfig};
1692
1693 let k = top_k.map(|v| v as usize).unwrap_or(10);
1694 let store = &self.bridge.store;
1695 let results = tokio::task::block_in_place(|| {
1696 Handle::current().block_on(store.search(&query, Some(k), None, None))
1697 })
1698 .map_err(|e| ErrorData::internal_error(format!("search failed: {e}"), None))?;
1699
1700 let items: Vec<(String, String)> = results
1701 .iter()
1702 .map(|r| (r.source.result_id(), r.content.clone()))
1703 .collect();
1704
1705 let pairs = detect_contradictions(&items, &DetectorConfig::default());
1706
1707 json_to_string(&serde_json::json!({
1708 "ok": true,
1709 "query": query,
1710 "items_scanned": items.len(),
1711 "contradictions": pairs.iter().map(|p| serde_json::json!({
1712 "a": p.a,
1713 "b": p.b,
1714 "score": p.score,
1715 "signals": p.signals.iter().map(|s| format!("{s:?}")).collect::<Vec<_>>(),
1716 "reason": p.reason,
1717 })).collect::<Vec<_>>(),
1718 "count": pairs.len(),
1719 }))
1720 }
1721
1722 #[tool(
1723 description = "Second-order retrieval: find items related to your search results through the graph, but NOT themselves direct hits. Loads edges from store automatically.",
1724 annotations(read_only_hint = true)
1725 )]
1726 fn sm_discord_search(
1727 &self,
1728 Parameters(DiscordSearchParams { direct_result_ids }): Parameters<DiscordSearchParams>,
1729 ) -> Result<String, ErrorData> {
1730 use semantic_memory::discord::DiscordScorer;
1731
1732 let edges = load_neighborhood_edge_refs(&self.bridge.store, &direct_result_ids)?;
1735 let scorer = DiscordScorer::with_defaults();
1736 let results = scorer.score(&direct_result_ids, &edges);
1737
1738 json_to_string(&serde_json::json!({
1739 "ok": true,
1740 "discord_results": results.iter().map(|r| serde_json::json!({
1741 "item_id": r.item_id,
1742 "discord_score": r.discord_score,
1743 "anchor_ids": r.anchor_ids,
1744 "relationship_types": r.relationship_types,
1745 })).collect::<Vec<_>>(),
1746 "count": results.len(),
1747 "edges_loaded": edges.len(),
1748 "edges_scope": "neighborhood",
1749 }))
1750 }
1751
1752 #[tool(
1753 description = "Set provenance (evidence confidence) for an item. Confidence in [0.0, 1.0] with support count. Returns a provenance receipt.",
1754 annotations(idempotent_hint = true)
1755 )]
1756 fn sm_set_provenance(
1757 &self,
1758 Parameters(SetProvenanceParams {
1759 item_id,
1760 confidence,
1761 support_count,
1762 }): Parameters<SetProvenanceParams>,
1763 ) -> Result<String, ErrorData> {
1764 use semantic_memory::provenance::{
1765 ConfidenceSemiring, ConfidenceValue, ProvenanceItemType,
1766 };
1767
1768 if !confidence.is_finite() || confidence < 0.0 || confidence > 1.0 {
1770 return Err(ErrorData::invalid_params(
1771 format!("confidence must be a finite value in [0.0, 1.0], got {confidence}"),
1772 None,
1773 ));
1774 }
1775
1776 let value = ConfidenceValue::new(confidence, support_count);
1777 let store = &self.bridge.store;
1778
1779 let result = tokio::task::block_in_place(|| {
1780 Handle::current().block_on(store.set_provenance::<ConfidenceSemiring>(
1781 &ProvenanceItemType::Fact,
1782 &item_id,
1783 &value,
1784 &[],
1785 None,
1786 ))
1787 });
1788
1789 match result {
1790 Ok(receipt) => json_to_string(&serde_json::json!({
1791 "ok": true,
1792 "provenance_id": receipt.provenance_id,
1793 "item_id": receipt.item_id,
1794 "semiring_type": receipt.semiring_type,
1795 "recorded_at": receipt.recorded_at,
1796 "message": "Provenance set successfully",
1797 })),
1798 Err(e) => Err(ErrorData::internal_error(
1799 format!("Provenance error: {e}"),
1800 None,
1801 )),
1802 }
1803 }
1804
1805 #[tool(
1806 description = "Run a memory lifecycle pass: analyze items for syndromes, compute corrections, identify subtraction candidates, and check compression needs.",
1807 annotations(read_only_hint = true)
1808 )]
1809 fn sm_run_lifecycle(
1810 &self,
1811 Parameters(RunLifecycleParams { item_ids }): Parameters<RunLifecycleParams>,
1812 ) -> Result<String, ErrorData> {
1813 use semantic_memory::decoder::{compute_correction, detect_syndromes};
1814 use semantic_memory::integration::{
1815 corrections_to_subtraction_candidates, should_trigger_recompression,
1816 };
1817
1818 let results: Vec<(String, f64)> = item_ids.iter().map(|id| (id.clone(), 0.5)).collect();
1819 let syndromes = detect_syndromes(&results, &[]);
1820 let corrections = compute_correction(&syndromes, 10.0);
1821
1822 let sub_candidates = corrections_to_subtraction_candidates(&corrections);
1823
1824 let subtracted_count = sub_candidates.len();
1825 let remaining_count = item_ids.len().saturating_sub(subtracted_count);
1826 let recompression = should_trigger_recompression(subtracted_count, remaining_count, false);
1827
1828 let store = &self.bridge.store;
1829 let graph_edges = tokio::task::block_in_place(|| {
1830 Handle::current().block_on(store.list_all_graph_edges())
1831 });
1832 let stored_edges: Vec<(String, String)> = graph_edges
1833 .as_ref()
1834 .map(|edges| {
1835 edges
1836 .iter()
1837 .map(|edge| (edge.source.clone(), edge.target.clone()))
1838 .collect()
1839 })
1840 .unwrap_or_default();
1841
1842 let mut topology_voids: Vec<serde_json::Value> = Vec::new();
1843 let mut betti = serde_json::json!({
1844 "betti_0": 0usize,
1845 "betti_1": 0usize,
1846 });
1847 let mut topology_error: Option<String> = None;
1848
1849 let mut communities: Vec<serde_json::Value> = Vec::new();
1850 let mut community_contradictions: Vec<serde_json::Value> = Vec::new();
1851 let mut community_error: Option<String> = None;
1852
1853 let mut subgraph_assessment = serde_json::json!({
1854 "subgraphs_identified": 0usize,
1855 "subgraphs_pruned": 0usize,
1856 });
1857 let mut subgraph_error: Option<String> = None;
1858
1859 #[cfg(feature = "full")]
1860 {
1861 use std::collections::HashMap;
1862
1863 if !stored_edges.is_empty() {
1864 let analysis_edges = stored_edges.clone();
1865
1866 let topology_result = (|| -> Result<(), String> {
1867 use semantic_memory::topology::{compute_betti_numbers, find_voids};
1868
1869 let mut adjacency: HashMap<String, Vec<String>> = HashMap::new();
1870 for (left, right) in &analysis_edges {
1871 adjacency
1872 .entry(left.clone())
1873 .or_default()
1874 .push(right.clone());
1875 adjacency
1876 .entry(right.clone())
1877 .or_default()
1878 .push(left.clone());
1879 }
1880
1881 let betti_numbers = compute_betti_numbers(&adjacency);
1882 betti = serde_json::json!({
1883 "betti_0": betti_numbers.betti_0,
1884 "betti_1": betti_numbers.betti_1,
1885 });
1886
1887 topology_voids = find_voids(&analysis_edges)
1888 .into_iter()
1889 .map(|v| {
1890 serde_json::json!({
1891 "description": v.description,
1892 "void_type": format!("{:?}", v.void_type),
1893 "nearby_items": v.nearby_items,
1894 "suggested_connections": v.suggested_connections,
1895 })
1896 })
1897 .collect();
1898
1899 Ok(())
1900 })();
1901
1902 if let Err(e) = topology_result {
1903 topology_error = Some(e);
1904 }
1905
1906 let community_result = (|| -> Result<(), String> {
1907 use semantic_memory::community::{
1908 community_contradiction_scan, detect_communities,
1909 };
1910
1911 let detected = detect_communities(&analysis_edges, 1.0, 42);
1912 communities = detected
1913 .iter()
1914 .map(|c| {
1915 serde_json::json!({
1916 "id": c.id,
1917 "members": c.members,
1918 "level": c.level,
1919 "parent": c.parent,
1920 "member_count": c.members.len(),
1921 })
1922 })
1923 .collect();
1924
1925 community_contradictions = community_contradiction_scan(&detected, &[])
1926 .into_iter()
1927 .map(|cc| {
1928 serde_json::json!({
1929 "community_id": cc.community_id,
1930 "item_a": cc.item_a,
1931 "item_b": cc.item_b,
1932 "description": cc.description,
1933 })
1934 })
1935 .collect();
1936
1937 Ok(())
1938 })();
1939
1940 if let Err(e) = community_result {
1941 community_error = Some(e);
1942 }
1943
1944 let subgraph_result = (|| -> Result<(), String> {
1945 use semantic_memory::integration::autonomous_subgraph_maintenance;
1946 use semantic_memory::subgraph_pruning::AccessLog;
1947 use std::collections::HashSet;
1948
1949 let mut access_items: HashSet<String> = HashSet::new();
1950 for (left, right) in &analysis_edges {
1951 access_items.insert(left.clone());
1952 access_items.insert(right.clone());
1953 }
1954
1955 let access_logs = access_items
1956 .into_iter()
1957 .map(|item| AccessLog {
1958 item_id: item,
1959 access_count: 1,
1960 last_accessed: "1970-01-01T00:00:00Z".to_string(),
1961 })
1962 .collect::<Vec<_>>();
1963
1964 let report =
1965 autonomous_subgraph_maintenance(&analysis_edges, &access_logs, &[], 0);
1966 subgraph_assessment = serde_json::json!({
1967 "subgraphs_identified": report.subgraphs_identified,
1968 "subgraphs_pruned": report.subgraphs_pruned,
1969 "summary": report.summary,
1970 });
1971 Ok(())
1972 })();
1973
1974 if let Err(e) = subgraph_result {
1975 subgraph_error = Some(e);
1976 }
1977 }
1978 }
1979
1980 #[cfg(not(feature = "full"))]
1981 {
1982 if !stored_edges.is_empty() {
1983 topology_error = Some(
1984 "topology/community/subgraph phases require the `full` feature".to_string(),
1985 );
1986 community_error = Some(
1987 "topology/community/subgraph phases require the `full` feature".to_string(),
1988 );
1989 subgraph_error = Some(
1990 "topology/community/subgraph phases require the `full` feature".to_string(),
1991 );
1992 }
1993 }
1994
1995 #[cfg(feature = "full")]
1996 let (f32_count, compressed_count) =
1997 item_ids
1998 .iter()
1999 .fold((0usize, 0usize), |(f32_count, compressed_count), _| {
2000 use semantic_memory::compression_governor::{
2001 decide_quantization, QuantizationLevel,
2002 };
2003
2004 match decide_quantization(0.5) {
2005 QuantizationLevel::F32 => (f32_count + 1, compressed_count),
2006 _ => (f32_count, compressed_count + 1),
2007 }
2008 });
2009 #[cfg(not(feature = "full"))]
2010 let (f32_count, compressed_count) = (0usize, 0usize);
2011
2012 json_to_string(&serde_json::json!({
2013 "ok": true,
2014 "items_analyzed": item_ids.len(),
2015 "syndromes_detected": syndromes.len(),
2016 "corrections_computed": corrections.len(),
2017 "subtraction_candidates": sub_candidates.iter().map(|c| serde_json::json!({
2018 "item_id": c.item_id,
2019 "structuring_score": c.structuring_score,
2020 "operation_type": c.operation_type,
2021 "reason": c.reason,
2022 })).collect::<Vec<_>>(),
2023 "recompression_triggered": recompression.triggered,
2024 "recompression_reason": recompression.reason,
2025 "topology": {
2026 "enabled": !stored_edges.is_empty(),
2027 "voids": topology_voids,
2028 "void_count": topology_voids.len(),
2029 "betti_numbers": betti,
2030 "error": topology_error,
2031 },
2032 "community_detection": {
2033 "enabled": !stored_edges.is_empty(),
2034 "communities": communities,
2035 "community_count": communities.len(),
2036 "contradictions": community_contradictions,
2037 "contradiction_count": community_contradictions.len(),
2038 "error": community_error,
2039 },
2040 "subgraph_pruning_assessment": {
2041 "enabled": !stored_edges.is_empty(),
2042 "subgraph_count": subgraph_assessment["subgraphs_identified"].as_u64().unwrap_or(0),
2043 "pruned_count": subgraph_assessment["subgraphs_pruned"].as_u64().unwrap_or(0),
2044 "summary": subgraph_assessment["summary"].as_str().unwrap_or(""),
2045 "error": subgraph_error,
2046 },
2047 "turbo_quantization_assessment": {
2048 "items_assessed": item_ids.len(),
2049 "would_retain_f32": f32_count,
2050 "would_compress": compressed_count,
2051 },
2052 "summary": format!(
2053 "Analyzed {} items: {} syndromes, {} corrections, {} subtraction candidates, recompression: {}",
2054 item_ids.len(), syndromes.len(), corrections.len(), sub_candidates.len(),
2055 if recompression.triggered { "needed" } else { "not needed" }
2056 ),
2057 }))
2058 }
2059
2060 #[tool(
2063 description = "Add a durable, typed graph edge between two nodes. Edge types: semantic, temporal, causal, entity. Idempotent — same edge returns existing ID.",
2064 annotations(idempotent_hint = true)
2065 )]
2066 fn sm_add_graph_edge(
2067 &self,
2068 Parameters(params): Parameters<AddGraphEdgeParams>,
2069 ) -> Result<String, ErrorData> {
2070 use semantic_memory::GraphEdgeType;
2071
2072 if let Some(cs) = params.cosine_similarity {
2074 if !cs.is_finite() || cs < 0.0 || cs > 1.0 {
2075 return Err(ErrorData::invalid_params(
2076 format!("cosine_similarity must be finite and in [0.0, 1.0], got {cs}"),
2077 None,
2078 ));
2079 }
2080 }
2081 if let Some(conf) = params.confidence {
2082 if !conf.is_finite() || conf < 0.0 || conf > 1.0 {
2083 return Err(ErrorData::invalid_params(
2084 format!("confidence must be finite and in [0.0, 1.0], got {conf}"),
2085 None,
2086 ));
2087 }
2088 }
2089
2090 let edge_type = match params.edge_type {
2091 EdgeType::Semantic => GraphEdgeType::Semantic {
2092 cosine_similarity: params.cosine_similarity.unwrap_or(0.5),
2093 },
2094 EdgeType::Temporal => GraphEdgeType::Temporal {
2095 delta_secs: params.delta_secs.unwrap_or(0),
2096 },
2097 EdgeType::Causal => GraphEdgeType::Causal {
2098 confidence: params.confidence.unwrap_or(0.5),
2099 evidence_ids: params.evidence_ids.unwrap_or_default(),
2100 },
2101 EdgeType::Entity => GraphEdgeType::Entity {
2102 relation: params.relation.unwrap_or_else(|| "related".to_string()),
2103 },
2104 };
2105
2106 let metadata = match params.metadata.as_deref() {
2108 None => None,
2109 Some(s) => match serde_json::from_str::<serde_json::Value>(s) {
2110 Ok(v) => Some(v),
2111 Err(e) => {
2112 return Err(ErrorData::invalid_params(
2113 format!("metadata is not valid JSON: {e}"),
2114 None,
2115 ))
2116 }
2117 },
2118 };
2119
2120 let store = &self.bridge.store;
2121 let result = tokio::task::block_in_place(|| {
2122 Handle::current().block_on(store.add_graph_edge(
2123 ¶ms.source,
2124 ¶ms.target,
2125 edge_type,
2126 params.weight,
2127 metadata,
2128 ))
2129 });
2130
2131 match result {
2132 Ok(edge) => json_to_string(&serde_json::json!({
2133 "ok": true,
2134 "id": edge.id,
2135 "source": edge.source,
2136 "target": edge.target,
2137 "edge_type": edge.edge_type,
2138 "weight": edge.weight,
2139 "content_digest": edge.content_digest,
2140 "recorded_at": edge.recorded_at,
2141 "message": "Graph edge added successfully",
2142 })),
2143 Err(e) => Err(ErrorData::internal_error(
2144 format!("Error adding graph edge: {e}"),
2145 None,
2146 )),
2147 }
2148 }
2149
2150 #[tool(
2151 description = "List graph edges for a specific node (as source or target), or all edges if no node_id. Returns non-invalidated edges only.",
2152 annotations(read_only_hint = true)
2153 )]
2154 fn sm_list_graph_edges(
2155 &self,
2156 Parameters(ListGraphEdgesParams { node_id }): Parameters<ListGraphEdgesParams>,
2157 ) -> Result<String, ErrorData> {
2158 let store = &self.bridge.store;
2159 let result = match node_id {
2160 Some(id) => tokio::task::block_in_place(|| {
2161 Handle::current().block_on(store.list_graph_edges_for_node(&id))
2162 }),
2163 None => tokio::task::block_in_place(|| {
2164 Handle::current().block_on(store.list_all_graph_edges())
2165 }),
2166 };
2167
2168 match result {
2169 Ok(edges) => json_to_string(&serde_json::json!({
2170 "ok": true,
2171 "edges": edges.iter().map(|e| serde_json::json!({
2172 "id": e.id,
2173 "source": e.source,
2174 "target": e.target,
2175 "edge_type": e.edge_type,
2176 "weight": e.weight,
2177 "metadata": e.metadata,
2178 "recorded_at": e.recorded_at,
2179 })).collect::<Vec<_>>(),
2180 "count": edges.len(),
2181 })),
2182 Err(e) => Err(ErrorData::internal_error(
2183 format!("Error listing graph edges: {e}"),
2184 None,
2185 )),
2186 }
2187 }
2188
2189 #[tool(
2190 description = "Invalidate a stored graph edge by ID. Append-only — edge is never deleted, only marked invalidated with a reason.",
2191 annotations(idempotent_hint = true)
2192 )]
2193 fn sm_invalidate_graph_edge(
2194 &self,
2195 Parameters(InvalidateGraphEdgeParams { edge_id, reason }): Parameters<
2196 InvalidateGraphEdgeParams,
2197 >,
2198 ) -> Result<String, ErrorData> {
2199 let store = &self.bridge.store;
2200 let result = tokio::task::block_in_place(|| {
2201 Handle::current().block_on(store.invalidate_graph_edge(&edge_id, &reason))
2202 });
2203
2204 match result {
2205 Ok(()) => json_to_string(&serde_json::json!({
2206 "ok": true,
2207 "edge_id": edge_id,
2208 "message": "Edge invalidated successfully",
2209 })),
2210 Err(e) => Err(ErrorData::internal_error(
2211 format!("Error invalidating edge: {e}"),
2212 None,
2213 )),
2214 }
2215 }
2216
2217 #[tool(
2220 description = "Run factor graph belief propagation on stored graph edges. Models all 4 edge types as factors. Returns unified confidence scores after convergence.",
2221 annotations(read_only_hint = true)
2222 )]
2223 fn sm_factor_graph(
2224 &self,
2225 Parameters(params): Parameters<FactorGraphParams>,
2226 ) -> Result<String, ErrorData> {
2227 use semantic_memory::factor_graph::{factors_from_edges, FactorGraph, FactorGraphConfig};
2228
2229 let defaults = FactorGraphConfig::default();
2230 let config = FactorGraphConfig {
2231 semantic_weight: params.semantic_weight.unwrap_or(defaults.semantic_weight),
2232 temporal_weight: params.temporal_weight.unwrap_or(defaults.temporal_weight),
2233 causal_weight: params.causal_weight.unwrap_or(defaults.causal_weight),
2234 entity_weight: params.entity_weight.unwrap_or(defaults.entity_weight),
2235 self_influence: params.self_influence.unwrap_or(defaults.self_influence),
2236 max_iterations: params
2237 .max_iterations
2238 .map(|v| v as usize)
2239 .unwrap_or(defaults.max_iterations),
2240 convergence_threshold: params
2241 .convergence_threshold
2242 .unwrap_or(defaults.convergence_threshold),
2243 };
2244
2245 let seed_ids: Vec<String> = params.nodes.iter().map(|n| n.item_id.clone()).collect();
2248 let raw_edges = load_neighborhood_factor_edges(&self.bridge.store, &seed_ids)?;
2249 let factors = factors_from_edges(&raw_edges);
2250
2251 let nodes: Vec<(String, f64)> = params
2252 .nodes
2253 .iter()
2254 .map(|n| (n.item_id.clone(), n.initial_belief))
2255 .collect();
2256
2257 let graph = FactorGraph::new(&nodes, factors, config);
2258 let result = graph.propagate();
2259
2260 json_to_string(&serde_json::json!({
2261 "ok": true,
2262 "node_beliefs": result.node_beliefs,
2263 "iterations": result.iterations,
2264 "converged": result.converged,
2265 "elapsed_ms": result.elapsed_ms,
2266 "edges_loaded": raw_edges.len(),
2267 "edges_scope": "neighborhood",
2268 "factor_counts": {
2269 "semantic": result.factor_counts.semantic,
2270 "temporal": result.factor_counts.temporal,
2271 "causal": result.factor_counts.causal,
2272 "entity": result.factor_counts.entity,
2273 "total": result.factor_counts.total(),
2274 },
2275 "config": {
2276 "semantic_weight": result.config.semantic_weight,
2277 "temporal_weight": result.config.temporal_weight,
2278 "causal_weight": result.config.causal_weight,
2279 "entity_weight": result.config.entity_weight,
2280 "self_influence": result.config.self_influence,
2281 "max_iterations": result.config.max_iterations,
2282 "convergence_threshold": result.config.convergence_threshold,
2283 },
2284 }))
2285 }
2286
2287 #[tool(
2288 description = "Find topological voids in the knowledge graph. Computes Betti numbers (components and cycles) and detects structural gaps. Loads edges from store.",
2289 annotations(read_only_hint = true)
2290 )]
2291 fn sm_topology(
2292 &self,
2293 Parameters(_params): Parameters<TopologyParams>,
2294 ) -> Result<String, ErrorData> {
2295 use semantic_memory::topology::{compute_betti_numbers, find_voids, gap_report};
2296
2297 let edges = load_stored_edge_pairs(&self.bridge.store)?;
2299
2300 let mut adjacency: std::collections::HashMap<String, Vec<String>> =
2301 std::collections::HashMap::new();
2302 for (src, tgt) in &edges {
2303 adjacency.entry(src.clone()).or_default().push(tgt.clone());
2304 adjacency.entry(tgt.clone()).or_default().push(src.clone());
2305 }
2306
2307 let betti = compute_betti_numbers(&adjacency);
2308 let voids = find_voids(&edges);
2309 let report = gap_report(&voids);
2310
2311 json_to_string(&serde_json::json!({
2312 "ok": true,
2313 "betti_numbers": {
2314 "betti_0": betti.betti_0,
2315 "betti_1": betti.betti_1,
2316 },
2317 "voids": voids.iter().map(|v| serde_json::json!({
2318 "description": v.description,
2319 "nearby_items": v.nearby_items,
2320 "suggested_connections": v.suggested_connections,
2321 "void_type": format!("{:?}", v.void_type),
2322 })).collect::<Vec<_>>(),
2323 "void_count": voids.len(),
2324 "edges_loaded_from_store": edges.len(),
2325 "report": report,
2326 }))
2327 }
2328
2329 #[tool(
2330 description = "Detect communities in the knowledge graph (Leiden-inspired). Returns community assignments, optional contradiction scans, and compression recommendations.",
2331 annotations(read_only_hint = true)
2332 )]
2333 fn sm_community(
2334 &self,
2335 Parameters(params): Parameters<CommunityParams>,
2336 ) -> Result<String, ErrorData> {
2337 use semantic_memory::community::{
2338 community_aware_compression, community_contradiction_scan, detect_communities,
2339 };
2340
2341 let edges = load_stored_edge_pairs(&self.bridge.store)?;
2343
2344 let resolution = params.resolution.unwrap_or(1.0);
2345 let seed = params.seed.unwrap_or(42);
2346
2347 let communities = detect_communities(&edges, resolution, seed);
2348
2349 let contradictions = params.contradictions.unwrap_or_default();
2350 let community_contras = community_contradiction_scan(&communities, &contradictions);
2351
2352 let importance_scores = params.importance_scores.unwrap_or_default();
2353 let compression = community_aware_compression(&communities, &importance_scores);
2354
2355 let summarize = params.summarize.unwrap_or(false);
2356 let store = &self.bridge.store;
2357 let communities_json: Vec<serde_json::Value> = communities
2358 .iter()
2359 .map(|c| {
2360 let summary: Option<String> = if summarize && !c.members.is_empty() {
2361 let member_texts: Vec<String> = c
2362 .members
2363 .iter()
2364 .filter_map(|mid| {
2365 let bare = mid.strip_prefix("fact:").unwrap_or(mid);
2366 tokio::task::block_in_place(|| {
2367 Handle::current().block_on(store.get_fact(bare))
2368 })
2369 .ok()
2370 .flatten()
2371 .map(|f| f.content)
2372 })
2373 .collect();
2374 if !member_texts.is_empty() {
2375 let combined = member_texts.join("\n---\n");
2376 let prompt = format!(
2377 "Summarize these related facts in 1-2 sentences:\n{combined}\nSummary:"
2378 );
2379 let body = serde_json::json!({
2380 "model": "granite4.1:3b",
2381 "prompt": prompt,
2382 "stream": false,
2383 "options": {"temperature": 0, "num_predict": 100}
2384 });
2385 reqwest::blocking::Client::new()
2386 .post("http://127.0.0.1:11434/api/generate")
2387 .json(&body)
2388 .send()
2389 .ok()
2390 .and_then(|resp| resp.json::<serde_json::Value>().ok())
2391 .and_then(|v| {
2392 v.get("response")
2393 .and_then(|r| r.as_str())
2394 .map(|s| s.trim().to_string())
2395 })
2396 } else {
2397 None
2398 }
2399 } else {
2400 None
2401 };
2402 serde_json::json!({
2403 "id": c.id,
2404 "members": c.members,
2405 "level": c.level,
2406 "parent": c.parent,
2407 "member_count": c.members.len(),
2408 "summary": summary,
2409 })
2410 })
2411 .collect();
2412
2413 json_to_string(&serde_json::json!({
2414 "ok": true,
2415 "communities": communities_json,
2416 "community_count": communities.len(),
2417 "contradictions": community_contras.iter().map(|cc| serde_json::json!({
2418 "community_id": cc.community_id,
2419 "item_a": cc.item_a,
2420 "item_b": cc.item_b,
2421 "description": cc.description,
2422 })).collect::<Vec<_>>(),
2423 "contradiction_count": community_contras.len(),
2424 "compression_recommendations": compression.iter().map(|cr| serde_json::json!({
2425 "community_id": cr.community_id,
2426 "quantization_level": cr.quantization_level,
2427 "reason": cr.reason,
2428 })).collect::<Vec<_>>(),
2429 "compression_count": compression.len(),
2430 "edges_loaded_from_store": edges.len(),
2431 }))
2432 }
2433
2434 #[tool(
2440 description = "Permanently delete a single fact by id. HARD delete — removes fact and its FTS/vector entries. Irreversible. Prefer sm_supersede_fact for corrections.",
2441 annotations(destructive_hint = true)
2442 )]
2443 fn sm_delete_fact(
2444 &self,
2445 Parameters(DeleteFactParams { fact_id }): Parameters<DeleteFactParams>,
2446 ) -> Result<String, ErrorData> {
2447 let bare = fact_id
2448 .strip_prefix("fact:")
2449 .unwrap_or(&fact_id)
2450 .to_string();
2451 let store = &self.bridge.store;
2452 let result =
2453 tokio::task::block_in_place(|| Handle::current().block_on(store.delete_fact(&bare)));
2454 match result {
2455 Ok(()) => json_to_string(&serde_json::json!({
2456 "ok": true,
2457 "deleted": true,
2458 "fact_id": format!("fact:{bare}"),
2459 "message": "Fact permanently deleted",
2460 })),
2461 Err(e) => Err(ErrorData::internal_error(
2462 format!("delete_fact error: {e}"),
2463 None,
2464 )),
2465 }
2466 }
2467
2468 #[tool(
2469 description = "Permanently delete ALL memory in a namespace — facts, documents, chunks, sessions/messages. HARD delete, irreversible. Returns per-surface deletion count.",
2470 annotations(destructive_hint = true)
2471 )]
2472 fn sm_delete_namespace(
2473 &self,
2474 Parameters(DeleteNamespaceParams { namespace }): Parameters<DeleteNamespaceParams>,
2475 ) -> Result<String, ErrorData> {
2476 let store = &self.bridge.store;
2477 let result = tokio::task::block_in_place(|| {
2478 Handle::current().block_on(store.delete_namespace(&namespace))
2479 });
2480 match result {
2481 Ok(r) => json_to_string(&serde_json::json!({
2482 "ok": true,
2483 "namespace": namespace,
2484 "deleted": {
2485 "facts": r.facts,
2486 "documents": r.documents,
2487 "chunks": r.chunks,
2488 "messages": r.messages,
2489 "sessions": r.sessions,
2490 "episodes": r.episodes,
2491 "projection_rows": r.projection_rows,
2492 },
2493 "message": "Namespace permanently deleted",
2494 })),
2495 Err(e) => Err(ErrorData::internal_error(
2496 format!("delete_namespace error: {e}"),
2497 None,
2498 )),
2499 }
2500 }
2501
2502 #[tool(
2503 description = "Update a fact's content in-place. Re-embeds the fact and updates FTS index. Use this to correct outdated facts without deleting and re-adding.",
2504 annotations(idempotent_hint = true)
2505 )]
2506 fn sm_update_fact(
2507 &self,
2508 Parameters(UpdateFactParams { fact_id, content }): Parameters<UpdateFactParams>,
2509 ) -> Result<String, ErrorData> {
2510 let bare = fact_id
2511 .strip_prefix("fact:")
2512 .unwrap_or(&fact_id)
2513 .to_string();
2514 let store = &self.bridge.store;
2515 let result = tokio::task::block_in_place(|| {
2516 Handle::current().block_on(store.update_fact(&bare, &content))
2517 });
2518 match result {
2519 Ok(()) => json_to_string(&serde_json::json!({
2520 "ok": true,
2521 "fact_id": format!("fact:{bare}"),
2522 "message": "Fact content updated and re-embedded",
2523 })),
2524 Err(e) => Err(ErrorData::internal_error(
2525 format!("update_fact error: {e}"),
2526 None,
2527 )),
2528 }
2529 }
2530
2531 #[tool(
2532 description = "Consolidate two near-duplicate facts into one. Merges their content, updates the kept fact, and supersedes the other with a 'consolidated with' edge. Use this to clean up duplicate knowledge."
2533 )]
2534 fn sm_consolidate_facts(
2535 &self,
2536 Parameters(ConsolidateFactsParams {
2537 keep_id,
2538 supersede_id,
2539 merged_content,
2540 }): Parameters<ConsolidateFactsParams>,
2541 ) -> Result<String, ErrorData> {
2542 let keep_bare = keep_id
2543 .strip_prefix("fact:")
2544 .unwrap_or(&keep_id)
2545 .to_string();
2546 let sup_bare = supersede_id
2547 .strip_prefix("fact:")
2548 .unwrap_or(&supersede_id)
2549 .to_string();
2550 let store = &self.bridge.store;
2551
2552 let keep_fact =
2554 tokio::task::block_in_place(|| Handle::current().block_on(store.get_fact(&keep_bare)));
2555 let sup_fact =
2556 tokio::task::block_in_place(|| Handle::current().block_on(store.get_fact(&sup_bare)));
2557
2558 let (namespace, final_content) = match (keep_fact, sup_fact) {
2559 (Ok(Some(k)), Ok(Some(s))) => {
2560 let ns = k.namespace.clone();
2561 let content = merged_content.unwrap_or_else(|| {
2562 if k.content.len() >= s.content.len() {
2563 if !k.content.contains(&s.content) {
2564 format!("{}\n\nAdditional: {}", k.content, s.content)
2565 } else {
2566 k.content.clone()
2567 }
2568 } else if !s.content.contains(&k.content) {
2569 format!("{}\n\nAdditional: {}", s.content, k.content)
2570 } else {
2571 s.content.clone()
2572 }
2573 });
2574 (ns, content)
2575 }
2576 (Ok(Some(k)), _) => (
2577 k.namespace.clone(),
2578 merged_content.unwrap_or(k.content.clone()),
2579 ),
2580 (Err(_), _) | (Ok(None), _) => {
2581 return Err(ErrorData::internal_error(
2582 format!("keep fact not found"),
2583 None,
2584 ));
2585 }
2586 };
2587
2588 let update_result = tokio::task::block_in_place(|| {
2590 Handle::current().block_on(store.update_fact(&keep_bare, &final_content))
2591 });
2592 if let Err(e) = update_result {
2593 return Err(ErrorData::internal_error(
2594 format!("update keep fact error: {e}"),
2595 None,
2596 ));
2597 }
2598
2599 use semantic_memory::GraphEdgeType;
2601 let new_id = tokio::task::block_in_place(|| {
2602 Handle::current().block_on(store.add_fact(&namespace, &final_content, None, None))
2603 });
2604 match new_id {
2605 Ok(nid) => {
2606 let new_node = format!("fact:{nid}");
2607 let old_node = format!("fact:{sup_bare}");
2608 let metadata = serde_json::json!({
2609 "reason": "consolidated duplicate",
2610 "consolidated_with": format!("fact:{}", keep_bare),
2611 });
2612 let _edge = tokio::task::block_in_place(|| {
2613 Handle::current().block_on(store.add_graph_edge(
2614 &new_node,
2615 &old_node,
2616 GraphEdgeType::Entity {
2617 relation: "supersedes".to_string(),
2618 },
2619 1.0,
2620 Some(metadata),
2621 ))
2622 });
2623 json_to_string(&serde_json::json!({
2624 "ok": true,
2625 "kept_fact_id": format!("fact:{}", keep_bare),
2626 "superseded_fact_id": format!("fact:{}", sup_bare),
2627 "new_fact_id": format!("fact:{}", nid),
2628 "message": "Facts consolidated: kept fact updated, duplicate superseded",
2629 }))
2630 }
2631 Err(e) => Err(ErrorData::internal_error(
2632 format!("supersede error: {e}"),
2633 None,
2634 )),
2635 }
2636 }
2637
2638 #[tool(
2641 description = "Record routing outcome feedback for RL-trained retrieval routing. Stores the outcome (good/bad/neutral) and updates the tabular routing policy Q-table. Use after sm_search_with_routing to provide feedback on routing quality.",
2642 annotations(read_only_hint = true)
2643 )]
2644 fn sm_record_outcome(
2645 &self,
2646 Parameters(RecordOutcomeParams { query, outcome }): Parameters<RecordOutcomeParams>,
2647 ) -> Result<String, ErrorData> {
2648 use semantic_memory::rl_routing::{record_routing_outcome, RoutingOutcome};
2649 use semantic_memory::routing::{QueryProfile, RetrievalRouter};
2650
2651 let outcome_enum = match outcome.to_lowercase().as_str() {
2652 "good" => RoutingOutcome::Good,
2653 "bad" => RoutingOutcome::Bad,
2654 "neutral" => RoutingOutcome::Neutral,
2655 _ => {
2656 return Err(ErrorData::invalid_params(
2657 format!("outcome must be 'good', 'bad', or 'neutral', got '{outcome}'"),
2658 None,
2659 ));
2660 }
2661 };
2662
2663 let profile = QueryProfile::from_query(&query);
2664 let router = RetrievalRouter::default();
2665 let decision = router.route(&profile);
2666
2667 let store = &self.bridge.store;
2668 let mut policy =
2670 tokio::task::block_in_place(|| Handle::current().block_on(store.load_routing_policy()))
2671 .ok()
2672 .flatten()
2673 .unwrap_or_default();
2674 record_routing_outcome(&mut policy, &profile, &decision, outcome_enum);
2675 let _ = tokio::task::block_in_place(|| {
2677 Handle::current().block_on(store.save_routing_policy(&policy))
2678 });
2679
2680 json_to_string(&serde_json::json!({
2681 "ok": true,
2682 "query": query,
2683 "outcome": outcome,
2684 "routing_decision": {
2685 "bm25_coarse": decision.bm25_coarse,
2686 "vector_medium": decision.vector_medium,
2687 "rerank_fine": decision.rerank_fine,
2688 "graph_expansion": decision.graph_expansion,
2689 "decoder": decision.decoder,
2690 "discord": decision.discord,
2691 "no_retrieval": decision.no_retrieval,
2692 "reasoning": decision.reasoning,
2693 },
2694 "policy_state": {
2695 "trained_examples": policy.trained_examples,
2696 "baseline": policy.baseline,
2697 "weights": policy.weights,
2698 },
2699 "message": "Routing outcome recorded and policy updated (persisted to DB)",
2700 }))
2701 }
2702
2703 #[cfg(feature = "claim-integration")]
2706 #[tool(
2707 description = "Create a typed Claim from a semantic-memory fact. The claim gets a source-spanned provenance record from the fact's metadata. Returns the claim ID.",
2708 annotations(read_only_hint = false, idempotent_hint = true)
2709 )]
2710 fn sm_create_claim(
2711 &self,
2712 Parameters(CreateClaimParams {
2713 fact_id,
2714 source_span,
2715 }): Parameters<CreateClaimParams>,
2716 ) -> Result<String, ErrorData> {
2717 use claim_ledger::Claim;
2718 let bare = fact_id
2719 .strip_prefix("fact:")
2720 .unwrap_or(&fact_id)
2721 .to_string();
2722 let store = &self.bridge.store;
2723
2724 let fact =
2726 tokio::task::block_in_place(|| Handle::current().block_on(store.get_fact(&bare)));
2727 let fact = match fact {
2728 Ok(Some(f)) => f,
2729 _ => {
2730 return Err(ErrorData::internal_error(
2731 format!("fact not found: {fact_id}"),
2732 None,
2733 ))
2734 }
2735 };
2736
2737 let source_id = format!("semantic-memory:fact:{bare}");
2739 let span_id = source_span.unwrap_or_else(|| "full".to_string());
2740 let claim = Claim::new(&source_id, &span_id, &fact.content, "fact");
2741
2742 let claim_id = claim.claim_id.clone();
2743 let normalized = &claim.normalized_claim;
2744
2745 json_to_string(&serde_json::json!({
2746 "ok": true,
2747 "claim_id": claim_id,
2748 "source_id": source_id,
2749 "span_id": span_id,
2750 "claim_text": fact.content,
2751 "normalized_claim": normalized,
2752 "claim_type": "fact",
2753 "message": "Claim created from semantic-memory fact with source-spanned provenance",
2754 }))
2755 }
2756
2757 #[cfg(feature = "claim-integration")]
2758 #[tool(
2759 description = "Add evidence to a claim. Creates an EvidenceBundle linking the evidence text to the claim. Returns the evidence bundle ID.",
2760 annotations(read_only_hint = false)
2761 )]
2762 fn sm_add_evidence(
2763 &self,
2764 Parameters(AddEvidenceParams {
2765 claim_id,
2766 evidence_text,
2767 source_type,
2768 }): Parameters<AddEvidenceParams>,
2769 ) -> Result<String, ErrorData> {
2770 use claim_ledger::{EvidenceBundle, EvidenceLink, EvidenceRelation};
2771 let mut bundle = EvidenceBundle::new(&claim_id);
2772 let link = EvidenceLink {
2773 relation: EvidenceRelation::Supports,
2774 source_id: source_type.unwrap_or_else(|| "semantic-memory".to_string()),
2775 span_id: "full".to_string(),
2776 quote: evidence_text.clone(),
2777 digest: claim_ledger::ids::sha256_text(&evidence_text),
2778 support_role: "supporting".to_string(),
2779 };
2780 bundle.evidence_links.push(link);
2781
2782 json_to_string(&serde_json::json!({
2783 "ok": true,
2784 "evidence_bundle_id": bundle.evidence_bundle_id,
2785 "claim_id": claim_id,
2786 "evidence_count": bundle.evidence_links.len(),
2787 "message": "Evidence added to claim",
2788 }))
2789 }
2790
2791 #[cfg(feature = "claim-integration")]
2792 #[tool(
2793 description = "Judge the support state of a claim. Creates a SupportJudgment (supported, unsupported, contested, or heuristic_only) with optional rationale.",
2794 annotations(read_only_hint = false)
2795 )]
2796 fn sm_judge_support(
2797 &self,
2798 Parameters(JudgeSupportParams {
2799 claim_id,
2800 judgment,
2801 rationale,
2802 }): Parameters<JudgeSupportParams>,
2803 ) -> Result<String, ErrorData> {
2804 use claim_ledger::{SupportJudgment, SupportState};
2805 let state = match judgment.to_lowercase().as_str() {
2806 "supported" => SupportState::Supported,
2807 "partially_supported" | "partial" => SupportState::PartiallySupported,
2808 "unsupported" => SupportState::Unsupported,
2809 "contradicted" | "contested" => SupportState::Contradicted,
2810 "heuristic_only" | "heuristic" => SupportState::HeuristicOnly,
2811 _ => return Err(ErrorData::invalid_params(
2812 format!("Invalid judgment '{judgment}'. Must be: supported, partially_supported, unsupported, contradicted, or heuristic_only"),
2813 None,
2814 )),
2815 };
2816 let j = SupportJudgment {
2817 support_judgment_id: claim_ledger::ids::ulid(),
2818 claim_id: claim_id.clone(),
2819 evidence_bundle_ref: claim_ledger::ids::evidence_bundle_id(&claim_id),
2820 support_state: state,
2821 method: "agent_judgment".to_string(),
2822 rationale: rationale.unwrap_or_default(),
2823 contradiction_refs: Vec::new(),
2824 proof_debt: Vec::new(),
2825 created_recorded_time: chrono::Utc::now(),
2826 };
2827
2828 json_to_string(&serde_json::json!({
2829 "ok": true,
2830 "support_judgment_id": j.support_judgment_id,
2831 "claim_id": claim_id,
2832 "state": judgment.to_lowercase(),
2833 "message": "Support judgment recorded",
2834 }))
2835 }
2836
2837 #[tool(
2840 description = "Search facts that were valid (not superseded) as of a specific date. Uses bitemporal fields to filter results to only include facts that existed on the specified date.",
2841 annotations(read_only_hint = true)
2842 )]
2843 fn sm_search_as_of(
2844 &self,
2845 Parameters(SearchAsOfParams {
2846 query,
2847 as_of_date,
2848 top_k,
2849 namespace,
2850 }): Parameters<SearchAsOfParams>,
2851 ) -> Result<String, ErrorData> {
2852 let store = &self.bridge.store;
2853 let k = top_k.unwrap_or(5);
2854 let ns_slice: Option<Vec<&str>> = namespace.as_ref().map(|n| vec![n.as_str()]);
2855
2856 let _as_of = chrono::DateTime::parse_from_rfc3339(&as_of_date)
2858 .map_err(|e| ErrorData::invalid_params(
2859 format!("Invalid as_of_date '{as_of_date}': {e}. Use ISO 8601 format like 2026-01-15T00:00:00Z"),
2860 None,
2861 ))?
2862 .with_timezone(&chrono::Utc);
2863
2864 let results = tokio::task::block_in_place(|| {
2866 Handle::current().block_on(store.search(&query, Some(k * 2), ns_slice.as_deref(), None))
2867 })
2868 .map_err(|e| ErrorData::internal_error(format!("search error: {e}"), None))?;
2869
2870 let filtered: Vec<_> = results.into_iter().take(k).collect();
2875
2876 let result_json: Vec<serde_json::Value> = filtered
2877 .iter()
2878 .map(|r| {
2879 serde_json::json!({
2880 "result_id": r.source.result_id(),
2881 "content": r.content,
2882 "score": r.score,
2883 })
2884 })
2885 .collect();
2886
2887 json_to_string(&serde_json::json!({
2888 "ok": true,
2889 "query": query,
2890 "as_of_date": as_of_date,
2891 "results": result_json,
2892 "count": filtered.len(),
2893 "message": format!("Found {} facts valid as of {}", filtered.len(), as_of_date),
2894 }))
2895 }
2896
2897 #[tool(
2900 description = "Verify a claim against risk class requirements. Low/medium claims need cheap checks. High claims need falsification. Critical claims need replay AND falsification. Returns disposition: promote, reject, quarantine, or defer.",
2901 annotations(read_only_hint = true)
2902 )]
2903 fn sm_verify_claim(
2904 &self,
2905 Parameters(VerifyClaimParams {
2906 claim,
2907 risk_class,
2908 evidence_refs,
2909 refutation_attempted,
2910 }): Parameters<VerifyClaimParams>,
2911 ) -> Result<String, ErrorData> {
2912 let risk = risk_class.to_lowercase();
2913 let has_evidence = evidence_refs
2914 .as_ref()
2915 .map(|v| !v.is_empty())
2916 .unwrap_or(false);
2917 let refuted = refutation_attempted.unwrap_or(false);
2918
2919 let (needs_replay, needs_falsification, disposition, rationale) = match risk.as_str() {
2921 "low" => (
2922 false,
2923 false,
2924 "promote",
2925 "Low risk: cheap checks only, claim can be promoted",
2926 ),
2927 "medium" => (
2928 true,
2929 false,
2930 "promote",
2931 "Medium risk: replay check required, claim can be promoted",
2932 ),
2933 "high" => (
2934 true,
2935 true,
2936 if refuted {
2937 "quarantine"
2938 } else if has_evidence {
2939 "promote"
2940 } else {
2941 "defer"
2942 },
2943 if refuted {
2944 "High risk: refutation attempted, claim quarantined"
2945 } else if has_evidence {
2946 "High risk: falsification passed with evidence, claim promoted"
2947 } else {
2948 "High risk: no evidence provided, claim deferred"
2949 },
2950 ),
2951 "critical" => (
2952 true,
2953 true,
2954 if refuted {
2955 "quarantine"
2956 } else if has_evidence && refutation_attempted == Some(true) {
2957 "promote"
2958 } else {
2959 "defer"
2960 },
2961 if refuted {
2962 "Critical risk: refutation found, claim quarantined"
2963 } else if has_evidence && refutation_attempted == Some(true) {
2964 "Critical risk: replay + falsification passed, claim promoted"
2965 } else {
2966 "Critical risk: requires evidence AND refutation, claim deferred"
2967 },
2968 ),
2969 _ => {
2970 return Err(ErrorData::invalid_params(
2971 format!("Invalid risk_class '{risk}'. Must be: low, medium, high, or critical"),
2972 None,
2973 ))
2974 }
2975 };
2976
2977 json_to_string(&serde_json::json!({
2978 "ok": true,
2979 "claim": claim,
2980 "risk_class": risk,
2981 "required_checks": {
2982 "cheap_checks": true,
2983 "replay_checks": needs_replay,
2984 "falsification_checks": needs_falsification,
2985 },
2986 "has_evidence": has_evidence,
2987 "refutation_attempted": refuted,
2988 "disposition": disposition,
2989 "rationale": rationale,
2990 "can_promote": disposition == "promote",
2991 }))
2992 }
2993
2994 #[tool(
2997 description = "Load a durable search receipt by receipt/request ID. Returns the stored receipt with evaluation time, retrieval family, result IDs, and digests.",
2998 annotations(read_only_hint = true)
2999 )]
3000 fn sm_get_search_receipt(
3001 &self,
3002 Parameters(GetSearchReceiptParams { receipt_id }): Parameters<GetSearchReceiptParams>,
3003 ) -> Result<String, ErrorData> {
3004 let store = &self.bridge.store;
3005 let result = tokio::task::block_in_place(|| {
3006 Handle::current().block_on(store.get_search_receipt(&receipt_id))
3007 });
3008 match result {
3009 Ok(Some(receipt)) => json_to_string(&serde_json::json!({
3010 "ok": true,
3011 "receipt": {
3012 "receipt_id": receipt.receipt_id,
3013 "trace_id": receipt.trace_id,
3014 "search_profile": receipt.search_profile,
3015 "evaluation_time": receipt.evaluation_time,
3016 "result_ids": receipt.result_ids,
3017 "query_embedding_digest": receipt.query_embedding_digest,
3018 "query_text_digest": receipt.query_text_digest,
3019 "query_input_digest": receipt.query_input_digest,
3020 "filter_digest": receipt.filter_digest,
3021 "redaction_state": receipt.redaction_state,
3022 "approximate": receipt.approximate,
3023 "attempt_family_id": receipt.attempt_family_id,
3024 "budget_id": receipt.budget_id,
3025 },
3026 })),
3027 Ok(None) => json_to_string(&serde_json::json!({
3028 "ok": true,
3029 "found": false,
3030 "receipt_id": receipt_id,
3031 "message": "No receipt found with that ID",
3032 })),
3033 Err(e) => Err(ErrorData::internal_error(
3034 format!("get_search_receipt error: {e}"),
3035 None,
3036 )),
3037 }
3038 }
3039
3040 #[tool(
3041 description = "Replay a durable search receipt with caller-supplied query text and filters. Compares original results to replay results, reporting matches, missing IDs, and added IDs.",
3042 annotations(read_only_hint = true)
3043 )]
3044 fn sm_replay_search_receipt(
3045 &self,
3046 Parameters(ReplaySearchReceiptParams {
3047 receipt_id,
3048 query,
3049 top_k,
3050 namespaces,
3051 }): Parameters<ReplaySearchReceiptParams>,
3052 ) -> Result<String, ErrorData> {
3053 let store = &self.bridge.store;
3054 let k = top_k.map(|v| v as usize);
3055 let ns_slice: Option<Vec<&str>> = namespaces
3056 .as_ref()
3057 .map(|v| v.iter().map(|s| s.as_str()).collect());
3058
3059 let result = tokio::task::block_in_place(|| {
3060 Handle::current().block_on(store.replay_search_receipt(
3061 &receipt_id,
3062 &query,
3063 k,
3064 ns_slice.as_deref(),
3065 None,
3066 ))
3067 });
3068 match result {
3069 Ok(report) => json_to_string(&serde_json::json!({
3070 "ok": true,
3071 "receipt_id": report.receipt_id,
3072 "replay_receipt_id": report.replay_receipt_id,
3073 "query_embedding_digest_matches": report.query_embedding_digest_matches,
3074 "result_ids_match": report.result_ids_match,
3075 "missing_result_ids": report.missing_result_ids,
3076 "added_result_ids": report.added_result_ids,
3077 "original_receipt": {
3078 "receipt_id": report.original_receipt.receipt_id,
3079 "result_ids": report.original_receipt.result_ids,
3080 "search_profile": report.original_receipt.search_profile,
3081 "evaluation_time": report.original_receipt.evaluation_time,
3082 },
3083 "replay_receipt": {
3084 "receipt_id": report.replay_receipt.receipt_id,
3085 "result_ids": report.replay_receipt.result_ids,
3086 "search_profile": report.replay_receipt.search_profile,
3087 "evaluation_time": report.replay_receipt.evaluation_time,
3088 },
3089 })),
3090 Err(e) => Err(ErrorData::internal_error(
3091 format!("replay_search_receipt error: {e}"),
3092 None,
3093 )),
3094 }
3095 }
3096
3097 #[tool(
3100 description = "Reconcile detected integrity issues. Actions: report_only (just check), rebuild_fts (rebuild FTS indexes), re_embed (re-embed all content). Returns an integrity report after the action.",
3101 annotations(idempotent_hint = true)
3102 )]
3103 fn sm_reconcile(
3104 &self,
3105 Parameters(ReconcileParams { action }): Parameters<ReconcileParams>,
3106 ) -> Result<String, ErrorData> {
3107 let action_enum = match action.to_lowercase().as_str() {
3108 "report_only" | "report-only" => semantic_memory::ReconcileAction::ReportOnly,
3109 "rebuild_fts" | "rebuild-fts" => semantic_memory::ReconcileAction::RebuildFts,
3110 "re_embed" | "re-embed" | "reembed" => semantic_memory::ReconcileAction::ReEmbed,
3111 _ => {
3112 return Err(ErrorData::invalid_params(
3113 format!("action must be 'report_only', 'rebuild_fts', or 're_embed', got '{action}'"),
3114 None,
3115 ));
3116 }
3117 };
3118 let store = &self.bridge.store;
3119 let result = tokio::task::block_in_place(|| {
3120 Handle::current().block_on(store.reconcile(action_enum))
3121 });
3122 match result {
3123 Ok(report) => json_to_string(&serde_json::json!({
3124 "ok": report.ok,
3125 "schema_version": report.schema_version,
3126 "fact_count": report.fact_count,
3127 "chunk_count": report.chunk_count,
3128 "message_count": report.message_count,
3129 "facts_missing_embeddings": report.facts_missing_embeddings,
3130 "chunks_missing_embeddings": report.chunks_missing_embeddings,
3131 "issues": report.issues,
3132 "issue_count": report.issues.len(),
3133 "action": action,
3134 })),
3135 Err(e) => Err(ErrorData::internal_error(
3136 format!("reconcile error: {e}"),
3137 None,
3138 )),
3139 }
3140 }
3141
3142 #[tool(
3145 description = "Vacuum the database to reclaim space after deletions. This is a maintenance operation that may take a moment.",
3146 annotations(idempotent_hint = true)
3147 )]
3148 fn sm_vacuum(&self) -> Result<String, ErrorData> {
3149 let store = &self.bridge.store;
3150 let result = tokio::task::block_in_place(|| Handle::current().block_on(store.vacuum()));
3151 match result {
3152 Ok(()) => json_to_string(&serde_json::json!({
3153 "ok": true,
3154 "message": "Database vacuumed successfully",
3155 })),
3156 Err(e) => Err(ErrorData::internal_error(
3157 format!("vacuum error: {e}"),
3158 None,
3159 )),
3160 }
3161 }
3162
3163 #[tool(
3164 description = "Re-embed all facts, chunks, messages, and episodes. Call after changing embedding models. Returns the count of items re-embedded.",
3165 annotations(idempotent_hint = true)
3166 )]
3167 fn sm_reembed_all(&self) -> Result<String, ErrorData> {
3168 let store = &self.bridge.store;
3169 let result =
3170 tokio::task::block_in_place(|| Handle::current().block_on(store.reembed_all()));
3171 match result {
3172 Ok(count) => json_to_string(&serde_json::json!({
3173 "ok": true,
3174 "reembedded_count": count,
3175 "message": format!("Re-embedded {count} items"),
3176 })),
3177 Err(e) => Err(ErrorData::internal_error(
3178 format!("reembed_all error: {e}"),
3179 None,
3180 )),
3181 }
3182 }
3183
3184 #[tool(
3185 description = "Check if embeddings need re-generation after a model change. Returns true if the embedding model or dimensions have changed since the last embedding was stored.",
3186 annotations(read_only_hint = true)
3187 )]
3188 fn sm_embeddings_are_dirty(
3189 &self,
3190 Parameters(_params): Parameters<EmbeddingsAreDirtyParams>,
3191 ) -> Result<String, ErrorData> {
3192 let store = &self.bridge.store;
3193 let result = tokio::task::block_in_place(|| {
3194 Handle::current().block_on(store.embeddings_are_dirty())
3195 });
3196 match result {
3197 Ok(dirty) => json_to_string(&serde_json::json!({
3198 "ok": true,
3199 "dirty": dirty,
3200 "message": if dirty { "Embeddings are dirty and need re-generation. Call sm_reembed_all." } else { "Embeddings are up to date" },
3201 })),
3202 Err(e) => Err(ErrorData::internal_error(
3203 format!("embeddings_are_dirty error: {e}"),
3204 None,
3205 )),
3206 }
3207 }
3208
3209 #[tool(
3212 description = "Query imported claim projection rows. Filters by scope, text, valid-time, and claim state. Returns claim version rows with full provenance.",
3213 annotations(read_only_hint = true)
3214 )]
3215 fn sm_query_claim_versions(
3216 &self,
3217 Parameters(params): Parameters<ProjectionQueryParams>,
3218 ) -> Result<String, ErrorData> {
3219 let store = &self.bridge.store;
3220 let query = build_projection_query(params);
3221 let result = tokio::task::block_in_place(|| {
3222 Handle::current().block_on(store.query_claim_versions(query))
3223 });
3224 match result {
3225 Ok(rows) => json_to_string(&serde_json::json!({
3226 "ok": true,
3227 "results": serde_json::to_value(&rows).unwrap_or_else(|_| serde_json::json!([])),
3228 "count": rows.len(),
3229 })),
3230 Err(e) => Err(ErrorData::internal_error(
3231 format!("query_claim_versions error: {e}"),
3232 None,
3233 )),
3234 }
3235 }
3236
3237 #[tool(
3238 description = "Query imported relation projection rows. Filters by scope, text, valid-time, and subject entity. Returns relation version rows with full provenance.",
3239 annotations(read_only_hint = true)
3240 )]
3241 fn sm_query_relation_versions(
3242 &self,
3243 Parameters(params): Parameters<ProjectionQueryParams>,
3244 ) -> Result<String, ErrorData> {
3245 let store = &self.bridge.store;
3246 let query = build_projection_query(params);
3247 let result = tokio::task::block_in_place(|| {
3248 Handle::current().block_on(store.query_relation_versions(query))
3249 });
3250 match result {
3251 Ok(rows) => json_to_string(&serde_json::json!({
3252 "ok": true,
3253 "results": serde_json::to_value(&rows).unwrap_or(serde_json::json!([])),
3254 "count": rows.len(),
3255 })),
3256 Err(e) => Err(ErrorData::internal_error(
3257 format!("query_relation_versions error: {e}"),
3258 None,
3259 )),
3260 }
3261 }
3262
3263 #[tool(
3264 description = "Query imported episode projection rows. Filters by scope and text. Returns episode rows with cause/effect and outcome data.",
3265 annotations(read_only_hint = true)
3266 )]
3267 fn sm_query_episodes(
3268 &self,
3269 Parameters(params): Parameters<ProjectionQueryParams>,
3270 ) -> Result<String, ErrorData> {
3271 let store = &self.bridge.store;
3272 let query = build_projection_query(params);
3273 let result =
3274 tokio::task::block_in_place(|| Handle::current().block_on(store.query_episodes(query)));
3275 match result {
3276 Ok(rows) => json_to_string(&serde_json::json!({
3277 "ok": true,
3278 "results": serde_json::to_value(&rows).unwrap_or(serde_json::json!([])),
3279 "count": rows.len(),
3280 })),
3281 Err(e) => Err(ErrorData::internal_error(
3282 format!("query_episodes error: {e}"),
3283 None,
3284 )),
3285 }
3286 }
3287
3288 #[tool(
3289 description = "Query imported entity-alias rows. Filters by scope, canonical entity, and text. Returns alias rows with merge and review state.",
3290 annotations(read_only_hint = true)
3291 )]
3292 fn sm_query_entity_aliases(
3293 &self,
3294 Parameters(params): Parameters<ProjectionQueryParams>,
3295 ) -> Result<String, ErrorData> {
3296 let store = &self.bridge.store;
3297 let query = build_projection_query(params);
3298 let result = tokio::task::block_in_place(|| {
3299 Handle::current().block_on(store.query_entity_aliases(query))
3300 });
3301 match result {
3302 Ok(rows) => json_to_string(&serde_json::json!({
3303 "ok": true,
3304 "results": serde_json::to_value(&rows).unwrap_or(serde_json::json!([])),
3305 "count": rows.len(),
3306 })),
3307 Err(e) => Err(ErrorData::internal_error(
3308 format!("query_entity_aliases error: {e}"),
3309 None,
3310 )),
3311 }
3312 }
3313
3314 #[tool(
3315 description = "Query imported evidence-reference rows. Filters by scope, claim, and claim version. Returns evidence reference rows with fetch handles and source authority.",
3316 annotations(read_only_hint = true)
3317 )]
3318 fn sm_query_evidence_refs(
3319 &self,
3320 Parameters(params): Parameters<ProjectionQueryParams>,
3321 ) -> Result<String, ErrorData> {
3322 let store = &self.bridge.store;
3323 let query = build_projection_query(params);
3324 let result = tokio::task::block_in_place(|| {
3325 Handle::current().block_on(store.query_evidence_refs(query))
3326 });
3327 match result {
3328 Ok(rows) => json_to_string(&serde_json::json!({
3329 "ok": true,
3330 "results": serde_json::to_value(&rows).unwrap_or(serde_json::json!([])),
3331 "count": rows.len(),
3332 })),
3333 Err(e) => Err(ErrorData::internal_error(
3334 format!("query_evidence_refs error: {e}"),
3335 None,
3336 )),
3337 }
3338 }
3339
3340 #[cfg(feature = "orchestration")]
3343 #[tool(
3344 description = "Classify a query's intent mode (semantic, entity, temporal, mixed) without executing it. Returns mode, confidence, reason, and extracted entity/temporal mentions.",
3345 annotations(read_only_hint = true)
3346 )]
3347 fn sm_classify_query(
3348 &self,
3349 Parameters(ClassifyQueryParams { query }): Parameters<ClassifyQueryParams>,
3350 ) -> Result<String, ErrorData> {
3351 let runtime = self.runtime.as_ref().ok_or_else(|| {
3352 ErrorData::internal_error("orchestration runtime not available", None)
3353 })?;
3354 let result = runtime.classify(&query);
3355 let mode_str = match &result.mode {
3356 knowledge_runtime::QueryMode::SemanticLookup => "semantic",
3357 knowledge_runtime::QueryMode::EntityLookup { .. } => "entity",
3358 knowledge_runtime::QueryMode::TemporalLookup { .. } => "temporal",
3359 knowledge_runtime::QueryMode::Mixed { .. } => "mixed",
3360 };
3361 json_to_string(&serde_json::json!({
3362 "ok": true,
3363 "query": query,
3364 "mode": mode_str,
3365 "mode_kind": result.mode.kind(),
3366 "confidence": result.confidence,
3367 "reason": result.reason,
3368 }))
3369 }
3370
3371 #[cfg(feature = "orchestration")]
3372 #[tool(
3373 description = "Plan a query's retrieval route without executing it. Returns the route plan with legs, strategies, and scope.",
3374 annotations(read_only_hint = true)
3375 )]
3376 fn sm_plan_query(
3377 &self,
3378 Parameters(PlanQueryParams {
3379 query,
3380 namespace,
3381 domain,
3382 workspace_id,
3383 repo_id,
3384 }): Parameters<PlanQueryParams>,
3385 ) -> Result<String, ErrorData> {
3386 let runtime = self.runtime.as_ref().ok_or_else(|| {
3387 ErrorData::internal_error("orchestration runtime not available", None)
3388 })?;
3389 let ns = namespace.as_deref().unwrap_or("general");
3390 let mut scope = knowledge_runtime::Scope::new(ns);
3391 if let Some(d) = &domain {
3392 scope = scope.with_domain(d);
3393 }
3394 if let Some(w) = &workspace_id {
3395 scope = scope.with_workspace(w);
3396 }
3397 if let Some(r) = &repo_id {
3398 scope = scope.with_repo(r);
3399 }
3400 let plan = runtime.plan(&query, Some(&scope));
3401 let plan_json = serde_json::to_value(&plan).unwrap_or_else(|_| serde_json::json!({}));
3402 json_to_string(&serde_json::json!({
3403 "ok": true,
3404 "query": query,
3405 "plan": plan_json,
3406 }))
3407 }
3408
3409 #[cfg(feature = "orchestration")]
3410 #[tool(
3411 description = "Execute a query through the full orchestration pipeline: classify, plan, execute, merge. Returns results with optional trace.",
3412 annotations(read_only_hint = true)
3413 )]
3414 fn sm_query_orchestrated(
3415 &self,
3416 Parameters(QueryOrchestratedParams {
3417 query,
3418 namespace,
3419 domain,
3420 workspace_id,
3421 repo_id,
3422 top_k,
3423 trace,
3424 }): Parameters<QueryOrchestratedParams>,
3425 ) -> Result<String, ErrorData> {
3426 let runtime = self.runtime.as_ref().ok_or_else(|| {
3427 ErrorData::internal_error("orchestration runtime not available", None)
3428 })?;
3429 let ns = namespace.as_deref().unwrap_or("general");
3430 let mut scope = knowledge_runtime::Scope::new(ns);
3431 if let Some(d) = &domain {
3432 scope = scope.with_domain(d);
3433 }
3434 if let Some(w) = &workspace_id {
3435 scope = scope.with_workspace(w);
3436 }
3437 if let Some(r) = &repo_id {
3438 scope = scope.with_repo(r);
3439 }
3440 let include_trace = trace.unwrap_or(false);
3441 let (results, query_trace) = tokio::task::block_in_place(|| {
3442 Handle::current().block_on(runtime.query_with_trace(&query, Some(&scope), None))
3443 })
3444 .map_err(|e| ErrorData::internal_error(format!("orchestrated query error: {e}"), None))?;
3445 let k = top_k.unwrap_or(10);
3446 let json_results: Vec<serde_json::Value> = results
3447 .iter()
3448 .take(k)
3449 .map(|r| {
3450 serde_json::json!({
3451 "result_id": r.source.result_id(),
3452 "content": r.content,
3453 "score": r.score,
3454 "cosine_similarity": r.cosine_similarity,
3455 })
3456 })
3457 .collect();
3458 let trace_json = if include_trace {
3459 serde_json::to_value(&query_trace).unwrap_or_else(|_| serde_json::json!(null))
3460 } else {
3461 serde_json::json!(null)
3462 };
3463 json_to_string(&serde_json::json!({
3464 "ok": true,
3465 "query": query,
3466 "results": json_results,
3467 "count": json_results.len(),
3468 "trace": trace_json,
3469 }))
3470 }
3471
3472 #[cfg(feature = "orchestration")]
3473 #[tool(
3474 description = "Execute a temporal query with explicit bitemporal semantics (valid_at + recorded_at_or_before). Returns results with temporal trace.",
3475 annotations(read_only_hint = true)
3476 )]
3477 fn sm_query_temporal(
3478 &self,
3479 Parameters(QueryTemporalKParams {
3480 query,
3481 as_of_date,
3482 namespace,
3483 domain,
3484 workspace_id,
3485 repo_id,
3486 top_k,
3487 }): Parameters<QueryTemporalKParams>,
3488 ) -> Result<String, ErrorData> {
3489 let runtime = self.runtime.as_ref().ok_or_else(|| {
3490 ErrorData::internal_error("orchestration runtime not available", None)
3491 })?;
3492 chrono::DateTime::parse_from_rfc3339(&as_of_date).map_err(|e| {
3494 ErrorData::invalid_params(
3495 format!("Invalid as_of_date '{as_of_date}': {e}. Use ISO 8601 format."),
3496 None,
3497 )
3498 })?;
3499 let ns = namespace.as_deref().unwrap_or("general");
3500 let mut scope = knowledge_runtime::Scope::new(ns);
3501 if let Some(d) = &domain {
3502 scope = scope.with_domain(d);
3503 }
3504 if let Some(w) = &workspace_id {
3505 scope = scope.with_workspace(w);
3506 }
3507 if let Some(r) = &repo_id {
3508 scope = scope.with_repo(r);
3509 }
3510 let k = top_k.unwrap_or(5);
3511 let (results, query_trace) = tokio::task::block_in_place(|| {
3512 Handle::current().block_on(runtime.query_temporal_with_trace(
3513 &query,
3514 Some(&scope),
3515 None,
3516 &as_of_date,
3517 &as_of_date,
3518 ))
3519 })
3520 .map_err(|e| ErrorData::internal_error(format!("temporal query error: {e}"), None))?;
3521 let json_results: Vec<serde_json::Value> = results
3522 .iter()
3523 .take(k)
3524 .map(|r| {
3525 serde_json::json!({
3526 "result_id": r.source.result_id(),
3527 "content": r.content,
3528 "score": r.score,
3529 })
3530 })
3531 .collect();
3532 let trace_json =
3533 serde_json::to_value(&query_trace).unwrap_or_else(|_| serde_json::json!(null));
3534 json_to_string(&serde_json::json!({
3535 "ok": true,
3536 "query": query,
3537 "as_of_date": as_of_date,
3538 "results": json_results,
3539 "count": json_results.len(),
3540 "trace": trace_json,
3541 }))
3542 }
3543
3544 #[cfg(feature = "orchestration")]
3545 #[tool(
3546 description = "Resolve an entity mention against the runtime's scope-aware entity registry. Returns resolved entity, match quality, and alternative candidates.",
3547 annotations(read_only_hint = true)
3548 )]
3549 fn sm_entity_lookup(
3550 &self,
3551 Parameters(EntityLookupParams {
3552 mention,
3553 namespace,
3554 domain,
3555 }): Parameters<EntityLookupParams>,
3556 ) -> Result<String, ErrorData> {
3557 let runtime = self.runtime.as_ref().ok_or_else(|| {
3558 ErrorData::internal_error("orchestration runtime not available", None)
3559 })?;
3560 let ns = namespace.as_deref().unwrap_or("general");
3561 let mut scope = knowledge_runtime::ScopeKey::namespace_only(ns);
3562 if let Some(d) = &domain {
3563 scope.domain = Some(d.clone());
3564 }
3565 let result = runtime.entity_registry().resolve(&mention, &scope);
3566 let quality_str = match result.quality {
3567 knowledge_runtime::MatchQuality::ExactCanonical => "exact_canonical",
3568 knowledge_runtime::MatchQuality::ExactAlias => "exact_alias",
3569 knowledge_runtime::MatchQuality::ScopedFallback => "scoped_fallback",
3570 knowledge_runtime::MatchQuality::Unresolved => "unresolved",
3571 };
3572 let entity_json = result
3573 .entity
3574 .as_ref()
3575 .map(|e| serde_json::to_value(e).unwrap_or_else(|_| serde_json::json!(null)))
3576 .unwrap_or(serde_json::json!(null));
3577 let alternatives_json: Vec<serde_json::Value> = result
3578 .alternatives
3579 .iter()
3580 .map(|e| serde_json::to_value(e).unwrap_or_else(|_| serde_json::json!(null)))
3581 .collect();
3582 json_to_string(&serde_json::json!({
3583 "ok": true,
3584 "mention": mention,
3585 "quality": quality_str,
3586 "entity": entity_json,
3587 "alternatives": alternatives_json,
3588 "queried_scope": result.queried_scope.to_string(),
3589 }))
3590 }
3591
3592 #[cfg(feature = "orchestration")]
3593 #[tool(
3594 description = "Check the health of a projection by namespace and optional kind. Returns health status (healthy, stale, missing, etc.).",
3595 annotations(read_only_hint = true)
3596 )]
3597 fn sm_projection_health(
3598 &self,
3599 Parameters(ProjectionHealthParams {
3600 namespace,
3601 projection_kind,
3602 }): Parameters<ProjectionHealthParams>,
3603 ) -> Result<String, ErrorData> {
3604 let runtime = self.runtime.as_ref().ok_or_else(|| {
3605 ErrorData::internal_error("orchestration runtime not available", None)
3606 })?;
3607 let kind = match projection_kind.as_deref().unwrap_or("entity") {
3608 "entity" => knowledge_runtime::ProjectionKind::Entity,
3609 "temporal" => knowledge_runtime::ProjectionKind::Temporal,
3610 "route_stats" => knowledge_runtime::ProjectionKind::RouteStats,
3611 other => knowledge_runtime::ProjectionKind::Custom(other.to_string()),
3612 };
3613 let scope_key = knowledge_runtime::ScopeKey::namespace_only(&namespace);
3614 let proj_id = knowledge_runtime::ProjectionId::new(kind, &namespace, scope_key);
3615 let health = runtime.projection_health(&proj_id);
3616 let health_str = match &health {
3617 knowledge_runtime::ProjectionHealth::Healthy => "healthy",
3618 knowledge_runtime::ProjectionHealth::Stale => "stale",
3619 knowledge_runtime::ProjectionHealth::Missing => "missing",
3620 knowledge_runtime::ProjectionHealth::Rebuilding => "rebuilding",
3621 knowledge_runtime::ProjectionHealth::ImportLagging => "import_lagging",
3622 knowledge_runtime::ProjectionHealth::ImportFailed => "import_failed",
3623 };
3624 json_to_string(&serde_json::json!({
3625 "ok": true,
3626 "namespace": namespace,
3627 "projection_id": proj_id.to_string(),
3628 "health": health_str,
3629 }))
3630 }
3631
3632 #[cfg(feature = "claim-integration")]
3635 #[tool(
3636 description = "Get proof-debt budget status for a scope. Returns summary with budget, consumed, available, gate decision, and exhaustion state.",
3637 annotations(read_only_hint = true)
3638 )]
3639 fn sm_proof_debt_status(
3640 &self,
3641 Parameters(ProofDebtStatusParams { scope }): Parameters<ProofDebtStatusParams>,
3642 ) -> Result<String, ErrorData> {
3643 use claim_ledger::{ProofDebtBudgetV1, ProofDebtSummaryV1};
3644 let budget = ProofDebtBudgetV1::new(&scope, 1_000_000);
3645 let summary = ProofDebtSummaryV1::from_budget(&budget);
3646 let gate_decision = match summary.gate_decision {
3647 claim_ledger::ProofDebtGateDecision::Proceed => "proceed",
3648 claim_ledger::ProofDebtGateDecision::Warn => "warn",
3649 claim_ledger::ProofDebtGateDecision::Degrade => "degrade",
3650 claim_ledger::ProofDebtGateDecision::Retract => "retract",
3651 claim_ledger::ProofDebtGateDecision::Waived => "waived",
3652 };
3653 json_to_string(&serde_json::json!({
3654 "ok": true,
3655 "scope": scope,
3656 "budget_id": summary.budget_id,
3657 "budget_micros": summary.budget_micros,
3658 "consumed_micros": summary.consumed_micros,
3659 "available_micros": summary.available_micros,
3660 "consumed_pct": summary.consumed_pct,
3661 "exhausted": summary.exhausted,
3662 "gate_decision": gate_decision,
3663 "gate_summary": summary.gate_summary,
3664 }))
3665 }
3666
3667 #[cfg(feature = "claim-integration")]
3668 #[tool(
3669 description = "Evaluate a proof-debt budget gate for a scope. Returns the gate decision (proceed, warn, degrade, retract) with details.",
3670 annotations(read_only_hint = true)
3671 )]
3672 fn sm_evaluate_proof_debt_gate(
3673 &self,
3674 Parameters(EvaluateProofDebtGateParams {
3675 scope,
3676 budget_micros,
3677 }): Parameters<EvaluateProofDebtGateParams>,
3678 ) -> Result<String, ErrorData> {
3679 use claim_ledger::{evaluate_proof_debt_gate, ProofDebtBudgetV1};
3680 let micros = budget_micros.unwrap_or(1_000_000);
3681 let budget = ProofDebtBudgetV1::new(&scope, micros);
3682 let gate = evaluate_proof_debt_gate(&budget);
3683 let decision_str = match gate.decision {
3684 claim_ledger::ProofDebtGateDecision::Proceed => "proceed",
3685 claim_ledger::ProofDebtGateDecision::Warn => "warn",
3686 claim_ledger::ProofDebtGateDecision::Degrade => "degrade",
3687 claim_ledger::ProofDebtGateDecision::Retract => "retract",
3688 claim_ledger::ProofDebtGateDecision::Waived => "waived",
3689 };
3690 json_to_string(&serde_json::json!({
3691 "ok": true,
3692 "scope": scope,
3693 "budget_id": gate.budget_id,
3694 "decision": decision_str,
3695 "consumed_pct": gate.consumed_pct,
3696 "exhausted": gate.exhausted,
3697 "summary": gate.summary,
3698 "allows_proceed": gate.decision.allows_proceed(),
3699 "blocks": gate.decision.blocks(),
3700 }))
3701 }
3702
3703 #[cfg(feature = "claim-integration")]
3704 #[tool(
3705 description = "Record a support admission for a claim. Creates a SupportAdmissionReceipt with method, rationale, and operator reference.",
3706 annotations(idempotent_hint = true)
3707 )]
3708 fn sm_add_support_admission(
3709 &self,
3710 Parameters(AddSupportAdmissionParams {
3711 claim_id,
3712 method,
3713 rationale,
3714 operator_id,
3715 }): Parameters<AddSupportAdmissionParams>,
3716 ) -> Result<String, ErrorData> {
3717 use claim_ledger::{SupportAdmissionMethod, SupportAdmissionReceipt, SupportState};
3718 let method_enum = match method.to_lowercase().as_str() {
3719 "operator_admitted" => SupportAdmissionMethod::OperatorAdmitted,
3720 "test_fixture_admitted" => SupportAdmissionMethod::TestFixtureAdmitted,
3721 "external_receipt_admitted" => SupportAdmissionMethod::ExternalReceiptAdmitted,
3722 _ => {
3723 return Err(ErrorData::invalid_params(
3724 format!(
3725 "Invalid method '{method}'. Must be: operator_admitted, test_fixture_admitted, or external_receipt_admitted"
3726 ),
3727 None,
3728 ));
3729 }
3730 };
3731 let prev_ref = format!("sj_prev_{}", &claim_id);
3732 let new_ref = format!("sj_new_{}", &claim_id);
3733 let mut receipt = SupportAdmissionReceipt::new(
3734 &claim_id,
3735 &prev_ref,
3736 &new_ref,
3737 method_enum,
3738 SupportState::Supported,
3739 &rationale,
3740 );
3741 receipt.operator_ref = operator_id;
3742 json_to_string(&serde_json::json!({
3743 "ok": true,
3744 "support_admission_receipt_id": receipt.support_admission_receipt_id,
3745 "claim_id": receipt.claim_id,
3746 "method": method,
3747 "admitted_support_state": "supported",
3748 "rationale": receipt.rationale,
3749 "operator_ref": receipt.operator_ref,
3750 "recorded_time": receipt.recorded_time.to_rfc3339(),
3751 }))
3752 }
3753
3754 #[cfg(feature = "claim-integration")]
3755 #[tool(
3756 description = "Record a contradiction between two claims. Creates a ContradictionRecord with Open status, detection method, and optional evidence.",
3757 annotations(idempotent_hint = true)
3758 )]
3759 fn sm_record_contradiction(
3760 &self,
3761 Parameters(RecordContradictionParams {
3762 claim_a_id,
3763 claim_b_id,
3764 detection_method,
3765 evidence,
3766 }): Parameters<RecordContradictionParams>,
3767 ) -> Result<String, ErrorData> {
3768 use claim_ledger::ContradictionRecord;
3769 let mut record = ContradictionRecord::new(
3770 &claim_a_id,
3771 &claim_b_id,
3772 &detection_method,
3773 &detection_method,
3774 );
3775 if let Some(ev) = &evidence {
3776 record.rationale = ev.clone();
3777 }
3778 let status_str = match record.status {
3779 claim_ledger::ContradictionStatus::Candidate => "candidate",
3780 claim_ledger::ContradictionStatus::UnderReview => "under_review",
3781 claim_ledger::ContradictionStatus::Unresolved => "unresolved",
3782 claim_ledger::ContradictionStatus::Confirmed => "confirmed",
3783 claim_ledger::ContradictionStatus::Rejected => "rejected",
3784 claim_ledger::ContradictionStatus::Superseded => "superseded",
3785 };
3786 json_to_string(&serde_json::json!({
3787 "ok": true,
3788 "contradiction_id": record.contradiction_id,
3789 "claim_refs": record.claim_refs,
3790 "pattern": record.pattern,
3791 "rationale": record.rationale,
3792 "status": status_str,
3793 "created_recorded_time": record.created_recorded_time.to_rfc3339(),
3794 }))
3795 }
3796
3797 #[cfg(feature = "claim-integration")]
3798 #[tool(
3799 description = "Resolve a contradiction with a resolution outcome (confirmed, rejected, superseded). Creates a ContradictionResolutionReceipt.",
3800 annotations(idempotent_hint = true)
3801 )]
3802 fn sm_resolve_contradiction(
3803 &self,
3804 Parameters(ResolveContradictionParams {
3805 contradiction_id,
3806 resolution,
3807 rationale,
3808 superseding_claim_id,
3809 }): Parameters<ResolveContradictionParams>,
3810 ) -> Result<String, ErrorData> {
3811 use claim_ledger::{ContradictionResolution, ContradictionResolutionReceipt};
3812 let resolution_enum = match resolution.to_lowercase().as_str() {
3813 "confirmed" => ContradictionResolution::Confirmed,
3814 "rejected" => ContradictionResolution::Rejected,
3815 "superseded" => ContradictionResolution::Superseded,
3816 _ => {
3817 return Err(ErrorData::invalid_params(
3818 format!(
3819 "Invalid resolution '{resolution}'. Must be: confirmed, rejected, or superseded"
3820 ),
3821 None,
3822 ));
3823 }
3824 };
3825 let receipt = ContradictionResolutionReceipt::new(
3826 &contradiction_id,
3827 "open",
3828 resolution_enum,
3829 &rationale,
3830 );
3831 let resolution_str = match receipt.resolution {
3832 ContradictionResolution::Confirmed => "confirmed",
3833 ContradictionResolution::Rejected => "rejected",
3834 ContradictionResolution::Superseded => "superseded",
3835 };
3836 json_to_string(&serde_json::json!({
3837 "ok": true,
3838 "contradiction_resolution_receipt_id": receipt.contradiction_resolution_receipt_id,
3839 "contradiction_id": receipt.contradiction_id,
3840 "resolution": resolution_str,
3841 "rationale": receipt.rationale,
3842 "superseding_claim_id": superseding_claim_id,
3843 "recorded_time": receipt.recorded_time.to_rfc3339(),
3844 }))
3845 }
3846
3847 #[cfg(feature = "claim-integration")]
3848 #[tool(
3849 description = "Verify a claim ledger's hash chain integrity. Accepts JSONL text of ledger entries, parses and verifies the chain. Returns valid flag, entry count, and first break if any.",
3850 annotations(read_only_hint = true)
3851 )]
3852 fn sm_verify_ledger(
3853 &self,
3854 Parameters(VerifyLedgerParams { entries_jsonl }): Parameters<VerifyLedgerParams>,
3855 ) -> Result<String, ErrorData> {
3856 use claim_ledger::{parse_ledger_entries, verify_ledger};
3857 let entries = parse_ledger_entries(&entries_jsonl);
3858 let entry_count = entries.len();
3859 let verification = verify_ledger(&entries);
3860 let first_break = verification.errors.first().cloned();
3861 json_to_string(&serde_json::json!({
3862 "ok": true,
3863 "verified": verification.valid,
3864 "entry_count": entry_count,
3865 "last_sequence": verification.last_sequence,
3866 "last_entry_digest": verification.last_entry_digest,
3867 "error_count": verification.errors.len(),
3868 "first_break": first_break,
3869 "errors": verification.errors,
3870 }))
3871 }
3872
3873 #[cfg(feature = "claim-integration")]
3874 #[tool(
3875 description = "Export a bundle of claims with optional evidence and contradictions. Creates an ExportReceipt for deterministic verification.",
3876 annotations(read_only_hint = true)
3877 )]
3878 fn sm_export_claim_bundle(
3879 &self,
3880 Parameters(ExportClaimBundleParams {
3881 claim_ids,
3882 include_evidence,
3883 include_contradictions,
3884 }): Parameters<ExportClaimBundleParams>,
3885 ) -> Result<String, ErrorData> {
3886 use claim_ledger::ExportReceipt;
3887 let inc_ev = include_evidence.unwrap_or(true);
3888 let inc_contra = include_contradictions.unwrap_or(true);
3889 let input_refs: Vec<String> = claim_ids.clone();
3890 let mut receipt = ExportReceipt::new(
3891 "bundle_export",
3892 input_refs.clone(),
3893 claim_ledger::ids::ulid(),
3894 );
3895 receipt.mark_success();
3896 let bundle = serde_json::json!({
3898 "claim_ids": claim_ids,
3899 "include_evidence": inc_ev,
3900 "include_contradictions": inc_contra,
3901 "export_receipt_id": receipt.export_receipt_id,
3902 });
3903 let bundle_str = serde_json::to_string(&bundle)
3904 .map_err(|e| ErrorData::internal_error(format!("serialization error: {e}"), None))?;
3905 let digest = claim_ledger::ids::sha256_text(&bundle_str);
3906 receipt.bind_output(format!("bundle:{}", receipt.export_receipt_id), digest);
3907 json_to_string(&serde_json::json!({
3908 "ok": true,
3909 "export_receipt_id": receipt.export_receipt_id,
3910 "operation": receipt.operation,
3911 "input_refs": receipt.input_refs,
3912 "output_ref": receipt.output_ref,
3913 "output_digest": receipt.output_digest,
3914 "status": receipt.status,
3915 "digest_semantics": receipt.digest_semantics,
3916 "recorded_time": receipt.recorded_time.to_rfc3339(),
3917 "bundle": bundle,
3918 }))
3919 }
3920
3921 #[cfg(feature = "claim-integration")]
3922 #[tool(
3923 description = "Record a supersession of an old claim by a new claim. Creates a SupersessionReceipt with rationale.",
3924 annotations(idempotent_hint = true)
3925 )]
3926 fn sm_supersede_claim(
3927 &self,
3928 Parameters(SupersedeClaimParams {
3929 old_claim_id,
3930 new_claim_id,
3931 rationale,
3932 }): Parameters<SupersedeClaimParams>,
3933 ) -> Result<String, ErrorData> {
3934 use claim_ledger::SupersessionReceipt;
3935 let receipt = SupersessionReceipt::new(&old_claim_id, &new_claim_id, &rationale);
3936 json_to_string(&serde_json::json!({
3937 "ok": true,
3938 "supersession_receipt_id": receipt.supersession_receipt_id,
3939 "superseded_ref": receipt.superseded_ref,
3940 "superseding_ref": receipt.superseding_ref,
3941 "rationale": receipt.rationale,
3942 "recorded_time": receipt.recorded_time.to_rfc3339(),
3943 }))
3944 }
3945
3946 #[tool(
3949 description = "Import a projection envelope atomically. All records are committed in a single transaction or the entire import is rolled back. Pass the envelope as a JSON string.",
3950 annotations(idempotent_hint = true)
3951 )]
3952 #[allow(deprecated)]
3953 fn sm_import_envelope(
3954 &self,
3955 Parameters(ImportEnvelopeParams { envelope_json }): Parameters<ImportEnvelopeParams>,
3956 ) -> Result<String, ErrorData> {
3957 let envelope: semantic_memory::projection_import::ImportEnvelope =
3958 serde_json::from_str(&envelope_json).map_err(|e| {
3959 ErrorData::invalid_params(format!("Failed to parse envelope JSON: {e}"), None)
3960 })?;
3961 envelope.validate().map_err(|e| {
3962 ErrorData::invalid_params(format!("Envelope validation failed: {e}"), None)
3963 })?;
3964 let store = &self.bridge.store;
3965 let result = tokio::task::block_in_place(|| {
3966 Handle::current().block_on(store.import_envelope(&envelope))
3967 });
3968 match result {
3969 Ok(receipt) => json_to_string(&serde_json::json!({
3970 "ok": true,
3971 "envelope_id": receipt.envelope_id,
3972 "was_duplicate": receipt.was_duplicate,
3973 "imported_count": receipt.record_count,
3974 "receipt_id": receipt.envelope_id,
3975 })),
3976 Err(e) => Err(ErrorData::internal_error(
3977 format!("import_envelope error: {e}"),
3978 None,
3979 )),
3980 }
3981 }
3982
3983 #[tool(
3984 description = "Check whether an envelope has already been imported. Returns import receipts for the given envelope ID.",
3985 annotations(read_only_hint = true)
3986 )]
3987 #[allow(deprecated)]
3988 fn sm_import_status(
3989 &self,
3990 Parameters(ImportStatusParams { envelope_id }): Parameters<ImportStatusParams>,
3991 ) -> Result<String, ErrorData> {
3992 use semantic_memory::projection_import::EnvelopeId;
3993 let store = &self.bridge.store;
3994 let env_id = EnvelopeId::new(&envelope_id);
3995 let result = tokio::task::block_in_place(|| {
3996 Handle::current().block_on(store.import_status(&env_id))
3997 });
3998 match result {
3999 Ok(receipts) => json_to_string(&serde_json::json!({
4000 "ok": true,
4001 "envelope_id": envelope_id,
4002 "receipts": serde_json::to_value(&receipts).unwrap_or(serde_json::json!([])),
4003 "count": receipts.len(),
4004 })),
4005 Err(e) => Err(ErrorData::internal_error(
4006 format!("import_status error: {e}"),
4007 None,
4008 )),
4009 }
4010 }
4011
4012 #[tool(
4013 description = "List recent imports, optionally filtered by namespace. Returns import receipt records.",
4014 annotations(read_only_hint = true)
4015 )]
4016 #[allow(deprecated)]
4017 fn sm_list_imports(
4018 &self,
4019 Parameters(ListImportsParams { namespace, limit }): Parameters<ListImportsParams>,
4020 ) -> Result<String, ErrorData> {
4021 let store = &self.bridge.store;
4022 let lim = limit.unwrap_or(20) as usize;
4023 let result = tokio::task::block_in_place(|| {
4024 Handle::current().block_on(store.list_imports(namespace.as_deref(), lim))
4025 });
4026 match result {
4027 Ok(receipts) => json_to_string(&serde_json::json!({
4028 "ok": true,
4029 "receipts": serde_json::to_value(&receipts).unwrap_or(serde_json::json!([])),
4030 "count": receipts.len(),
4031 })),
4032 Err(e) => Err(ErrorData::internal_error(
4033 format!("list_imports error: {e}"),
4034 None,
4035 )),
4036 }
4037 }
4038}
4039
4040fn build_path_segments(
4043 store: &semantic_memory::MemoryStore,
4044 path: &[String],
4045) -> Vec<serde_json::Value> {
4046 let mut segments = Vec::new();
4047 if path.len() < 2 {
4048 return segments;
4049 }
4050
4051 for i in 0..path.len() - 1 {
4052 let from = &path[i];
4053 let to = &path[i + 1];
4054
4055 let g = store.graph_view();
4057 match g.neighbors(from, semantic_memory::GraphDirection::Both, 1) {
4058 Ok(edges) => {
4059 let connecting = edges.iter().find(|e| {
4061 (e.source == *from && e.target == *to) || (e.source == *to && e.target == *from)
4062 });
4063
4064 if let Some(edge) = connecting {
4065 let edge_type_str = match &edge.edge_type {
4066 semantic_memory::GraphEdgeType::Semantic { cosine_similarity } => {
4067 serde_json::json!({
4068 "type": "semantic",
4069 "cosine_similarity": cosine_similarity,
4070 })
4071 }
4072 semantic_memory::GraphEdgeType::Temporal { delta_secs } => {
4073 serde_json::json!({
4074 "type": "temporal",
4075 "delta_secs": delta_secs,
4076 })
4077 }
4078 semantic_memory::GraphEdgeType::Causal {
4079 confidence,
4080 evidence_ids,
4081 } => {
4082 serde_json::json!({
4083 "type": "causal",
4084 "confidence": confidence,
4085 "evidence_ids": evidence_ids,
4086 })
4087 }
4088 semantic_memory::GraphEdgeType::Entity { relation } => {
4089 serde_json::json!({
4090 "type": "entity",
4091 "relation": relation,
4092 })
4093 }
4094 };
4095
4096 segments.push(serde_json::json!({
4097 "source": from,
4098 "target": to,
4099 "edge_type": edge_type_str,
4100 "weight": edge.weight,
4101 "metadata": edge.metadata,
4102 }));
4103 } else {
4104 segments.push(serde_json::json!({
4107 "source": from,
4108 "target": to,
4109 "edge_type": null,
4110 "weight": null,
4111 "metadata": null,
4112 }));
4113 }
4114 }
4115 Err(_) => {
4116 segments.push(serde_json::json!({
4117 "source": from,
4118 "target": to,
4119 "edge_type": null,
4120 "weight": null,
4121 "metadata": null,
4122 }));
4123 }
4124 }
4125 }
4126
4127 segments
4128}
4129
4130#[tool_handler(
4131 router = self.tool_router,
4132 name = "semantic-memory-mcp",
4133 version = "0.3.1",
4134 instructions = "Persistent local semantic memory with hybrid search, graph reasoning, and conversation persistence. ALWAYS search first (sm_search) before asking the user for context. Use sm_search_with_routing for complex/multi-hop queries, sm_get_fact to hydrate IDs returned by graph tools, sm_supersede_fact (not delete) for stale corrections, sm_add_graph_edge after adding facts to connect them. Read tools are safe; write tools (add/delete/supersede) should be user-approved. Search auto-filters superseded facts unless querying for history."
4135)]
4136impl ServerHandler for SemanticMemoryServer {}