1use std::io::{BufRead, BufReader, Read, Write};
15use std::net::TcpListener;
16use tokio::runtime::Handle;
17use tokio::task::block_in_place;
18
19use crate::bridge::MemoryBridge;
20
21fn parse_body_json(body: &str) -> serde_json::Value {
26 #[cfg(feature = "llm-parser")]
27 {
28 llm_output_parser::parse_json_value(body)
29 .unwrap_or_else(|_| serde_json::from_str(body).unwrap_or(serde_json::Value::Null))
30 }
31 #[cfg(not(feature = "llm-parser"))]
32 {
33 serde_json::from_str(body).unwrap_or(serde_json::Value::Null)
34 }
35}
36
37fn rerank_results(
40 query: &str,
41 results: &[serde_json::Value],
42 model: &str,
43) -> Vec<serde_json::Value> {
44 let client = reqwest::blocking::Client::new();
45 let mut scored: Vec<(f64, serde_json::Value)> = results
46 .iter()
47 .map(|r| {
48 let content = r.get("content").and_then(|v| v.as_str()).unwrap_or("");
49 let truncated: String = content.chars().take(500).collect();
50 let prompt = format!(
51 "Rate the relevance of this document to the query on a scale of 1-5. Reply with ONLY the number.\nQuery: {query}\nDocument: {truncated}\nRating:"
52 );
53 let body = serde_json::json!({
54 "model": model,
55 "prompt": prompt,
56 "stream": false,
57 "options": {"temperature": 0, "num_predict": 1}
58 });
59 let rating = client
60 .post("http://127.0.0.1:11434/api/generate")
61 .json(&body)
62 .send()
63 .ok()
64 .and_then(|resp| resp.json::<serde_json::Value>().ok())
65 .and_then(|v| {
66 v.get("response")
67 .and_then(|r| r.as_str())
68 .and_then(|s| s.trim().chars().next())
69 .and_then(|c| c.to_digit(10))
70 .map(|d| d as f64)
71 })
72 .unwrap_or(1.0);
73 (rating, r.clone())
74 })
75 .collect();
76 scored.sort_by(|a, b| b.0.partial_cmp(&a.0).unwrap_or(std::cmp::Ordering::Equal));
77 scored
78 .into_iter()
79 .map(|(score, mut r)| {
80 if let Some(obj) = r.as_object_mut() {
81 obj.insert("rerank_score".to_string(), serde_json::json!(score));
82 }
83 r
84 })
85 .collect()
86}
87
88pub fn start_http_server(port: u16, bridge: MemoryBridge, handle: Handle) {
89 std::thread::spawn(move || {
90 let listener = match TcpListener::bind(("127.0.0.1", port)) {
91 Ok(l) => {
92 eprintln!("HTTP search server listening on 127.0.0.1:{}", port);
93 l
94 }
95 Err(e) => {
96 eprintln!("Failed to bind HTTP port {}: {}", port, e);
97 return;
98 }
99 };
100
101 for stream in listener.incoming() {
102 let stream = match stream {
103 Ok(s) => s,
104 Err(_) => continue,
105 };
106
107 let bridge = bridge.clone();
108 let h = handle.clone();
109 std::thread::spawn(move || {
110 handle_connection(stream, bridge, h);
111 });
112 }
113 });
114}
115
116fn handle_connection(mut stream: std::net::TcpStream, bridge: MemoryBridge, handle: Handle) {
117 let mut reader = BufReader::new(stream.try_clone().expect("clone"));
118 let mut request_line = String::new();
119 if reader.read_line(&mut request_line).is_err() {
120 return;
121 }
122
123 let parts: Vec<&str> = request_line.split_whitespace().collect();
124 if parts.len() < 2 {
125 return;
126 }
127 let method = parts[0];
128 let path = parts[1];
129
130 let mut content_length = 0;
131 loop {
132 let mut header = String::new();
133 if reader.read_line(&mut header).is_err() {
134 return;
135 }
136 if header.trim().is_empty() {
137 break;
138 }
139 if let Some(len_str) = header
140 .strip_prefix("Content-Length:")
141 .or_else(|| header.strip_prefix("content-length:"))
142 {
143 content_length = len_str.trim().parse().unwrap_or(0);
144 }
145 }
146
147 let mut body = vec![0u8; content_length];
148 if content_length > 0 && reader.read_exact(&mut body).is_err() {
149 return;
150 }
151 let body_str = String::from_utf8_lossy(&body);
152
153 let (status, response) = match (method, path) {
154 ("GET", "/health") => (
155 "200 OK",
156 serde_json::json!({"ok": true, "service": "semantic-memory-mcp"}),
157 ),
158 ("POST", "/search") => handle_search(&body_str, &bridge, &handle),
159 ("POST", "/search-routed") => handle_search_routed(&body_str, &bridge, &handle),
160 #[cfg(feature = "orchestration")]
161 ("POST", "/query-orchestrated") => handle_query_orchestrated(&body_str, &bridge, &handle),
162 ("POST", "/rerank") => handle_rerank(&body_str),
163 ("POST", "/stats") => handle_stats(&bridge, &handle),
164 ("POST", "/add") => handle_add_fact(&body_str, &bridge, &handle),
165 ("POST", "/add-edge") => handle_add_edge(&body_str, &bridge, &handle),
166 ("POST", "/delete-fact") => handle_delete_fact(&body_str, &bridge, &handle),
167 ("POST", "/record-outcome") => handle_record_outcome(&body_str, &bridge, &handle),
168 ("GET", "/verify-integrity") => handle_verify_integrity(&bridge, &handle),
169 ("POST", "/discord") => handle_discord(&body_str, &bridge, &handle),
170 ("POST", "/maintenance/check") => handle_maintenance_check(&bridge, &handle),
171 ("POST", "/maintenance/vacuum") => handle_maintenance_vacuum(&bridge, &handle),
172 ("POST", "/maintenance/reembed") => handle_maintenance_reembed(&bridge, &handle),
173 ("POST", "/maintenance/reconcile") => {
174 handle_maintenance_reconcile(&body_str, &bridge, &handle)
175 }
176 ("POST", "/maintenance/compact-hnsw") => handle_maintenance_compact_hnsw(&bridge, &handle),
177 ("POST", "/maintenance/auto-edge") => {
178 handle_maintenance_auto_edge(&body_str, &bridge, &handle)
179 }
180 _ => (
181 "404 Not Found",
182 serde_json::json!({"error": "not found", "path": path}),
183 ),
184 };
185
186 let response_str = serde_json::to_string(&response).unwrap_or_default();
187 let response_bytes = response_str.as_bytes();
188 let http_response = format!(
189 "HTTP/1.1 {}\r\nContent-Type: application/json\r\nContent-Length: {}\r\nConnection: close\r\n\r\n",
190 status,
191 response_bytes.len()
192 );
193
194 let _ = stream.write_all(http_response.as_bytes());
195 let _ = stream.write_all(response_bytes);
196 let _ = stream.flush();
197}
198
199fn handle_search(
200 body: &str,
201 bridge: &MemoryBridge,
202 handle: &Handle,
203) -> (&'static str, serde_json::Value) {
204 let params: serde_json::Value = parse_body_json(body);
205 if params.is_null() {
206 return (
207 "400 Bad Request",
208 serde_json::json!({"ok": false, "error": "invalid JSON body"}),
209 );
210 }
211
212 let query = params.get("query").and_then(|v| v.as_str()).unwrap_or("");
213 let top_k = params.get("top_k").and_then(|v| v.as_u64()).unwrap_or(5) as usize;
214 let namespaces: Option<Vec<String>> = params
215 .get("namespaces")
216 .and_then(|v| serde_json::from_value(v.clone()).ok());
217 let do_rerank = params
218 .get("rerank")
219 .and_then(|v| v.as_bool())
220 .unwrap_or(false);
221
222 if query.is_empty() {
223 return (
224 "400 Bad Request",
225 serde_json::json!({"ok": false, "error": "missing 'query' field"}),
226 );
227 }
228
229 let store = &bridge.store;
230 let ns_slice: Option<Vec<&str>> = namespaces
231 .as_ref()
232 .map(|v| v.iter().map(|s| s.as_str()).collect());
233 let fetch_k = if do_rerank { top_k * 2 } else { top_k };
235 let result = block_in_place(|| {
236 handle.block_on(store.search(query, Some(fetch_k), ns_slice.as_deref(), None))
237 });
238
239 match result {
240 Ok(results) => {
241 let json_results: Vec<serde_json::Value> = results
242 .iter()
243 .map(|r| {
244 let namespace = match &r.source {
245 semantic_memory::SearchSource::Fact { namespace, .. } => namespace.clone(),
246 semantic_memory::SearchSource::Chunk { document_title, .. } => {
247 document_title.clone()
248 }
249 _ => String::new(),
250 };
251 serde_json::json!({
252 "result_id": r.source.result_id(),
253 "content": r.content,
254 "score": r.score,
255 "cosine_similarity": r.cosine_similarity,
256 "namespace": namespace,
257 })
258 })
259 .collect();
260
261 let final_results: Vec<serde_json::Value> = if do_rerank && !json_results.is_empty() {
262 rerank_results(query, &json_results, "granite4.1:3b")
263 .into_iter()
264 .take(top_k)
265 .collect()
266 } else {
267 json_results
268 };
269
270 let count = final_results.len();
271 let provenance = serde_json::json!({
272 "stages_fired": {
273 "bm25": true,
274 "vector": true,
275 "late_interaction": false,
276 "rerank": do_rerank,
277 },
278 "result_count": count,
279 "view": "semantic",
280 "widening_occurred": false,
281 "widening_reason": null,
282 "verification_status": "verified",
283 });
284 (
285 "200 OK",
286 serde_json::json!({
287 "ok": true,
288 "query": query,
289 "top_k": top_k,
290 "results": final_results,
291 "count": count,
292 "reranked": do_rerank,
293 "provenance": provenance,
294 }),
295 )
296 }
297 Err(e) => (
298 "500 Internal Server Error",
299 serde_json::json!({"ok": false, "error": format!("search error: {e}")}),
300 ),
301 }
302}
303
304fn handle_search_routed(
312 body: &str,
313 bridge: &MemoryBridge,
314 handle: &Handle,
315) -> (&'static str, serde_json::Value) {
316 use semantic_memory::integration::plan_execution;
317 use semantic_memory::routing::RetrievalRouter;
318
319 let params: serde_json::Value = parse_body_json(body);
320 if params.is_null() {
321 return (
322 "400 Bad Request",
323 serde_json::json!({"ok": false, "error": "invalid JSON body"}),
324 );
325 }
326
327 let query = params.get("query").and_then(|v| v.as_str()).unwrap_or("");
328 let base_top_k = params.get("top_k").and_then(|v| v.as_u64()).unwrap_or(12) as usize;
329 let query_class = params
330 .get("query_class")
331 .and_then(|v| v.as_str())
332 .unwrap_or("A");
333 let namespaces: Option<Vec<String>> = params
334 .get("namespaces")
335 .and_then(|v| serde_json::from_value(v.clone()).ok());
336 let contradictions: Vec<(String, String)> = params
337 .get("contradictions")
338 .and_then(|v| serde_json::from_value(v.clone()).ok())
339 .unwrap_or_default();
340 let group_by_community = params
341 .get("group_by_community")
342 .and_then(|v| v.as_bool())
343 .unwrap_or(false);
344
345 if query.is_empty() {
346 return (
347 "400 Bad Request",
348 serde_json::json!({"ok": false, "error": "missing 'query' field"}),
349 );
350 }
351
352 let router = RetrievalRouter {
354 decoder_enabled: true,
355 discord_enabled: true,
356 corpus_density: 0.5,
357 ..Default::default()
358 };
359 let decision = router.route_query(query);
360 let contras = contradictions.clone();
361 let plan = plan_execution(&decision, contras.clone());
362
363 let top_k = if query_class == "D" {
365 (base_top_k * 2).min(20)
366 } else {
367 base_top_k
368 };
369
370 let store = &bridge.store;
371 let ns_slice: Option<Vec<&str>> = namespaces
372 .as_ref()
373 .map(|v| v.iter().map(|s| s.as_str()).collect());
374
375 let result = if query_class == "C" {
377 use semantic_memory::{ExactnessProfile, SearchContext};
378 let mut ctx = SearchContext::default_now();
379 ctx.exactness_profile = ExactnessProfile::PreferExact;
380 block_in_place(|| {
381 handle.block_on(store.search_with_context(
382 query,
383 Some(top_k),
384 ns_slice.as_deref(),
385 None,
386 ctx,
387 ))
388 })
389 .map(|r| r.results)
390 } else {
391 block_in_place(|| {
392 handle.block_on(store.search(query, Some(top_k), ns_slice.as_deref(), None))
393 })
394 };
395
396 match result {
397 Ok(results) => {
398 let json_results: Vec<serde_json::Value> = results
399 .iter()
400 .map(|r| {
401 let namespace = match &r.source {
402 semantic_memory::SearchSource::Fact { namespace, .. } => namespace.clone(),
403 semantic_memory::SearchSource::Chunk { document_title, .. } => {
404 document_title.clone()
405 }
406 _ => String::new(),
407 };
408 serde_json::json!({
409 "result_id": r.source.result_id(),
410 "content": r.content,
411 "score": r.score,
412 "cosine_similarity": r.cosine_similarity,
413 "namespace": namespace,
414 "source_type": match &r.source {
415 semantic_memory::SearchSource::Fact { .. } => "fact",
416 semantic_memory::SearchSource::Chunk { .. } => "chunk",
417 semantic_memory::SearchSource::Message { .. } => "message",
418 _ => "unknown",
419 },
420 })
421 })
422 .collect();
423
424 let mut factor_graph_payload = serde_json::json!({"enabled": false});
425 let mut decoder_executed = false;
426 let mut discord_executed = false;
427 let mut discord_results_payload: Vec<serde_json::Value> = Vec::new();
428
429 if decision.decoder {
431 #[cfg(feature = "full")]
432 {
433 use semantic_memory::factor_graph::{
434 factors_from_edges, FactorGraph, FactorGraphConfig,
435 };
436
437 let graph_edges =
438 block_in_place(|| handle.block_on(store.list_all_graph_edges()));
439
440 if let Ok(edges) = graph_edges {
441 let raw_edges: Vec<(
442 String,
443 String,
444 semantic_memory::GraphEdgeType,
445 f64,
446 Option<String>,
447 )> = edges
448 .iter()
449 .map(|edge| {
450 let parsed_type = edge
451 .edge_type_parsed
452 .clone()
453 .or_else(|| serde_json::from_str(&edge.edge_type).ok())
454 .unwrap_or(semantic_memory::GraphEdgeType::Entity {
455 relation: "unknown".to_string(),
456 });
457 (
458 edge.source.clone(),
459 edge.target.clone(),
460 parsed_type,
461 edge.weight,
462 edge.metadata.clone(),
463 )
464 })
465 .collect();
466
467 let nodes: Vec<(String, f64)> = results
468 .iter()
469 .map(|r| (r.source.result_id(), r.score))
470 .collect();
471 let factors = factors_from_edges(&raw_edges);
472 let graph = FactorGraph::new(&nodes, factors, FactorGraphConfig::default());
473 let propagated = graph.propagate();
474 let top_beliefs = propagated.top_k(top_k);
475
476 factor_graph_payload = serde_json::json!({
477 "enabled": true,
478 "top_k_beliefs": top_beliefs
479 .into_iter()
480 .map(|(item_id, belief)| serde_json::json!({
481 "item_id": item_id,
482 "belief": belief,
483 }))
484 .collect::<Vec<_>>(),
485 "iterations": propagated.iterations,
486 "converged": propagated.converged,
487 "elapsed_ms": propagated.elapsed_ms,
488 "factor_counts": {
489 "semantic": propagated.factor_counts.semantic,
490 "temporal": propagated.factor_counts.temporal,
491 "causal": propagated.factor_counts.causal,
492 "entity": propagated.factor_counts.entity,
493 "total": propagated.factor_counts.total(),
494 },
495 });
496 decoder_executed = true;
497 }
498 }
499
500 if !plan.contradictions.is_empty() {
502 use semantic_memory::decoder::{compute_correction, detect_syndromes};
503 let result_scores: Vec<(String, f64)> = results
504 .iter()
505 .map(|r| (r.source.result_id(), r.score))
506 .collect();
507 let syndromes = detect_syndromes(&result_scores, &plan.contradictions);
508 let _ = compute_correction(&syndromes, 10.0);
509 decoder_executed = true;
510 }
511 }
512
513 if plan.use_discord {
515 use semantic_memory::discord::DiscordScorer;
516 let direct_ids: Vec<String> =
517 results.iter().map(|r| r.source.result_id()).collect();
518 let existing_ids: std::collections::HashSet<String> =
519 direct_ids.iter().cloned().collect();
520 let edges_result = block_in_place(|| {
521 handle.block_on(store.list_graph_edges_for_neighborhood(
522 direct_ids.clone(),
523 2,
524 200,
525 ))
526 });
527 if let Ok(raw_edges) = edges_result {
528 let edge_refs: Vec<semantic_memory::discord::GraphEdgeRef> = raw_edges
529 .iter()
530 .map(|edge| {
531 let parsed_type = edge
532 .edge_type_parsed
533 .clone()
534 .or_else(|| serde_json::from_str(&edge.edge_type).ok())
535 .unwrap_or(semantic_memory::GraphEdgeType::Entity {
536 relation: "unknown".to_string(),
537 });
538 let type_str = match parsed_type {
539 semantic_memory::GraphEdgeType::Semantic { .. } => "semantic",
540 semantic_memory::GraphEdgeType::Temporal { .. } => "temporal",
541 semantic_memory::GraphEdgeType::Causal { .. } => "causal",
542 semantic_memory::GraphEdgeType::Entity { .. } => "entity",
543 };
544 semantic_memory::discord::GraphEdgeRef {
545 source: edge.source.clone(),
546 target: edge.target.clone(),
547 edge_type: type_str.to_string(),
548 weight: edge.weight,
549 }
550 })
551 .collect();
552 let scorer = DiscordScorer::with_defaults();
553 let discord_hits = scorer.score(&direct_ids, &edge_refs);
554 for hit in &discord_hits {
555 if !existing_ids.contains(&hit.item_id) {
556 let bare_id = hit.item_id.strip_prefix("fact:").unwrap_or(&hit.item_id);
559 let (content, namespace) = {
560 let fact_result = handle.block_on(store.get_fact(bare_id));
561 match fact_result {
562 Ok(Some(fact)) => (fact.content, fact.namespace),
563 Ok(None) => {
564 eprintln!(
565 "[discord] get_fact returned None for id={}",
566 bare_id
567 );
568 (String::new(), String::new())
569 }
570 Err(e) => {
571 eprintln!(
572 "[discord] get_fact error for id={}: {}",
573 bare_id, e
574 );
575 (String::new(), String::new())
576 }
577 }
578 };
579 discord_results_payload.push(serde_json::json!({
580 "result_id": hit.item_id,
581 "content": content,
582 "namespace": namespace,
583 "discord_score": hit.discord_score,
584 "anchor_ids": hit.anchor_ids,
585 "relationship_types": hit.relationship_types,
586 }));
587 }
588 }
589 discord_executed = true;
590 }
591 }
592
593 let grouped_results_payload: serde_json::Value = if group_by_community {
595 let seed_ids: Vec<String> = results.iter().map(|r| r.source.result_id()).collect();
596 let edges_result = block_in_place(|| {
597 handle.block_on(store.list_graph_edges_for_neighborhood(
598 seed_ids.clone(),
599 2,
600 200,
601 ))
602 });
603 let edges: Vec<(String, String)> = match edges_result {
604 Ok(raw_edges) => raw_edges
605 .iter()
606 .map(|edge| (edge.source.clone(), edge.target.clone()))
607 .collect(),
608 Err(_) => Vec::new(),
609 };
610 if !edges.is_empty() {
611 use semantic_memory::community::detect_communities;
612 let communities = detect_communities(&edges, 1.0, 42);
613 let mut member_to_comm: std::collections::HashMap<String, String> =
614 std::collections::HashMap::new();
615 for c in &communities {
616 for m in &c.members {
617 member_to_comm.insert(m.clone(), c.id.clone());
618 }
619 }
620 let mut groups: std::collections::HashMap<String, Vec<serde_json::Value>> =
621 std::collections::HashMap::new();
622 let mut ungrouped: Vec<serde_json::Value> = Vec::new();
623 for r in &json_results {
624 if let Some(rid) = r.get("result_id").and_then(|v| v.as_str()) {
625 match member_to_comm.get(rid).cloned() {
626 Some(cid) => groups.entry(cid).or_default().push(r.clone()),
627 None => ungrouped.push(r.clone()),
628 }
629 }
630 }
631 let mut map = serde_json::Map::new();
632 for (cid, items) in groups {
633 map.insert(format!("community_{cid}"), serde_json::json!(items));
634 }
635 if !ungrouped.is_empty() {
636 map.insert("ungrouped".to_string(), serde_json::json!(ungrouped));
637 }
638 serde_json::Value::Object(map)
639 } else {
640 serde_json::Value::Null
641 }
642 } else {
643 serde_json::Value::Null
644 };
645
646 let provenance = serde_json::json!({
648 "stages_fired": {
649 "bm25": results.iter().any(|r| r.bm25_rank.is_some()),
650 "vector": results.iter().any(|r| r.vector_rank.is_some()),
651 "late_interaction": true,
652 "discord": discord_executed,
653 "decoder": decoder_executed,
654 },
655 "result_count": results.len(),
656 "view": "routed",
657 "query_class": query_class,
658 "widening_occurred": false,
659 "widening_reason": null,
660 "verification_status": "verified",
661 });
662
663 (
664 "200 OK",
665 serde_json::json!({
666 "ok": true,
667 "query": query,
668 "top_k": base_top_k,
669 "results": json_results,
670 "provenance": provenance,
671 "query_class": query_class,
672 "routed": true,
673 "routing_decision": {
674 "bm25_coarse": decision.bm25_coarse,
675 "vector_medium": decision.vector_medium,
676 "rerank_fine": decision.rerank_fine,
677 "graph_expansion": decision.graph_expansion,
678 "decoder": decision.decoder,
679 "discord": decision.discord,
680 "no_retrieval": decision.no_retrieval,
681 "reasoning": decision.reasoning,
682 },
683 "decoder_planned": plan.use_decoder,
684 "decoder_executed": decoder_executed,
685 "discord_planned": plan.use_discord,
686 "discord_executed": discord_executed,
687 "discord_results": discord_results_payload,
688 "factor_graph": factor_graph_payload,
689 "grouped_results": grouped_results_payload,
690 }),
691 )
692 }
693 Err(e) => (
694 "500 Internal Server Error",
695 serde_json::json!({"ok": false, "error": format!("search error: {e}")}),
696 ),
697 }
698}
699
700#[cfg(feature = "orchestration")]
704fn handle_query_orchestrated(
705 body: &str,
706 bridge: &MemoryBridge,
707 handle: &Handle,
708) -> (&'static str, serde_json::Value) {
709 use knowledge_runtime::adapters::semantic_memory::SemanticMemoryAdapter;
710 use knowledge_runtime::config::RuntimeConfig;
711 use knowledge_runtime::KnowledgeRuntime;
712 use stack_ids::Scope;
713
714 let params: serde_json::Value = parse_body_json(body);
715 if params.is_null() {
716 return (
717 "400 Bad Request",
718 serde_json::json!({"ok": false, "error": "invalid JSON body"}),
719 );
720 }
721
722 let query = params.get("query").and_then(|v| v.as_str()).unwrap_or("");
723 let top_k = params.get("top_k").and_then(|v| v.as_u64()).unwrap_or(12) as usize;
724 let namespace = params
725 .get("namespace")
726 .and_then(|v| v.as_str())
727 .unwrap_or("general");
728 let domain = params.get("domain").and_then(|v| v.as_str());
729 let include_trace = params
730 .get("trace")
731 .and_then(|v| v.as_bool())
732 .unwrap_or(true);
733
734 if query.is_empty() {
735 return (
736 "400 Bad Request",
737 serde_json::json!({"ok": false, "error": "missing 'query' field"}),
738 );
739 }
740
741 let adapter = SemanticMemoryAdapter::new(bridge.store.clone());
744 let scope = Scope::new(namespace);
745 let scope = if let Some(d) = domain {
746 scope.with_domain(d)
747 } else {
748 scope
749 };
750 let config = RuntimeConfig {
751 default_scope: Scope::new(namespace),
752 query: knowledge_runtime::config::QueryConfig {
753 max_results_per_leg: top_k,
754 max_route_legs: 4,
755 default_limit: top_k,
756 ..Default::default()
757 },
758 entity: knowledge_runtime::config::EntityConfig::default(),
759 projection: knowledge_runtime::config::ProjectionConfig::default(),
760 strict_temporal: false,
761 strict_scope: false,
762 };
763
764 let runtime = match KnowledgeRuntime::new(config, adapter) {
765 Ok(rt) => rt,
766 Err(e) => {
767 return (
768 "500 Internal Server Error",
769 serde_json::json!({"ok": false, "error": format!("runtime init: {e}")}),
770 )
771 }
772 };
773
774 let classification = runtime.classify(query);
776
777 let route_plan = runtime.plan(query, Some(&scope));
779
780 let result =
782 block_in_place(|| handle.block_on(runtime.query_with_trace(query, Some(&scope), None)));
783
784 match result {
785 Ok((results, trace)) => {
786 let json_results: Vec<serde_json::Value> = results
787 .iter()
788 .map(|r| {
789 let ns = match &r.source {
790 semantic_memory::SearchSource::Fact { namespace, .. } => namespace.clone(),
791 semantic_memory::SearchSource::Chunk { document_title, .. } => {
792 document_title.clone()
793 }
794 _ => String::new(),
795 };
796 serde_json::json!({
797 "result_id": r.source.result_id(),
798 "content": r.content,
799 "score": r.score,
800 "cosine_similarity": r.cosine_similarity,
801 "namespace": ns,
802 "source_type": match &r.source {
803 semantic_memory::SearchSource::Fact { .. } => "fact",
804 semantic_memory::SearchSource::Chunk { .. } => "chunk",
805 semantic_memory::SearchSource::Message { .. } => "message",
806 _ => "unknown",
807 },
808 })
809 })
810 .collect();
811
812 let trace_json = if include_trace {
813 serde_json::to_value(&trace).unwrap_or(serde_json::Value::Null)
814 } else {
815 serde_json::Value::Null
816 };
817
818 let classification_json = serde_json::json!({
819 "mode": classification.mode.kind(),
820 "confidence": classification.confidence,
821 "reason": classification.reason,
822 });
823
824 let plan_json = serde_json::to_value(&route_plan).unwrap_or(serde_json::Value::Null);
825
826 (
827 "200 OK",
828 serde_json::json!({
829 "ok": true,
830 "query": query,
831 "top_k": top_k,
832 "results": json_results,
833 "classification": classification_json,
834 "route_plan": plan_json,
835 "trace": trace_json,
836 "routed": true,
837 "orchestrated": true,
838 }),
839 )
840 }
841 Err(e) => (
842 "500 Internal Server Error",
843 serde_json::json!({"ok": false, "error": format!("query error: {e}")}),
844 ),
845 }
846}
847
848fn handle_stats(bridge: &MemoryBridge, handle: &Handle) -> (&'static str, serde_json::Value) {
849 let store = &bridge.store;
850 let result = block_in_place(|| handle.block_on(store.stats()));
851 match result {
852 Ok(stats) => (
853 "200 OK",
854 serde_json::json!({
855 "ok": true,
856 "facts": stats.total_facts,
857 "documents": stats.total_documents,
858 "chunks": stats.total_chunks,
859 "db_size_mb": (stats.database_size_bytes as f64) / (1024.0 * 1024.0),
860 }),
861 ),
862 Err(e) => (
863 "500 Internal Server Error",
864 serde_json::json!({"ok": false, "error": format!("{e}")}),
865 ),
866 }
867}
868
869fn handle_rerank(body: &str) -> (&'static str, serde_json::Value) {
870 let params: serde_json::Value = parse_body_json(body);
871 if params.is_null() {
872 return (
873 "400 Bad Request",
874 serde_json::json!({"ok": false, "error": "invalid JSON body"}),
875 );
876 }
877
878 let query = params.get("query").and_then(|v| v.as_str()).unwrap_or("");
879 let model = params
880 .get("model")
881 .and_then(|v| v.as_str())
882 .unwrap_or("granite4.1:3b");
883 let results = match params.get("results").and_then(|v| v.as_array()) {
884 Some(r) => r.clone(),
885 None => {
886 return (
887 "400 Bad Request",
888 serde_json::json!({"ok": false, "error": "missing 'results' array"}),
889 )
890 }
891 };
892
893 if query.is_empty() {
894 return (
895 "400 Bad Request",
896 serde_json::json!({"ok": false, "error": "missing 'query' field"}),
897 );
898 }
899
900 let reranked = rerank_results(query, &results, model);
901 let count = reranked.len();
902 (
903 "200 OK",
904 serde_json::json!({
905 "ok": true,
906 "results": reranked,
907 "count": count,
908 }),
909 )
910}
911
912fn handle_add_fact(
913 body: &str,
914 bridge: &MemoryBridge,
915 handle: &Handle,
916) -> (&'static str, serde_json::Value) {
917 let params: serde_json::Value = parse_body_json(body);
918 if params.is_null() {
919 return (
920 "400 Bad Request",
921 serde_json::json!({"ok": false, "error": "invalid JSON body"}),
922 );
923 }
924
925 let content = params.get("content").and_then(|v| v.as_str()).unwrap_or("");
926 let namespace = params
927 .get("namespace")
928 .and_then(|v| v.as_str())
929 .unwrap_or("general");
930 let source = params.get("source").and_then(|v| v.as_str());
931
932 if content.is_empty() {
933 return (
934 "400 Bad Request",
935 serde_json::json!({"ok": false, "error": "missing 'content' field"}),
936 );
937 }
938
939 let store = &bridge.store;
940 let result =
941 block_in_place(|| handle.block_on(store.add_fact(namespace, content, source, None)));
942
943 match result {
944 Ok(fact_id) => (
945 "200 OK",
946 serde_json::json!({"ok": true, "fact_id": fact_id}),
947 ),
948 Err(e) => (
949 "500 Internal Server Error",
950 serde_json::json!({"ok": false, "error": format!("{e}")}),
951 ),
952 }
953}
954
955fn handle_add_edge(
957 body: &str,
958 bridge: &MemoryBridge,
959 handle: &Handle,
960) -> (&'static str, serde_json::Value) {
961 let params: serde_json::Value = parse_body_json(body);
962 if params.is_null() {
963 return (
964 "400 Bad Request",
965 serde_json::json!({"ok": false, "error": "invalid JSON body"}),
966 );
967 }
968
969 let source = params.get("source").and_then(|v| v.as_str()).unwrap_or("");
970 let target = params.get("target").and_then(|v| v.as_str()).unwrap_or("");
971 let edge_type_str = params
972 .get("edge_type")
973 .and_then(|v| v.as_str())
974 .unwrap_or("semantic");
975 let weight = params.get("weight").and_then(|v| v.as_f64()).unwrap_or(1.0);
976 let cosine_similarity = params.get("cosine_similarity").and_then(|v| v.as_f64());
977 let relation = params.get("relation").and_then(|v| v.as_str());
978
979 if source.is_empty() || target.is_empty() {
980 return (
981 "400 Bad Request",
982 serde_json::json!({"ok": false, "error": "missing 'source' or 'target'"}),
983 );
984 }
985
986 let edge_type = match edge_type_str {
987 "semantic" => semantic_memory::GraphEdgeType::Semantic {
988 cosine_similarity: cosine_similarity.unwrap_or(weight) as f32,
989 },
990 "temporal" => semantic_memory::GraphEdgeType::Temporal {
991 delta_secs: params
992 .get("delta_secs")
993 .and_then(|v| v.as_u64())
994 .unwrap_or(0),
995 },
996 "causal" => semantic_memory::GraphEdgeType::Causal {
997 confidence: cosine_similarity.unwrap_or(weight) as f32,
998 evidence_ids: Vec::new(),
999 },
1000 "entity" => semantic_memory::GraphEdgeType::Entity {
1001 relation: relation.unwrap_or("mentions").to_string(),
1002 },
1003 _ => semantic_memory::GraphEdgeType::Semantic {
1004 cosine_similarity: cosine_similarity.unwrap_or(weight) as f32,
1005 },
1006 };
1007
1008 let store = &bridge.store;
1009 let result = block_in_place(|| {
1010 handle.block_on(store.add_graph_edge(source, target, edge_type, weight, None))
1011 });
1012
1013 match result {
1014 Ok(edge) => (
1015 "200 OK",
1016 serde_json::json!({"ok": true, "edge_id": edge.id}),
1017 ),
1018 Err(e) => (
1019 "500 Internal Server Error",
1020 serde_json::json!({"ok": false, "error": format!("{e}")}),
1021 ),
1022 }
1023}
1024
1025fn handle_delete_fact(
1027 body: &str,
1028 bridge: &MemoryBridge,
1029 handle: &Handle,
1030) -> (&'static str, serde_json::Value) {
1031 let params: serde_json::Value = parse_body_json(body);
1032 if params.is_null() {
1033 return (
1034 "400 Bad Request",
1035 serde_json::json!({"ok": false, "error": "invalid JSON body"}),
1036 );
1037 }
1038
1039 let fact_id = params.get("fact_id").and_then(|v| v.as_str()).unwrap_or("");
1040 if fact_id.is_empty() {
1041 return (
1042 "400 Bad Request",
1043 serde_json::json!({"ok": false, "error": "missing 'fact_id'"}),
1044 );
1045 }
1046
1047 let bare_id = fact_id.strip_prefix("fact:").unwrap_or(fact_id);
1048 let store = &bridge.store;
1049 let result = block_in_place(|| handle.block_on(store.delete_fact(bare_id)));
1050
1051 match result {
1052 Ok(()) => ("200 OK", serde_json::json!({"ok": true, "deleted": true})),
1053 Err(e) => (
1054 "500 Internal Server Error",
1055 serde_json::json!({"ok": false, "error": format!("{e}")}),
1056 ),
1057 }
1058}
1059
1060fn handle_record_outcome(
1062 body: &str,
1063 bridge: &MemoryBridge,
1064 handle: &Handle,
1065) -> (&'static str, serde_json::Value) {
1066 use semantic_memory::rl_routing::{record_routing_outcome, RoutingOutcome};
1067 use semantic_memory::routing::{QueryProfile, RetrievalRouter};
1068
1069 let params: serde_json::Value = parse_body_json(body);
1070 if params.is_null() {
1071 return (
1072 "400 Bad Request",
1073 serde_json::json!({"ok": false, "error": "invalid JSON body"}),
1074 );
1075 }
1076
1077 let query = params.get("query").and_then(|v| v.as_str()).unwrap_or("");
1078 let outcome = params
1079 .get("outcome")
1080 .and_then(|v| v.as_str())
1081 .unwrap_or("neutral");
1082 let _query_class = params
1083 .get("query_class")
1084 .and_then(|v| v.as_str())
1085 .unwrap_or("A");
1086
1087 if query.is_empty() {
1088 return (
1089 "400 Bad Request",
1090 serde_json::json!({"ok": false, "error": "missing 'query' field"}),
1091 );
1092 }
1093
1094 let outcome_enum = match outcome.to_lowercase().as_str() {
1095 "good" => RoutingOutcome::Good,
1096 "bad" => RoutingOutcome::Bad,
1097 "neutral" => RoutingOutcome::Neutral,
1098 _ => {
1099 return (
1100 "400 Bad Request",
1101 serde_json::json!({"ok": false, "error": "outcome must be 'good', 'bad', or 'neutral'"}),
1102 )
1103 }
1104 };
1105
1106 let profile = QueryProfile::from_query(query);
1107 let router = RetrievalRouter::default();
1108 let decision = router.route(&profile);
1109
1110 let store = &bridge.store;
1111 let mut policy = block_in_place(|| handle.block_on(store.load_routing_policy()))
1113 .ok()
1114 .flatten()
1115 .unwrap_or_default();
1116 record_routing_outcome(&mut policy, &profile, &decision, outcome_enum);
1117 let _ = block_in_place(|| handle.block_on(store.save_routing_policy(&policy)));
1119
1120 (
1121 "200 OK",
1122 serde_json::json!({
1123 "ok": true,
1124 "recorded": true,
1125 "outcome": outcome,
1126 "routing_decision": {
1127 "bm25_coarse": decision.bm25_coarse,
1128 "vector_medium": decision.vector_medium,
1129 "rerank_fine": decision.rerank_fine,
1130 "graph_expansion": decision.graph_expansion,
1131 "decoder": decision.decoder,
1132 "discord": decision.discord,
1133 "no_retrieval": decision.no_retrieval,
1134 "reasoning": decision.reasoning,
1135 },
1136 "policy_state": {
1137 "trained_examples": policy.trained_examples,
1138 "baseline": policy.baseline,
1139 },
1140 }),
1141 )
1142}
1143
1144fn handle_verify_integrity(
1146 bridge: &MemoryBridge,
1147 handle: &Handle,
1148) -> (&'static str, serde_json::Value) {
1149 let store = &bridge.store;
1150 let result = block_in_place(|| {
1151 handle.block_on(store.verify_integrity(semantic_memory::VerifyMode::Quick))
1152 });
1153
1154 match result {
1155 Ok(report) => (
1156 "200 OK",
1157 serde_json::json!({
1158 "ok": report.ok,
1159 "integrity": report.ok,
1160 "schema_version": report.schema_version,
1161 "fact_count": report.fact_count,
1162 "chunk_count": report.chunk_count,
1163 "message_count": report.message_count,
1164 "facts_missing_embeddings": report.facts_missing_embeddings,
1165 "chunks_missing_embeddings": report.chunks_missing_embeddings,
1166 "issues": report.issues,
1167 "issue_count": report.issues.len(),
1168 "message": if report.ok { "All integrity checks passed".to_string() } else { format!("{} integrity issues found", report.issues.len()) },
1169 }),
1170 ),
1171 Err(e) => (
1172 "500 Internal Server Error",
1173 serde_json::json!({"ok": false, "integrity": false, "error": format!("verify_integrity error: {e}")}),
1174 ),
1175 }
1176}
1177
1178fn handle_discord(
1183 body: &str,
1184 bridge: &MemoryBridge,
1185 handle: &Handle,
1186) -> (&'static str, serde_json::Value) {
1187 use semantic_memory::discord::DiscordScorer;
1188
1189 let params: serde_json::Value = parse_body_json(body);
1190 if params.is_null() {
1191 return (
1192 "400 Bad Request",
1193 serde_json::json!({"ok": false, "error": "invalid JSON body"}),
1194 );
1195 }
1196
1197 let query = params.get("query").and_then(|v| v.as_str()).unwrap_or("");
1198 let top_k = params.get("top_k").and_then(|v| v.as_u64()).unwrap_or(5) as usize;
1199
1200 let direct_ids: Vec<String> = match params.get("direct_ids").and_then(|v| v.as_array()) {
1202 Some(arr) => arr
1203 .iter()
1204 .filter_map(|v| v.as_str().map(|s| s.to_string()))
1205 .collect(),
1206 None => {
1207 if query.is_empty() {
1209 return (
1210 "400 Bad Request",
1211 serde_json::json!({"ok": false, "error": "either 'direct_ids' or 'query' must be provided"}),
1212 );
1213 }
1214 let store = &bridge.store;
1215 let search_result =
1216 block_in_place(|| handle.block_on(store.search(query, Some(top_k), None, None)));
1217 match search_result {
1218 Ok(results) => results.iter().map(|r| r.source.result_id()).collect(),
1219 Err(e) => {
1220 return (
1221 "500 Internal Server Error",
1222 serde_json::json!({"ok": false, "error": format!("search error: {e}")}),
1223 )
1224 }
1225 }
1226 }
1227 };
1228
1229 if direct_ids.is_empty() {
1230 return (
1231 "200 OK",
1232 serde_json::json!({"ok": true, "discord_results": [], "count": 0, "edges_loaded": 0}),
1233 );
1234 }
1235
1236 let store = &bridge.store;
1237 let edges_result = block_in_place(|| {
1239 handle.block_on(store.list_graph_edges_for_neighborhood(direct_ids.clone(), 2, 200))
1240 });
1241
1242 let edges: Vec<semantic_memory::discord::GraphEdgeRef> = match edges_result {
1243 Ok(raw_edges) => raw_edges
1244 .iter()
1245 .map(|edge| {
1246 let parsed_type = edge
1247 .edge_type_parsed
1248 .clone()
1249 .or_else(|| serde_json::from_str(&edge.edge_type).ok())
1250 .unwrap_or(semantic_memory::GraphEdgeType::Entity {
1251 relation: "unknown".to_string(),
1252 });
1253 let type_str = match parsed_type {
1254 semantic_memory::GraphEdgeType::Semantic { .. } => "semantic",
1255 semantic_memory::GraphEdgeType::Temporal { .. } => "temporal",
1256 semantic_memory::GraphEdgeType::Causal { .. } => "causal",
1257 semantic_memory::GraphEdgeType::Entity { .. } => "entity",
1258 };
1259 semantic_memory::discord::GraphEdgeRef {
1260 source: edge.source.clone(),
1261 target: edge.target.clone(),
1262 edge_type: type_str.to_string(),
1263 weight: edge.weight,
1264 }
1265 })
1266 .collect(),
1267 Err(e) => {
1268 return (
1269 "500 Internal Server Error",
1270 serde_json::json!({"ok": false, "error": format!("failed to load graph edges: {e}")}),
1271 )
1272 }
1273 };
1274
1275 let edges_loaded = edges.len();
1276 let scorer = DiscordScorer::with_defaults();
1277 let discord_hits = scorer.score(&direct_ids, &edges);
1278
1279 let existing: std::collections::HashSet<String> = direct_ids.iter().cloned().collect();
1281 let filtered_hits: Vec<serde_json::Value> = discord_hits
1282 .iter()
1283 .filter(|hit| !existing.contains(&hit.item_id))
1284 .map(|hit| {
1285 let bare_id = hit.item_id.strip_prefix("fact:").unwrap_or(&hit.item_id);
1287 let (content, namespace) = {
1288 let fact_result = handle.block_on(store.get_fact(bare_id));
1289 match fact_result {
1290 Ok(Some(fact)) => (fact.content, fact.namespace),
1291 _ => (String::new(), String::new()),
1292 }
1293 };
1294 serde_json::json!({
1295 "result_id": hit.item_id,
1296 "content": content,
1297 "namespace": namespace,
1298 "discord_score": hit.discord_score,
1299 "anchor_ids": hit.anchor_ids,
1300 "relationship_types": hit.relationship_types,
1301 })
1302 })
1303 .collect();
1304
1305 (
1306 "200 OK",
1307 serde_json::json!({
1308 "ok": true,
1309 "discord_results": filtered_hits,
1310 "count": filtered_hits.len(),
1311 "edges_loaded": edges_loaded,
1312 "direct_ids": direct_ids,
1313 }),
1314 )
1315}
1316
1317fn handle_maintenance_check(
1325 bridge: &MemoryBridge,
1326 handle: &Handle,
1327) -> (&'static str, serde_json::Value) {
1328 let store = &bridge.store;
1329 let embeddings_dirty = block_in_place(|| handle.block_on(store.embeddings_are_dirty()));
1330 let embeddings_dirty = match embeddings_dirty {
1331 Ok(v) => v,
1332 Err(e) => {
1333 return (
1334 "500 Internal Server Error",
1335 serde_json::json!({"ok": false, "error": format!("embeddings_are_dirty error: {e}")}),
1336 )
1337 }
1338 };
1339 let integrity_result = block_in_place(|| {
1340 handle.block_on(store.verify_integrity(semantic_memory::VerifyMode::Quick))
1341 });
1342
1343 match integrity_result {
1344 Ok(report) => (
1345 "200 OK",
1346 serde_json::json!({
1347 "ok": report.ok,
1348 "embeddings_are_dirty": embeddings_dirty,
1349 "integrity": {
1350 "ok": report.ok,
1351 "schema_version": report.schema_version,
1352 "fact_count": report.fact_count,
1353 "chunk_count": report.chunk_count,
1354 "message_count": report.message_count,
1355 "facts_missing_embeddings": report.facts_missing_embeddings,
1356 "chunks_missing_embeddings": report.chunks_missing_embeddings,
1357 "issues": report.issues,
1358 "issue_count": report.issues.len(),
1359 },
1360 "message": if report.ok && !embeddings_dirty {
1361 "All checks passed".to_string()
1362 } else if report.ok && embeddings_dirty {
1363 "Integrity OK but embeddings need re-embedding".to_string()
1364 } else {
1365 format!("{} integrity issues found", report.issues.len())
1366 },
1367 }),
1368 ),
1369 Err(e) => (
1370 "500 Internal Server Error",
1371 serde_json::json!({
1372 "ok": false,
1373 "embeddings_are_dirty": embeddings_dirty,
1374 "error": format!("verify_integrity error: {e}"),
1375 }),
1376 ),
1377 }
1378}
1379
1380fn handle_maintenance_vacuum(
1382 bridge: &MemoryBridge,
1383 handle: &Handle,
1384) -> (&'static str, serde_json::Value) {
1385 let store = &bridge.store;
1386 let result = block_in_place(|| handle.block_on(store.vacuum()));
1387
1388 match result {
1389 Ok(()) => (
1390 "200 OK",
1391 serde_json::json!({"ok": true, "action": "vacuum", "message": "Database vacuumed successfully"}),
1392 ),
1393 Err(e) => (
1394 "500 Internal Server Error",
1395 serde_json::json!({"ok": false, "error": format!("vacuum error: {e}")}),
1396 ),
1397 }
1398}
1399
1400fn handle_maintenance_reembed(
1403 bridge: &MemoryBridge,
1404 handle: &Handle,
1405) -> (&'static str, serde_json::Value) {
1406 let store = &bridge.store;
1407 let result = block_in_place(|| handle.block_on(store.reembed_all()));
1408
1409 match result {
1410 Ok(count) => (
1411 "200 OK",
1412 serde_json::json!({"ok": true, "action": "reembed", "reembedded_count": count, "message": format!("Re-embedded {count} items")}),
1413 ),
1414 Err(e) => (
1415 "500 Internal Server Error",
1416 serde_json::json!({"ok": false, "error": format!("reembed_all error: {e}")}),
1417 ),
1418 }
1419}
1420
1421fn handle_maintenance_reconcile(
1424 body: &str,
1425 bridge: &MemoryBridge,
1426 handle: &Handle,
1427) -> (&'static str, serde_json::Value) {
1428 let params: serde_json::Value = parse_body_json(body);
1429 if params.is_null() {
1430 return (
1431 "400 Bad Request",
1432 serde_json::json!({"ok": false, "error": "invalid JSON body"}),
1433 );
1434 }
1435
1436 let action_str = params
1437 .get("action")
1438 .and_then(|v| v.as_str())
1439 .unwrap_or("ReportOnly");
1440
1441 let action = match action_str {
1442 "RebuildFts" => semantic_memory::ReconcileAction::RebuildFts,
1443 "ReEmbed" => semantic_memory::ReconcileAction::ReEmbed,
1444 _ => semantic_memory::ReconcileAction::ReportOnly,
1445 };
1446
1447 let store = &bridge.store;
1448 let result = block_in_place(|| handle.block_on(store.reconcile(action)));
1449
1450 match result {
1451 Ok(report) => (
1452 "200 OK",
1453 serde_json::json!({
1454 "ok": report.ok,
1455 "action": "reconcile",
1456 "reconcile_action": action_str,
1457 "integrity": {
1458 "ok": report.ok,
1459 "schema_version": report.schema_version,
1460 "fact_count": report.fact_count,
1461 "chunk_count": report.chunk_count,
1462 "message_count": report.message_count,
1463 "facts_missing_embeddings": report.facts_missing_embeddings,
1464 "chunks_missing_embeddings": report.chunks_missing_embeddings,
1465 "issues": report.issues,
1466 "issue_count": report.issues.len(),
1467 },
1468 "message": if report.ok {
1469 "Reconciliation completed, no issues found".to_string()
1470 } else {
1471 format!("Reconciliation completed with {} issues", report.issues.len())
1472 },
1473 }),
1474 ),
1475 Err(e) => (
1476 "500 Internal Server Error",
1477 serde_json::json!({"ok": false, "error": format!("reconcile error: {e}")}),
1478 ),
1479 }
1480}
1481
1482fn handle_maintenance_compact_hnsw(
1487 bridge: &MemoryBridge,
1488 handle: &Handle,
1489) -> (&'static str, serde_json::Value) {
1490 #[cfg(feature = "hnsw")]
1491 {
1492 let store = &bridge.store;
1493 let result = block_in_place(|| handle.block_on(store.compact_hnsw()));
1494
1495 return match result {
1496 Ok(()) => (
1497 "200 OK",
1498 serde_json::json!({"ok": true, "action": "compact-hnsw", "message": "HNSW index compacted successfully"}),
1499 ),
1500 Err(e) => (
1501 "500 Internal Server Error",
1502 serde_json::json!({"ok": false, "error": format!("compact_hnsw error: {e}")}),
1503 ),
1504 };
1505 }
1506
1507 #[cfg(not(feature = "hnsw"))]
1508 {
1509 let _ = bridge;
1510 let _ = handle;
1511 (
1512 "200 OK",
1513 serde_json::json!({
1514 "ok": true,
1515 "action": "compact-hnsw",
1516 "message": "HNSW compaction not applicable — usearch backend does not require compaction",
1517 "skipped": true,
1518 }),
1519 )
1520 }
1521}
1522
1523fn handle_maintenance_auto_edge(
1533 body: &str,
1534 bridge: &MemoryBridge,
1535 handle: &Handle,
1536) -> (&'static str, serde_json::Value) {
1537 let start = std::time::Instant::now();
1538
1539 let params: serde_json::Value = parse_body_json(body);
1540 let batch_size = params
1541 .get("batch_size")
1542 .and_then(|v| v.as_u64())
1543 .unwrap_or(500) as usize;
1544 let max_edges_per_fact = params
1545 .get("max_edges_per_fact")
1546 .and_then(|v| v.as_u64())
1547 .unwrap_or(20) as usize;
1548 let dry_run = params
1549 .get("dry_run")
1550 .and_then(|v| v.as_bool())
1551 .unwrap_or(false);
1552 let rebuild = params
1553 .get("rebuild")
1554 .and_then(|v| v.as_bool())
1555 .unwrap_or(false);
1556
1557 let store = &bridge.store;
1558
1559 let mut edges_invalidated: u64 = 0;
1561 if rebuild && !dry_run {
1562 let all_edges = block_in_place(|| handle.block_on(store.list_all_graph_edges()));
1563 if let Ok(edges) = all_edges {
1564 for edge in &edges {
1565 let parsed = edge.edge_type_parsed.clone().or_else(|| {
1566 serde_json::from_str::<semantic_memory::GraphEdgeType>(&edge.edge_type).ok()
1567 });
1568 if matches!(parsed, Some(semantic_memory::GraphEdgeType::Entity { .. })) {
1569 let _ = block_in_place(|| {
1570 handle.block_on(store.invalidate_graph_edge(&edge.id, "auto-edge rebuild"))
1571 });
1572 edges_invalidated += 1;
1573 }
1574 }
1575 }
1576 eprintln!("[auto-edge] rebuild: invalidated {edges_invalidated} existing entity edges");
1577 }
1578
1579 const SKIP_NAMESPACES: &[&str] =
1581 &["mixed", "chatgpt", "twitter", "tool-receipts", "agentguard"];
1582 let skip_ns: std::collections::HashSet<&str> = SKIP_NAMESPACES.iter().copied().collect();
1583
1584 const STOPWORDS: &[&str] = &[
1586 "the",
1587 "and",
1588 "for",
1589 "are",
1590 "but",
1591 "not",
1592 "you",
1593 "all",
1594 "can",
1595 "her",
1596 "was",
1597 "one",
1598 "our",
1599 "out",
1600 "has",
1601 "have",
1602 "had",
1603 "his",
1604 "how",
1605 "its",
1606 "may",
1607 "new",
1608 "now",
1609 "old",
1610 "see",
1611 "him",
1612 "way",
1613 "who",
1614 "did",
1615 "yes",
1616 "yet",
1617 "say",
1618 "she",
1619 "too",
1620 "use",
1621 "via",
1622 "any",
1623 "few",
1624 "get",
1625 "got",
1626 "let",
1627 "put",
1628 "run",
1629 "set",
1630 "try",
1631 "two",
1632 "bad",
1633 "big",
1634 "far",
1635 "off",
1636 "own",
1637 "per",
1638 "sub",
1639 "top",
1640 "end",
1641 "add",
1642 "also",
1643 "been",
1644 "from",
1645 "into",
1646 "that",
1647 "this",
1648 "with",
1649 "will",
1650 "your",
1651 "they",
1652 "them",
1653 "were",
1654 "what",
1655 "when",
1656 "where",
1657 "which",
1658 "while",
1659 "there",
1660 "their",
1661 "about",
1662 "after",
1663 "before",
1664 "between",
1665 "because",
1666 "being",
1667 "would",
1668 "could",
1669 "should",
1670 "than",
1671 "then",
1672 "these",
1673 "those",
1674 "only",
1675 "over",
1676 "under",
1677 "again",
1678 "more",
1679 "most",
1680 "some",
1681 "such",
1682 "very",
1683 "just",
1684 "like",
1685 "even",
1686 "back",
1687 "both",
1688 "down",
1689 "here",
1690 "make",
1691 "made",
1692 "each",
1693 "want",
1694 "need",
1695 "know",
1696 "same",
1697 "other",
1698 "many",
1699 "much",
1700 "last",
1701 "first",
1702 "third",
1703 "next",
1704 "best",
1705 "main",
1706 "full",
1707 "upon",
1708 "within",
1709 "without",
1710 "through",
1711 "during",
1712 "above",
1713 "below",
1714 "against",
1715 "among",
1716 "across",
1717 "behind",
1718 "beside",
1719 "beyond",
1720 "project",
1722 "work",
1723 "thing",
1724 "things",
1725 "fix",
1726 "fixed",
1727 "build",
1728 "built",
1729 "code",
1730 "data",
1731 "system",
1732 "update",
1733 "updated",
1734 "check",
1735 "checked",
1736 "test",
1737 "tested",
1738 "error",
1739 "issue",
1740 "problem",
1741 "result",
1742 "results",
1743 "status",
1744 "state",
1745 "info",
1746 "note",
1747 "notes",
1748 "list",
1749 "item",
1750 "items",
1751 "type",
1752 "types",
1753 "field",
1754 "fields",
1755 "name",
1756 "names",
1757 "value",
1758 "values",
1759 "line",
1760 "lines",
1761 "file",
1762 "files",
1763 "part",
1764 "parts",
1765 "section",
1766 "sections",
1767 "step",
1768 "steps",
1769 "task",
1770 "tasks",
1771 "goal",
1772 "goals",
1773 "plan",
1774 "plans",
1775 "done",
1776 "open",
1777 "close",
1778 "closed",
1779 "start",
1780 "started",
1781 "stop",
1782 "stopped",
1783 "change",
1784 "changed",
1785 "changing",
1786 "good",
1787 "bad",
1788 "right",
1789 "wrong",
1790 "true",
1791 "false",
1792 "yes",
1793 "no",
1794 "ok",
1795 "lots",
1796 "lot",
1797 "really",
1798 "actually",
1799 "basically",
1800 "probably",
1801 "maybe",
1802 "going",
1803 "getting",
1804 "looking",
1805 "trying",
1806 "working",
1807 "something",
1808 "anything",
1809 "everything",
1810 "nothing",
1811 "someone",
1812 "anyone",
1813 "everyone",
1814 "still",
1815 "always",
1816 "never",
1817 "sometimes",
1818 "usually",
1819 "often",
1820 "since",
1821 "until",
1822 "though",
1823 "although",
1824 "however",
1825 "therefore",
1826 "either",
1827 "neither",
1828 "both",
1829 "all",
1830 "any",
1831 "some",
1832 "none",
1833 "version",
1834 "config",
1835 "setup",
1836 "install",
1837 "installed",
1838 "running",
1839 "report",
1840 "reports",
1841 "summary",
1842 "detail",
1843 "details",
1844 "feature",
1845 "features",
1846 "function",
1847 "functions",
1848 "method",
1849 "methods",
1850 "class",
1851 "classes",
1852 "module",
1853 "modules",
1854 "crate",
1855 "crates",
1856 "package",
1857 "packages",
1858 "library",
1859 "libraries",
1860 "app",
1861 "apps",
1862 "web",
1863 "page",
1864 "pages",
1865 "site",
1866 "sites",
1867 "link",
1868 "links",
1869 "post",
1870 "posts",
1871 "reply",
1872 "replies",
1873 "comment",
1874 "comments",
1875 "read",
1876 "write",
1877 "call",
1878 "called",
1879 "calling",
1880 "return",
1881 "returns",
1882 "input",
1883 "output",
1884 "source",
1885 "target",
1886 "source_id",
1887 "target_id",
1888 "create",
1889 "created",
1890 "creating",
1891 "delete",
1892 "deleted",
1893 "removing",
1894 "find",
1895 "found",
1896 "finding",
1897 "search",
1898 "searching",
1899 "replace",
1900 "replaced",
1901 "show",
1902 "shown",
1903 "showing",
1904 "hide",
1905 "hidden",
1906 "display",
1907 "rendered",
1908 "enable",
1909 "enabled",
1910 "disable",
1911 "disabled",
1912 "allow",
1913 "allowed",
1914 "require",
1915 "required",
1916 "requiring",
1917 "support",
1918 "supported",
1919 "default",
1920 "custom",
1921 "general",
1922 "specific",
1923 "standard",
1924 "current",
1925 "latest",
1926 "previous",
1927 "old",
1928 "new",
1929 "future",
1930 "high",
1931 "low",
1932 "medium",
1933 "critical",
1934 "normal",
1935 "minor",
1936 "major",
1937 "single",
1938 "multiple",
1939 "total",
1940 "count",
1941 "number",
1942 "description",
1943 "summary",
1944 "overview",
1945 "introduction",
1946 "conclusion",
1947 "todo",
1948 "fixme",
1949 "wip",
1950 "draft",
1951 "final",
1952 "complete",
1953 "completed",
1954 "include",
1955 "includes",
1956 "included",
1957 "exclude",
1958 "excludes",
1959 "excluded",
1960 "require",
1961 "requires",
1962 "required",
1963 "optional",
1964 "important",
1965 "urgent",
1966 "priority",
1967 "blocking",
1968 "question",
1969 "answer",
1970 "response",
1971 "request",
1972 ];
1973 let stopword_set: std::collections::HashSet<&str> = STOPWORDS.iter().copied().collect();
1974
1975 fn extract_terms(
1979 text: &str,
1980 stopwords: &std::collections::HashSet<&str>,
1981 ) -> std::collections::HashSet<String> {
1982 let mut terms = std::collections::HashSet::new();
1983 for word in text.split(|c: char| !c.is_alphanumeric() && c != '_' && c != '-') {
1984 let w = word.trim_matches(|c: char| c == '-' || c == '_');
1985 if w.len() < 3 {
1986 continue;
1987 }
1988 let lower = w.to_lowercase();
1989 if stopwords.contains(lower.as_str()) {
1990 continue;
1991 }
1992 if lower.chars().all(|c| c.is_ascii_digit()) {
1994 continue;
1995 }
1996 let has_separator = w.contains('_') || w.contains('-');
2001 let has_upper = w.chars().any(|c| c.is_uppercase());
2002 let starts_upper = w.chars().next().map(|c| c.is_uppercase()).unwrap_or(false);
2003 let long_enough = lower.len() >= 5;
2004
2005 if has_separator || has_upper || starts_upper || long_enough {
2006 terms.insert(lower);
2007 }
2008 }
2009 terms
2010 }
2011
2012 let namespaces = match block_in_place(|| handle.block_on(store.list_fact_namespaces())) {
2014 Ok(ns) => ns,
2015 Err(e) => {
2016 return (
2017 "500 Internal Server Error",
2018 serde_json::json!({"ok": false, "error": format!("list_fact_namespaces error: {e}")}),
2019 )
2020 }
2021 };
2022
2023 let namespaces: Vec<String> = namespaces
2025 .into_iter()
2026 .filter(|ns| !skip_ns.contains(ns.as_str()))
2027 .collect();
2028
2029 struct FactInfo {
2031 id: String,
2032 namespace: String,
2033 terms: std::collections::HashSet<String>,
2034 }
2035
2036 let mut all_facts: Vec<FactInfo> = Vec::new();
2037 for ns in &namespaces {
2038 let mut offset = 0usize;
2039 loop {
2040 let batch = match block_in_place(|| {
2041 handle.block_on(store.list_facts(ns, batch_size, offset))
2042 }) {
2043 Ok(facts) => facts,
2044 Err(e) => {
2045 eprintln!("[auto-edge] list_facts error for ns={ns} offset={offset}: {e}");
2046 break;
2047 }
2048 };
2049 if batch.is_empty() {
2050 break;
2051 }
2052 for fact in &batch {
2053 let terms = extract_terms(&fact.content, &stopword_set);
2054 if terms.len() >= 3 {
2056 all_facts.push(FactInfo {
2057 id: fact.id.clone(),
2058 namespace: fact.namespace.clone(),
2059 terms,
2060 });
2061 }
2062 }
2063 offset += batch.len();
2064 if batch.len() < batch_size {
2065 break;
2066 }
2067 }
2068 }
2069
2070 let facts_processed = all_facts.len();
2071
2072 let existing_edges: std::collections::HashSet<(String, String)> =
2074 match block_in_place(|| handle.block_on(store.list_all_graph_edges())) {
2075 Ok(edges) => edges
2076 .iter()
2077 .filter_map(|e| {
2078 let parsed = e.edge_type_parsed.clone().or_else(|| {
2079 serde_json::from_str::<semantic_memory::GraphEdgeType>(&e.edge_type).ok()
2080 });
2081 match parsed {
2082 Some(semantic_memory::GraphEdgeType::Entity { .. }) => {
2083 Some((e.source.clone(), e.target.clone()))
2084 }
2085 _ => None,
2086 }
2087 })
2088 .collect(),
2089 Err(_) => std::collections::HashSet::new(),
2090 };
2091
2092 let mut edges_created: u64 = 0;
2095 let mut edges_skipped: u64 = 0;
2096 let mut edge_counts: std::collections::HashMap<String, u64> = std::collections::HashMap::new();
2097
2098 for i in 0..all_facts.len() {
2099 let fact_i = &all_facts[i];
2100 let src_id = format!("fact:{}", fact_i.id);
2101
2102 if *edge_counts.get(&src_id).unwrap_or(&0) >= max_edges_per_fact as u64 {
2104 continue;
2105 }
2106
2107 for j in (i + 1)..all_facts.len() {
2108 let fact_j = &all_facts[j];
2109
2110 if fact_i.namespace == fact_j.namespace {
2112 continue;
2113 }
2114
2115 let shared: usize = fact_i.terms.intersection(&fact_j.terms).count();
2117 if shared < 3 {
2118 continue;
2119 }
2120
2121 let tgt_id = format!("fact:{}", fact_j.id);
2122
2123 let (edge_source, edge_target) = if src_id <= tgt_id {
2125 (src_id.clone(), tgt_id.clone())
2126 } else {
2127 (tgt_id.clone(), src_id.clone())
2128 };
2129
2130 if existing_edges.contains(&(edge_source.clone(), edge_target.clone())) {
2132 edges_skipped += 1;
2133 continue;
2134 }
2135
2136 let count_src = *edge_counts.get(&src_id).unwrap_or(&0);
2138 let count_tgt = *edge_counts.get(&tgt_id).unwrap_or(&0);
2139 if count_src >= max_edges_per_fact as u64 || count_tgt >= max_edges_per_fact as u64 {
2140 continue;
2141 }
2142
2143 if dry_run {
2144 edges_created += 1;
2145 continue;
2146 }
2147
2148 let relation = format!("shared_terms:{}", shared);
2149 let edge_type = semantic_memory::GraphEdgeType::Entity {
2150 relation: relation.clone(),
2151 };
2152 let weight = (shared.min(10) as f64) / 10.0;
2154
2155 let result = block_in_place(|| {
2156 handle.block_on(store.add_graph_edge(
2157 &edge_source,
2158 &edge_target,
2159 edge_type,
2160 weight,
2161 None,
2162 ))
2163 });
2164
2165 match result {
2166 Ok(_) => {
2167 edges_created += 1;
2168 *edge_counts.entry(src_id.clone()).or_insert(0) += 1;
2169 *edge_counts.entry(tgt_id.clone()).or_insert(0) += 1;
2170 }
2171 Err(e) => {
2172 eprintln!(
2173 "[auto-edge] add_graph_edge error: {e} for {edge_source} -> {edge_target}"
2174 );
2175 }
2176 }
2177
2178 if *edge_counts.get(&src_id).unwrap_or(&0) >= max_edges_per_fact as u64 {
2180 break;
2181 }
2182 }
2183 }
2184
2185 let elapsed = start.elapsed();
2186
2187 (
2188 "200 OK",
2189 serde_json::json!({
2190 "ok": true,
2191 "action": "auto-edge",
2192 "dry_run": dry_run,
2193 "facts_processed": facts_processed,
2194 "namespaces_scanned": namespaces.len(),
2195 "edges_created": edges_created,
2196 "edges_skipped": edges_skipped,
2197 "max_edges_per_fact": max_edges_per_fact,
2198 "time_elapsed_ms": elapsed.as_millis(),
2199 "time_elapsed_secs": (elapsed.as_millis() as f64) / 1000.0,
2200 }),
2201 )
2202}