Skip to main content

smriti/services/
scanner.rs

1//! Directory scanning service
2//!
3//! Recursively scans directories to find supported image files,
4//! extracts EXIF metadata, and stores everything in the database.
5//! Uses rayon for parallel hash computation and EXIF extraction.
6
7use std::path::{Path, PathBuf};
8use std::sync::atomic::{AtomicBool, Ordering};
9use std::sync::Arc;
10use std::time::Instant;
11
12use async_channel::{bounded, Receiver, Sender};
13use sha2::{Digest, Sha256};
14use walkdir::{DirEntry, WalkDir};
15
16use crate::db::photo_repo::{PhotoInsert, PhotoRepo};
17use crate::db::Database;
18use crate::models::MediaType;
19use crate::services::exclusions::ExclusionMatcher;
20
21/// Supported image extensions.
22///
23/// Two groups:
24///   1. Stills the `image` crate (or our HEIC arm) handles directly.
25///   2. TIFF-based RAWs, decoded via the camera's embedded full-res
26///      JPEG preview — see `services::image_io::open_image` and
27///      `services::raw_preview`. Files in this group don't need any
28///      special handling here at scan time; they go through the same
29///      hash + EXIF + thumbnail pipeline as ordinary JPEGs.
30const SUPPORTED_EXTENSIONS: &[&str] = &[
31    // Stills
32    "jpg", "jpeg", "png", "heic", "heif", "webp", "tif", "tiff", "avif", "bmp", "gif",
33    // RAWs (TIFF-based, embedded-JPEG path)
34    "nef", "cr2", "cr3", "arw", "dng", "orf", "rw2", "pef", "rwl", "srw", "raf",
35    // Videos
36    "mp4", "m4v", "mov", "webm", "mkv", "avi", "3gp", "3g2", "mts", "m2ts",
37];
38
39/// Directories to skip during scanning
40const SKIP_DIRECTORIES: &[&str] = &[
41    ".photovault",
42    ".Trash",
43    "$RECYCLE.BIN",
44    "System Volume Information",
45    ".DS_Store",
46    "Thumbs.db",
47    ".thumbnails",
48    "@eaDir", // Synology thumbnails
49];
50
51/// Minimum file size to consider (10KB)
52const MIN_FILE_SIZE: u64 = 10 * 1024;
53
54/// Batch size for database inserts
55const DB_BATCH_SIZE: usize = 100;
56
57/// How many file paths we buffer between the walker and the stub-writer.
58/// 1000 keeps memory under ~200 KB even on libraries that average 200-byte paths.
59const WALKER_CHANNEL_DEPTH: usize = 1000;
60
61/// Bytes used for the fast-hash prefix. 64 KB is enough that JPEG/HEIC
62/// markers + first scanline diverge for distinct photos, while staying
63/// well within typical disk read-ahead (so it's effectively free).
64const FAST_HASH_PREFIX_BYTES: usize = 64 * 1024;
65
66/// How often the stub writer emits a progress event (in files written).
67const STUB_PROGRESS_EVERY: u64 = 200;
68
69/// Scan progress information
70#[derive(Debug, Clone)]
71pub struct ScanProgress {
72    pub files_found: u64,
73    pub files_processed: u64,
74    pub bytes_processed: u64,
75    pub current_directory: String,
76    pub current_file: String,
77    pub errors: Vec<String>,
78    pub is_complete: bool,
79    pub elapsed_seconds: f64,
80}
81
82impl Default for ScanProgress {
83    fn default() -> Self {
84        Self {
85            files_found: 0,
86            files_processed: 0,
87            bytes_processed: 0,
88            current_directory: String::new(),
89            current_file: "Preparing...".to_string(),
90            errors: Vec::new(),
91            is_complete: false,
92            elapsed_seconds: 0.0,
93        }
94    }
95}
96
97/// Result of a completed scan -- returns the database back to the caller
98pub struct ScanResult {
99    pub database: Database,
100    pub final_progress: ScanProgress,
101}
102
103#[derive(Debug, Clone)]
104pub struct ScanReport {
105    pub files_inserted: u64,
106    pub errors: Vec<String>,
107    pub elapsed_seconds: f64,
108}
109
110/// A candidate file discovered during directory walk
111struct FileCandidate {
112    path: PathBuf,
113    relative_path: String,
114    file_name: String,
115    file_size: i64,
116    mtime: Option<i64>,
117    media_type: MediaType,
118}
119
120/// Start a scan in a background task.
121///
122/// Borrows the database via `Arc<Mutex<Database>>` so the library
123/// stays open and queryable while the scan runs. The caller provides
124/// the cancel flag (usually `job.cancel`) so pause/resume is unified
125/// across all background jobs.
126pub fn start_scan(
127    root_path: PathBuf,
128    db: Arc<tokio::sync::Mutex<Database>>,
129    cancel: Arc<AtomicBool>,
130    scan_hidden_folders: bool,
131) -> (Receiver<ScanProgress>, tokio::task::JoinHandle<ScanReport>) {
132    let (progress_tx, progress_rx) = bounded::<ScanProgress>(64);
133    let handle = tokio::task::spawn(async move {
134        let result =
135            run_scan_streaming(root_path, db, cancel, scan_hidden_folders, progress_tx).await;
136        result.unwrap_or_else(|e| ScanReport {
137            files_inserted: 0,
138            errors: vec![format!("scan failed: {e}")],
139            elapsed_seconds: 0.0,
140        })
141    });
142    (progress_rx, handle)
143}
144
145/// Cheap identity fingerprint for the scan stage.
146///
147/// Reads the first 64 KB of the file, hashes (size, mtime, prefix bytes)
148/// with SHA-256. The hex digest is stored directly in `photos.file_hash`;
149/// downstream code (thumbnails, duplicates) uses the first two hex chars
150/// as a shard subdirectory, so no prefix is added — a prefix like `fh:`
151/// would collide with that path scheme and produce invalid Windows
152/// filenames (`:` is reserved).
153///
154/// A future deferred Stage 5 worker can overwrite the column with a true
155/// full-file SHA-256 if duplicate detection ever needs byte-level
156/// certainty; both are hex digests of the same length, so no consumer
157/// has to care.
158///
159/// Collisions: in practice essentially zero for organic photo libraries.
160/// Two photos would have to be byte-identical in their first 64 KB AND
161/// have the same file size AND the same mtime. Even bursts from the same
162/// camera diverge in the first few EXIF bytes (timestamp).
163pub(crate) fn calculate_fast_hash<P: AsRef<Path>>(
164    path: P,
165    file_size: u64,
166    mtime: Option<i64>,
167) -> std::io::Result<String> {
168    use std::io::Read;
169
170    let mut file = std::fs::File::open(&path)?;
171    let mut buf = vec![0u8; FAST_HASH_PREFIX_BYTES];
172    let n = file.read(&mut buf)?;
173
174    let mut hasher = Sha256::new();
175    hasher.update(file_size.to_le_bytes());
176    hasher.update(mtime.unwrap_or(0).to_le_bytes());
177    hasher.update(&buf[..n]);
178    Ok(format!("{:x}", hasher.finalize()))
179}
180
181async fn run_scan_streaming(
182    root_path: PathBuf,
183    db: Arc<tokio::sync::Mutex<Database>>,
184    cancel: Arc<AtomicBool>,
185    scan_hidden_folders: bool,
186    progress_tx: Sender<ScanProgress>,
187) -> Result<ScanReport, String> {
188    let start = Instant::now();
189    let (paths_tx, paths_rx) = bounded::<FileCandidate>(WALKER_CHANNEL_DEPTH);
190    let exclusions = {
191        let guard = db.lock().await;
192        ExclusionMatcher::from_db(&guard.conn).unwrap_or_else(|e| {
193            tracing::warn!("scan: failed to load folder exclusions: {e}");
194            ExclusionMatcher::empty()
195        })
196    };
197
198    // ----- Producer: walker thread -----
199    let walker_cancel = cancel.clone();
200    let walker_root = root_path.clone();
201    let walker = tokio::task::spawn_blocking(move || -> u64 {
202        let mut count: u64 = 0;
203        let walker = WalkDir::new(&walker_root)
204            .follow_links(false)
205            .into_iter()
206            .filter_entry(|e| !should_skip(e, &walker_root, scan_hidden_folders, &exclusions));
207
208        for entry in walker {
209            if walker_cancel.load(Ordering::Relaxed) {
210                tracing::info!("Scan cancelled during walk");
211                break;
212            }
213            let Ok(entry) = entry else { continue };
214            if !entry.file_type().is_file() {
215                continue;
216            }
217            if !is_supported_file(&entry) {
218                continue;
219            }
220            let media_type = media_type_for_path(entry.path()).unwrap_or_default();
221            let Ok(metadata) = entry.metadata() else {
222                continue;
223            };
224            if metadata.len() < MIN_FILE_SIZE {
225                continue;
226            }
227
228            let relative_path = entry
229                .path()
230                .strip_prefix(&walker_root)
231                .map(crate::services::path_util::relative_path_for_storage)
232                .unwrap_or_else(|_| {
233                    crate::services::path_util::relative_path_for_storage(entry.path())
234                });
235
236            let mtime = metadata
237                .modified()
238                .ok()
239                .and_then(|t| t.duration_since(std::time::UNIX_EPOCH).ok())
240                .map(|d| d.as_secs() as i64);
241
242            let candidate = FileCandidate {
243                path: entry.path().to_path_buf(),
244                relative_path,
245                file_name: entry.file_name().to_string_lossy().to_string(),
246                file_size: metadata.len() as i64,
247                mtime,
248                media_type,
249            };
250
251            // send_blocking is OK here — this whole closure is on a blocking thread.
252            if paths_tx.send_blocking(candidate).is_err() {
253                break; // receiver dropped (consumer ended early)
254            }
255            count += 1;
256        }
257        drop(paths_tx);
258        count
259    });
260
261    // ----- Consumer: stub-row writer + fast-hash batcher -----
262    let writer_cancel = cancel.clone();
263    let writer_db = db.clone();
264    let writer_progress = progress_tx.clone();
265    let writer = tokio::task::spawn(async move {
266        let mut errors: Vec<String> = Vec::new();
267        let mut buf: Vec<PhotoInsert> = Vec::with_capacity(DB_BATCH_SIZE);
268        let mut files_inserted: u64 = 0;
269
270        // Channel receives are async; the actual fast-hash + insert is blocking.
271        // We do them in chunks of DB_BATCH_SIZE so the SQLite writer is hit
272        // in bursts, not one row at a time.
273        loop {
274            if writer_cancel.load(Ordering::Relaxed) {
275                break;
276            }
277
278            let Ok(candidate) = paths_rx.recv().await else {
279                break;
280            };
281
282            // Fast-hash on this async task is fine — it's tiny I/O (64 KB).
283            // Doing it on a tokio worker thread would be marginally faster
284            // but adds spawn overhead.
285            let hash = match calculate_fast_hash(
286                &candidate.path,
287                candidate.file_size as u64,
288                candidate.mtime,
289            ) {
290                Ok(h) => h,
291                Err(e) => {
292                    errors.push(format!("fast-hash {}: {e}", candidate.path.display()));
293                    continue;
294                }
295            };
296
297            buf.push(PhotoInsert {
298                relative_path: candidate.relative_path,
299                file_name: candidate.file_name,
300                file_hash: hash,
301                file_size: candidate.file_size,
302                file_mtime: candidate.mtime,
303                // Everything else stays None / default. Phase 2 fills it in.
304                date_taken: None,
305                date_taken_source: None,
306                gps_latitude: None,
307                gps_longitude: None,
308                location_city: None,
309                location_country: None,
310                camera_make: None,
311                camera_model: None,
312                iso: None,
313                aperture: None,
314                shutter_speed: None,
315                focal_length: None,
316                lens_model: None,
317                flash: None,
318                gps_altitude: None,
319                width: None,
320                height: None,
321                orientation: 1,
322                media_type: candidate.media_type,
323                duration_ms: None,
324                video_codec: None,
325                audio_codec: None,
326                frame_rate: None,
327                bitrate: None,
328                has_audio: false,
329            });
330
331            if buf.len() >= DB_BATCH_SIZE {
332                let inserted = flush_stub_batch(&writer_db, &mut buf, &mut errors).await;
333                files_inserted += inserted;
334
335                if files_inserted.is_multiple_of(STUB_PROGRESS_EVERY) {
336                    let _ = writer_progress.try_send(ScanProgress {
337                        files_found: files_inserted,
338                        files_processed: files_inserted,
339                        bytes_processed: 0,
340                        current_directory: String::new(),
341                        current_file: format!("Indexed {files_inserted} files…"),
342                        errors: errors.clone(),
343                        is_complete: false,
344                        elapsed_seconds: start.elapsed().as_secs_f64(),
345                    });
346                }
347            }
348        }
349
350        // Drain final batch.
351        if !buf.is_empty() {
352            files_inserted += flush_stub_batch(&writer_db, &mut buf, &mut errors).await;
353        }
354        (files_inserted, errors)
355    });
356
357    let total_walked = walker.await.map_err(|e| format!("walker join: {e}"))?;
358    let (files_inserted, errors) = writer.await.map_err(|e| format!("writer join: {e}"))?;
359    tracing::info!("Phase 1 done: walked {total_walked}, inserted {files_inserted}");
360
361    let report = ScanReport {
362        files_inserted,
363        errors,
364        elapsed_seconds: start.elapsed().as_secs_f64(),
365    };
366
367    // Emit final progress; the IPC wrapper interprets is_complete=true.
368    let _ = progress_tx
369        .send(ScanProgress {
370            files_found: report.files_inserted,
371            files_processed: report.files_inserted,
372            bytes_processed: 0,
373            current_directory: String::new(),
374            current_file: format!("Indexed {} files", report.files_inserted),
375            errors: report.errors.clone(),
376            is_complete: true,
377            elapsed_seconds: report.elapsed_seconds,
378        })
379        .await;
380
381    Ok(report)
382}
383
384async fn flush_stub_batch(
385    db: &Arc<tokio::sync::Mutex<Database>>,
386    buf: &mut Vec<PhotoInsert>,
387    errors: &mut Vec<String>,
388) -> u64 {
389    let guard = db.lock().await;
390    let repo = PhotoRepo::new(&guard.conn);
391    let inserted = match repo.insert_batch_stub(buf) {
392        Ok(n) => n as u64,
393        Err(e) => {
394            errors.push(format!("stub batch insert: {e}"));
395            0
396        }
397    };
398    buf.clear();
399    inserted
400}
401
402/// Check if a directory entry should be skipped
403fn should_skip(
404    entry: &DirEntry,
405    root: &Path,
406    scan_hidden_folders: bool,
407    exclusions: &ExclusionMatcher,
408) -> bool {
409    // The selected library root itself is always traversed. A perfectly
410    // valid library can live in a dot-prefixed parent/test directory; hidden
411    // filtering applies to descendants, not to the root the user selected.
412    if entry.path() == root {
413        return false;
414    }
415    let file_name = entry.file_name().to_string_lossy();
416
417    if !scan_hidden_folders && file_name.starts_with('.') {
418        return true;
419    }
420
421    for skip in SKIP_DIRECTORIES {
422        if file_name == *skip {
423            return true;
424        }
425    }
426
427    exclusions.should_skip_path(root, entry.path())
428}
429
430/// Check if a file has a supported extension
431fn is_supported_file(entry: &DirEntry) -> bool {
432    media_type_for_path(entry.path()).is_some()
433}
434
435pub fn media_type_for_path(path: &Path) -> Option<MediaType> {
436    let lower = path.extension()?.to_str()?.to_lowercase();
437    if !SUPPORTED_EXTENSIONS.contains(&lower.as_str()) {
438        return None;
439    }
440    match lower.as_str() {
441        "mp4" | "m4v" | "mov" | "webm" | "mkv" | "avi" | "3gp" | "3g2" | "mts" | "m2ts" => {
442            Some(MediaType::Video)
443        }
444        _ => Some(MediaType::Photo),
445    }
446}
447
448/// Calculate SHA256 hash of a file (full streaming hash).
449pub fn calculate_hash<P: AsRef<Path>>(path: P) -> std::io::Result<String> {
450    use std::io::Read;
451
452    let mut file = std::fs::File::open(path)?;
453    let mut hasher = Sha256::new();
454    let mut buffer = [0u8; 65536]; // 64KB buffer
455
456    loop {
457        let bytes_read = file.read(&mut buffer)?;
458        if bytes_read == 0 {
459            break;
460        }
461        hasher.update(&buffer[..bytes_read]);
462    }
463
464    let result = hasher.finalize();
465    Ok(format!("{:x}", result))
466}
467
468#[cfg(test)]
469mod tests {
470    use super::*;
471    use std::fs::File;
472    use std::io::Write;
473    use tempfile::tempdir;
474
475    #[test]
476    fn test_should_skip_hidden() {
477        let temp = tempdir().unwrap();
478        let hidden = temp.path().join(".hidden");
479        std::fs::create_dir(&hidden).unwrap();
480        let exclusions = ExclusionMatcher::empty();
481
482        let walker = WalkDir::new(temp.path()).into_iter();
483        for entry in walker {
484            let entry = entry.unwrap();
485            if entry.file_name().to_string_lossy() == ".hidden" {
486                assert!(should_skip(&entry, temp.path(), false, &exclusions));
487                assert!(!should_skip(&entry, temp.path(), true, &exclusions));
488            }
489        }
490    }
491
492    #[test]
493    fn test_should_skip_excluded_folder_descendants() {
494        let temp = tempdir().unwrap();
495        let excluded = temp.path().join("Trips").join("Goa");
496        let similar = temp.path().join("Trips").join("Goa2");
497        std::fs::create_dir_all(&excluded).unwrap();
498        std::fs::create_dir_all(&similar).unwrap();
499        let exclusions = ExclusionMatcher::new(vec!["Trips/Goa".into()]);
500
501        let mut saw_excluded = false;
502        let mut saw_similar = false;
503        for entry in WalkDir::new(temp.path()).into_iter().filter_map(Result::ok) {
504            let name = entry.path();
505            if name == excluded {
506                saw_excluded = true;
507                assert!(should_skip(&entry, temp.path(), false, &exclusions));
508            }
509            if name == similar {
510                saw_similar = true;
511                assert!(!should_skip(&entry, temp.path(), false, &exclusions));
512            }
513        }
514        assert!(saw_excluded);
515        assert!(saw_similar);
516    }
517
518    #[test]
519    fn test_is_supported_file() {
520        let temp = tempdir().unwrap();
521
522        let jpg = temp.path().join("test.jpg");
523        File::create(&jpg).unwrap();
524
525        let txt = temp.path().join("test.txt");
526        File::create(&txt).unwrap();
527
528        let png = temp.path().join("test.PNG");
529        File::create(&png).unwrap();
530
531        let mp4 = temp.path().join("clip.MP4");
532        File::create(&mp4).unwrap();
533
534        let walker = WalkDir::new(temp.path()).into_iter();
535        for entry in walker {
536            let entry = entry.unwrap();
537            if !entry.file_type().is_file() {
538                continue;
539            }
540            let name = entry.file_name().to_string_lossy().to_string();
541            match name.as_str() {
542                "test.jpg" => assert!(is_supported_file(&entry)),
543                "test.txt" => assert!(!is_supported_file(&entry)),
544                "test.PNG" => assert!(is_supported_file(&entry)),
545                "clip.MP4" => {
546                    assert!(is_supported_file(&entry));
547                    assert_eq!(media_type_for_path(entry.path()), Some(MediaType::Video));
548                }
549                _ => {}
550            }
551        }
552    }
553
554    #[test]
555    fn test_calculate_hash() {
556        let temp = tempdir().unwrap();
557        let file_path = temp.path().join("test.bin");
558        let mut f = File::create(&file_path).unwrap();
559        f.write_all(b"hello world").unwrap();
560
561        let hash = calculate_hash(&file_path).unwrap();
562        assert_eq!(
563            hash,
564            "b94d27b9934d3e08a52e52d7da7dabfac484efe37a5380ee9088f7ace2efcde9"
565        );
566    }
567
568    #[test]
569    fn test_min_file_size_filter() {
570        const _: () = assert!(MIN_FILE_SIZE == 10 * 1024);
571    }
572
573    #[test]
574    fn test_fast_hash_changes_with_mtime() {
575        let temp = tempdir().unwrap();
576        let file_path = temp.path().join("test_mtime.bin");
577        std::fs::write(&file_path, b"some content here").unwrap();
578
579        let hash1 = calculate_fast_hash(&file_path, 100, Some(1700000000)).unwrap();
580        let hash2 = calculate_fast_hash(&file_path, 100, Some(1700000001)).unwrap();
581        assert_ne!(hash1, hash2, "fast hash must change when mtime changes");
582    }
583
584    #[test]
585    fn test_fast_hash_changes_with_size() {
586        let temp = tempdir().unwrap();
587        let file_path = temp.path().join("test_size.bin");
588        std::fs::write(&file_path, b"same prefix data").unwrap();
589
590        let hash1 = calculate_fast_hash(&file_path, 100, Some(1700000000)).unwrap();
591        let hash2 = calculate_fast_hash(&file_path, 200, Some(1700000000)).unwrap();
592        assert_ne!(hash1, hash2, "fast hash must change when file size changes");
593    }
594
595    #[test]
596    fn test_fast_hash_stable_on_repeat() {
597        let temp = tempdir().unwrap();
598        let file_path = temp.path().join("test_stable.bin");
599        std::fs::write(&file_path, b"repeatable content").unwrap();
600
601        let hash1 = calculate_fast_hash(&file_path, 500, Some(1700000000)).unwrap();
602        let hash2 = calculate_fast_hash(&file_path, 500, Some(1700000000)).unwrap();
603        assert_eq!(hash1, hash2, "fast hash must be stable on repeat calls");
604    }
605
606    #[test]
607    fn test_fast_hash_is_plain_hex_digest() {
608        // No `fh:` prefix — the hash is consumed directly as a path
609        // shard by thumbnail + duplicate code (first two chars as
610        // subdirectory). A prefix with `:` would also be illegal on
611        // Windows filenames.
612        let temp = tempdir().unwrap();
613        let file_path = temp.path().join("test_prefix.bin");
614        std::fs::write(&file_path, b"check prefix").unwrap();
615
616        let hash = calculate_fast_hash(&file_path, 12, None).unwrap();
617        assert_eq!(hash.len(), 64, "SHA-256 hex digest is 64 chars");
618        assert!(
619            hash.chars().all(|c| c.is_ascii_hexdigit()),
620            "hash must be pure ASCII hex"
621        );
622    }
623}