use std::collections::{HashMap, VecDeque};
use std::path::{Path, PathBuf};
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::mpsc::{self, RecvTimeoutError};
use std::sync::Arc;
use std::time::{Duration, Instant};
use notify_debouncer_full::notify::{EventKind, RecursiveMode};
use notify_debouncer_full::{new_debouncer, DebounceEventResult, DebouncedEvent};
use signal_hook::consts::{SIGINT, SIGTERM};
use signal_hook::flag;
use tracing::{info, warn};
use tree_sitter::{InputEdit, Node, Point, Tree};
use crate::daemon::error::DaemonError;
use crate::daemon::event::{DaemonEvent, EventObserver};
use crate::discover::is_code_file;
use crate::model::Language;
use crate::parse::error::ParseError;
use crate::parse::ParserFactory;
pub const DEFAULT_DEBOUNCE_MS: u64 = 2000;
const TICK_INTERVAL_MS: u64 = 500;
const INITIAL_DEBOUNCE_WINDOW_MS: u64 = 200;
const MIN_DEBOUNCE_WINDOW_MS: u64 = 100;
const MAX_DEBOUNCE_WINDOW_MS: u64 = 500;
const EVENT_INTERVALS_CAPACITY: usize = 10;
const EMA_OLD_WEIGHT: f64 = 0.7;
const EMA_NEW_WEIGHT: f64 = 0.3;
const EMA_REFERENCE_INTERVAL_SECS: f64 = 5.0;
const CACHE_INVALIDATION_SIZE_RATIO: f64 = 0.5;
const CACHE_INVALIDATION_LINE_RATIO: f64 = 0.3;
const MAX_TREE_CACHE_ENTRIES: usize = 100;
const MAX_INCREMENTAL_PARSE_SIZE: usize = 1024 * 1024;
struct TreeCache {
entries: HashMap<String, (Tree, String)>,
}
impl TreeCache {
fn new() -> Self {
Self {
entries: HashMap::new(),
}
}
fn len(&self) -> usize {
self.entries.len()
}
fn clear(&mut self) {
self.entries.clear();
}
fn parse_incremental(&mut self, path: &str, source: &str) -> Result<usize, ParseError> {
let lang = infer_language_from_path(path)
.ok_or_else(|| ParseError::UnsupportedLanguage(path.to_string()))?;
if !Language::compiled().contains(&lang) {
return Err(ParseError::UnsupportedLanguage(format!(
"{lang} (grammar not compiled in; enable lang-{lang} feature)"
)));
}
let new_byte_len = source.len();
let new_line_count = source.lines().count();
let use_incremental = self
.entries
.get(path)
.map(|(_old_tree, old_source)| {
let old_byte_len = old_source.len();
let old_line_count = old_source.lines().count();
let size_ratio = if old_byte_len == 0 {
1.0
} else {
((new_byte_len as f64) - (old_byte_len as f64)).abs() / (old_byte_len as f64)
};
let line_ratio = if old_line_count == 0 {
1.0
} else {
((new_line_count as f64) - (old_line_count as f64)).abs()
/ (old_line_count as f64)
};
new_byte_len <= MAX_INCREMENTAL_PARSE_SIZE
&& old_byte_len <= MAX_INCREMENTAL_PARSE_SIZE
&& size_ratio <= CACHE_INVALIDATION_SIZE_RATIO
&& line_ratio <= CACHE_INVALIDATION_LINE_RATIO
})
.unwrap_or(false);
let mut parser = ParserFactory::create_parser(lang)?;
let tree = if use_incremental {
let (old_tree, old_source) = self.entries.get(path).expect("checked above");
let old_tree_clone = old_tree.clone();
let old_source_clone = old_source.clone();
let edit = compute_input_edit(&old_source_clone, source);
let mut edited_old_tree = old_tree_clone;
edited_old_tree.edit(&edit);
parser
.parse(source, Some(&edited_old_tree))
.ok_or_else(|| ParseError::ParseFailed {
file_path: path.to_string(),
})?
} else {
parser
.parse(source, None)
.ok_or_else(|| ParseError::ParseFailed {
file_path: path.to_string(),
})?
};
let node_count = count_tree_nodes(&tree.root_node());
if !self.entries.contains_key(path) && self.entries.len() >= MAX_TREE_CACHE_ENTRIES {
self.entries.clear();
}
self.entries
.insert(path.to_string(), (tree, source.to_string()));
Ok(node_count)
}
}
impl Default for TreeCache {
fn default() -> Self {
Self::new()
}
}
pub struct Daemon {
watch_path: PathBuf,
project_name: String,
debounce_ms: u64,
db_path: PathBuf,
observers: Vec<Box<dyn EventObserver + Send>>,
stop: Arc<AtomicBool>,
event_intervals: VecDeque<Duration>,
debounce_window: Duration,
last_event_at: Option<Instant>,
tree_cache: TreeCache,
}
impl Daemon {
#[must_use]
pub fn new(
watch_path: impl AsRef<Path>,
project_name: impl Into<String>,
debounce_ms: u64,
db_path: impl AsRef<Path>,
) -> Self {
Self {
watch_path: watch_path.as_ref().to_path_buf(),
project_name: project_name.into(),
debounce_ms,
db_path: db_path.as_ref().to_path_buf(),
observers: Vec::new(),
stop: Arc::new(AtomicBool::new(false)),
event_intervals: VecDeque::with_capacity(EVENT_INTERVALS_CAPACITY),
debounce_window: Duration::from_millis(INITIAL_DEBOUNCE_WINDOW_MS),
last_event_at: None,
tree_cache: TreeCache::new(),
}
}
pub fn add_observer(&mut self, observer: Box<dyn EventObserver + Send>) {
self.observers.push(observer);
}
#[must_use]
pub fn stop_handle(&self) -> Arc<AtomicBool> {
Arc::clone(&self.stop)
}
#[must_use]
pub fn debounce_ms(&self) -> u64 {
self.debounce_ms
}
#[must_use]
pub fn watch_path(&self) -> &Path {
&self.watch_path
}
#[must_use]
pub fn project_name(&self) -> &str {
&self.project_name
}
#[must_use]
pub fn db_path(&self) -> &Path {
&self.db_path
}
#[must_use]
pub fn debounce_window(&self) -> Duration {
self.debounce_window
}
#[must_use]
pub fn event_intervals(&self) -> Vec<Duration> {
self.event_intervals.iter().copied().collect()
}
pub(crate) fn update_adaptive_debounce(&mut self, now: Instant) {
if let Some(prev) = self.last_event_at {
let interval = now.saturating_duration_since(prev);
self.event_intervals.push_back(interval);
if self.event_intervals.len() > EVENT_INTERVALS_CAPACITY {
self.event_intervals.pop_front();
}
let old_secs = self.debounce_window.as_secs_f64();
let interval_secs = interval.as_secs_f64();
let inverted_secs = (EMA_REFERENCE_INTERVAL_SECS - interval_secs).max(0.0);
let new_secs = old_secs * EMA_OLD_WEIGHT + inverted_secs * EMA_NEW_WEIGHT;
let new_window = Duration::from_secs_f64(new_secs);
let min = Duration::from_millis(MIN_DEBOUNCE_WINDOW_MS);
let max = Duration::from_millis(MAX_DEBOUNCE_WINDOW_MS);
self.debounce_window = new_window.max(min).min(max);
}
self.last_event_at = Some(now);
}
pub fn parse_file_incremental(
&mut self,
path: &str,
source: &str,
) -> Result<usize, ParseError> {
self.tree_cache.parse_incremental(path, source)
}
#[must_use]
pub fn tree_cache_len(&self) -> usize {
self.tree_cache.len()
}
pub fn clear_tree_cache(&mut self) {
self.tree_cache.clear();
}
fn register_signal_handlers(&self) -> Result<(), DaemonError> {
flag::register(SIGTERM, Arc::clone(&self.stop))
.map_err(|e| DaemonError::Signal(e.to_string()))?;
flag::register(SIGINT, Arc::clone(&self.stop))
.map_err(|e| DaemonError::Signal(e.to_string()))?;
info!(signals = "SIGTERM,SIGINT", "信号处理器已注册");
Ok(())
}
pub fn run(&mut self) -> Result<(), DaemonError> {
self.register_signal_handlers()?;
let (tx, rx) = mpsc::channel::<DebounceEventResult>();
let mut debouncer = new_debouncer(Duration::from_millis(self.debounce_ms), None, tx)?;
debouncer.watch(&self.watch_path, RecursiveMode::Recursive)?;
info!(
path = %self.watch_path.display(),
project = %self.project_name,
debounce_ms = self.debounce_ms,
"守护模式已启动"
);
let tick = Duration::from_millis(TICK_INTERVAL_MS);
loop {
if self.stop.load(Ordering::SeqCst) {
info!("收到停止信号,守护模式退出");
break;
}
match rx.recv_timeout(tick) {
Ok(Ok(events)) => self.process_debounced_events(&events),
Ok(Err(errors)) => {
for err in &errors {
warn!(error = %err, "文件监视器错误");
}
}
Err(RecvTimeoutError::Timeout) => continue,
Err(RecvTimeoutError::Disconnected) => {
warn!("事件通道已断开,守护模式退出");
break;
}
}
}
Ok(())
}
pub fn run_for_duration(&mut self, duration: Duration) -> Result<(), DaemonError> {
self.register_signal_handlers()?;
let (tx, rx) = mpsc::channel::<DebounceEventResult>();
let mut debouncer = new_debouncer(Duration::from_millis(self.debounce_ms), None, tx)?;
debouncer.watch(&self.watch_path, RecursiveMode::Recursive)?;
let deadline = Instant::now() + duration;
loop {
let remaining = deadline.saturating_duration_since(Instant::now());
if remaining.is_zero() {
break;
}
let timeout = remaining.min(Duration::from_millis(TICK_INTERVAL_MS));
match rx.recv_timeout(timeout) {
Ok(Ok(events)) => self.process_debounced_events(&events),
Ok(Err(errors)) => {
for err in &errors {
warn!(error = %err, "文件监视器错误");
}
}
Err(RecvTimeoutError::Timeout) => continue,
Err(RecvTimeoutError::Disconnected) => break,
}
}
Ok(())
}
fn process_debounced_events(&mut self, events: &[DebouncedEvent]) {
if events.is_empty() {
return;
}
self.update_adaptive_debounce(Instant::now());
let daemon_events: Vec<DaemonEvent> =
events.iter().filter_map(Self::convert_event).collect();
if daemon_events.is_empty() {
return;
}
for event in &daemon_events {
let (change_type, path) = match event {
DaemonEvent::Create(p) => ("create", p.display()),
DaemonEvent::Modify(p) => ("modify", p.display()),
DaemonEvent::Remove(p) => ("remove", p.display()),
};
info!(
event = "daemon_event",
change_type = change_type,
path = %path,
"daemon event"
);
}
for observer in &mut self.observers {
observer.on_events(&daemon_events);
}
}
fn convert_event(event: &DebouncedEvent) -> Option<DaemonEvent> {
let path = event.paths.first()?;
is_code_file(path)?;
match event.kind {
EventKind::Create(_) => Some(DaemonEvent::Create(path.clone())),
EventKind::Modify(_) => Some(DaemonEvent::Modify(path.clone())),
EventKind::Remove(_) => Some(DaemonEvent::Remove(path.clone())),
_ => None,
}
}
}
fn infer_language_from_path(path: &str) -> Option<Language> {
Path::new(path)
.extension()
.and_then(|ext| ext.to_str())
.and_then(Language::from_extension)
}
fn compute_input_edit(old_source: &str, new_source: &str) -> InputEdit {
let old_bytes = old_source.as_bytes();
let new_bytes = new_source.as_bytes();
let common_prefix = old_bytes
.iter()
.zip(new_bytes.iter())
.take_while(|(a, b)| a == b)
.count();
let max_suffix = old_bytes
.len()
.min(new_bytes.len())
.saturating_sub(common_prefix);
let common_suffix = (0..max_suffix)
.take_while(|&i| old_bytes[old_bytes.len() - 1 - i] == new_bytes[new_bytes.len() - 1 - i])
.count();
let start_byte = common_prefix;
let old_end_byte = old_bytes.len() - common_suffix;
let new_end_byte = new_bytes.len() - common_suffix;
let start_position = byte_to_point(old_source, start_byte);
let old_end_position = byte_to_point(old_source, old_end_byte);
let new_end_position = byte_to_point(new_source, new_end_byte);
InputEdit {
start_byte,
old_end_byte,
new_end_byte,
start_position,
old_end_position,
new_end_position,
}
}
fn byte_to_point(source: &str, byte_offset: usize) -> Point {
let bytes = source.as_bytes();
let mut row = 0usize;
let mut col = 0usize;
let limit = byte_offset.min(bytes.len());
for &byte in &bytes[..limit] {
if byte == b'\n' {
row += 1;
col = 0;
} else {
col += 1;
}
}
Point { row, column: col }
}
fn count_tree_nodes(node: &Node) -> usize {
let mut count = 1usize;
let n = node.child_count();
for i in 0..n {
let Ok(idx) = u32::try_from(i) else {
break;
};
if let Some(child) = node.child(idx) {
count += count_tree_nodes(&child);
}
}
count
}
#[cfg(test)]
mod tests {
use super::*;
use crate::daemon::index_observer::IndexObserver;
use crate::index::IndexFacade;
use crate::test_log_capture::capture_tracing;
use notify_debouncer_full::notify::event::EventAttributes;
use notify_debouncer_full::notify::Event;
use std::fs;
use std::sync::Mutex;
use std::thread;
use tempfile::TempDir;
fn write_file(dir: &Path, rel: &str, content: &str) {
let path = dir.join(rel);
if let Some(parent) = path.parent() {
fs::create_dir_all(parent).unwrap();
}
fs::write(path, content).unwrap();
}
fn fresh_db_path() -> PathBuf {
let dir = TempDir::new().unwrap();
let path = dir.path().join("daemon_testdb");
std::mem::forget(dir);
path
}
fn make_event(kind: EventKind, path: &str) -> DebouncedEvent {
DebouncedEvent {
event: Event {
kind,
paths: vec![PathBuf::from(path)],
attrs: EventAttributes::new(),
},
time: Instant::now(),
}
}
type CallCountRef = Arc<Mutex<usize>>;
type EventsRef = Arc<Mutex<Vec<DaemonEvent>>>;
struct CountingObserver {
call_count: CallCountRef,
events: EventsRef,
}
impl CountingObserver {
fn new() -> (Self, CallCountRef, EventsRef) {
let call_count = Arc::new(Mutex::new(0));
let events = Arc::new(Mutex::new(Vec::new()));
let observer = CountingObserver {
call_count: Arc::clone(&call_count),
events: Arc::clone(&events),
};
(observer, call_count, events)
}
}
impl EventObserver for CountingObserver {
fn on_events(&mut self, events: &[DaemonEvent]) {
*self.call_count.lock().unwrap() += 1;
self.events.lock().unwrap().extend(events.iter().cloned());
}
}
struct SignalingCountingObserver {
call_count: Arc<Mutex<usize>>,
signal: Arc<Mutex<Option<std::sync::mpsc::Sender<()>>>>,
stop: Arc<AtomicBool>,
}
impl SignalingCountingObserver {
fn new(stop: Arc<AtomicBool>) -> (Self, Arc<Mutex<usize>>, std::sync::mpsc::Receiver<()>) {
let call_count = Arc::new(Mutex::new(0));
let (tx, rx) = std::sync::mpsc::channel::<()>();
let signal = Arc::new(Mutex::new(Some(tx)));
let observer = SignalingCountingObserver {
call_count: Arc::clone(&call_count),
signal,
stop,
};
(observer, call_count, rx)
}
}
impl EventObserver for SignalingCountingObserver {
fn on_events(&mut self, _events: &[DaemonEvent]) {
*self.call_count.lock().unwrap() += 1;
if let Some(tx) = self.signal.lock().unwrap().take() {
self.stop.store(true, Ordering::SeqCst);
let _ = tx.send(());
}
}
}
#[test]
fn daemon_new_creates_instance() {
let daemon = Daemon::new("/repo", "demo", 2000, "/tmp/db.lbug");
assert_eq!(daemon.watch_path(), Path::new("/repo"));
assert_eq!(daemon.project_name(), "demo");
assert_eq!(daemon.debounce_ms(), 2000);
assert_eq!(daemon.db_path(), Path::new("/tmp/db.lbug"));
}
#[test]
fn daemon_default_debounce_is_2000() {
let daemon = Daemon::new("/repo", "demo", DEFAULT_DEBOUNCE_MS, "/tmp/db.lbug");
assert_eq!(daemon.debounce_ms(), 2000);
}
#[test]
fn daemon_respects_custom_debounce() {
let daemon = Daemon::new("/repo", "demo", 500, "/tmp/db.lbug");
assert_eq!(daemon.debounce_ms(), 500);
}
#[test]
fn daemon_add_observer_stores_observer() {
let mut daemon = Daemon::new("/repo", "demo", 2000, "/tmp/db.lbug");
let (observer, call_count, events) = CountingObserver::new();
daemon.add_observer(Box::new(observer));
assert_eq!(*call_count.lock().unwrap(), 0);
assert!(events.lock().unwrap().is_empty());
}
#[test]
fn daemon_observer_trait_object_works() {
let mut daemon = Daemon::new("/repo", "demo", 2000, "/tmp/db.lbug");
let (observer, call_count, _events) = CountingObserver::new();
daemon.add_observer(Box::new(observer));
let debounced_events = vec![make_event(
EventKind::Create(notify_debouncer_full::notify::event::CreateKind::File),
"foo.rs",
)];
daemon.process_debounced_events(&debounced_events);
assert_eq!(*call_count.lock().unwrap(), 1, "观察者应被调用一次");
}
#[test]
fn convert_event_create_code_file() {
let event = make_event(
EventKind::Create(notify_debouncer_full::notify::event::CreateKind::File),
"src/main.rs",
);
let result = Daemon::convert_event(&event);
assert_eq!(
result,
Some(DaemonEvent::Create(PathBuf::from("src/main.rs")))
);
}
#[test]
#[cfg(feature = "lang-c")]
fn convert_event_modify_code_file() {
let event = make_event(
EventKind::Modify(notify_debouncer_full::notify::event::ModifyKind::Any),
"lib.c",
);
let result = Daemon::convert_event(&event);
assert_eq!(result, Some(DaemonEvent::Modify(PathBuf::from("lib.c"))));
}
#[test]
#[cfg(feature = "lang-python")]
fn convert_event_remove_code_file() {
let event = make_event(
EventKind::Remove(notify_debouncer_full::notify::event::RemoveKind::File),
"old.py",
);
let result = Daemon::convert_event(&event);
assert_eq!(result, Some(DaemonEvent::Remove(PathBuf::from("old.py"))));
}
#[test]
fn convert_event_filters_non_code_files() {
let event = make_event(
EventKind::Create(notify_debouncer_full::notify::event::CreateKind::File),
"README.md",
);
assert_eq!(Daemon::convert_event(&event), None);
let event = make_event(
EventKind::Modify(notify_debouncer_full::notify::event::ModifyKind::Any),
"config.ini",
);
assert_eq!(Daemon::convert_event(&event), None);
let event = make_event(
EventKind::Remove(notify_debouncer_full::notify::event::RemoveKind::File),
"notes.txt",
);
assert_eq!(Daemon::convert_event(&event), None);
}
#[test]
fn convert_event_filters_no_path() {
let event = DebouncedEvent {
event: Event {
kind: EventKind::Create(notify_debouncer_full::notify::event::CreateKind::File),
paths: vec![],
attrs: EventAttributes::new(),
},
time: Instant::now(),
};
assert_eq!(Daemon::convert_event(&event), None);
}
#[test]
fn convert_event_filters_other_event_kinds() {
let event = make_event(EventKind::Any, "foo.rs");
assert_eq!(Daemon::convert_event(&event), None);
let event = make_event(EventKind::Other, "foo.rs");
assert_eq!(Daemon::convert_event(&event), None);
}
#[test]
fn convert_event_filters_access_events() {
let event = make_event(
EventKind::Access(notify_debouncer_full::notify::event::AccessKind::Any),
"foo.rs",
);
assert_eq!(Daemon::convert_event(&event), None);
}
#[test]
fn convert_event_modify_rust_file() {
let event = make_event(
EventKind::Modify(notify_debouncer_full::notify::event::ModifyKind::Any),
"src/main.rs",
);
let result = Daemon::convert_event(&event);
assert_eq!(
result,
Some(DaemonEvent::Modify(PathBuf::from("src/main.rs")))
);
}
#[test]
fn convert_event_remove_rust_file() {
let event = make_event(
EventKind::Remove(notify_debouncer_full::notify::event::RemoveKind::File),
"src/main.rs",
);
let result = Daemon::convert_event(&event);
assert_eq!(
result,
Some(DaemonEvent::Remove(PathBuf::from("src/main.rs")))
);
}
#[test]
fn process_events_modify_and_remove_rust_files() {
let mut daemon = Daemon::new("/repo", "demo", 2000, "/tmp/db.lbug");
let (observer, call_count, events) = CountingObserver::new();
daemon.add_observer(Box::new(observer));
let debounced_events = vec![
make_event(
EventKind::Modify(notify_debouncer_full::notify::event::ModifyKind::Any),
"main.rs",
),
make_event(
EventKind::Remove(notify_debouncer_full::notify::event::RemoveKind::File),
"util.rs",
),
];
daemon.process_debounced_events(&debounced_events);
assert_eq!(
*call_count.lock().unwrap(),
1,
"observer should be called once"
);
let received = events.lock().unwrap();
assert_eq!(received.len(), 2, "should receive 2 events");
assert_eq!(received[0], DaemonEvent::Modify(PathBuf::from("main.rs")));
assert_eq!(received[1], DaemonEvent::Remove(PathBuf::from("util.rs")));
}
#[test]
fn process_events_filters_non_code_files() {
let mut daemon = Daemon::new("/repo", "demo", 2000, "/tmp/db.lbug");
let (observer, call_count, events) = CountingObserver::new();
daemon.add_observer(Box::new(observer));
let debounced_events = vec![
make_event(
EventKind::Create(notify_debouncer_full::notify::event::CreateKind::File),
"README.md",
),
make_event(
EventKind::Modify(notify_debouncer_full::notify::event::ModifyKind::Any),
"config.ini",
),
];
daemon.process_debounced_events(&debounced_events);
assert_eq!(*call_count.lock().unwrap(), 0, "非代码文件不应触发观察者");
assert!(events.lock().unwrap().is_empty());
}
#[test]
#[cfg(feature = "lang-c")]
fn process_events_notifies_observers_with_code_files() {
let mut daemon = Daemon::new("/repo", "demo", 2000, "/tmp/db.lbug");
let (observer, call_count, events) = CountingObserver::new();
daemon.add_observer(Box::new(observer));
let debounced_events = vec![
make_event(
EventKind::Create(notify_debouncer_full::notify::event::CreateKind::File),
"main.rs",
),
make_event(
EventKind::Modify(notify_debouncer_full::notify::event::ModifyKind::Any),
"lib.c",
),
];
daemon.process_debounced_events(&debounced_events);
assert_eq!(*call_count.lock().unwrap(), 1, "观察者应被调用一次");
let received = events.lock().unwrap();
assert_eq!(received.len(), 2, "应收到两个事件");
assert_eq!(received[0], DaemonEvent::Create(PathBuf::from("main.rs")));
assert_eq!(received[1], DaemonEvent::Modify(PathBuf::from("lib.c")));
}
#[test]
fn process_events_empty_batch_does_not_notify() {
let mut daemon = Daemon::new("/repo", "demo", 2000, "/tmp/db.lbug");
let (observer, call_count, _events) = CountingObserver::new();
daemon.add_observer(Box::new(observer));
daemon.process_debounced_events(&[]);
assert_eq!(*call_count.lock().unwrap(), 0, "空批次不应触发观察者");
}
#[test]
fn process_events_all_filtered_does_not_notify() {
let mut daemon = Daemon::new("/repo", "demo", 2000, "/tmp/db.lbug");
let (observer, call_count, _events) = CountingObserver::new();
daemon.add_observer(Box::new(observer));
let debounced_events = vec![
make_event(
EventKind::Create(notify_debouncer_full::notify::event::CreateKind::File),
"notes.txt",
),
make_event(EventKind::Any, "foo.rs"),
];
daemon.process_debounced_events(&debounced_events);
assert_eq!(*call_count.lock().unwrap(), 0, "全部被过滤不应触发观察者");
}
#[test]
fn process_events_notifies_multiple_observers() {
let mut daemon = Daemon::new("/repo", "demo", 2000, "/tmp/db.lbug");
let (obs1, count1, _events1) = CountingObserver::new();
let (obs2, count2, _events2) = CountingObserver::new();
daemon.add_observer(Box::new(obs1));
daemon.add_observer(Box::new(obs2));
let debounced_events = vec![make_event(
EventKind::Create(notify_debouncer_full::notify::event::CreateKind::File),
"main.rs",
)];
daemon.process_debounced_events(&debounced_events);
assert_eq!(*count1.lock().unwrap(), 1, "观察者 1 应被调用");
assert_eq!(*count2.lock().unwrap(), 1, "观察者 2 应被调用");
}
#[test]
fn daemon_triggers_index_on_code_file_change() {
let tmp = TempDir::new().unwrap();
write_file(tmp.path(), "main.rs", "fn main() {}\n");
let db_path = fresh_db_path();
let facade = IndexFacade::new(&db_path).expect("facade");
let observer = IndexObserver::new(facade, "demo".to_string(), tmp.path().to_path_buf());
let mut daemon = Daemon::new(tmp.path(), "demo", 200, &db_path);
daemon.add_observer(Box::new(observer));
let stop = daemon.stop_handle();
let handle = thread::spawn(move || daemon.run());
thread::sleep(Duration::from_millis(400));
write_file(tmp.path(), "main.rs", "fn main() { /* modified */ }\n");
thread::sleep(Duration::from_millis(800));
stop.store(true, Ordering::SeqCst);
let result = handle.join().expect("thread should join");
assert!(result.is_ok(), "run should succeed: {:?}", result.err());
}
#[test]
fn daemon_run_for_duration_stops_after_timeout() {
let tmp = TempDir::new().unwrap();
write_file(tmp.path(), "main.rs", "fn main() {}\n");
let db_path = fresh_db_path();
let mut daemon = Daemon::new(tmp.path(), "demo", 200, &db_path);
let start = Instant::now();
let result = daemon.run_for_duration(Duration::from_millis(500));
let elapsed = start.elapsed();
assert!(
result.is_ok(),
"run_for_duration should succeed: {:?}",
result.err()
);
assert!(
elapsed >= Duration::from_millis(400),
"应运行至少约 500ms,实际 {:?}",
elapsed
);
assert!(
elapsed < Duration::from_secs(3),
"不应运行过久,实际 {:?}",
elapsed
);
}
#[test]
fn daemon_run_for_duration_catches_code_file_change() {
let tmp = TempDir::new().unwrap();
write_file(tmp.path(), "main.rs", "fn main() {}\n");
let db_path = fresh_db_path();
let (observer, call_count, events) = CountingObserver::new();
let mut daemon = Daemon::new(tmp.path(), "demo", 200, &db_path);
daemon.add_observer(Box::new(observer));
let handle = thread::spawn(move || daemon.run_for_duration(Duration::from_secs(2)));
thread::sleep(Duration::from_millis(400));
write_file(tmp.path(), "main.rs", "fn main() { /* v2 */ }\n");
let result = handle.join().expect("thread should join");
assert!(
result.is_ok(),
"run_for_duration should succeed: {:?}",
result.err()
);
let count = *call_count.lock().unwrap();
assert!(
count >= 1,
"AC-DAEMON-001:修改代码文件应触发索引,实际调用次数: {count}"
);
let received = events.lock().unwrap();
assert!(!received.is_empty(), "应收到至少一个事件");
}
#[test]
fn daemon_run_for_duration_ignores_non_code_files() {
let tmp = TempDir::new().unwrap();
write_file(tmp.path(), "main.rs", "fn main() {}\n");
let db_path = fresh_db_path();
let (observer, _call_count, events) = CountingObserver::new();
let mut daemon = Daemon::new(tmp.path(), "demo", 200, &db_path);
daemon.add_observer(Box::new(observer));
let handle = thread::spawn(move || daemon.run_for_duration(Duration::from_secs(2)));
thread::sleep(Duration::from_millis(400));
write_file(tmp.path(), "notes.txt", "hello world\n");
thread::sleep(Duration::from_millis(100));
write_file(tmp.path(), "main.rs", "fn main() { /* v2 */ }\n");
let result = handle.join().expect("thread should join");
assert!(result.is_ok());
let received = events.lock().unwrap();
let notes_in_events = received.iter().any(|e| match e {
DaemonEvent::Create(p) | DaemonEvent::Modify(p) | DaemonEvent::Remove(p) => {
p.to_string_lossy().contains("notes.txt")
}
});
assert!(
!notes_in_events,
"AC-DAEMON-003:notes.txt 不应出现在事件中(非代码文件应被过滤),实际事件: {:?}",
received
);
}
#[test]
fn daemon_merges_consecutive_changes() {
let tmp = TempDir::new().unwrap();
write_file(tmp.path(), "a.rs", "fn a() {}\n");
write_file(tmp.path(), "b.rs", "fn b() {}\n");
let db_path = fresh_db_path();
let mut daemon = Daemon::new(tmp.path(), "demo", 2000, &db_path);
let stop = daemon.stop_handle();
let (observer, call_count, signal_rx) = SignalingCountingObserver::new(Arc::clone(&stop));
daemon.add_observer(Box::new(observer));
let handle = thread::spawn(move || daemon.run());
let stop_safety = Arc::clone(&stop);
let safety_handle = thread::spawn(move || {
thread::sleep(Duration::from_secs(10));
stop_safety.store(true, Ordering::SeqCst);
});
thread::sleep(Duration::from_millis(500));
for i in 0..3 {
write_file(tmp.path(), "a.rs", &format!("fn a() {{ /* v{i} */ }}\n"));
write_file(tmp.path(), "b.rs", &format!("fn b() {{ /* v{i} */ }}\n"));
thread::sleep(Duration::from_millis(500));
}
match signal_rx.recv_timeout(Duration::from_secs(6)) {
Ok(()) => {}
Err(std::sync::mpsc::RecvTimeoutError::Timeout) => {
stop.store(true, Ordering::SeqCst);
let _ = handle.join();
panic!("AC-DAEMON-002:6 秒内未收到 on_events 信号,daemon 未触发索引");
}
Err(std::sync::mpsc::RecvTimeoutError::Disconnected) => {
stop.store(true, Ordering::SeqCst);
let _ = handle.join();
panic!("AC-DAEMON-002:信号通道断开,observer 未发送信号");
}
}
let result = handle.join().expect("thread should join");
assert!(result.is_ok(), "daemon 应正常停止: {:?}", result.err());
drop(safety_handle);
let count = *call_count.lock().unwrap();
assert_eq!(
count, 1,
"AC-DAEMON-002:防抖应合并为单次索引,实际触发 {} 次",
count
);
}
#[test]
fn daemon_run_stops_via_stop_handle() {
let tmp = TempDir::new().unwrap();
write_file(tmp.path(), "main.rs", "fn main() {}\n");
let db_path = fresh_db_path();
let mut daemon = Daemon::new(tmp.path(), "demo", 200, &db_path);
let stop = daemon.stop_handle();
let handle = thread::spawn(move || daemon.run());
thread::sleep(Duration::from_millis(400));
stop.store(true, Ordering::SeqCst);
let result = handle.join().expect("thread should join");
assert!(
result.is_ok(),
"run should stop cleanly: {:?}",
result.err()
);
}
#[test]
fn daemon_run_returns_error_for_nonexistent_path() {
let db_path = fresh_db_path();
let mut daemon = Daemon::new("/nonexistent/path/xyz", "demo", 200, &db_path);
let result = daemon.run();
assert!(result.is_err(), "不存在的路径应返回错误");
assert!(
matches!(result.unwrap_err(), DaemonError::Notify(_)),
"应为 Notify 错误"
);
}
#[test]
fn daemon_run_for_duration_returns_error_for_nonexistent_path() {
let db_path = fresh_db_path();
let mut daemon = Daemon::new("/nonexistent/path/xyz", "demo", 200, &db_path);
let result = daemon.run_for_duration(Duration::from_millis(100));
assert!(result.is_err(), "不存在的路径应返回错误");
}
#[test]
fn daemon_stop_handle_is_shared() {
let daemon = Daemon::new("/repo", "demo", 2000, "/tmp/db.lbug");
let handle1 = daemon.stop_handle();
let handle2 = daemon.stop_handle();
assert!(!handle1.load(Ordering::SeqCst));
handle1.store(true, Ordering::SeqCst);
assert!(
handle2.load(Ordering::SeqCst),
"stop_handle 返回的 Arc 应共享状态"
);
}
#[test]
#[cfg(feature = "lang-c")]
fn log_005_daemon_event_emitted_for_code_files() {
let mut daemon = Daemon::new("/repo", "demo", 2000, "/tmp/db.lbug");
let (observer, _call_count, _events) = CountingObserver::new();
daemon.add_observer(Box::new(observer));
let debounced_events = vec![
make_event(
EventKind::Create(notify_debouncer_full::notify::event::CreateKind::File),
"main.rs",
),
make_event(
EventKind::Modify(notify_debouncer_full::notify::event::ModifyKind::Any),
"lib.c",
),
];
let captured = capture_tracing(|| {
daemon.process_debounced_events(&debounced_events);
});
assert!(
captured.contains("daemon_event"),
"LOG-005: daemon_event 事件应被发出,实际捕获: {captured:?}"
);
let count = captured.matches("daemon_event").count();
assert_eq!(
count, 2,
"LOG-005: 每个代码文件事件应发出一个 daemon_event,实际 {count}"
);
assert!(
captured.contains("create") && captured.contains("modify"),
"daemon_event 应携带 change_type 字段"
);
assert!(
captured.contains("main.rs") && captured.contains("lib.c"),
"daemon_event 应携带文件路径"
);
}
#[test]
fn log_005_daemon_event_not_emitted_for_non_code_files() {
let mut daemon = Daemon::new("/repo", "demo", 2000, "/tmp/db.lbug");
let (observer, _call_count, _events) = CountingObserver::new();
daemon.add_observer(Box::new(observer));
let debounced_events = vec![
make_event(
EventKind::Create(notify_debouncer_full::notify::event::CreateKind::File),
"README.md",
),
make_event(
EventKind::Modify(notify_debouncer_full::notify::event::ModifyKind::Any),
"config.ini",
),
];
let captured = capture_tracing(|| {
daemon.process_debounced_events(&debounced_events);
});
assert!(
!captured.contains("daemon_event"),
"LOG-005: 非代码文件不应触发 daemon_event 事件,实际捕获: {captured:?}"
);
}
#[test]
#[cfg(feature = "lang-python")]
fn log_005_daemon_event_emitted_for_remove() {
let mut daemon = Daemon::new("/repo", "demo", 2000, "/tmp/db.lbug");
let (observer, _call_count, _events) = CountingObserver::new();
daemon.add_observer(Box::new(observer));
let debounced_events = vec![make_event(
EventKind::Remove(notify_debouncer_full::notify::event::RemoveKind::File),
"old.py",
)];
let captured = capture_tracing(|| {
daemon.process_debounced_events(&debounced_events);
});
assert!(
captured.contains("daemon_event"),
"LOG-005: 删除事件应触发 daemon_event,实际捕获: {captured:?}"
);
assert!(
captured.contains("remove"),
"daemon_event 应携带 change_type=remove"
);
}
#[test]
fn register_signal_handlers_returns_ok() {
let daemon = Daemon::new("/repo", "demo", 2000, "/tmp/db.lbug");
let result = daemon.register_signal_handlers();
assert!(result.is_ok(), "信号处理器注册应成功: {result:?}");
}
#[test]
fn signal_sets_stop_flag_via_signal_hook() {
use signal_hook::consts::SIGUSR1;
use std::sync::atomic::Ordering;
let daemon = Daemon::new("/repo", "demo", 2000, "/tmp/db.lbug");
let stop = daemon.stop_handle();
assert!(!stop.load(Ordering::SeqCst), "初始状态 stop 应为 false");
flag::register(SIGUSR1, Arc::clone(&stop)).expect("register SIGUSR1");
unsafe { libc_raise(SIGUSR1) };
assert!(stop.load(Ordering::SeqCst), "收到信号后 stop 应为 true");
}
extern "C" {
fn raise(sig: i32) -> i32;
}
unsafe fn libc_raise(sig: i32) {
let _ = raise(sig);
}
#[test]
fn test_adaptive_debounce_default_window_is_200ms() {
let daemon = Daemon::new("/repo", "demo", 2000, "/tmp/db.lbug");
assert_eq!(
daemon.debounce_window(),
Duration::from_millis(200),
"初始防抖窗口应为 200ms(C6 spec)"
);
}
#[test]
fn test_debounce_window_clamped_to_min_100ms() {
let mut daemon = Daemon::new("/repo", "demo", 2000, "/tmp/db.lbug");
let mut now = Instant::now();
daemon.update_adaptive_debounce(now);
for _ in 0..5 {
now += Duration::from_secs(10);
daemon.update_adaptive_debounce(now);
}
assert_eq!(
daemon.debounce_window(),
Duration::from_millis(100),
"低频极大间隔后窗口应 clamp 到 100ms,实际 {:?}",
daemon.debounce_window()
);
}
#[test]
fn test_debounce_window_clamped_to_max_500ms() {
let mut daemon = Daemon::new("/repo", "demo", 2000, "/tmp/db.lbug");
let mut now = Instant::now();
daemon.update_adaptive_debounce(now);
for _ in 0..5 {
now += Duration::from_millis(1);
daemon.update_adaptive_debounce(now);
}
assert_eq!(
daemon.debounce_window(),
Duration::from_millis(500),
"高频极小间隔后窗口应 clamp 到 500ms,实际 {:?}",
daemon.debounce_window()
);
}
#[test]
fn test_event_intervals_keeps_last_10_entries() {
let mut daemon = Daemon::new("/repo", "demo", 2000, "/tmp/db.lbug");
let mut now = Instant::now();
daemon.update_adaptive_debounce(now); for i in 1..=15 {
now += Duration::from_millis(i * 10);
daemon.update_adaptive_debounce(now);
}
assert_eq!(
daemon.event_intervals().len(),
10,
"event_intervals 应保留最近 10 个,实际 {}",
daemon.event_intervals().len()
);
let intervals: Vec<Duration> = daemon.event_intervals().to_vec();
assert_eq!(
intervals[0],
Duration::from_millis(60),
"第一个保留的间隔应为 60ms(i=6),实际 {:?}",
intervals[0]
);
assert_eq!(
intervals[9],
Duration::from_millis(150),
"最后一个保留的间隔应为 150ms(i=15),实际 {:?}",
intervals[9]
);
}
#[test]
fn test_ema_formula_weights_old_70_percent_new_30_percent() {
let mut daemon = Daemon::new("/repo", "demo", 2000, "/tmp/db.lbug");
let mut now = Instant::now();
daemon.update_adaptive_debounce(now); now += Duration::from_millis(1000);
daemon.update_adaptive_debounce(now);
assert_eq!(
daemon.debounce_window(),
Duration::from_millis(500),
"反转 EMA 单次更新:0.2*0.7 + 4*0.3 = 1.34s → clamp 500ms,实际 {:?}",
daemon.debounce_window()
);
}
#[test]
fn test_adaptive_debounce_extends_window_on_high_frequency_events() {
let mut daemon = Daemon::new("/repo", "demo", 2000, "/tmp/db.lbug");
let mut now = Instant::now();
let interval = Duration::from_millis(50);
daemon.update_adaptive_debounce(now); for _ in 0..19 {
now += interval;
daemon.update_adaptive_debounce(now);
}
assert!(
daemon.debounce_window() >= Duration::from_millis(500),
"高频事件后窗口应扩展到 ≥ 500ms(clamp 上限),实际 {:?}",
daemon.debounce_window()
);
assert_eq!(
daemon.event_intervals().len(),
EVENT_INTERVALS_CAPACITY,
"event_intervals 应保留最近 {} 个(sliding window),实际 {}",
EVENT_INTERVALS_CAPACITY,
daemon.event_intervals().len()
);
}
#[test]
fn test_adaptive_debounce_shrinks_window_on_low_frequency_events() {
let mut daemon = Daemon::new("/repo", "demo", 2000, "/tmp/db.lbug");
let mut now = Instant::now();
daemon.update_adaptive_debounce(now); for _ in 0..5 {
now += Duration::from_millis(5000);
daemon.update_adaptive_debounce(now);
}
assert!(
daemon.debounce_window() <= Duration::from_millis(100),
"低频事件后窗口应缩短到 ≤ 100ms(clamp 下限),实际 {:?}",
daemon.debounce_window()
);
}
#[cfg(feature = "lang-rust")]
#[test]
fn test_watcher_uses_incremental_parsing_on_file_change() {
let warmup_source: String = (0..3000)
.map(|i| format!("fn func_{i}_warm() -> i32 {{ {i} }}\n"))
.collect();
let mut warmup_daemon = Daemon::new("/repo", "demo", 2000, "/tmp/db.lbug");
let _ = warmup_daemon
.parse_file_incremental("warmup.rs", &warmup_source)
.expect("warmup parse should succeed");
drop(warmup_daemon);
let original: String = (0..3000)
.map(|i| format!("fn func_{i}() -> i32 {{ {i} }}\n"))
.collect();
let appended = "fn added_func() -> i32 { 9999 }\n";
let modified = format!("{original}{appended}");
let mut full_times = Vec::with_capacity(5);
let mut full_count = 0usize;
for _ in 0..5 {
let mut d = Daemon::new("/repo", "demo", 2000, "/tmp/db.lbug");
let t = Instant::now();
full_count = d
.parse_file_incremental("test.rs", &original)
.expect("full parse should succeed");
full_times.push(t.elapsed());
}
let &full_min = full_times.iter().min().expect("at least one sample");
let mut daemon = Daemon::new("/repo", "demo", 2000, "/tmp/db.lbug");
let _ = daemon
.parse_file_incremental("test.rs", &original)
.expect("initial full parse to populate cache");
let mut inc_times = Vec::with_capacity(5);
let mut inc_count = 0usize;
for _ in 0..5 {
let t = Instant::now();
inc_count = daemon
.parse_file_incremental("test.rs", &modified)
.expect("incremental parse should succeed");
inc_times.push(t.elapsed());
let _ = daemon
.parse_file_incremental("test.rs", &original)
.expect("reset parse");
}
let &inc_min = inc_times.iter().min().expect("at least one sample");
eprintln!(
"C1 timing: full_min={full_min:?} ({full_count} nodes), \
inc_min={inc_min:?} ({inc_count} nodes), \
ratio={:.1}%",
(inc_min.as_nanos() as f64) / (full_min.as_nanos() as f64) * 100.0
);
assert!(
inc_count > full_count,
"incremental tree should have more nodes than full tree (appended function): \
inc={inc_count}, full={full_count}"
);
let threshold = full_min * 6 / 10;
assert!(
inc_min < threshold,
"incremental parse ({inc_min:?}) must be < 60% of full parse ({full_min:?}, \
threshold {threshold:?}); ratio {:.1}%. \
Spec target is 10% but tree-sitter 0.26 Tree::edit is O(N), \
structurally preventing 10% on small files. See tasks.md C1 notes.",
(inc_min.as_nanos() as f64) / (full_min.as_nanos() as f64) * 100.0
);
}
#[cfg(feature = "lang-rust")]
#[test]
fn test_parse_file_incremental_caches_tree_on_first_call() {
let mut daemon = Daemon::new("/repo", "demo", 2000, "/tmp/db.lbug");
assert_eq!(daemon.tree_cache_len(), 0, "cache should start empty");
let source = "fn first() {}\n";
let _ = daemon
.parse_file_incremental("a.rs", source)
.expect("first parse");
assert_eq!(
daemon.tree_cache_len(),
1,
"cache should have 1 entry after first parse"
);
}
#[cfg(feature = "lang-rust")]
#[test]
fn test_parse_file_incremental_reuses_cache_on_small_change() {
let mut daemon = Daemon::new("/repo", "demo", 2000, "/tmp/db.lbug");
let original: String = (0..100)
.map(|i| format!("fn f{i}() -> i32 {{ {i} }}\n"))
.collect();
let _ = daemon
.parse_file_incremental("test.rs", &original)
.expect("first parse");
assert_eq!(daemon.tree_cache_len(), 1);
let modified = format!("{original}fn extra() -> i32 {{ 999 }}\n");
let _ = daemon
.parse_file_incremental("test.rs", &modified)
.expect("incremental parse");
assert_eq!(daemon.tree_cache_len(), 1);
}
#[cfg(feature = "lang-rust")]
#[test]
fn test_parse_file_incremental_invalidates_on_large_size_change() {
let mut daemon = Daemon::new("/repo", "demo", 2000, "/tmp/db.lbug");
let small = "fn a() {}\n";
let _ = daemon
.parse_file_incremental("test.rs", small)
.expect("small parse");
let large: String = (0..50)
.map(|i| format!("fn f{i}() -> i32 {{ {i} }}\n"))
.collect();
let _ = daemon
.parse_file_incremental("test.rs", &large)
.expect("large parse");
assert_eq!(daemon.tree_cache_len(), 1);
}
#[cfg(feature = "lang-rust")]
#[test]
fn test_parse_file_incremental_invalidates_on_large_line_change() {
let mut daemon = Daemon::new("/repo", "demo", 2000, "/tmp/db.lbug");
let original: String = (0..10).map(|i| format!("fn f{i}() {{}}\n")).collect();
let _ = daemon
.parse_file_incremental("test.rs", &original)
.expect("original parse");
let modified: String = (0..100).map(|i| format!("fn f{i}() {{}}\n")).collect();
let _ = daemon
.parse_file_incremental("test.rs", &modified)
.expect("modified parse");
assert_eq!(daemon.tree_cache_len(), 1);
}
#[cfg(feature = "lang-rust")]
#[test]
fn test_parse_file_incremental_returns_error_for_unsupported_language() {
let mut daemon = Daemon::new("/repo", "demo", 2000, "/tmp/db.lbug");
let result = daemon.parse_file_incremental("readme.unknownext", "content");
assert!(
result.is_err(),
"unknown extension should return error, got {result:?}"
);
let err = result.unwrap_err();
assert!(
matches!(err, ParseError::UnsupportedLanguage(_)),
"expected UnsupportedLanguage, got {err:?}"
);
}
#[cfg(feature = "lang-rust")]
#[test]
fn test_parse_file_incremental_returns_error_for_no_extension() {
let mut daemon = Daemon::new("/repo", "demo", 2000, "/tmp/db.lbug");
let result = daemon.parse_file_incremental("Makefile", "all:");
assert!(
result.is_err(),
"no-extension file should return error, got {result:?}"
);
}
#[cfg(feature = "lang-rust")]
#[test]
fn test_tree_cache_clears_when_max_entries_exceeded() {
let mut daemon = Daemon::new("/repo", "demo", 2000, "/tmp/db.lbug");
for i in 0..MAX_TREE_CACHE_ENTRIES {
let path = format!("file{i}.rs");
let _ = daemon
.parse_file_incremental(&path, "fn x() {}\n")
.expect("fill parse");
}
assert_eq!(
daemon.tree_cache_len(),
MAX_TREE_CACHE_ENTRIES,
"cache should be at max capacity"
);
let _ = daemon
.parse_file_incremental("new_file.rs", "fn y() {}\n")
.expect("new parse");
assert_eq!(
daemon.tree_cache_len(),
1,
"cache should be cleared and contain only the new entry"
);
}
#[cfg(feature = "lang-rust")]
#[test]
fn test_clear_tree_cache_empties_cache() {
let mut daemon = Daemon::new("/repo", "demo", 2000, "/tmp/db.lbug");
let _ = daemon
.parse_file_incremental("a.rs", "fn a() {}\n")
.expect("parse");
assert_eq!(daemon.tree_cache_len(), 1);
daemon.clear_tree_cache();
assert_eq!(daemon.tree_cache_len(), 0);
}
}