Skip to main content

ragfs_index/
indexer.rs

1//! Main indexing service.
2
3use chrono::Utc;
4use ragfs_chunker::ChunkerRegistry;
5use ragfs_core::{
6    Chunk, ChunkConfig, ChunkMetadata, ChunkOutput, ContentType, DirectoryScope, EmbeddingConfig,
7    Error, FileEvent, FileRecord, FileStatus, IndexStats, Indexer, Result, VectorStore,
8};
9use ragfs_embed::EmbedderPool;
10use ragfs_extract::ExtractorRegistry;
11use std::path::{Path, PathBuf};
12use std::sync::Arc;
13use std::sync::atomic::{AtomicUsize, Ordering};
14use tokio::sync::{Notify, RwLock, broadcast, mpsc};
15use tracing::{debug, error, info, warn};
16use uuid::Uuid;
17
18use crate::watcher::FileWatcher;
19
20/// Index update events.
21#[derive(Debug, Clone)]
22pub enum IndexUpdate {
23    FileIndexed { path: PathBuf, chunk_count: u32 },
24    FileRemoved { path: PathBuf },
25    FileError { path: PathBuf, error: String },
26    IndexingStarted { path: PathBuf },
27}
28
29/// Configuration for the indexer.
30#[derive(Debug, Clone)]
31pub struct IndexerConfig {
32    /// Chunk configuration
33    pub chunk_config: ChunkConfig,
34    /// Embedding configuration
35    pub embed_config: EmbeddingConfig,
36    /// Include patterns (glob)
37    pub include_patterns: Vec<String>,
38    /// Exclude patterns (glob)
39    pub exclude_patterns: Vec<String>,
40    /// Debounce duration for the file watcher (milliseconds)
41    pub debounce_ms: u64,
42    /// Maximum file size to index (bytes)
43    pub max_file_size: u64,
44    /// Skip the content-hash short-circuit and reindex every eligible file
45    pub force: bool,
46}
47
48impl Default for IndexerConfig {
49    fn default() -> Self {
50        Self {
51            chunk_config: ChunkConfig::default(),
52            embed_config: EmbeddingConfig::default(),
53            include_patterns: vec!["**/*".to_string()],
54            exclude_patterns: vec![
55                "**/.*".to_string(),
56                "**/.git/**".to_string(),
57                "**/node_modules/**".to_string(),
58                "**/target/**".to_string(),
59                "**/__pycache__/**".to_string(),
60                "**/*.lock".to_string(),
61            ],
62            debounce_ms: 500,
63            max_file_size: 52_428_800,
64            force: false,
65        }
66    }
67}
68
69impl IndexerConfig {
70    /// Whether `path` is eligible for indexing given include/exclude patterns.
71    #[must_use]
72    pub fn should_process_path(&self, path: &Path) -> bool {
73        let path_str = path.to_string_lossy();
74        if self
75            .exclude_patterns
76            .iter()
77            .any(|pattern| path_matches_pattern(&path_str, pattern))
78        {
79            return false;
80        }
81        if self.include_patterns.is_empty() {
82            return true;
83        }
84        self.include_patterns
85            .iter()
86            .any(|pattern| path_matches_pattern(&path_str, pattern))
87    }
88}
89
90/// Match a path against a glob-style include/exclude pattern.
91///
92/// Supports `**` (any directories), `*` / `?` within a path segment, and
93/// basename-only patterns such as `test_*.rs` (matched against any component).
94/// Patterns that do not start with `**/` also match as a suffix of an absolute
95/// path (`src/**/*.rs` matches `/proj/src/lib.rs`).
96#[must_use]
97pub fn path_matches_pattern(path: &str, pattern: &str) -> bool {
98    let path = path.replace('\\', "/");
99    let pattern = pattern.replace('\\', "/");
100
101    if pattern.is_empty() {
102        return false;
103    }
104    if pattern == "**/*" || pattern == "**" || pattern == "*" {
105        return true;
106    }
107    if pattern == "**/.*" || pattern == ".*" {
108        // Hidden *files* only — do not treat `/tmp/.tmpXXXX/file.txt` as hidden.
109        return path
110            .rsplit('/')
111            .next()
112            .is_some_and(|name| !name.is_empty() && name.starts_with('.'));
113    }
114
115    if glob_match_path(&path, &pattern) {
116        return true;
117    }
118    // Absolute indexed paths should still match repo-relative globs.
119    if !pattern.starts_with("**/") && !pattern.starts_with('/') {
120        return glob_match_path(&path, &format!("**/{pattern}"));
121    }
122    false
123}
124
125fn glob_match_path(path: &str, pattern: &str) -> bool {
126    let path: Vec<&str> = path.split('/').filter(|s| !s.is_empty()).collect();
127    let pat: Vec<&str> = pattern.split('/').filter(|s| !s.is_empty()).collect();
128    glob_match_segments(&path, &pat)
129}
130
131fn glob_match_segments(path: &[&str], pat: &[&str]) -> bool {
132    match pat.split_first() {
133        None => path.is_empty(),
134        Some((&"**", rest)) => {
135            if rest.is_empty() {
136                return true;
137            }
138            glob_match_segments(path, rest)
139                || (!path.is_empty() && glob_match_segments(&path[1..], pat))
140        }
141        Some((seg, rest)) => path.split_first().is_some_and(|(head, tail)| {
142            glob_match_segment(head, seg) && glob_match_segments(tail, rest)
143        }),
144    }
145}
146
147fn glob_match_segment(name: &str, pat: &str) -> bool {
148    glob_match_bytes(name.as_bytes(), pat.as_bytes())
149}
150
151fn glob_match_bytes(name: &[u8], pat: &[u8]) -> bool {
152    match pat.split_first() {
153        None => name.is_empty(),
154        Some((b'*', rest)) => (0..=name.len()).any(|i| glob_match_bytes(&name[i..], rest)),
155        Some((b'?', rest)) => !name.is_empty() && glob_match_bytes(&name[1..], rest),
156        Some((c, rest)) => name.first() == Some(c) && glob_match_bytes(&name[1..], rest),
157    }
158}
159
160/// Main indexing service.
161pub struct IndexerService {
162    /// Root path being indexed
163    root: PathBuf,
164    /// Vector store
165    store: Arc<dyn VectorStore>,
166    /// Extractor registry
167    extractors: Arc<ExtractorRegistry>,
168    /// Chunker registry
169    chunkers: Arc<ChunkerRegistry>,
170    /// Embedder pool
171    embedder: Arc<EmbedderPool>,
172    /// Configuration
173    config: IndexerConfig,
174    /// Current stats
175    stats: Arc<RwLock<IndexStats>>,
176    /// Event sender for file watcher
177    event_tx: mpsc::Sender<FileEvent>,
178    /// Event receiver
179    event_rx: Arc<RwLock<mpsc::Receiver<FileEvent>>>,
180    /// Update broadcast
181    update_tx: broadcast::Sender<IndexUpdate>,
182    /// File watcher (if active)
183    watcher: Arc<RwLock<Option<FileWatcher>>>,
184    /// Running flag
185    running: Arc<RwLock<bool>>,
186    /// Outstanding file events (queued or in flight)
187    pending: Arc<AtomicUsize>,
188    /// Signaled when [`Self::pending`] reaches zero
189    idle_notify: Arc<Notify>,
190}
191
192impl IndexerService {
193    /// Create a new indexer service.
194    pub fn new(
195        root: PathBuf,
196        store: Arc<dyn VectorStore>,
197        extractors: Arc<ExtractorRegistry>,
198        chunkers: Arc<ChunkerRegistry>,
199        embedder: Arc<EmbedderPool>,
200        config: IndexerConfig,
201    ) -> Self {
202        let (event_tx, event_rx) = mpsc::channel(1024);
203        let (update_tx, _) = broadcast::channel(256);
204
205        Self {
206            root,
207            store,
208            extractors,
209            chunkers,
210            embedder,
211            config,
212            stats: Arc::new(RwLock::new(IndexStats::default())),
213            event_tx,
214            event_rx: Arc::new(RwLock::new(event_rx)),
215            update_tx,
216            watcher: Arc::new(RwLock::new(None)),
217            running: Arc::new(RwLock::new(false)),
218            pending: Arc::new(AtomicUsize::new(0)),
219            idle_notify: Arc::new(Notify::new()),
220        }
221    }
222
223    /// Configuration used by this indexer (include/exclude, chunk, force, …).
224    #[must_use]
225    pub fn config(&self) -> &IndexerConfig {
226        &self.config
227    }
228
229    /// Wait until every queued file event has been processed.
230    ///
231    /// Call this after [`Self::start`] for a one-shot index so the CLI does not
232    /// return while the background worker is still embedding files.
233    pub async fn wait_until_idle(&self) {
234        loop {
235            if self.pending.load(Ordering::SeqCst) == 0 {
236                tokio::task::yield_now().await;
237                if self.pending.load(Ordering::SeqCst) == 0 {
238                    return;
239                }
240            }
241            tokio::select! {
242                () = self.idle_notify.notified() => {}
243                () = tokio::time::sleep(std::time::Duration::from_millis(20)) => {}
244            }
245        }
246    }
247
248    /// Subscribe to index updates.
249    #[must_use]
250    pub fn subscribe(&self) -> broadcast::Receiver<IndexUpdate> {
251        self.update_tx.subscribe()
252    }
253
254    /// Get the root path.
255    #[must_use]
256    pub fn root(&self) -> &Path {
257        &self.root
258    }
259
260    /// Start the indexer background task.
261    pub async fn start(&self) -> Result<()> {
262        let mut running = self.running.write().await;
263        if *running {
264            return Ok(());
265        }
266        *running = true;
267        drop(running);
268
269        info!("Starting indexer for {:?}", self.root);
270
271        // Initialize store
272        self.store.init().await.map_err(Error::Store)?;
273
274        // Start file watcher (debounce from config)
275        let watcher = FileWatcher::with_pending(
276            self.event_tx.clone(),
277            std::time::Duration::from_millis(self.config.debounce_ms),
278            Some(Arc::clone(&self.pending)),
279        )
280        .map_err(|e| Error::Other(format!("watcher error: {e}")))?;
281        {
282            let mut w = self.watcher.write().await;
283            *w = Some(watcher);
284        }
285
286        // Watch root directory
287        {
288            let mut w = self.watcher.write().await;
289            if let Some(ref mut watcher) = *w {
290                watcher
291                    .watch(&self.root)
292                    .map_err(|e| Error::Other(format!("watch error: {e}")))?;
293            }
294        }
295
296        // Clone what we need for the background task
297        let event_rx = Arc::clone(&self.event_rx);
298        let update_tx = self.update_tx.clone();
299        let running = Arc::clone(&self.running);
300        let store = Arc::clone(&self.store);
301        let extractors = Arc::clone(&self.extractors);
302        let chunkers = Arc::clone(&self.chunkers);
303        let embedder = Arc::clone(&self.embedder);
304        let config = self.config.clone();
305        let stats = Arc::clone(&self.stats);
306        let pending = Arc::clone(&self.pending);
307        let idle_notify = Arc::clone(&self.idle_notify);
308        let root = self.root.clone();
309
310        // Spawn event processing task
311        tokio::spawn(async move {
312            let mut rx = event_rx.write().await;
313            while *running.read().await {
314                match rx.recv().await {
315                    Some(event) => {
316                        debug!("Received file event: {:?}", event);
317                        match &event {
318                            FileEvent::Created(path) | FileEvent::Modified(path) => {
319                                let _ = update_tx
320                                    .send(IndexUpdate::IndexingStarted { path: path.clone() });
321
322                                match process_file(
323                                    path,
324                                    &root,
325                                    &store,
326                                    &extractors,
327                                    &chunkers,
328                                    &embedder,
329                                    &config,
330                                )
331                                .await
332                                {
333                                    Ok(chunk_count) => {
334                                        info!("Indexed {:?} ({} chunks)", path, chunk_count);
335                                        let _ = update_tx.send(IndexUpdate::FileIndexed {
336                                            path: path.clone(),
337                                            chunk_count,
338                                        });
339
340                                        // Update stats
341                                        let mut s = stats.write().await;
342                                        s.indexed_files += 1;
343                                        s.total_chunks += u64::from(chunk_count);
344                                        s.last_update = Some(Utc::now());
345                                    }
346                                    Err(e) => {
347                                        error!("Failed to index {:?}: {}", path, e);
348                                        let _ = update_tx.send(IndexUpdate::FileError {
349                                            path: path.clone(),
350                                            error: e.to_string(),
351                                        });
352
353                                        // Update stats
354                                        let mut s = stats.write().await;
355                                        s.error_files += 1;
356                                    }
357                                }
358                            }
359                            FileEvent::Deleted(path) => {
360                                if let Err(e) = store.delete_by_file_path(path).await {
361                                    error!("Failed to delete {:?}: {}", path, e);
362                                } else {
363                                    let _ = update_tx
364                                        .send(IndexUpdate::FileRemoved { path: path.clone() });
365                                }
366                            }
367                            FileEvent::Renamed { from, to } => {
368                                // Delete old, index new
369                                if let Err(e) = store.delete_by_file_path(from).await {
370                                    warn!("Failed to delete old path {:?}: {}", from, e);
371                                }
372                                let _ = update_tx
373                                    .send(IndexUpdate::IndexingStarted { path: to.clone() });
374
375                                match process_file(
376                                    to,
377                                    &root,
378                                    &store,
379                                    &extractors,
380                                    &chunkers,
381                                    &embedder,
382                                    &config,
383                                )
384                                .await
385                                {
386                                    Ok(chunk_count) => {
387                                        let _ = update_tx.send(IndexUpdate::FileIndexed {
388                                            path: to.clone(),
389                                            chunk_count,
390                                        });
391                                    }
392                                    Err(e) => {
393                                        let _ = update_tx.send(IndexUpdate::FileError {
394                                            path: to.clone(),
395                                            error: e.to_string(),
396                                        });
397                                    }
398                                }
399                            }
400                        }
401                        finish_pending(&pending, &idle_notify);
402                    }
403                    None => break,
404                }
405            }
406        });
407
408        // Initial scan
409        self.scan().await?;
410
411        Ok(())
412    }
413
414    /// Perform initial scan of the root directory.
415    async fn scan(&self) -> Result<()> {
416        info!("Scanning {:?}", self.root);
417
418        let root = self.root.clone();
419        let event_tx = self.event_tx.clone();
420        let exclude_patterns = self.config.exclude_patterns.clone();
421        let include_patterns = self.config.include_patterns.clone();
422        let pending = Arc::clone(&self.pending);
423
424        // Walk directory in background thread (blocking I/O)
425        tokio::task::spawn_blocking(move || {
426            scan_directory(
427                &root,
428                &event_tx,
429                &exclude_patterns,
430                &include_patterns,
431                &pending,
432            );
433        })
434        .await
435        .map_err(|e| Error::Other(format!("scan task failed: {e}")))?;
436
437        Ok(())
438    }
439
440    /// Process a single file through the pipeline.
441    pub async fn process_single(&self, path: &Path) -> Result<u32> {
442        process_file(
443            path,
444            &self.root,
445            &self.store,
446            &self.extractors,
447            &self.chunkers,
448            &self.embedder,
449            &self.config,
450        )
451        .await
452    }
453
454    /// Reindex a path (file or directory).
455    ///
456    /// If the path is a file, it will be reindexed (existing chunks deleted first).
457    /// If the path is a directory, all files in it will be reindexed recursively.
458    pub async fn reindex_path(&self, path: &Path) -> Result<()> {
459        if !path.exists() {
460            return Err(Error::Other(format!(
461                "Path does not exist: {}",
462                path.display()
463            )));
464        }
465
466        if path.is_file() {
467            // Delete existing chunks and reindex
468            let _ = self.store.delete_by_file_path(path).await;
469
470            let _ = self.update_tx.send(IndexUpdate::IndexingStarted {
471                path: path.to_path_buf(),
472            });
473
474            match process_file(
475                path,
476                &self.root,
477                &self.store,
478                &self.extractors,
479                &self.chunkers,
480                &self.embedder,
481                &self.config,
482            )
483            .await
484            {
485                Ok(chunk_count) => {
486                    info!("Reindexed {:?} ({} chunks)", path, chunk_count);
487                    let _ = self.update_tx.send(IndexUpdate::FileIndexed {
488                        path: path.to_path_buf(),
489                        chunk_count,
490                    });
491                }
492                Err(e) => {
493                    error!("Failed to reindex {:?}: {}", path, e);
494                    let _ = self.update_tx.send(IndexUpdate::FileError {
495                        path: path.to_path_buf(),
496                        error: e.to_string(),
497                    });
498                    return Err(e);
499                }
500            }
501        } else if path.is_dir() {
502            // Reindex all files in directory recursively
503            info!("Reindexing directory {:?}", path);
504            self.reindex_directory(path).await?;
505        }
506
507        Ok(())
508    }
509
510    /// Recursively reindex all files in a directory.
511    async fn reindex_directory(&self, dir: &Path) -> Result<()> {
512        let entries = tokio::fs::read_dir(dir)
513            .await
514            .map_err(|e| Error::Other(format!("Failed to read directory: {e}")))?;
515
516        let mut entries_stream = tokio_stream::wrappers::ReadDirStream::new(entries);
517
518        use tokio_stream::StreamExt;
519        while let Some(entry) = entries_stream.next().await {
520            let entry = entry.map_err(|e| Error::Other(format!("Failed to read entry: {e}")))?;
521            let path = entry.path();
522
523            if path.is_dir() {
524                if self
525                    .config
526                    .exclude_patterns
527                    .iter()
528                    .any(|pattern| path_matches_pattern(&path.to_string_lossy(), pattern))
529                {
530                    continue;
531                }
532                // Recurse into subdirectory using a boxed future to avoid infinite recursion
533                Box::pin(self.reindex_directory(&path)).await?;
534            } else if path.is_file() {
535                if !self.config.should_process_path(&path) {
536                    continue;
537                }
538                // Reindex file - use process_file directly to avoid recursion
539                let _ = self.store.delete_by_file_path(&path).await;
540
541                let _ = self
542                    .update_tx
543                    .send(IndexUpdate::IndexingStarted { path: path.clone() });
544
545                match process_file(
546                    &path,
547                    &self.root,
548                    &self.store,
549                    &self.extractors,
550                    &self.chunkers,
551                    &self.embedder,
552                    &self.config,
553                )
554                .await
555                {
556                    Ok(chunk_count) => {
557                        info!("Reindexed {:?} ({} chunks)", path, chunk_count);
558                        let _ = self.update_tx.send(IndexUpdate::FileIndexed {
559                            path: path.clone(),
560                            chunk_count,
561                        });
562                    }
563                    Err(e) => {
564                        warn!("Failed to reindex {:?}: {}", path, e);
565                        let _ = self.update_tx.send(IndexUpdate::FileError {
566                            path: path.clone(),
567                            error: e.to_string(),
568                        });
569                        // Continue with other files
570                    }
571                }
572            }
573        }
574
575        Ok(())
576    }
577}
578
579/// Scan a directory and send file events.
580fn scan_directory(
581    root: &Path,
582    event_tx: &mpsc::Sender<FileEvent>,
583    exclude_patterns: &[String],
584    include_patterns: &[String],
585    pending: &AtomicUsize,
586) {
587    use std::fs;
588
589    fn visit_dir(
590        dir: &Path,
591        event_tx: &mpsc::Sender<FileEvent>,
592        exclude_patterns: &[String],
593        include_patterns: &[String],
594        pending: &AtomicUsize,
595    ) {
596        let entries = match fs::read_dir(dir) {
597            Ok(e) => e,
598            Err(e) => {
599                warn!("Cannot read directory {:?}: {}", dir, e);
600                return;
601            }
602        };
603
604        for entry in entries.flatten() {
605            let path = entry.path();
606            let path_str = path.to_string_lossy();
607
608            if exclude_patterns
609                .iter()
610                .any(|pattern| path_matches_pattern(&path_str, pattern))
611            {
612                continue;
613            }
614
615            if path.is_dir() {
616                visit_dir(&path, event_tx, exclude_patterns, include_patterns, pending);
617            } else if path.is_file() {
618                let included = include_patterns.is_empty()
619                    || include_patterns
620                        .iter()
621                        .any(|pattern| path_matches_pattern(&path_str, pattern));
622                if !included {
623                    continue;
624                }
625                queue_event_blocking(event_tx, pending, FileEvent::Created(path));
626            }
627        }
628    }
629
630    visit_dir(root, event_tx, exclude_patterns, include_patterns, pending);
631}
632
633fn queue_event_blocking(
634    event_tx: &mpsc::Sender<FileEvent>,
635    pending: &AtomicUsize,
636    event: FileEvent,
637) {
638    pending.fetch_add(1, Ordering::SeqCst);
639    if let Err(e) = event_tx.blocking_send(event) {
640        pending.fetch_sub(1, Ordering::SeqCst);
641        warn!("Failed to queue file event: {e}");
642    }
643}
644
645fn finish_pending(pending: &AtomicUsize, idle_notify: &Notify) {
646    let prev = pending.fetch_sub(1, Ordering::SeqCst);
647    if prev <= 1 {
648        idle_notify.notify_waiters();
649    }
650}
651
652/// Process a file through the full pipeline: extract → chunk → embed → store.
653async fn process_file(
654    path: &Path,
655    root: &Path,
656    store: &Arc<dyn VectorStore>,
657    extractors: &Arc<ExtractorRegistry>,
658    chunkers: &Arc<ChunkerRegistry>,
659    embedder: &Arc<EmbedderPool>,
660    config: &IndexerConfig,
661) -> Result<u32> {
662    // Get file metadata
663    let metadata = tokio::fs::metadata(path)
664        .await
665        .map_err(|e| Error::Other(format!("Failed to get metadata: {e}")))?;
666
667    if !metadata.is_file() {
668        return Ok(0);
669    }
670
671    if !config.should_process_path(path) {
672        debug!("Skipping {:?} (include/exclude patterns)", path);
673        return Ok(0);
674    }
675
676    if metadata.len() > config.max_file_size {
677        debug!(
678            "Skipping {:?} ({} bytes > max_file_size {})",
679            path,
680            metadata.len(),
681            config.max_file_size
682        );
683        let _ = store.delete_by_file_path(path).await;
684        return Ok(0);
685    }
686
687    // Compute content hash
688    let content_hash = compute_hash(path).await?;
689
690    // Check if already indexed with same hash (unless --force)
691    if !config.force
692        && let Ok(Some(existing)) = store.get_file(path).await
693        && existing.content_hash == content_hash
694        && existing.status == FileStatus::Indexed
695    {
696        debug!("File {:?} already indexed, skipping", path);
697        return Ok(existing.chunk_count);
698    }
699
700    // Determine MIME type
701    let mime_type = mime_guess::from_path(path)
702        .first_or_text_plain()
703        .to_string();
704
705    // Extract content
706    let content = extractors
707        .extract(path, &mime_type)
708        .await
709        .map_err(Error::Extraction)?;
710
711    if content.text.is_empty() {
712        debug!("Empty content for {:?}, skipping", path);
713        return Ok(0);
714    }
715
716    // Determine content type
717    let content_type = determine_content_type(path, &mime_type, &content);
718
719    // Chunk content
720    let chunk_outputs = chunkers
721        .chunk(&content, &content_type, &config.chunk_config)
722        .await
723        .map_err(Error::Chunking)?;
724
725    if chunk_outputs.is_empty() {
726        return Ok(0);
727    }
728
729    // Prepare texts for embedding
730    let texts: Vec<&str> = chunk_outputs.iter().map(|c| c.content.as_str()).collect();
731
732    // Generate embeddings
733    let embeddings = embedder
734        .embed_batch(&texts, &config.embed_config)
735        .await
736        .map_err(Error::Embedding)?;
737
738    // Create Chunk objects
739    let file_id = Uuid::new_v4();
740    let now = Utc::now();
741    let model_name = embedder.model_name().to_string();
742
743    let chunks: Vec<Chunk> = chunk_outputs
744        .into_iter()
745        .zip(embeddings)
746        .enumerate()
747        .map(|(idx, (output, emb_output))| {
748            build_chunk(
749                file_id,
750                path,
751                root,
752                idx as u32,
753                output,
754                emb_output.embedding,
755                &content_type,
756                &mime_type,
757                &model_name,
758                now,
759            )
760        })
761        .collect();
762
763    let chunk_count = chunks.len() as u32;
764
765    // Delete old chunks for this file
766    let _ = store.delete_by_file_path(path).await;
767
768    // Store chunks
769    store.upsert_chunks(&chunks).await.map_err(Error::Store)?;
770
771    // Store file record
772    let file_record = FileRecord {
773        id: file_id,
774        path: path.to_path_buf(),
775        size_bytes: metadata.len(),
776        mime_type,
777        content_hash,
778        modified_at: metadata
779            .modified()
780            .map_or_else(|_| now, chrono::DateTime::<Utc>::from),
781        indexed_at: Some(now),
782        chunk_count,
783        status: FileStatus::Indexed,
784        error_message: None,
785    };
786
787    store
788        .upsert_file(&file_record)
789        .await
790        .map_err(Error::Store)?;
791
792    Ok(chunk_count)
793}
794
795/// Compute blake3 hash of file content.
796async fn compute_hash(path: &Path) -> Result<String> {
797    let content = tokio::fs::read(path)
798        .await
799        .map_err(|e| Error::Other(format!("Failed to read file: {e}")))?;
800
801    let hash = blake3::hash(&content);
802    Ok(hash.to_hex().to_string())
803}
804
805/// Determine content type from path, MIME type, and extracted content.
806fn determine_content_type(
807    path: &Path,
808    mime_type: &str,
809    content: &ragfs_core::ExtractedContent,
810) -> ContentType {
811    // Check for code files
812    let code_extensions = [
813        ("rs", "rust"),
814        ("py", "python"),
815        ("js", "javascript"),
816        ("ts", "typescript"),
817        ("tsx", "typescript"),
818        ("jsx", "javascript"),
819        ("java", "java"),
820        ("go", "go"),
821        ("c", "c"),
822        ("cpp", "cpp"),
823        ("h", "c"),
824        ("hpp", "cpp"),
825        ("rb", "ruby"),
826        ("php", "php"),
827        ("swift", "swift"),
828        ("kt", "kotlin"),
829        ("scala", "scala"),
830        ("ex", "elixir"),
831        ("hs", "haskell"),
832        ("ml", "ocaml"),
833        ("lua", "lua"),
834        ("sh", "bash"),
835        ("sql", "sql"),
836    ];
837
838    if let Some(ext) = path.extension().and_then(|e| e.to_str()) {
839        let ext_lower = ext.to_lowercase();
840        for (code_ext, lang) in &code_extensions {
841            if ext_lower == *code_ext {
842                return ContentType::Code {
843                    language: (*lang).to_string(),
844                    symbol: None,
845                };
846            }
847        }
848    }
849
850    // Check for markdown
851    if mime_type.contains("markdown")
852        || path
853            .extension()
854            .is_some_and(|e| e == "md" || e == "markdown")
855    {
856        return ContentType::Markdown;
857    }
858
859    // Check language from extraction metadata
860    if let Some(ref lang) = content.metadata.language
861        && !lang.is_empty()
862    {
863        return ContentType::Code {
864            language: lang.clone(),
865            symbol: None,
866        };
867    }
868
869    ContentType::Text
870}
871
872/// Build a Chunk from chunk output and embedding.
873fn build_chunk(
874    file_id: Uuid,
875    file_path: &Path,
876    root: &Path,
877    chunk_index: u32,
878    output: ChunkOutput,
879    embedding: Vec<f32>,
880    content_type: &ContentType,
881    mime_type: &str,
882    model_name: &str,
883    now: chrono::DateTime<Utc>,
884) -> Chunk {
885    let scope = DirectoryScope::from_paths(file_path, Some(root));
886    Chunk {
887        id: Uuid::new_v4(),
888        file_id,
889        file_path: file_path.to_path_buf(),
890        content: output.content,
891        content_type: content_type.clone(),
892        mime_type: Some(mime_type.to_string()),
893        chunk_index,
894        byte_range: output.byte_range,
895        line_range: output.line_range,
896        parent_chunk_id: None,
897        depth: output.depth,
898        embedding: Some(embedding),
899        dir_path: scope.dir_path,
900        dir_depth: scope.dir_depth,
901        path_components: scope.path_components,
902        metadata: ChunkMetadata {
903            embedding_model: Some(model_name.to_string()),
904            indexed_at: Some(now),
905            token_count: None,
906            extra: Default::default(),
907        },
908    }
909}
910
911#[async_trait::async_trait]
912impl Indexer for IndexerService {
913    async fn watch(&self, path: &Path) -> Result<()> {
914        let mut w = self.watcher.write().await;
915        if let Some(ref mut watcher) = *w {
916            watcher
917                .watch(path)
918                .map_err(|e| Error::Other(format!("watch error: {e}")))?;
919        }
920        Ok(())
921    }
922
923    async fn stop(&self) -> Result<()> {
924        let mut running = self.running.write().await;
925        *running = false;
926        info!("Indexer stopped");
927        Ok(())
928    }
929
930    async fn index(&self, path: &Path, force: bool) -> Result<()> {
931        if force {
932            // Delete existing and reindex
933            let _ = self.store.delete_by_file_path(path).await;
934        }
935
936        self.pending.fetch_add(1, Ordering::SeqCst);
937        if let Err(e) = self
938            .event_tx
939            .send(FileEvent::Modified(path.to_path_buf()))
940            .await
941        {
942            self.pending.fetch_sub(1, Ordering::SeqCst);
943            return Err(Error::Other(format!("send error: {e}")));
944        }
945        Ok(())
946    }
947
948    async fn stats(&self) -> Result<IndexStats> {
949        Ok(self.stats.read().await.clone())
950    }
951
952    async fn needs_reindex(&self, path: &Path) -> Result<bool> {
953        // Check if file exists in store and compare hash
954        match self.store.get_file(path).await {
955            Ok(Some(record)) => {
956                let current_hash = compute_hash(path).await?;
957                Ok(record.content_hash != current_hash)
958            }
959            Ok(None) => Ok(true),
960            Err(e) => Err(Error::Store(e)),
961        }
962    }
963}
964
965#[cfg(test)]
966mod tests {
967    use super::*;
968    use ragfs_chunker::{ChunkerRegistry, FixedSizeChunker};
969    use ragfs_core::{
970        ContentMetadataInfo, EmbedError, Embedder, EmbeddingConfig, EmbeddingOutput, Modality,
971        SearchQuery, SearchResult, StoreError, StoreStats,
972    };
973    use ragfs_embed::EmbedderPool;
974    use ragfs_extract::ExtractorRegistry;
975    use std::collections::HashMap;
976    use tempfile::tempdir;
977
978    const TEST_DIM: usize = 384;
979
980    // ==================== Mock Embedder ====================
981
982    struct MockEmbedder {
983        dimension: usize,
984    }
985
986    impl MockEmbedder {
987        fn new(dimension: usize) -> Self {
988            Self { dimension }
989        }
990    }
991
992    #[async_trait::async_trait]
993    impl Embedder for MockEmbedder {
994        fn model_name(&self) -> &str {
995            "mock-embedder"
996        }
997
998        fn dimension(&self) -> usize {
999            self.dimension
1000        }
1001
1002        fn max_tokens(&self) -> usize {
1003            512
1004        }
1005
1006        fn modalities(&self) -> &[Modality] {
1007            &[Modality::Text]
1008        }
1009
1010        async fn embed_text(
1011            &self,
1012            texts: &[&str],
1013            _config: &EmbeddingConfig,
1014        ) -> std::result::Result<Vec<EmbeddingOutput>, EmbedError> {
1015            Ok(texts
1016                .iter()
1017                .map(|_| EmbeddingOutput {
1018                    embedding: vec![0.1; self.dimension],
1019                    token_count: 10,
1020                })
1021                .collect())
1022        }
1023
1024        async fn embed_query(
1025            &self,
1026            _query: &str,
1027            _config: &EmbeddingConfig,
1028        ) -> std::result::Result<EmbeddingOutput, EmbedError> {
1029            Ok(EmbeddingOutput {
1030                embedding: vec![0.1; self.dimension],
1031                token_count: 10,
1032            })
1033        }
1034    }
1035
1036    // ==================== Mock VectorStore ====================
1037
1038    struct MockStore {
1039        chunks: Arc<RwLock<Vec<Chunk>>>,
1040        files: Arc<RwLock<HashMap<PathBuf, FileRecord>>>,
1041    }
1042
1043    impl MockStore {
1044        fn new() -> Self {
1045            Self {
1046                chunks: Arc::new(RwLock::new(Vec::new())),
1047                files: Arc::new(RwLock::new(HashMap::new())),
1048            }
1049        }
1050    }
1051
1052    #[async_trait::async_trait]
1053    impl VectorStore for MockStore {
1054        async fn init(&self) -> std::result::Result<(), StoreError> {
1055            Ok(())
1056        }
1057
1058        async fn upsert_chunks(&self, chunks: &[Chunk]) -> std::result::Result<(), StoreError> {
1059            let mut store = self.chunks.write().await;
1060            for chunk in chunks {
1061                store.push(chunk.clone());
1062            }
1063            Ok(())
1064        }
1065
1066        async fn search(
1067            &self,
1068            _query: SearchQuery,
1069        ) -> std::result::Result<Vec<SearchResult>, StoreError> {
1070            Ok(vec![])
1071        }
1072
1073        async fn hybrid_search(
1074            &self,
1075            _query: SearchQuery,
1076        ) -> std::result::Result<Vec<SearchResult>, StoreError> {
1077            Ok(vec![])
1078        }
1079
1080        async fn delete_by_file_path(&self, path: &Path) -> std::result::Result<u64, StoreError> {
1081            let mut chunks = self.chunks.write().await;
1082            let initial_len = chunks.len();
1083            chunks.retain(|c| c.file_path != path);
1084            let deleted = initial_len - chunks.len();
1085            let mut files = self.files.write().await;
1086            files.remove(path);
1087            Ok(deleted as u64)
1088        }
1089
1090        async fn get_file(
1091            &self,
1092            path: &Path,
1093        ) -> std::result::Result<Option<FileRecord>, StoreError> {
1094            let files = self.files.read().await;
1095            Ok(files.get(path).cloned())
1096        }
1097
1098        async fn upsert_file(&self, record: &FileRecord) -> std::result::Result<(), StoreError> {
1099            let mut files = self.files.write().await;
1100            files.insert(record.path.clone(), record.clone());
1101            Ok(())
1102        }
1103
1104        async fn stats(&self) -> std::result::Result<StoreStats, StoreError> {
1105            let chunks = self.chunks.read().await;
1106            let files = self.files.read().await;
1107            Ok(StoreStats {
1108                total_chunks: chunks.len() as u64,
1109                total_files: files.len() as u64,
1110                index_size_bytes: 0,
1111                last_updated: None,
1112            })
1113        }
1114
1115        async fn update_file_path(
1116            &self,
1117            _old_path: &Path,
1118            _new_path: &Path,
1119        ) -> std::result::Result<u64, StoreError> {
1120            Ok(0)
1121        }
1122
1123        async fn get_chunks_for_file(
1124            &self,
1125            path: &Path,
1126        ) -> std::result::Result<Vec<Chunk>, StoreError> {
1127            let chunks = self.chunks.read().await;
1128            Ok(chunks
1129                .iter()
1130                .filter(|c| c.file_path == path)
1131                .cloned()
1132                .collect())
1133        }
1134
1135        async fn get_all_chunks(&self) -> std::result::Result<Vec<Chunk>, StoreError> {
1136            let chunks = self.chunks.read().await;
1137            Ok(chunks.clone())
1138        }
1139
1140        async fn get_all_files(&self) -> std::result::Result<Vec<FileRecord>, StoreError> {
1141            let files = self.files.read().await;
1142            Ok(files.values().cloned().collect())
1143        }
1144    }
1145
1146    // ==================== Helper function tests ====================
1147
1148    #[test]
1149    fn test_determine_content_type_rust() {
1150        let path = PathBuf::from("/test/file.rs");
1151        let content = ragfs_core::ExtractedContent {
1152            text: "fn main() {}".to_string(),
1153            elements: vec![],
1154            images: vec![],
1155            metadata: ContentMetadataInfo::default(),
1156        };
1157
1158        let result = determine_content_type(&path, "text/x-rust", &content);
1159        match result {
1160            ContentType::Code { language, .. } => assert_eq!(language, "rust"),
1161            _ => panic!("Expected Code content type"),
1162        }
1163    }
1164
1165    #[test]
1166    fn test_determine_content_type_python() {
1167        let path = PathBuf::from("/test/script.py");
1168        let content = ragfs_core::ExtractedContent {
1169            text: "def main(): pass".to_string(),
1170            elements: vec![],
1171            images: vec![],
1172            metadata: ContentMetadataInfo::default(),
1173        };
1174
1175        let result = determine_content_type(&path, "text/x-python", &content);
1176        match result {
1177            ContentType::Code { language, .. } => assert_eq!(language, "python"),
1178            _ => panic!("Expected Code content type"),
1179        }
1180    }
1181
1182    #[test]
1183    fn test_determine_content_type_markdown() {
1184        let path = PathBuf::from("/test/readme.md");
1185        let content = ragfs_core::ExtractedContent {
1186            text: "# Hello".to_string(),
1187            elements: vec![],
1188            images: vec![],
1189            metadata: ContentMetadataInfo::default(),
1190        };
1191
1192        let result = determine_content_type(&path, "text/markdown", &content);
1193        assert!(matches!(result, ContentType::Markdown));
1194    }
1195
1196    #[test]
1197    fn test_determine_content_type_text() {
1198        let path = PathBuf::from("/test/notes.txt");
1199        let content = ragfs_core::ExtractedContent {
1200            text: "Some text content".to_string(),
1201            elements: vec![],
1202            images: vec![],
1203            metadata: ContentMetadataInfo::default(),
1204        };
1205
1206        let result = determine_content_type(&path, "text/plain", &content);
1207        assert!(matches!(result, ContentType::Text));
1208    }
1209
1210    #[test]
1211    fn test_determine_content_type_javascript() {
1212        let path = PathBuf::from("/test/app.js");
1213        let content = ragfs_core::ExtractedContent {
1214            text: "const x = 1;".to_string(),
1215            elements: vec![],
1216            images: vec![],
1217            metadata: ContentMetadataInfo::default(),
1218        };
1219
1220        let result = determine_content_type(&path, "application/javascript", &content);
1221        match result {
1222            ContentType::Code { language, .. } => assert_eq!(language, "javascript"),
1223            _ => panic!("Expected Code content type"),
1224        }
1225    }
1226
1227    #[tokio::test]
1228    async fn test_compute_hash() {
1229        let temp_dir = tempdir().unwrap();
1230        let file_path = temp_dir.path().join("test.txt");
1231        std::fs::write(&file_path, "test content").unwrap();
1232
1233        let hash = compute_hash(&file_path).await.unwrap();
1234        assert!(!hash.is_empty());
1235        assert_eq!(hash.len(), 64); // blake3 hex is 64 chars
1236
1237        // Same content should produce same hash
1238        let hash2 = compute_hash(&file_path).await.unwrap();
1239        assert_eq!(hash, hash2);
1240    }
1241
1242    #[tokio::test]
1243    async fn test_compute_hash_different_content() {
1244        let temp_dir = tempdir().unwrap();
1245
1246        let file1 = temp_dir.path().join("file1.txt");
1247        let file2 = temp_dir.path().join("file2.txt");
1248
1249        std::fs::write(&file1, "content 1").unwrap();
1250        std::fs::write(&file2, "content 2").unwrap();
1251
1252        let hash1 = compute_hash(&file1).await.unwrap();
1253        let hash2 = compute_hash(&file2).await.unwrap();
1254
1255        assert_ne!(hash1, hash2);
1256    }
1257
1258    #[test]
1259    fn test_build_chunk() {
1260        use ragfs_core::ChunkOutputMetadata;
1261
1262        let file_id = Uuid::new_v4();
1263        let root = PathBuf::from("/project");
1264        let file_path = PathBuf::from("/project/src/auth/login.rs");
1265        let chunk_output = ChunkOutput {
1266            content: "Test chunk content".to_string(),
1267            byte_range: 0..18,
1268            line_range: Some(0..1),
1269            parent_index: None,
1270            depth: 0,
1271            metadata: ChunkOutputMetadata::default(),
1272        };
1273        let embedding = vec![0.1; TEST_DIM];
1274        let content_type = ContentType::Text;
1275        let mime_type = "text/plain";
1276        let model_name = "test-model";
1277        let now = Utc::now();
1278
1279        let chunk = build_chunk(
1280            file_id,
1281            &file_path,
1282            &root,
1283            0,
1284            chunk_output,
1285            embedding.clone(),
1286            &content_type,
1287            mime_type,
1288            model_name,
1289            now,
1290        );
1291
1292        assert_eq!(chunk.file_id, file_id);
1293        assert_eq!(chunk.file_path, file_path);
1294        assert_eq!(chunk.chunk_index, 0);
1295        assert_eq!(chunk.content, "Test chunk content");
1296        assert_eq!(chunk.embedding, Some(embedding));
1297        assert_eq!(chunk.mime_type, Some("text/plain".to_string()));
1298        assert!(matches!(chunk.content_type, ContentType::Text));
1299        assert_eq!(chunk.dir_path, "src/auth");
1300        assert_eq!(chunk.dir_depth, 2);
1301        assert_eq!(chunk.path_components, "src,auth,login.rs");
1302        assert!(!chunk.dir_path.starts_with('/'));
1303    }
1304
1305    // ==================== IndexerService tests ====================
1306
1307    fn create_test_indexer(store: Arc<dyn VectorStore>) -> IndexerService {
1308        use ragfs_extract::TextExtractor;
1309
1310        let mut extractors = ExtractorRegistry::new();
1311        extractors.register("text", TextExtractor::new());
1312        let extractors = Arc::new(extractors);
1313
1314        let mut chunkers = ChunkerRegistry::new();
1315        chunkers.register("fixed", FixedSizeChunker::new());
1316        chunkers.set_default("fixed");
1317        let chunkers = Arc::new(chunkers);
1318
1319        let embedder = Arc::new(MockEmbedder::new(TEST_DIM));
1320        let embedder_pool = Arc::new(EmbedderPool::new(embedder, 1));
1321
1322        let config = IndexerConfig::default();
1323
1324        IndexerService::new(
1325            PathBuf::from("/tmp"),
1326            store,
1327            extractors,
1328            chunkers,
1329            embedder_pool,
1330            config,
1331        )
1332    }
1333
1334    #[tokio::test]
1335    async fn test_process_single_file() {
1336        let temp_dir = tempdir().unwrap();
1337        let file_path = temp_dir.path().join("test.txt");
1338        std::fs::write(&file_path, "This is test content for indexing.").unwrap();
1339
1340        let store = Arc::new(MockStore::new());
1341        let indexer = create_test_indexer(Arc::clone(&store) as Arc<dyn VectorStore>);
1342
1343        let chunk_count = indexer.process_single(&file_path).await.unwrap();
1344
1345        assert!(chunk_count > 0, "Should have created at least one chunk");
1346
1347        // Verify chunks were stored
1348        let stored_chunks = store.chunks.read().await;
1349        assert!(!stored_chunks.is_empty());
1350        assert!(stored_chunks.iter().all(|c| c.file_path == file_path));
1351
1352        // Verify file record was stored
1353        let files = store.files.read().await;
1354        assert!(files.contains_key(&file_path));
1355        let file_record = files.get(&file_path).unwrap();
1356        assert_eq!(file_record.chunk_count, chunk_count);
1357        assert_eq!(file_record.status, FileStatus::Indexed);
1358    }
1359
1360    #[tokio::test]
1361    async fn test_process_single_skip_already_indexed() {
1362        let temp_dir = tempdir().unwrap();
1363        let file_path = temp_dir.path().join("test.txt");
1364        std::fs::write(&file_path, "Test content").unwrap();
1365
1366        let store = Arc::new(MockStore::new());
1367        let indexer = create_test_indexer(Arc::clone(&store) as Arc<dyn VectorStore>);
1368
1369        // First indexing
1370        let chunk_count1 = indexer.process_single(&file_path).await.unwrap();
1371
1372        // Count chunks after first indexing
1373        let _chunks_after_first = store.chunks.read().await.len();
1374
1375        // Second indexing (should skip because content hasn't changed)
1376        let chunk_count2 = indexer.process_single(&file_path).await.unwrap();
1377
1378        assert_eq!(chunk_count1, chunk_count2);
1379
1380        // Chunk count should be same (old chunks deleted and new ones added if reindexed,
1381        // or no change if skipped)
1382        let chunks_after_second = store.chunks.read().await.len();
1383        // Note: process_file deletes old chunks before adding new ones, so count may vary
1384        assert!(chunks_after_second > 0);
1385    }
1386
1387    #[tokio::test]
1388    async fn test_needs_reindex_new_file() {
1389        let temp_dir = tempdir().unwrap();
1390        let file_path = temp_dir.path().join("new_file.txt");
1391        std::fs::write(&file_path, "New content").unwrap();
1392
1393        let store = Arc::new(MockStore::new());
1394        let indexer = create_test_indexer(Arc::clone(&store) as Arc<dyn VectorStore>);
1395
1396        // File not in store should need reindex
1397        let needs = indexer.needs_reindex(&file_path).await.unwrap();
1398        assert!(needs);
1399    }
1400
1401    #[tokio::test]
1402    async fn test_needs_reindex_unchanged_file() {
1403        let temp_dir = tempdir().unwrap();
1404        let file_path = temp_dir.path().join("unchanged.txt");
1405        std::fs::write(&file_path, "Unchanged content").unwrap();
1406
1407        let store = Arc::new(MockStore::new());
1408        let indexer = create_test_indexer(Arc::clone(&store) as Arc<dyn VectorStore>);
1409
1410        // Index the file first
1411        indexer.process_single(&file_path).await.unwrap();
1412
1413        // Same content should not need reindex
1414        let needs = indexer.needs_reindex(&file_path).await.unwrap();
1415        assert!(!needs);
1416    }
1417
1418    #[tokio::test]
1419    async fn test_needs_reindex_modified_file() {
1420        let temp_dir = tempdir().unwrap();
1421        let file_path = temp_dir.path().join("modified.txt");
1422        std::fs::write(&file_path, "Original content").unwrap();
1423
1424        let store = Arc::new(MockStore::new());
1425        let indexer = create_test_indexer(Arc::clone(&store) as Arc<dyn VectorStore>);
1426
1427        // Index the file
1428        indexer.process_single(&file_path).await.unwrap();
1429
1430        // Modify the file
1431        std::fs::write(&file_path, "Modified content - different!").unwrap();
1432
1433        // Modified file should need reindex
1434        let needs = indexer.needs_reindex(&file_path).await.unwrap();
1435        assert!(needs);
1436    }
1437
1438    #[tokio::test]
1439    async fn test_reindex_path_file() {
1440        let temp_dir = tempdir().unwrap();
1441        let file_path = temp_dir.path().join("reindex_test.txt");
1442        std::fs::write(&file_path, "Original content for reindex test").unwrap();
1443
1444        let store = Arc::new(MockStore::new());
1445        let indexer = create_test_indexer(Arc::clone(&store) as Arc<dyn VectorStore>);
1446
1447        // Index initially
1448        indexer.process_single(&file_path).await.unwrap();
1449
1450        let initial_chunks = store.chunks.read().await.len();
1451        assert!(initial_chunks > 0);
1452
1453        // Modify content
1454        std::fs::write(&file_path, "New content after modification").unwrap();
1455
1456        // Reindex
1457        indexer.reindex_path(&file_path).await.unwrap();
1458
1459        // Verify file was reindexed
1460        let files = store.files.read().await;
1461        assert!(files.contains_key(&file_path));
1462    }
1463
1464    #[tokio::test]
1465    async fn test_empty_file_returns_zero_chunks() {
1466        let temp_dir = tempdir().unwrap();
1467        let file_path = temp_dir.path().join("empty.txt");
1468        std::fs::write(&file_path, "").unwrap();
1469
1470        let store = Arc::new(MockStore::new());
1471        let indexer = create_test_indexer(Arc::clone(&store) as Arc<dyn VectorStore>);
1472
1473        let chunk_count = indexer.process_single(&file_path).await.unwrap();
1474
1475        assert_eq!(chunk_count, 0, "Empty file should produce zero chunks");
1476    }
1477
1478    #[test]
1479    fn test_indexer_config_default() {
1480        let config = IndexerConfig::default();
1481
1482        assert!(!config.include_patterns.is_empty());
1483        assert!(config.include_patterns.contains(&"**/*".to_string()));
1484
1485        assert!(!config.exclude_patterns.is_empty());
1486        assert!(config.exclude_patterns.contains(&"**/.git/**".to_string()));
1487        assert!(
1488            config
1489                .exclude_patterns
1490                .contains(&"**/node_modules/**".to_string())
1491        );
1492        assert!(
1493            config
1494                .exclude_patterns
1495                .contains(&"**/target/**".to_string())
1496        );
1497        assert_eq!(config.debounce_ms, 500);
1498        assert_eq!(config.max_file_size, 52_428_800);
1499        assert!(!config.force);
1500    }
1501
1502    #[test]
1503    fn test_path_matches_exclude_and_include_patterns() {
1504        assert!(path_matches_pattern("/proj/src/main.rs", "**/*"));
1505        assert!(path_matches_pattern(
1506            "/proj/node_modules/pkg/index.js",
1507            "**/node_modules/**"
1508        ));
1509        assert!(path_matches_pattern("/proj/.git/config", "**/.git/**"));
1510        assert!(path_matches_pattern("/proj/foo.pyc", "**/*.pyc"));
1511        assert!(path_matches_pattern("/proj/.env", "**/.env"));
1512        assert!(path_matches_pattern("/proj/.env", "**/.*"));
1513        assert!(
1514            !path_matches_pattern("/tmp/.tmpABC/file.txt", "**/.*"),
1515            "hidden parent dirs must not exclude a visible file"
1516        );
1517        assert!(!path_matches_pattern(
1518            "/proj/src/main.rs",
1519            "**/node_modules/**"
1520        ));
1521        assert!(path_matches_pattern("/proj/src/main.rs", "**/*.rs"));
1522        assert!(!path_matches_pattern("/proj/readme.md", "**/*.rs"));
1523        assert!(
1524            path_matches_pattern("src/lib.rs", "src/**/*.rs"),
1525            "nested ** must match a file directly under the prefix"
1526        );
1527        assert!(path_matches_pattern("/proj/src/lib.rs", "src/**/*.rs"));
1528        assert!(path_matches_pattern(
1529            "/proj/src/nested/lib.rs",
1530            "src/**/*.rs"
1531        ));
1532        assert!(path_matches_pattern("test_unit.rs", "test_*.rs"));
1533        assert!(path_matches_pattern(
1534            "/proj/tests/test_unit.rs",
1535            "test_*.rs"
1536        ));
1537        assert!(!path_matches_pattern("/proj/src/lib.rs", "test_*.rs"));
1538
1539        let config = IndexerConfig {
1540            include_patterns: vec!["**/*.rs".to_string()],
1541            exclude_patterns: vec!["**/target/**".to_string()],
1542            ..Default::default()
1543        };
1544        assert!(config.should_process_path(Path::new("/proj/src/lib.rs")));
1545        assert!(!config.should_process_path(Path::new("/proj/src/lib.md")));
1546        assert!(!config.should_process_path(Path::new("/proj/target/debug/lib.rs")));
1547    }
1548
1549    fn create_test_indexer_with_config(
1550        store: Arc<dyn VectorStore>,
1551        config: IndexerConfig,
1552    ) -> IndexerService {
1553        use ragfs_extract::TextExtractor;
1554
1555        let mut extractors = ExtractorRegistry::new();
1556        extractors.register("text", TextExtractor::new());
1557        let extractors = Arc::new(extractors);
1558
1559        let mut chunkers = ChunkerRegistry::new();
1560        chunkers.register("fixed", FixedSizeChunker::new());
1561        chunkers.set_default("fixed");
1562        let chunkers = Arc::new(chunkers);
1563
1564        let embedder = Arc::new(MockEmbedder::new(TEST_DIM));
1565        let embedder_pool = Arc::new(EmbedderPool::new(embedder, 1));
1566
1567        IndexerService::new(
1568            PathBuf::from("/tmp"),
1569            store,
1570            extractors,
1571            chunkers,
1572            embedder_pool,
1573            config,
1574        )
1575    }
1576
1577    #[tokio::test]
1578    async fn test_max_file_size_skips_large_files() {
1579        let temp_dir = tempdir().unwrap();
1580        let file_path = temp_dir.path().join("big.txt");
1581        std::fs::write(&file_path, "this file is definitely larger than 8 bytes").unwrap();
1582
1583        let store = Arc::new(MockStore::new());
1584        let config = IndexerConfig {
1585            max_file_size: 8,
1586            ..Default::default()
1587        };
1588        let indexer =
1589            create_test_indexer_with_config(Arc::clone(&store) as Arc<dyn VectorStore>, config);
1590
1591        let chunk_count = indexer.process_single(&file_path).await.unwrap();
1592        assert_eq!(chunk_count, 0);
1593        assert!(store.files.read().await.is_empty());
1594        assert!(store.chunks.read().await.is_empty());
1595    }
1596
1597    #[tokio::test]
1598    async fn test_oversized_previously_indexed_file_is_removed() {
1599        let temp_dir = tempdir().unwrap();
1600        let file_path = temp_dir.path().join("grow.txt");
1601        std::fs::write(&file_path, "small").unwrap();
1602
1603        let store = Arc::new(MockStore::new());
1604        let indexer = create_test_indexer(Arc::clone(&store) as Arc<dyn VectorStore>);
1605        let chunk_count = indexer.process_single(&file_path).await.unwrap();
1606        assert!(chunk_count > 0);
1607        assert!(store.files.read().await.contains_key(&file_path));
1608        assert!(!store.chunks.read().await.is_empty());
1609
1610        std::fs::write(&file_path, "this file is now much larger than eight bytes").unwrap();
1611
1612        let limited = IndexerConfig {
1613            max_file_size: 8,
1614            ..Default::default()
1615        };
1616        let limited_indexer =
1617            create_test_indexer_with_config(Arc::clone(&store) as Arc<dyn VectorStore>, limited);
1618        let skipped = limited_indexer.process_single(&file_path).await.unwrap();
1619        assert_eq!(skipped, 0);
1620        assert!(
1621            store.files.read().await.is_empty(),
1622            "growing past max_file_size must drop the file record"
1623        );
1624        assert!(
1625            store.chunks.read().await.is_empty(),
1626            "growing past max_file_size must drop indexed chunks"
1627        );
1628    }
1629
1630    #[tokio::test]
1631    async fn test_force_reindexes_unchanged_file() {
1632        let temp_dir = tempdir().unwrap();
1633        let file_path = temp_dir.path().join("stable.txt");
1634        std::fs::write(&file_path, "unchanged payload").unwrap();
1635
1636        let store = Arc::new(MockStore::new());
1637        let indexer = create_test_indexer(Arc::clone(&store) as Arc<dyn VectorStore>);
1638        indexer.process_single(&file_path).await.unwrap();
1639        let first_id = store.files.read().await.get(&file_path).unwrap().id;
1640
1641        // Incremental path skips when the hash matches.
1642        indexer.process_single(&file_path).await.unwrap();
1643        let skipped_id = store.files.read().await.get(&file_path).unwrap().id;
1644        assert_eq!(first_id, skipped_id);
1645
1646        let force_config = IndexerConfig {
1647            force: true,
1648            ..Default::default()
1649        };
1650        let force_indexer = create_test_indexer_with_config(
1651            Arc::clone(&store) as Arc<dyn VectorStore>,
1652            force_config,
1653        );
1654        force_indexer.process_single(&file_path).await.unwrap();
1655        let forced_id = store.files.read().await.get(&file_path).unwrap().id;
1656        assert_ne!(
1657            first_id, forced_id,
1658            "--force must replace the existing file record"
1659        );
1660    }
1661
1662    #[tokio::test]
1663    async fn test_wait_until_idle_indexes_all_files() {
1664        let temp_dir = tempdir().unwrap();
1665        std::fs::write(temp_dir.path().join("a.txt"), "alpha file contents").unwrap();
1666        std::fs::write(temp_dir.path().join("b.txt"), "bravo file contents").unwrap();
1667        std::fs::create_dir(temp_dir.path().join("skip_me")).unwrap();
1668        std::fs::write(temp_dir.path().join("skip_me").join("c.bin"), "nope").unwrap();
1669
1670        let store = Arc::new(MockStore::new());
1671        let mut extractors = ExtractorRegistry::new();
1672        extractors.register("text", ragfs_extract::TextExtractor::new());
1673        let mut chunkers = ChunkerRegistry::new();
1674        chunkers.register("fixed", FixedSizeChunker::new());
1675        chunkers.set_default("fixed");
1676        let embedder_pool = Arc::new(EmbedderPool::new(Arc::new(MockEmbedder::new(TEST_DIM)), 1));
1677
1678        let config = IndexerConfig {
1679            include_patterns: vec!["**/*.txt".to_string()],
1680            exclude_patterns: vec!["**/skip_me/**".to_string()],
1681            debounce_ms: 50,
1682            ..Default::default()
1683        };
1684
1685        let indexer = IndexerService::new(
1686            temp_dir.path().to_path_buf(),
1687            Arc::clone(&store) as Arc<dyn VectorStore>,
1688            Arc::new(extractors),
1689            Arc::new(chunkers),
1690            embedder_pool,
1691            config,
1692        );
1693
1694        indexer.start().await.unwrap();
1695        tokio::time::timeout(
1696            std::time::Duration::from_secs(10),
1697            indexer.wait_until_idle(),
1698        )
1699        .await
1700        .expect("indexer should become idle");
1701
1702        let files = store.get_all_files().await.unwrap();
1703        assert_eq!(files.len(), 2, "only included .txt files should be indexed");
1704        indexer.stop().await.unwrap();
1705    }
1706
1707    #[tokio::test]
1708    async fn test_chunk_config_from_indexer_is_used() {
1709        let temp_dir = tempdir().unwrap();
1710        let file_path = temp_dir.path().join("chunky.txt");
1711        let body = "word ".repeat(80);
1712        std::fs::write(&file_path, body).unwrap();
1713
1714        let store = Arc::new(MockStore::new());
1715        let config = IndexerConfig {
1716            chunk_config: ChunkConfig {
1717                target_size: 8,
1718                max_size: 16,
1719                overlap: 0,
1720                ..Default::default()
1721            },
1722            ..Default::default()
1723        };
1724        let indexer =
1725            create_test_indexer_with_config(Arc::clone(&store) as Arc<dyn VectorStore>, config);
1726
1727        let chunk_count = indexer.process_single(&file_path).await.unwrap();
1728        assert!(
1729            chunk_count > 1,
1730            "small chunk target should produce multiple chunks, got {chunk_count}"
1731        );
1732    }
1733
1734    #[tokio::test]
1735    async fn test_subscribe_receives_updates() {
1736        let store = Arc::new(MockStore::new());
1737        let indexer = create_test_indexer(Arc::clone(&store) as Arc<dyn VectorStore>);
1738
1739        let mut receiver = indexer.subscribe();
1740
1741        // Note: This test is limited because we can't easily trigger updates
1742        // without starting the full indexer service. The subscribe mechanism
1743        // is tested indirectly through the update_tx.send() calls in process_file.
1744
1745        // Verify we can create a receiver without panicking
1746        assert!(receiver.try_recv().is_err()); // No messages yet
1747    }
1748
1749    #[test]
1750    fn test_index_update_variants() {
1751        // Test that all IndexUpdate variants can be created
1752        let _indexed = IndexUpdate::FileIndexed {
1753            path: PathBuf::from("/test"),
1754            chunk_count: 5,
1755        };
1756
1757        let _removed = IndexUpdate::FileRemoved {
1758            path: PathBuf::from("/test"),
1759        };
1760
1761        let _error = IndexUpdate::FileError {
1762            path: PathBuf::from("/test"),
1763            error: "test error".to_string(),
1764        };
1765
1766        let _started = IndexUpdate::IndexingStarted {
1767            path: PathBuf::from("/test"),
1768        };
1769    }
1770}