Files
mcp-memory/server/src/watcher.rs
T

224 lines
9.7 KiB
Rust

use crate::state::MemoryState;
use notify::{Config, Event, RecommendedWatcher, RecursiveMode, Watcher};
use std::path::Path;
use std::sync::Arc;
use tracing::{error, info};
pub fn spawn_watcher(state: Arc<MemoryState>) {
let watch_path = std::env::current_dir().unwrap_or_else(|_| std::path::PathBuf::from("."));
info!("Spawning proactive daemon watcher on {:?}", watch_path);
tokio::spawn(async move {
let (tx, mut rx) = tokio::sync::mpsc::channel::<notify::Result<Event>>(500);
let mut watcher = match RecommendedWatcher::new(
move |res| {
if let Err(e) = tx.try_send(res) {
tracing::warn!("Watcher event dropped due to channel backpressure: {}", e);
}
},
Config::default(),
) {
Ok(w) => w,
Err(e) => {
error!("Failed to create watcher: {}", e);
return;
}
};
if let Err(e) = watcher.watch(&watch_path, RecursiveMode::Recursive) {
error!("Failed to watch path: {}", e);
return;
}
let mut last_processed: std::collections::HashMap<std::path::PathBuf, std::time::Instant> =
std::collections::HashMap::new();
loop {
tokio::select! {
_ = state.shutdown_notify.notified() => {
info!("File watcher received shutdown notification; terminating cleanly.");
break;
}
res = rx.recv() => {
let Some(res) = res else {
break;
};
match res {
Ok(event) => {
if event.kind.is_modify() {
let now = std::time::Instant::now();
if last_processed.len() > 1000 {
let ten_mins = std::time::Duration::from_secs(600);
last_processed.retain(|_, last_time| now.duration_since(*last_time) < ten_mins);
}
for path in event.paths {
if should_review(&path) {
// 250ms debouncing window per file path
if let Some(last) = last_processed.get(&path)
&& now.duration_since(*last) < std::time::Duration::from_millis(250)
{
continue;
}
last_processed.insert(path.clone(), now);
info!("Proactive Daemon Hooks: File modified: {:?}", path);
trigger_autonomous_review(&path, Arc::clone(&state)).await;
}
}
}
}
Err(e) => error!("Watch error: {}", e),
}
}
}
}
});
}
fn should_review(path: &Path) -> bool {
let path_str = path.to_string_lossy();
if path_str.contains(".git") || path_str.contains("target") || path_str.contains(".gemini") || path_str.contains("node_modules") {
return false;
}
path.extension()
.and_then(|ext| ext.to_str())
.is_some_and(|ext| matches!(ext, "rs" | "md" | "toml" | "lua"))
}
async fn trigger_autonomous_review(path: &Path, state: Arc<MemoryState>) {
info!("Triggering autonomous review & incremental AST index for {:?}", path);
state.broadcast_activity("AUTONOMOUS", &format!("Modified: {:?}", path.file_name().unwrap_or_default()));
// 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();
let chunks_count = chunks.len();
let file_str_clone = file_str.clone();
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);
}
}
});
let _ = state.event_bus_tx.send(crate::state::GenericEvent {
topic: "ast:symbol_updated".to_string(),
session_id: None,
payload: serde_json::json!({
"file": file_str_clone,
"extension": ext,
"symbols_count": chunks_count,
}),
});
let _ = state.event_bus_tx.send(crate::state::GenericEvent {
topic: "resource:updated".to_string(),
session_id: None,
payload: serde_json::json!({
"uri": "memory://graph"
}),
});
}
}
}
}
}
info!("Autonomous review & incremental AST index complete for {:?}", path);
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn test_should_review() {
assert!(should_review(Path::new("src/lib.rs")));
assert!(should_review(Path::new("README.md")));
assert!(should_review(Path::new("Cargo.toml")));
assert!(should_review(Path::new("init.lua")));
assert!(!should_review(Path::new("target/debug/app.exe")));
assert!(!should_review(Path::new(".git/HEAD")));
assert!(!should_review(Path::new("data.json")));
assert!(!should_review(Path::new("image.png")));
}
#[tokio::test]
async fn test_spawn_watcher_lifecycle() {
let temp_dir = tempfile::tempdir().unwrap();
let state = Arc::new(MemoryState::new(temp_dir.path().to_str().unwrap()));
spawn_watcher(state);
}
#[tokio::test]
async fn test_spawn_watcher_invalid_path() {
let state = Arc::new(MemoryState::new("/nonexistent/path"));
spawn_watcher(state);
}
}