1use 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#[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#[derive(Debug, Clone)]
31pub struct IndexerConfig {
32 pub chunk_config: ChunkConfig,
34 pub embed_config: EmbeddingConfig,
36 pub include_patterns: Vec<String>,
38 pub exclude_patterns: Vec<String>,
40 pub debounce_ms: u64,
42 pub max_file_size: u64,
44 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 #[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#[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 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 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
160pub struct IndexerService {
162 root: PathBuf,
164 store: Arc<dyn VectorStore>,
166 extractors: Arc<ExtractorRegistry>,
168 chunkers: Arc<ChunkerRegistry>,
170 embedder: Arc<EmbedderPool>,
172 config: IndexerConfig,
174 stats: Arc<RwLock<IndexStats>>,
176 event_tx: mpsc::Sender<FileEvent>,
178 event_rx: Arc<RwLock<mpsc::Receiver<FileEvent>>>,
180 update_tx: broadcast::Sender<IndexUpdate>,
182 watcher: Arc<RwLock<Option<FileWatcher>>>,
184 running: Arc<RwLock<bool>>,
186 pending: Arc<AtomicUsize>,
188 idle_notify: Arc<Notify>,
190}
191
192impl IndexerService {
193 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 #[must_use]
225 pub fn config(&self) -> &IndexerConfig {
226 &self.config
227 }
228
229 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 #[must_use]
250 pub fn subscribe(&self) -> broadcast::Receiver<IndexUpdate> {
251 self.update_tx.subscribe()
252 }
253
254 #[must_use]
256 pub fn root(&self) -> &Path {
257 &self.root
258 }
259
260 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 self.store.init().await.map_err(Error::Store)?;
273
274 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 {
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 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 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 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 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 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 self.scan().await?;
410
411 Ok(())
412 }
413
414 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 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 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 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 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 info!("Reindexing directory {:?}", path);
504 self.reindex_directory(path).await?;
505 }
506
507 Ok(())
508 }
509
510 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 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 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 }
571 }
572 }
573 }
574
575 Ok(())
576 }
577}
578
579fn 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
652async 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 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 let content_hash = compute_hash(path).await?;
689
690 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 let mime_type = mime_guess::from_path(path)
702 .first_or_text_plain()
703 .to_string();
704
705 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 let content_type = determine_content_type(path, &mime_type, &content);
718
719 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 let texts: Vec<&str> = chunk_outputs.iter().map(|c| c.content.as_str()).collect();
731
732 let embeddings = embedder
734 .embed_batch(&texts, &config.embed_config)
735 .await
736 .map_err(Error::Embedding)?;
737
738 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 let _ = store.delete_by_file_path(path).await;
767
768 store.upsert_chunks(&chunks).await.map_err(Error::Store)?;
770
771 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
795async 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
805fn determine_content_type(
807 path: &Path,
808 mime_type: &str,
809 content: &ragfs_core::ExtractedContent,
810) -> ContentType {
811 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 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 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
872fn 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 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 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 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 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 #[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); 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 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 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 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 let chunk_count1 = indexer.process_single(&file_path).await.unwrap();
1371
1372 let _chunks_after_first = store.chunks.read().await.len();
1374
1375 let chunk_count2 = indexer.process_single(&file_path).await.unwrap();
1377
1378 assert_eq!(chunk_count1, chunk_count2);
1379
1380 let chunks_after_second = store.chunks.read().await.len();
1383 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 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 indexer.process_single(&file_path).await.unwrap();
1412
1413 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 indexer.process_single(&file_path).await.unwrap();
1429
1430 std::fs::write(&file_path, "Modified content - different!").unwrap();
1432
1433 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 indexer.process_single(&file_path).await.unwrap();
1449
1450 let initial_chunks = store.chunks.read().await.len();
1451 assert!(initial_chunks > 0);
1452
1453 std::fs::write(&file_path, "New content after modification").unwrap();
1455
1456 indexer.reindex_path(&file_path).await.unwrap();
1458
1459 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 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 assert!(receiver.try_recv().is_err()); }
1748
1749 #[test]
1750 fn test_index_update_variants() {
1751 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}