Skip to main content

smriti/services/
metadata_processor.rs

1//! Background EXIF + geocoding pass.
2//!
3//! Reads photos with `metadata_extracted = 0`, runs `ExifExtractor::extract`
4//! on each (header read only, ~10 KB per file), reverse-geocodes any GPS
5//! coordinates, and updates the row in place. Idempotent and resumable.
6
7use std::io::{Read, Seek, SeekFrom};
8use std::path::{Path, PathBuf};
9use std::sync::atomic::{AtomicBool, Ordering};
10use std::sync::Arc;
11use std::time::Instant;
12
13use async_channel::{bounded, Receiver, Sender};
14use rayon::prelude::*;
15use regex::Regex;
16use rusqlite::params;
17
18use crate::db::Database;
19use crate::models::MediaType;
20use crate::services::exif_extractor::ExifExtractor;
21use crate::services::scanner::media_type_for_path;
22use crate::services::GeocodingService;
23
24const METADATA_CHUNK_SIZE: usize = 50;
25
26#[derive(Debug, Clone)]
27pub struct MetadataProgress {
28    pub total: u64,
29    pub done: u64,
30    pub is_complete: bool,
31    pub elapsed_seconds: f64,
32    pub stage: Option<String>,
33    pub message: Option<String>,
34}
35
36pub fn start_metadata_job(
37    drive_root: PathBuf,
38    db: Arc<tokio::sync::Mutex<Database>>,
39    cancel: Arc<AtomicBool>,
40) -> (Receiver<MetadataProgress>, tokio::task::JoinHandle<()>) {
41    let (progress_tx, progress_rx) = bounded::<MetadataProgress>(32);
42    let handle = tokio::spawn(async move {
43        run_metadata_job(drive_root, db, cancel, progress_tx).await;
44    });
45    (progress_rx, handle)
46}
47
48async fn run_metadata_job(
49    drive_root: PathBuf,
50    db: Arc<tokio::sync::Mutex<Database>>,
51    cancel: Arc<AtomicBool>,
52    progress_tx: Sender<MetadataProgress>,
53) {
54    let start = Instant::now();
55
56    let total: u64 = {
57        let guard = db.lock().await;
58        guard.conn.query_row(
59            "SELECT COUNT(*) FROM photos WHERE metadata_extracted = FALSE AND is_trashed = FALSE",
60            [],
61            |row| row.get(0),
62        ).unwrap_or(0)
63    };
64    let total = total as u64;
65    let mut done = 0u64;
66    let mut hard_error: Option<String> = None;
67    let mut failed_ids: std::collections::HashSet<i64> = std::collections::HashSet::new();
68
69    let geonames_path = crate::db::geonames::geonames_db_path();
70    let geocoder = if geonames_path.exists() {
71        GeocodingService::new(&geonames_path).ok()
72    } else {
73        None
74    };
75
76    loop {
77        if cancel.load(Ordering::Relaxed) {
78            break;
79        }
80
81        let chunk: Vec<(i64, String, String)> = {
82            let guard = db.lock().await;
83            match load_metadata_chunk(&guard.conn, &failed_ids, METADATA_CHUNK_SIZE) {
84                Ok(rows) => rows,
85                Err(e) => {
86                    hard_error = Some(format!("metadata query failed: {e}"));
87                    Vec::new()
88                }
89            }
90        };
91
92        if hard_error.is_some() || chunk.is_empty() {
93            break;
94        }
95
96        let root = drive_root.clone();
97        let extracted: Vec<(i64, String, ExtractedMetadata)> = chunk
98            .par_iter()
99            .map(|(id, rel_path, file_hash)| {
100                let abs = match crate::services::path_util::safe_join_relative(&root, rel_path) {
101                    Ok(abs) => abs,
102                    Err(e) => {
103                        tracing::debug!("metadata skipped unsafe path for photo_id={id}: {e}");
104                        return (*id, file_hash.clone(), ExtractedMetadata::Skipped);
105                    }
106                };
107                let meta = match media_type_for_path(&abs).unwrap_or_default() {
108                    MediaType::Video => ExtractedMetadata::Video(VideoMetadata::from_path(&abs)),
109                    MediaType::Photo => {
110                        ExtractedMetadata::Photo(Box::new(ExifExtractor::extract(&abs)))
111                    }
112                };
113                (*id, file_hash.clone(), meta)
114            })
115            .collect();
116
117        let mut errors_this_chunk = 0usize;
118        {
119            let mut guard = db.lock().await;
120            let tx = match guard.conn.transaction() {
121                Ok(t) => t,
122                Err(e) => {
123                    tracing::error!("metadata tx start: {e}");
124                    hard_error = Some(format!("metadata database transaction failed: {e}"));
125                    break;
126                }
127            };
128
129            for (id, file_hash, meta) in &extracted {
130                if matches!(meta, ExtractedMetadata::Skipped) {
131                    errors_this_chunk += 1;
132                    failed_ids.insert(*id);
133                    continue;
134                }
135                let gps_latitude = meta.gps_latitude();
136                let gps_longitude = meta.gps_longitude();
137                let (city, country): (Option<String>, Option<String>) =
138                    match (gps_latitude, gps_longitude, &geocoder) {
139                        (Some(lat), Some(lon), Some(g)) => g
140                            .reverse_geocode(lat, lon)
141                            .map(|r| (Some(r.city), Some(r.country)))
142                            .unwrap_or((None, None)),
143                        _ => (None, None),
144                    };
145
146                let res = tx.execute(
147                    "UPDATE photos SET
148                        date_taken = ?,
149                        date_taken_source = ?,
150                        gps_latitude = ?,
151                        gps_longitude = ?,
152                        location_city = ?,
153                        location_country = ?,
154                        camera_make = ?,
155                        camera_model = ?,
156                        iso = ?,
157                        aperture = ?,
158                        shutter_speed = ?,
159                        focal_length = ?,
160                        lens_model = ?,
161                        flash = ?,
162                        gps_altitude = ?,
163                        width = ?,
164                        height = ?,
165                        orientation = ?,
166                        media_type = ?,
167                        duration_ms = ?,
168                        video_codec = ?,
169                        audio_codec = ?,
170                        frame_rate = ?,
171                        bitrate = ?,
172                        has_audio = ?,
173                        metadata_extracted = TRUE
174                     WHERE id = ? AND file_hash = ?",
175                    params![
176                        meta.date_taken()
177                            .map(|d: chrono::DateTime<chrono::Utc>| d.to_rfc3339()),
178                        meta.date_taken_source(),
179                        gps_latitude,
180                        gps_longitude,
181                        city,
182                        country,
183                        meta.camera_make(),
184                        meta.camera_model(),
185                        meta.iso(),
186                        meta.aperture(),
187                        meta.shutter_speed(),
188                        meta.focal_length(),
189                        meta.lens_model(),
190                        meta.flash(),
191                        meta.gps_altitude(),
192                        meta.width().map(|v| v as i64),
193                        meta.height().map(|v| v as i64),
194                        meta.orientation() as i64,
195                        meta.media_type().as_str(),
196                        meta.duration_ms(),
197                        meta.video_codec(),
198                        meta.audio_codec(),
199                        meta.frame_rate(),
200                        meta.bitrate(),
201                        meta.has_audio(),
202                        id,
203                        file_hash,
204                    ],
205                );
206                match res {
207                    Ok(0) => {
208                        errors_this_chunk += 1;
209                        failed_ids.insert(*id);
210                    }
211                    Ok(_) => {}
212                    Err(_) => {
213                        errors_this_chunk += 1;
214                        failed_ids.insert(*id);
215                    }
216                }
217            }
218            if let Err(e) = tx.commit() {
219                tracing::error!("metadata tx commit: {e}");
220                hard_error = Some(format!("metadata database commit failed: {e}"));
221                break;
222            }
223        }
224
225        done += (chunk.len() - errors_this_chunk) as u64;
226        let _ = progress_tx.try_send(MetadataProgress {
227            total,
228            done,
229            is_complete: false,
230            elapsed_seconds: start.elapsed().as_secs_f64(),
231            stage: Some("extract".to_string()),
232            message: None,
233        });
234    }
235
236    let _ = progress_tx
237        .send(MetadataProgress {
238            total,
239            done,
240            is_complete: hard_error.is_none(),
241            elapsed_seconds: start.elapsed().as_secs_f64(),
242            stage: hard_error
243                .as_ref()
244                .map(|_| "error".to_string())
245                .or_else(|| Some("complete".to_string())),
246            message: hard_error,
247        })
248        .await;
249}
250
251enum ExtractedMetadata {
252    Photo(Box<crate::services::exif_extractor::ImageMetadata>),
253    Video(VideoMetadata),
254    Skipped,
255}
256
257fn load_metadata_chunk(
258    conn: &rusqlite::Connection,
259    failed_ids: &std::collections::HashSet<i64>,
260    limit: usize,
261) -> rusqlite::Result<Vec<(i64, String, String)>> {
262    let mut sql = String::from(
263        "SELECT id, file_path, file_hash FROM photos \
264         WHERE metadata_extracted = FALSE AND is_trashed = FALSE",
265    );
266    let mut params: Vec<i64> = Vec::with_capacity(failed_ids.len() + 1);
267    if !failed_ids.is_empty() {
268        sql.push_str(" AND id NOT IN (");
269        for idx in 0..failed_ids.len() {
270            if idx > 0 {
271                sql.push_str(", ");
272            }
273            sql.push('?');
274        }
275        sql.push(')');
276        params.extend(failed_ids.iter().copied());
277    }
278    sql.push_str(" ORDER BY id ASC LIMIT ?");
279    params.push(limit as i64);
280
281    let mut stmt = conn.prepare(&sql)?;
282    let rows = stmt.query_map(rusqlite::params_from_iter(params), |row| {
283        Ok((row.get(0)?, row.get(1)?, row.get(2)?))
284    })?;
285    rows.collect()
286}
287
288#[derive(Debug, Clone, Default)]
289struct VideoMetadata {
290    date_taken: Option<chrono::DateTime<chrono::Utc>>,
291    date_taken_source: Option<String>,
292}
293
294impl VideoMetadata {
295    fn from_path(path: &Path) -> Self {
296        if let Some(date) = video_capture_date(path) {
297            return Self {
298                date_taken: Some(date),
299                date_taken_source: Some("video_metadata".into()),
300            };
301        }
302        if let Some(date) = ExifExtractor::parse_date_from_filename(path) {
303            return Self {
304                date_taken: Some(date),
305                date_taken_source: Some("filename".into()),
306            };
307        }
308        Self {
309            date_taken: file_mtime(path),
310            date_taken_source: Some("mtime".into()),
311        }
312    }
313}
314
315fn video_capture_date(path: &Path) -> Option<chrono::DateTime<chrono::Utc>> {
316    let bytes = read_video_probe_bytes(path).ok()?;
317    embedded_video_date_string(&bytes).or_else(|| quicktime_epoch_date(&bytes))
318}
319
320fn read_video_probe_bytes(path: &Path) -> std::io::Result<Vec<u8>> {
321    const PROBE_CHUNK: usize = 8 * 1024 * 1024;
322    let mut file = std::fs::File::open(path)?;
323    let len = file.metadata()?.len();
324    if len <= (PROBE_CHUNK * 2) as u64 {
325        let mut bytes = Vec::with_capacity(len as usize);
326        file.read_to_end(&mut bytes)?;
327        return Ok(bytes);
328    }
329
330    let mut bytes = vec![0u8; PROBE_CHUNK];
331    file.read_exact(&mut bytes)?;
332    file.seek(SeekFrom::End(-(PROBE_CHUNK as i64)))?;
333    let mut tail = vec![0u8; PROBE_CHUNK];
334    file.read_exact(&mut tail)?;
335    bytes.extend_from_slice(&tail);
336    Ok(bytes)
337}
338
339fn quicktime_epoch_date(bytes: &[u8]) -> Option<chrono::DateTime<chrono::Utc>> {
340    find_quicktime_atom_date(bytes, b"mvhd").or_else(|| find_quicktime_atom_date(bytes, b"mdhd"))
341}
342
343fn find_quicktime_atom_date(bytes: &[u8], atom: &[u8; 4]) -> Option<chrono::DateTime<chrono::Utc>> {
344    const QT_TO_UNIX_SECONDS: i64 = 2_082_844_800;
345    for i in 0..bytes.len().saturating_sub(16) {
346        if &bytes[i..i + 4] != atom {
347            continue;
348        }
349        let version = bytes.get(i + 4).copied()?;
350        let seconds = match version {
351            0 => {
352                let raw = u32::from_be_bytes(bytes.get(i + 8..i + 12)?.try_into().ok()?);
353                i64::from(raw)
354            }
355            1 => {
356                let raw = u64::from_be_bytes(bytes.get(i + 8..i + 16)?.try_into().ok()?);
357                i64::try_from(raw).ok()?
358            }
359            _ => continue,
360        };
361        if seconds <= QT_TO_UNIX_SECONDS {
362            continue;
363        }
364        let unix = seconds - QT_TO_UNIX_SECONDS;
365        let dt = chrono::DateTime::from_timestamp(unix, 0)?;
366        ExifExtractor::plausible_datetime(dt.naive_utc())?;
367        return Some(dt);
368    }
369    None
370}
371
372fn embedded_video_date_string(bytes: &[u8]) -> Option<chrono::DateTime<chrono::Utc>> {
373    lazy_static::lazy_static! {
374        static ref VIDEO_DATE: Regex = Regex::new(
375            r"(?x)
376            (\d{4}[:-]\d{2}[:-]\d{2}
377            [ T]
378            \d{2}:\d{2}(?::\d{2}(?:\.\d+)? )?
379            (?:Z|[+-]\d{2}:?\d{2})?)
380            "
381        ).unwrap();
382    }
383    let text = String::from_utf8_lossy(bytes);
384    for caps in VIDEO_DATE.captures_iter(&text) {
385        let candidate = caps.get(1)?.as_str();
386        if let Some(dt) = ExifExtractor::parse_exif_date(candidate) {
387            return Some(dt);
388        }
389    }
390    None
391}
392
393fn file_mtime(path: &Path) -> Option<chrono::DateTime<chrono::Utc>> {
394    let metadata = std::fs::metadata(path).ok()?;
395    let mtime = metadata.modified().ok()?;
396    let duration = mtime.duration_since(std::time::UNIX_EPOCH).ok()?;
397    chrono::DateTime::from_timestamp(duration.as_secs() as i64, 0)
398}
399
400impl ExtractedMetadata {
401    fn media_type(&self) -> MediaType {
402        match self {
403            Self::Photo(_) => MediaType::Photo,
404            Self::Video(_) => MediaType::Video,
405            Self::Skipped => MediaType::Photo,
406        }
407    }
408
409    fn date_taken(&self) -> Option<chrono::DateTime<chrono::Utc>> {
410        match self {
411            Self::Photo(m) => m.date_taken,
412            Self::Video(m) => m.date_taken,
413            Self::Skipped => None,
414        }
415    }
416
417    fn date_taken_source(&self) -> Option<String> {
418        match self {
419            Self::Photo(m) => m.date_taken_source.clone(),
420            Self::Video(m) => m.date_taken_source.clone(),
421            Self::Skipped => None,
422        }
423    }
424
425    fn gps_latitude(&self) -> Option<f64> {
426        match self {
427            Self::Photo(m) => m.gps_latitude,
428            Self::Video(_) => None,
429            Self::Skipped => None,
430        }
431    }
432
433    fn gps_longitude(&self) -> Option<f64> {
434        match self {
435            Self::Photo(m) => m.gps_longitude,
436            Self::Video(_) => None,
437            Self::Skipped => None,
438        }
439    }
440
441    fn camera_make(&self) -> Option<String> {
442        match self {
443            Self::Photo(m) => m.camera_make.clone(),
444            Self::Video(_) => None,
445            Self::Skipped => None,
446        }
447    }
448
449    fn camera_model(&self) -> Option<String> {
450        match self {
451            Self::Photo(m) => m.camera_model.clone(),
452            Self::Video(_) => None,
453            Self::Skipped => None,
454        }
455    }
456
457    fn iso(&self) -> Option<i32> {
458        match self {
459            Self::Photo(m) => m.iso,
460            Self::Video(_) => None,
461            Self::Skipped => None,
462        }
463    }
464
465    fn aperture(&self) -> Option<String> {
466        match self {
467            Self::Photo(m) => m.aperture.clone(),
468            Self::Video(_) => None,
469            Self::Skipped => None,
470        }
471    }
472
473    fn shutter_speed(&self) -> Option<String> {
474        match self {
475            Self::Photo(m) => m.shutter_speed.clone(),
476            Self::Video(_) => None,
477            Self::Skipped => None,
478        }
479    }
480
481    fn focal_length(&self) -> Option<String> {
482        match self {
483            Self::Photo(m) => m.focal_length.clone(),
484            Self::Video(_) => None,
485            Self::Skipped => None,
486        }
487    }
488
489    fn lens_model(&self) -> Option<String> {
490        match self {
491            Self::Photo(m) => m.lens_model.clone(),
492            Self::Video(_) => None,
493            Self::Skipped => None,
494        }
495    }
496
497    fn flash(&self) -> Option<String> {
498        match self {
499            Self::Photo(m) => m.flash.clone(),
500            Self::Video(_) => None,
501            Self::Skipped => None,
502        }
503    }
504
505    fn gps_altitude(&self) -> Option<f64> {
506        match self {
507            Self::Photo(m) => m.gps_altitude,
508            Self::Video(_) => None,
509            Self::Skipped => None,
510        }
511    }
512
513    fn width(&self) -> Option<u32> {
514        match self {
515            Self::Photo(m) => m.width,
516            Self::Video(_) => None,
517            Self::Skipped => None,
518        }
519    }
520
521    fn height(&self) -> Option<u32> {
522        match self {
523            Self::Photo(m) => m.height,
524            Self::Video(_) => None,
525            Self::Skipped => None,
526        }
527    }
528
529    fn orientation(&self) -> u16 {
530        match self {
531            Self::Photo(m) => m.orientation.unwrap_or(1),
532            Self::Video(_) => 1,
533            Self::Skipped => 1,
534        }
535    }
536
537    fn duration_ms(&self) -> Option<i64> {
538        match self {
539            Self::Photo(_) => None,
540            Self::Video(_) => None,
541            Self::Skipped => None,
542        }
543    }
544
545    fn video_codec(&self) -> Option<String> {
546        match self {
547            Self::Photo(_) => None,
548            Self::Video(_) => None,
549            Self::Skipped => None,
550        }
551    }
552
553    fn audio_codec(&self) -> Option<String> {
554        match self {
555            Self::Photo(_) => None,
556            Self::Video(_) => None,
557            Self::Skipped => None,
558        }
559    }
560
561    fn frame_rate(&self) -> Option<f32> {
562        match self {
563            Self::Photo(_) => None,
564            Self::Video(_) => None,
565            Self::Skipped => None,
566        }
567    }
568
569    fn bitrate(&self) -> Option<i64> {
570        match self {
571            Self::Photo(_) => None,
572            Self::Video(_) => None,
573            Self::Skipped => None,
574        }
575    }
576
577    fn has_audio(&self) -> bool {
578        match self {
579            Self::Photo(_) => false,
580            Self::Video(_) => false,
581            Self::Skipped => false,
582        }
583    }
584}
585
586#[cfg(test)]
587mod tests {
588    use super::*;
589    use chrono::{Datelike, TimeZone, Timelike, Utc};
590    use std::collections::HashSet;
591
592    #[test]
593    fn quicktime_epoch_date_reads_mvhd_creation_time() {
594        const QT_TO_UNIX_SECONDS: i64 = 2_082_844_800;
595        let unix = Utc
596            .with_ymd_and_hms(2024, 1, 15, 10, 16, 0)
597            .single()
598            .unwrap()
599            .timestamp();
600        let qt_seconds = (unix + QT_TO_UNIX_SECONDS) as u32;
601        let mut bytes = Vec::new();
602        bytes.extend_from_slice(&108u32.to_be_bytes());
603        bytes.extend_from_slice(b"mvhd");
604        bytes.push(0);
605        bytes.extend_from_slice(&[0, 0, 0]);
606        bytes.extend_from_slice(&qt_seconds.to_be_bytes());
607        bytes.resize(108, 0);
608
609        let dt = quicktime_epoch_date(&bytes).unwrap();
610        assert_eq!(dt.year(), 2024);
611        assert_eq!(dt.hour(), 10);
612        assert_eq!(dt.minute(), 16);
613    }
614
615    #[test]
616    fn embedded_video_date_string_reads_quicktime_creationdate() {
617        let bytes = b"com.apple.quicktime.creationdate\0\x002024-01-15T10:16:00+05:30";
618        let dt = embedded_video_date_string(bytes).unwrap();
619        assert_eq!(dt.year(), 2024);
620        assert_eq!(dt.hour(), 10);
621        assert_eq!(dt.minute(), 16);
622    }
623
624    #[test]
625    fn metadata_chunk_skips_current_run_failures_without_hiding_pending_rows() {
626        let conn = rusqlite::Connection::open_in_memory().unwrap();
627        conn.execute_batch(
628            r#"
629            CREATE TABLE photos (
630                id INTEGER PRIMARY KEY,
631                file_path TEXT NOT NULL,
632                file_hash TEXT NOT NULL,
633                metadata_extracted BOOLEAN NOT NULL DEFAULT FALSE,
634                is_trashed BOOLEAN NOT NULL DEFAULT FALSE
635            );
636            INSERT INTO photos (id, file_path, file_hash, metadata_extracted, is_trashed)
637            VALUES
638                (1, 'bad.jpg', 'aa111', FALSE, FALSE),
639                (2, 'good.jpg', 'bb222', FALSE, FALSE),
640                (3, 'done.jpg', 'cc333', TRUE, FALSE),
641                (4, 'trashed.jpg', 'dd444', FALSE, TRUE);
642            "#,
643        )
644        .unwrap();
645        let mut failed = HashSet::new();
646        failed.insert(1);
647
648        let rows = load_metadata_chunk(&conn, &failed, 20).unwrap();
649
650        assert_eq!(rows, vec![(2, "good.jpg".into(), "bb222".into())]);
651    }
652}