Skip to main content

semantic_memory_mcp/
http_server.rs

1//! HTTP search server for semantic-memory-mcp.
2//!
3//! A minimal HTTP server that exposes the most-used semantic-memory
4//! operations over a local TCP port. Runs alongside the stdio MCP
5//! transport so the same warm process serves both MCP clients and
6//! HTTP clients (hooks, benchmarks, scripts).
7//!
8//! Endpoints:
9//!   POST /search   {"query": "...", "top_k": 10} -> search results
10//!   POST /stats    {} -> DB stats
11//!   POST /add      {"content": "...", "namespace": "..."} -> fact_id
12//!   GET  /health   -> {"ok": true}
13
14use 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
21/// Parse an HTTP request body as a JSON Value, using llm-output-parser
22/// when the llm-parser feature is enabled (for robust handling of
23/// potentially malformed agent-submitted input), falling back to
24/// serde_json::from_str.
25fn 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
37/// Call Ollama to rate each result's relevance to the query (1-5) and sort descending.
38/// Returns a new vec with a `rerank_score` field added to each result object.
39fn 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    // Fetch top_k * 2 candidates when reranking so the LLM has a richer pool to sort.
234    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
304/// Handle /search-routed: routing-aware search with full pipeline.
305///
306/// Uses the library's routing system to profile the query and decide which
307/// retrieval stages to activate. For class C/D queries with contradictions,
308/// runs factor graph belief propagation and decoder syndrome detection.
309/// When discord is enabled, runs second-order retrieval via graph neighborhood.
310/// Optionally groups results by community.
311fn 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    // Use the routing system to profile the query
353    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    // Class D (SYNTHESIS): retrieve more candidates to support comprehensive answers
364    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    // Class C (CONTRADICTION): use ExactSearch context for higher-fidelity results
376    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            // Factor graph belief propagation for class C/D with contradictions
430            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                // Decoder syndrome detection for contradictions
501                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            // Discord second-order retrieval
514            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                            // Fetch the fact's content directly from the DB.
557                            // get_fact expects a bare UUID (without "fact:" prefix).
558                            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            // Community grouping (opt-in)
594            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            // Query provenance: declare which retrieval stages contributed
647            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/// Handle /query-orchestrated: knowledge-runtime powered query with intent
701/// classification, multi-leg route planning, provenance-aware merge, and
702/// degradation reporting. Available when `orchestration` feature is enabled.
703#[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    // Construct KnowledgeRuntime on-demand from the bridge's store.
742    // MemoryStore is Clone (Arc internals), so this is cheap.
743    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    // Step 1: Classify
775    let classification = runtime.classify(query);
776
777    // Step 2: Plan
778    let route_plan = runtime.plan(query, Some(&scope));
779
780    // Step 3: Execute full pipeline
781    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
955/// Handle /add-edge: add a graph edge between two facts.
956fn 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
1025/// Handle /delete-fact: hard-delete a fact by ID.
1026fn 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
1060/// Handle /record-outcome: record a search outcome for RL routing feedback.
1061fn 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    // Load persisted policy (or default if none saved yet)
1112    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    // Save updated policy
1118    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
1144/// Handle GET /verify-integrity: check DB integrity using real library checks.
1145fn 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
1178/// Handle POST /discord: second-order retrieval via graph neighborhood.
1179///
1180/// Accepts {"query": "...", "top_k": 5, "direct_ids": ["fact:uuid1", ...]}.
1181/// If direct_ids not provided, runs a search first to get top_k results.
1182fn 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    // Get direct_ids from params, or run a search to get them
1201    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            // Need a query to search
1208            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    // Load graph edges for the neighborhood
1238    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    // Filter out items already in direct_ids
1280    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            // Fetch the fact's content directly from the DB.
1286            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
1317// ---------------------------------------------------------------------------
1318// Maintenance endpoints — for hooks and cron jobs to trigger auto-management
1319// without needing MCP tools to be visible (they're hidden in lean profile).
1320// ---------------------------------------------------------------------------
1321
1322/// Handle POST /maintenance/check: returns embeddings_are_dirty + verify_integrity(Quick)
1323/// in one call. This is the "health check" for auto-management.
1324fn 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
1380/// Handle POST /maintenance/vacuum: calls store.vacuum(). Returns ok.
1381fn 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
1400/// Handle POST /maintenance/reembed: calls store.reembed_all(). Returns count.
1401/// This is expensive so the handler just calls it and returns the count.
1402fn 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
1421/// Handle POST /maintenance/reconcile: accepts {"action": "ReportOnly"|"RebuildFts"|"ReEmbed"}.
1422/// Calls store.reconcile(action). Returns the IntegrityReport.
1423fn 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
1482/// Handle POST /maintenance/compact-hnsw: calls store.compact_hnsw(). Returns ok.
1483///
1484/// Only available when the `hnsw` feature is enabled. The default backend is
1485/// usearch, so this endpoint returns a not-applicable response without the feature.
1486fn 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
1523/// Handle POST /maintenance/auto-edge: automatically create entity edges
1524/// between facts in DIFFERENT namespaces that share 2+ meaningful terms.
1525///
1526/// Accepts optional JSON body:
1527///   - batch_size (default 500): facts per page when iterating
1528///   - max_edges_per_fact (default 20): cap edges per source fact to prevent explosion
1529///   - dry_run (default false): if true, report what would be created without creating edges
1530///
1531/// Returns a report with facts_processed, edges_created, edges_skipped, time_elapsed.
1532fn 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    // If rebuild mode, invalidate ALL existing entity edges first.
1560    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    // Namespaces to SKIP entirely — social media noise, not technical knowledge.
1580    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    // Large stopword list — common English words that should never create edges.
1585    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        // Common technical/process words that create false connections
1721        "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    /// Extract HIGH-QUALITY terms from text: only proper nouns, camelCase/snake_case
1976    /// identifiers, and words 5+ chars that aren't stopwords. This filters out
1977    /// common English words that create false connections.
1978    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            // Skip pure numbers
1993            if lower.chars().all(|c| c.is_ascii_digit()) {
1994                continue;
1995            }
1996            // Include if:
1997            // - camelCase or snake_case (has internal uppercase or _ or -)
1998            // - all lowercase and 5+ chars (filters short common words)
1999            // - starts with uppercase (proper noun)
2000            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    // 1. Get all namespaces that contain facts
2013    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    // Filter out skip namespaces
2024    let namespaces: Vec<String> = namespaces
2025        .into_iter()
2026        .filter(|ns| !skip_ns.contains(ns.as_str()))
2027        .collect();
2028
2029    // 2. Page through all facts in each namespace, extract terms, store (id, namespace, terms)
2030    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                // Skip facts with fewer than 3 quality terms — they can't form meaningful edges
2055                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    // 3. Load existing edges to skip pairs that already have edges.
2073    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    // 4. For each pair of facts in DIFFERENT namespaces sharing 3+ terms, create entity edge.
2093    // Higher threshold (3 instead of 2) + quality term extraction = much fewer, better edges.
2094    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        // Check if this fact has already hit its edge cap
2103        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            // Only create edges between DIFFERENT namespaces
2111            if fact_i.namespace == fact_j.namespace {
2112                continue;
2113            }
2114
2115            // Check shared terms (need 3+ for high quality)
2116            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            // Determine direction (smaller id first for consistency)
2124            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            // Skip if edge already exists
2131            if existing_edges.contains(&(edge_source.clone(), edge_target.clone())) {
2132                edges_skipped += 1;
2133                continue;
2134            }
2135
2136            // Check edge cap for both facts
2137            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            // Weight based on overlap quality: 3 terms = 0.3, up to 1.0 for 10+ terms
2153            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            // Re-check cap after creating edge
2179            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}