1use 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}