diff --git a/server/src/api/events.rs b/server/src/api/events.rs index 583b539..fe218d0 100644 --- a/server/src/api/events.rs +++ b/server/src/api/events.rs @@ -2,8 +2,11 @@ use crate::AppState; use crate::state::GenericEvent; use axum::extract::{Query, State}; use axum::response::IntoResponse; +use axum::response::sse::{Event, KeepAlive, Sse}; use std::collections::HashMap; use std::sync::Arc; +use tokio_stream::StreamExt; +use tokio_stream::wrappers::BroadcastStream; pub async fn wait_for_event_handler( State(state): State>, @@ -49,6 +52,36 @@ pub async fn post_event_handler( axum::Json(serde_json::json!({"status": "ok"})) } +/// Real-time Server-Sent Events (SSE) stream for agent execution events and observability (ADR-0110) +pub async fn sse_events_handler( + State(state): State>, + Query(params): Query>, +) -> impl IntoResponse { + let topic_filter = params.get("topic").cloned(); + let session_filter = params.get("session_id").cloned(); + let rx = state.handler.state.event_bus_tx.subscribe(); + + let stream = BroadcastStream::new(rx).filter_map(move |msg| match msg { + Ok(event) => { + let topic_matches = topic_filter.as_ref().is_none_or(|t| t == &event.topic); + let session_matches = session_filter + .as_ref() + .is_none_or(|s| Some(s) == event.session_id.as_ref()); + if topic_matches && session_matches { + let json_data = serde_json::to_string(&event).unwrap_or_default(); + Some(Ok::<_, std::convert::Infallible>( + Event::default().event(event.topic).data(json_data), + )) + } else { + None + } + } + Err(_) => None, + }); + + Sse::new(stream).keep_alive(KeepAlive::default()) +} + #[cfg(test)] mod tests { use super::*; diff --git a/server/src/api/setup.rs b/server/src/api/setup.rs index 4f9ff93..9e645f2 100644 --- a/server/src/api/setup.rs +++ b/server/src/api/setup.rs @@ -68,6 +68,8 @@ pub fn create_router(app_state: Arc) -> Router { .route("/terminal/telemetry", post(crate::api::telemetry::terminal_telemetry_handler)) .route("/events/wait", get(crate::api::events::wait_for_event_handler)) .route("/events", post(crate::api::events::post_event_handler)) + .route("/api/events/stream", get(crate::api::events::sse_events_handler)) + .route("/events/stream", get(crate::api::events::sse_events_handler)) .route( "/api/activity", get({ diff --git a/server/src/handlers/graph.rs b/server/src/handlers/graph.rs index 25785f0..68e5077 100644 --- a/server/src/handlers/graph.rs +++ b/server/src/handlers/graph.rs @@ -1044,9 +1044,11 @@ impl McpTool for SweepGraphHealthHandler { async fn execute(&self, args: Value, state: Arc) -> crate::error::Result { let req: SweepGraphHealthTool = serde_json::from_value(args).map_err(|e| e.to_string())?; let auto_prune = req.auto_prune_orphans.unwrap_or(false); + let auto_prune_stale = req.auto_prune_stale_files.unwrap_or(false); let mut orphans = Vec::new(); let mut duplicates = Vec::new(); + let mut stale_entities = Vec::new(); state.modify_graph(|g| { // 1. Identify Orphans @@ -1062,12 +1064,31 @@ impl McpTool for SweepGraphHealthHandler { } } + // ADR-0111: Automated Stale Symbol Pruning & Graph Tombstoning + for (name, entity) in g.entities.iter() { + if entity.entity_type == "File" || name.contains(".rs") || name.contains(".ts") || name.contains(".js") || name.contains(".py") { + let check_path = entity.file_path.as_deref().unwrap_or(name.as_str()); + // Strip any symbol qualifiers like File::symbol + let base_file = check_path.split("::").next().unwrap_or(check_path); + if (base_file.contains('/') || base_file.contains('\\') || base_file.ends_with(".rs") || base_file.ends_with(".ts")) && !std::path::Path::new(base_file).exists() { + stale_entities.push(name.clone()); + } + } + } + if auto_prune { for orphan in &orphans { g.entities.remove(orphan); } } + if auto_prune_stale { + for stale in &stale_entities { + g.entities.remove(stale); + g.relations.retain(|r| &r.from != stale && &r.to != stale); + } + } + // 2. Compute similarity pairs for duplicate detection using pre-computed lowercase names let names: Vec<_> = g.entities.keys().cloned().collect(); let lower_names: Vec = names.iter().map(|n| n.to_lowercase()).collect(); @@ -1092,8 +1113,10 @@ impl McpTool for SweepGraphHealthHandler { let report = serde_json::json!({ "orphaned_entities": orphans, "orphans_pruned": auto_prune, + "stale_entities": stale_entities, + "stale_pruned": auto_prune_stale, "potential_duplicates": duplicates, - "health_score": if orphans.is_empty() && duplicates.is_empty() { "100%" } else { "Needs Maintenance" } + "health_score": if orphans.is_empty() && duplicates.is_empty() && stale_entities.is_empty() { "100%" } else { "Needs Maintenance" } }); Ok(serde_json::to_string_pretty(&report)?) diff --git a/server/src/indexer.rs b/server/src/indexer.rs index 1d5cd41..58e5c69 100644 --- a/server/src/indexer.rs +++ b/server/src/indexer.rs @@ -143,7 +143,12 @@ pub async fn start_background_indexer(state: Arc) { }); } -fn extract_chunks(node: Node, code: &str, chunks: &mut Vec<(String, String, String)>, ext: &str) { +pub fn extract_chunks( + node: Node, + code: &str, + chunks: &mut Vec<(String, String, String)>, + ext: &str, +) { extract_chunks_with_parent(node, code, chunks, ext, None, 0); } diff --git a/server/src/tools.rs b/server/src/tools.rs index 436a9fb..ced298c 100644 --- a/server/src/tools.rs +++ b/server/src/tools.rs @@ -151,7 +151,6 @@ pub struct VisualizeGraphTool { pub namespace: Option, } - /// Condense or summarize an entity's observations to reduce size. #[derive(Debug, Deserialize, Serialize, JsonSchema)] pub struct CondenseEntityTool { @@ -327,13 +326,15 @@ pub struct VerifyAcceptanceCriteriaTool { pub proof: String, } -/// Audit the knowledge graph to detect orphaned entities, compute name similarity for potential duplicate merges, and optionally auto-prune orphans. +/// Audit the knowledge graph to detect orphaned entities, compute name similarity for potential duplicate merges, and optionally auto-prune orphans and stale file/symbol tombstones (ADR-0111). #[derive(Debug, Deserialize, Serialize, JsonSchema)] pub struct SweepGraphHealthTool { /// Optional flag to automatically prune orphaned nodes with 0 relations. Defaults to false. pub auto_prune_orphans: Option, /// Minimum string similarity threshold (0.0 to 1.0) to report duplicate entity pairs. Defaults to 0.8. pub similarity_threshold: Option, + /// Optional flag to prune stale entities whose file paths or symbols no longer exist on disk (ADR-0111). Defaults to false. + pub auto_prune_stale_files: Option, } /// Trace the causal provenance and historical lineage linking a task, ADR, git commit, code change, or error fix. diff --git a/server/src/watcher.rs b/server/src/watcher.rs index a49b9b7..3ad8ef7 100644 --- a/server/src/watcher.rs +++ b/server/src/watcher.rs @@ -89,11 +89,86 @@ fn should_review(path: &Path) -> bool { } async fn trigger_autonomous_review(path: &Path, state: Arc) { - info!("Triggering autonomous review for {:?}", path); + info!("Triggering autonomous review & incremental AST index for {:?}", path); state.broadcast_activity("AUTONOMOUS", &format!("Modified: {:?}", path.file_name().unwrap_or_default())); - info!("Autonomous review complete for {:?}", path); -} + // ADR-0109: Incremental Background AST Indexing & Differential Graph Updates + let ext = path.extension().and_then(|e| e.to_str()).unwrap_or(""); + if matches!(ext, "rs" | "ts" | "js" | "py" | "go" | "java" | "c" | "cpp") { + if let Ok(content) = std::fs::read_to_string(path) { + let language = match ext { + "rs" => Some(tree_sitter_rust::LANGUAGE), + "ts" | "js" => Some(tree_sitter_typescript::LANGUAGE_TYPESCRIPT), + "py" => Some(tree_sitter_python::LANGUAGE), + "java" => Some(tree_sitter_java::LANGUAGE), + "c" => Some(tree_sitter_c::LANGUAGE), + "cpp" => Some(tree_sitter_cpp::LANGUAGE), + "go" => Some(tree_sitter_go::LANGUAGE), + _ => None, + }; + + if let Some(lang) = language { + let mut parser = tree_sitter::Parser::new(); + if parser.set_language(&lang.into()).is_ok() { + if let Some(tree) = parser.parse(&content, None) { + let mut chunks = Vec::new(); + crate::indexer::extract_chunks(tree.root_node(), &content, &mut chunks, ext); + let file_str = path.to_string_lossy().to_string(); + let now = std::time::SystemTime::now().duration_since(std::time::UNIX_EPOCH).unwrap_or_default().as_secs(); + + state.modify_graph(|g| { + // Ensure File entity exists + g.entities.entry(file_str.clone()).or_insert_with(|| { + crate::models::Entity { + name: file_str.clone(), + entity_type: "File".to_string(), + namespace: "global".to_string(), + file_path: Some(file_str.clone()), + created_at: Some(now), + updated_at: Some(now), + ..Default::default() + } + }); + + for (chunk_name, chunk_code, chunk_desc) in chunks { + let symbol_name = format!("{}::{}", file_str, chunk_name); + let symbol_type = if chunk_desc.contains("struct") { + "DataStructure".to_string() + } else { + "McpTool".to_string() + }; + + g.entities.insert(symbol_name.clone(), crate::models::Entity { + name: symbol_name.clone(), + entity_type: symbol_type, + observations: vec![format!("AST definition: {} chars", chunk_code.len())], + namespace: "global".to_string(), + file_path: Some(file_str.clone()), + created_at: Some(now), + updated_at: Some(now), + ..Default::default() + }); + + let rel = crate::models::Relation { + from: file_str.clone(), + to: symbol_name, + relation_type: "declares".to_string(), + namespace: "global".to_string(), + ..Default::default() + }; + if !g.relations.contains(&rel) { + g.relations.push(rel); + } + } + }); + } + } + } + } + } + + info!("Autonomous review & incremental AST index complete for {:?}", path); +} #[cfg(test)] mod tests { use super::*;