Skip to main content

smriti/services/
takeout_import.rs

1//! Google Photos Takeout ZIP importer.
2//!
3//! The importer treats every selected ZIP as one logical export, so sidecars
4//! and repeated album copies can be reconciled across split Takeout parts.
5//! Media is streamed through a temporary file, hashed, and atomically moved
6//! into the library. The database ledger makes reruns idempotent and retains
7//! Google-only metadata for future refreshes.
8
9use std::collections::{BTreeSet, HashMap, HashSet};
10use std::fs::{self, File};
11use std::io::{BufReader, BufWriter, Read, Write};
12use std::path::{Path, PathBuf};
13use std::sync::atomic::{AtomicBool, Ordering};
14use std::time::Instant;
15
16use chrono::{DateTime, TimeZone, Utc};
17use rusqlite::{params, Connection, OptionalExtension};
18use serde::{Deserialize, Serialize};
19use sha2::{Digest, Sha256};
20use zip::read::ZipArchive;
21
22use crate::db::{AlbumRepo, TakeoutImportRepo, TakeoutLedgerItem};
23use crate::services::scanner::{calculate_hash, media_type_for_path};
24
25const IMPORT_DIR: &str = "Imported from Google Photos";
26const STAGING_DIR: &str = "takeout-import";
27const MAX_JSON_BYTES: u64 = 8 * 1024 * 1024;
28const MAX_ARCHIVE_ENTRIES: usize = 2_000_000;
29const MIN_MEDIA_BYTES: u64 = 10 * 1024;
30const SPACE_MARGIN_BYTES: u64 = 256 * 1024 * 1024;
31
32#[derive(Debug, Clone, Serialize, Deserialize, Default, PartialEq)]
33pub struct TakeoutMetadata {
34    pub title: Option<String>,
35    pub description: Option<String>,
36    pub taken_at_unix: Option<i64>,
37    pub latitude: Option<f64>,
38    pub longitude: Option<f64>,
39    pub altitude: Option<f64>,
40    pub favorited: bool,
41}
42
43impl TakeoutMetadata {
44    fn merge_from(&mut self, other: &Self) {
45        if self.title.is_none() {
46            self.title.clone_from(&other.title);
47        }
48        if self.description.is_none() {
49            self.description.clone_from(&other.description);
50        }
51        if self.taken_at_unix.is_none() {
52            self.taken_at_unix = other.taken_at_unix;
53        }
54        if self.latitude.is_none() {
55            self.latitude = other.latitude;
56        }
57        if self.longitude.is_none() {
58            self.longitude = other.longitude;
59        }
60        if self.altitude.is_none() {
61            self.altitude = other.altitude;
62        }
63        self.favorited |= other.favorited;
64    }
65
66    pub fn taken_at_rfc3339(&self) -> Option<String> {
67        self.taken_at_unix
68            .and_then(|ts| Utc.timestamp_opt(ts, 0).single())
69            .map(|dt| dt.to_rfc3339())
70    }
71}
72
73#[derive(Debug, Clone, Serialize)]
74pub struct TakeoutProgress {
75    pub stage: &'static str,
76    pub processed: u64,
77    pub total: u64,
78    pub message: String,
79    pub elapsed_seconds: f64,
80}
81
82#[derive(Debug, Clone, Serialize, Default)]
83pub struct TakeoutImportReport {
84    pub archives: u64,
85    pub media_found: u64,
86    pub imported: u64,
87    pub reused_existing: u64,
88    pub duplicates_collapsed: u64,
89    pub unsupported_or_small: u64,
90    pub unmatched_sidecars: u64,
91    pub albums_restored: u64,
92    pub metadata_restored: u64,
93    pub errors: Vec<String>,
94    pub cancelled: bool,
95}
96
97#[derive(Debug, Clone)]
98struct MediaEntry {
99    archive_index: usize,
100    entry_index: usize,
101    logical_path: String,
102    size: u64,
103    sidecar_key: String,
104    folder_key: String,
105    file_name: String,
106}
107
108#[derive(Debug, Default)]
109struct ImportPlan {
110    archives: Vec<PathBuf>,
111    media: Vec<MediaEntry>,
112    sidecars: HashMap<String, TakeoutMetadata>,
113    album_dirs: HashMap<String, String>,
114    json_count: u64,
115    ignored_count: u64,
116    estimated_unique_bytes: u64,
117}
118
119#[derive(Debug)]
120struct ImportedAggregate {
121    relative_path: String,
122    metadata: TakeoutMetadata,
123    albums: BTreeSet<String>,
124}
125
126#[derive(Debug, Deserialize)]
127struct RawTime {
128    timestamp: Option<String>,
129}
130
131#[derive(Debug, Deserialize, Default)]
132#[serde(rename_all = "camelCase")]
133struct RawGeo {
134    latitude: Option<f64>,
135    longitude: Option<f64>,
136    altitude: Option<f64>,
137}
138
139#[derive(Debug, Deserialize, Default)]
140#[serde(rename_all = "camelCase")]
141struct RawSidecar {
142    title: Option<String>,
143    description: Option<String>,
144    photo_taken_time: Option<RawTime>,
145    creation_time: Option<RawTime>,
146    geo_data_exif: Option<RawGeo>,
147    geo_data: Option<RawGeo>,
148    favorited: Option<bool>,
149}
150
151pub fn import_google_takeout(
152    archive_paths: &[PathBuf],
153    library_root: &Path,
154    conn: &Connection,
155    cancel: &AtomicBool,
156    mut progress: impl FnMut(TakeoutProgress),
157) -> Result<TakeoutImportReport, String> {
158    let started = Instant::now();
159    validate_library_root(library_root)?;
160    let mut report = TakeoutImportReport {
161        archives: archive_paths.len() as u64,
162        ..Default::default()
163    };
164    progress(TakeoutProgress {
165        stage: "inspect",
166        processed: 0,
167        total: archive_paths.len() as u64,
168        message: "Checking Takeout archives".into(),
169        elapsed_seconds: 0.0,
170    });
171    let plan = inspect_archives(archive_paths, cancel, |processed, total, message| {
172        progress(TakeoutProgress {
173            stage: "inspect",
174            processed,
175            total,
176            message,
177            elapsed_seconds: started.elapsed().as_secs_f64(),
178        });
179    })?;
180    report.media_found = plan.media.len() as u64;
181    report.unsupported_or_small = plan.ignored_count;
182    if cancel.load(Ordering::Relaxed) {
183        report.cancelled = true;
184        return Ok(report);
185    }
186    if plan.media.is_empty() {
187        return Err("No supported Google Photos media was found in the selected ZIP files".into());
188    }
189    check_available_space(library_root, plan.estimated_unique_bytes)?;
190
191    let candidate_sizes: HashSet<i64> = plan
192        .media
193        .iter()
194        .filter_map(|m| i64::try_from(m.size).ok())
195        .collect();
196    let existing_hashes = hash_existing_candidates(conn, library_root, &candidate_sizes, cancel);
197    let import_root = library_root.join(IMPORT_DIR);
198    fs::create_dir_all(&import_root)
199        .map_err(|e| format!("Could not create {}: {e}", import_root.display()))?;
200    let staging = crate::db::connection::library_metadata_dir(library_root).join(STAGING_DIR);
201    prepare_staging(&staging)?;
202
203    let mut aggregates: HashMap<String, ImportedAggregate> = HashMap::new();
204    let mut used_sidecars: HashSet<String> = HashSet::new();
205    let mut current_archive_index = usize::MAX;
206    let mut current_archive: Option<ZipArchive<BufReader<File>>> = None;
207    let repo = TakeoutImportRepo::new(conn);
208
209    for (ordinal, media) in plan.media.iter().enumerate() {
210        if cancel.load(Ordering::Relaxed) {
211            report.cancelled = true;
212            break;
213        }
214        if current_archive_index != media.archive_index {
215            let file = File::open(&plan.archives[media.archive_index]).map_err(|e| {
216                format!(
217                    "Could not reopen {}: {e}",
218                    plan.archives[media.archive_index].display()
219                )
220            })?;
221            current_archive = Some(
222                ZipArchive::new(BufReader::new(file))
223                    .map_err(|e| format!("Invalid ZIP during import: {e}"))?,
224            );
225            current_archive_index = media.archive_index;
226        }
227
228        let temp_path = staging.join(format!("{}.partial", ordinal));
229        if temp_path.exists() {
230            let _ = fs::remove_file(&temp_path);
231        }
232        let archive = current_archive.as_mut().expect("archive opened above");
233        let mut entry = match archive.by_index(media.entry_index) {
234            Ok(entry) => entry,
235            Err(error) => {
236                report
237                    .errors
238                    .push(format!("Could not read {}: {error}", media.logical_path));
239                continue;
240            }
241        };
242        let content_hash = match stream_entry_to_temp(&mut entry, &temp_path) {
243            Ok(hash) => hash,
244            Err(error) => {
245                let _ = fs::remove_file(&temp_path);
246                report
247                    .errors
248                    .push(format!("Could not import {}: {error}", media.logical_path));
249                continue;
250            }
251        };
252
253        let metadata = plan
254            .sidecars
255            .get(&media.sidecar_key)
256            .cloned()
257            .unwrap_or_default();
258        if plan.sidecars.contains_key(&media.sidecar_key) {
259            used_sidecars.insert(media.sidecar_key.clone());
260        }
261        let mut albums = BTreeSet::new();
262        if let Some(album) = plan.album_dirs.get(&media.folder_key) {
263            albums.insert(album.clone());
264        }
265
266        if let Some(existing) = aggregates.get_mut(&content_hash) {
267            existing.metadata.merge_from(&metadata);
268            existing.albums.extend(albums);
269            report.duplicates_collapsed += 1;
270            let _ = fs::remove_file(&temp_path);
271        } else {
272            let prior_takeout = repo
273                .path_for_hash(&content_hash)
274                .map_err(|e| e.to_string())?;
275            let existing_path = prior_takeout
276                .filter(|path| library_root.join(path).is_file())
277                .or_else(|| existing_hashes.get(&content_hash).cloned());
278            let relative_path = if let Some(path) = existing_path {
279                report.reused_existing += 1;
280                let _ = fs::remove_file(&temp_path);
281                path
282            } else {
283                let destination = destination_for(&import_root, media, &metadata, &content_hash)?;
284                if destination.is_file()
285                    && calculate_hash(&destination).ok().as_deref() == Some(&content_hash)
286                {
287                    let _ = fs::remove_file(&temp_path);
288                    report.reused_existing += 1;
289                } else {
290                    atomic_move(&temp_path, &destination)?;
291                    report.imported += 1;
292                }
293                crate::services::path_util::relative_path_for_storage(
294                    destination
295                        .strip_prefix(library_root)
296                        .map_err(|_| "Import destination escaped the library root".to_string())?,
297                )
298            };
299            aggregates.insert(
300                content_hash,
301                ImportedAggregate {
302                    relative_path,
303                    metadata,
304                    albums,
305                },
306            );
307        }
308
309        progress(TakeoutProgress {
310            stage: "extract",
311            processed: (ordinal + 1) as u64,
312            total: plan.media.len() as u64,
313            message: media.file_name.clone(),
314            elapsed_seconds: started.elapsed().as_secs_f64(),
315        });
316    }
317
318    let ledger_items: Vec<TakeoutLedgerItem> = aggregates
319        .iter()
320        .map(|(hash, aggregate)| TakeoutLedgerItem {
321            content_hash: hash.clone(),
322            file_path: aggregate.relative_path.clone(),
323            metadata_json: (aggregate.metadata != TakeoutMetadata::default())
324                .then(|| serde_json::to_string(&aggregate.metadata).ok())
325                .flatten(),
326            albums: aggregate.albums.clone(),
327        })
328        .collect();
329    repo.upsert_items(&ledger_items)
330        .map_err(|e| format!("Could not save Takeout import state: {e}"))?;
331    report.unmatched_sidecars = plan
332        .sidecars
333        .keys()
334        .filter(|key| !used_sidecars.contains(*key))
335        .count() as u64;
336    let _ = fs::remove_dir_all(&staging);
337    Ok(report)
338}
339
340pub fn apply_takeout_metadata_and_albums(conn: &Connection) -> Result<(u64, u64), String> {
341    let geocoder =
342        crate::services::geocoding::GeocodingService::new(crate::db::geonames::geonames_db_path())
343            .ok();
344    let mut stmt = conn
345        .prepare(
346            r#"
347            SELECT i.content_hash, i.metadata_json, p.id
348              FROM google_takeout_items i
349              JOIN photos p ON p.file_path = i.file_path
350             WHERE p.is_trashed = FALSE
351            "#,
352        )
353        .map_err(|e| e.to_string())?;
354    let rows: Vec<(String, Option<String>, i64)> = stmt
355        .query_map([], |r| Ok((r.get(0)?, r.get(1)?, r.get(2)?)))
356        .map_err(|e| e.to_string())?
357        .collect::<Result<_, _>>()
358        .map_err(|e| e.to_string())?;
359    drop(stmt);
360
361    let mut parsed_rows = Vec::with_capacity(rows.len());
362    let mut album_photos: HashMap<String, Vec<i64>> = HashMap::new();
363    for (hash, metadata_json, photo_id) in rows {
364        let metadata = metadata_json
365            .as_deref()
366            .and_then(|json| serde_json::from_str::<TakeoutMetadata>(json).ok());
367        parsed_rows.push((photo_id, metadata));
368        let mut album_stmt = conn
369            .prepare("SELECT album_name FROM google_takeout_albums WHERE content_hash = ?1")
370            .map_err(|e| e.to_string())?;
371        let albums: Vec<String> = album_stmt
372            .query_map(params![hash], |r| r.get(0))
373            .map_err(|e| e.to_string())?
374            .collect::<Result<_, _>>()
375            .map_err(|e| e.to_string())?;
376        for album in albums {
377            album_photos.entry(album).or_default().push(photo_id);
378        }
379    }
380
381    let tx = conn.unchecked_transaction().map_err(|e| e.to_string())?;
382    let mut metadata_restored = 0u64;
383    for (photo_id, metadata) in parsed_rows {
384        if let Some(metadata) = metadata {
385            let taken = metadata.taken_at_rfc3339();
386            let (city, country) = match (metadata.latitude, metadata.longitude, geocoder.as_ref()) {
387                (Some(lat), Some(lon), Some(g)) => g
388                    .reverse_geocode(lat, lon)
389                    .map(|place| (Some(place.city), Some(place.country)))
390                    .unwrap_or((None, None)),
391                _ => (None, None),
392            };
393            tx.execute(
394                    r#"
395                    UPDATE photos SET
396                        date_taken = COALESCE(?1, date_taken),
397                        date_taken_source = CASE WHEN ?1 IS NOT NULL THEN 'google_takeout' ELSE date_taken_source END,
398                        gps_latitude = COALESCE(?2, gps_latitude),
399                        gps_longitude = COALESCE(?3, gps_longitude),
400                        gps_altitude = COALESCE(?4, gps_altitude),
401                        location_city = COALESCE(?5, location_city),
402                        location_country = COALESCE(?6, location_country),
403                        is_favorite = CASE WHEN ?7 THEN TRUE ELSE is_favorite END,
404                        updated_at = CURRENT_TIMESTAMP
405                    WHERE id = ?8
406                    "#,
407                    params![
408                        taken,
409                        metadata.latitude,
410                        metadata.longitude,
411                        metadata.altitude,
412                        city,
413                        country,
414                        metadata.favorited,
415                        photo_id
416                    ],
417                )
418                .map_err(|e| e.to_string())?;
419            metadata_restored += 1;
420        }
421    }
422    tx.commit().map_err(|e| e.to_string())?;
423
424    let mut album_ids = HashSet::new();
425    for (album_name, photo_ids) in album_photos {
426        let existing: Option<i64> = conn
427            .query_row(
428                "SELECT id FROM albums WHERE name = ?1 ORDER BY id LIMIT 1",
429                params![album_name],
430                |r| r.get(0),
431            )
432            .optional()
433            .map_err(|e| e.to_string())?;
434        let album_id = match existing {
435            Some(id) => id,
436            None => AlbumRepo::new(conn)
437                .create(&album_name)
438                .map_err(|e| e.to_string())?,
439        };
440        AlbumRepo::new(conn)
441            .add_photos(album_id, &photo_ids)
442            .map_err(|e| e.to_string())?;
443        album_ids.insert(album_id);
444    }
445    Ok((metadata_restored, album_ids.len() as u64))
446}
447
448fn inspect_archives(
449    archive_paths: &[PathBuf],
450    cancel: &AtomicBool,
451    mut progress: impl FnMut(u64, u64, String),
452) -> Result<ImportPlan, String> {
453    if archive_paths.is_empty() {
454        return Err("Select at least one Google Takeout ZIP file".into());
455    }
456    let mut plan = ImportPlan {
457        archives: archive_paths.to_vec(),
458        ..Default::default()
459    };
460    let mut unique_size_crc = HashSet::new();
461    for (archive_index, path) in archive_paths.iter().enumerate() {
462        if cancel.load(Ordering::Relaxed) {
463            break;
464        }
465        if path
466            .extension()
467            .and_then(|s| s.to_str())
468            .map(str::to_ascii_lowercase)
469            .as_deref()
470            != Some("zip")
471        {
472            return Err(format!("{} is not a ZIP file", path.display()));
473        }
474        let file =
475            File::open(path).map_err(|e| format!("Could not open {}: {e}", path.display()))?;
476        let mut archive = ZipArchive::new(BufReader::new(file))
477            .map_err(|e| format!("Invalid ZIP {}: {e}", path.display()))?;
478        if archive.len() > MAX_ARCHIVE_ENTRIES {
479            return Err(format!(
480                "{} contains too many entries ({})",
481                path.display(),
482                archive.len()
483            ));
484        }
485        for entry_index in 0..archive.len() {
486            let mut entry = archive.by_index(entry_index).map_err(|e| {
487                format!(
488                    "Could not inspect {} entry {entry_index}: {e}",
489                    path.display()
490                )
491            })?;
492            if entry.is_dir() {
493                continue;
494            }
495            let enclosed = entry
496                .enclosed_name()
497                .ok_or_else(|| format!("Unsafe path in {}: {}", path.display(), entry.name()))?;
498            let logical_path = normalize_zip_path(&enclosed);
499            let Some(relative_google_path) = google_photos_relative(&logical_path) else {
500                continue;
501            };
502            let lower = relative_google_path.to_ascii_lowercase();
503            if lower.ends_with(".json") {
504                if entry.size() > MAX_JSON_BYTES {
505                    plan.ignored_count += 1;
506                    continue;
507                }
508                let mut json = String::new();
509                entry
510                    .read_to_string(&mut json)
511                    .map_err(|e| format!("Could not read sidecar {logical_path}: {e}"))?;
512                if file_name_of(&relative_google_path).eq_ignore_ascii_case("metadata.json") {
513                    if let Ok(value) = serde_json::from_str::<serde_json::Value>(&json) {
514                        if let Some(title) = value.get("title").and_then(|v| v.as_str()) {
515                            let folder = folder_of(&relative_google_path);
516                            if !is_year_bucket(&folder) && !title.trim().is_empty() {
517                                plan.album_dirs.insert(folder, title.trim().to_string());
518                            }
519                        }
520                    }
521                } else {
522                    match parse_sidecar(&relative_google_path, &json) {
523                        Some((key, metadata)) => {
524                            plan.sidecars
525                                .entry(key)
526                                .and_modify(|m| m.merge_from(&metadata))
527                                .or_insert(metadata);
528                            plan.json_count += 1;
529                        }
530                        None => plan.ignored_count += 1,
531                    }
532                }
533                continue;
534            }
535            if media_type_for_path(Path::new(&relative_google_path)).is_none()
536                || entry.size() < MIN_MEDIA_BYTES
537            {
538                plan.ignored_count += 1;
539                continue;
540            }
541            let file_name = file_name_of(&relative_google_path).to_string();
542            let folder_key = folder_of(&relative_google_path);
543            let sidecar_key = sidecar_key(&folder_key, &file_name);
544            let size_crc = (entry.size(), entry.crc32());
545            if unique_size_crc.insert(size_crc) {
546                plan.estimated_unique_bytes =
547                    plan.estimated_unique_bytes.saturating_add(entry.size());
548            }
549            plan.media.push(MediaEntry {
550                archive_index,
551                entry_index,
552                logical_path,
553                size: entry.size(),
554                sidecar_key,
555                folder_key,
556                file_name,
557            });
558        }
559        progress(
560            (archive_index + 1) as u64,
561            archive_paths.len() as u64,
562            path.file_name()
563                .and_then(|s| s.to_str())
564                .unwrap_or("Takeout archive")
565                .to_string(),
566        );
567    }
568
569    Ok(plan)
570}
571
572fn parse_sidecar(relative_path: &str, json: &str) -> Option<(String, TakeoutMetadata)> {
573    let raw: RawSidecar = serde_json::from_str(json).ok()?;
574    let folder = folder_of(relative_path);
575    let fallback = sidecar_media_name(file_name_of(relative_path));
576    let title = raw
577        .title
578        .as_deref()
579        .filter(|s| !s.trim().is_empty())
580        .unwrap_or(&fallback)
581        .to_string();
582    if title.is_empty() {
583        return None;
584    }
585    let taken_at_unix = raw
586        .photo_taken_time
587        .or(raw.creation_time)
588        .and_then(|t| t.timestamp)
589        .and_then(|s| s.parse::<i64>().ok());
590    let geo = valid_geo(raw.geo_data_exif).or_else(|| valid_geo(raw.geo_data));
591    let metadata = TakeoutMetadata {
592        title: Some(title.clone()),
593        description: raw.description.filter(|s| !s.is_empty()),
594        taken_at_unix,
595        latitude: geo.as_ref().and_then(|g| g.latitude),
596        longitude: geo.as_ref().and_then(|g| g.longitude),
597        altitude: geo.and_then(|g| g.altitude),
598        favorited: raw.favorited.unwrap_or(false),
599    };
600    Some((sidecar_key(&folder, &title), metadata))
601}
602
603fn valid_geo(geo: Option<RawGeo>) -> Option<RawGeo> {
604    let geo = geo?;
605    let lat = geo.latitude?;
606    let lon = geo.longitude?;
607    if !(-90.0..=90.0).contains(&lat)
608        || !(-180.0..=180.0).contains(&lon)
609        || (lat == 0.0 && lon == 0.0)
610    {
611        return None;
612    }
613    Some(geo)
614}
615
616fn hash_existing_candidates(
617    conn: &Connection,
618    library_root: &Path,
619    sizes: &HashSet<i64>,
620    cancel: &AtomicBool,
621) -> HashMap<String, String> {
622    let repo = TakeoutImportRepo::new(conn);
623    let mut result = HashMap::new();
624    let candidates = repo.candidate_existing_files(sizes).unwrap_or_default();
625    for (relative, _) in candidates {
626        if cancel.load(Ordering::Relaxed) {
627            break;
628        }
629        let Ok(path) = crate::services::path_util::safe_join_relative(library_root, &relative)
630        else {
631            continue;
632        };
633        if let Ok(hash) = calculate_hash(path) {
634            result.entry(hash).or_insert(relative);
635        }
636    }
637    result
638}
639
640fn stream_entry_to_temp<R: Read>(reader: &mut R, temp_path: &Path) -> Result<String, String> {
641    let file = File::create(temp_path)
642        .map_err(|e| format!("Could not create temporary import file: {e}"))?;
643    let mut writer = BufWriter::new(file);
644    let mut hasher = Sha256::new();
645    let mut buffer = vec![0u8; 1024 * 1024];
646    loop {
647        let count = reader
648            .read(&mut buffer)
649            .map_err(|e| format!("Could not decompress Takeout media: {e}"))?;
650        if count == 0 {
651            break;
652        }
653        writer
654            .write_all(&buffer[..count])
655            .map_err(|e| format!("Could not write imported media: {e}"))?;
656        hasher.update(&buffer[..count]);
657    }
658    writer
659        .flush()
660        .map_err(|e| format!("Could not finish imported media: {e}"))?;
661    Ok(format!("{:x}", hasher.finalize()))
662}
663
664fn destination_for(
665    import_root: &Path,
666    media: &MediaEntry,
667    metadata: &TakeoutMetadata,
668    hash: &str,
669) -> Result<PathBuf, String> {
670    let year = metadata
671        .taken_at_unix
672        .and_then(|ts| Utc.timestamp_opt(ts, 0).single())
673        .map(|d: DateTime<Utc>| d.format("%Y").to_string())
674        .unwrap_or_else(|| "Undated".to_string());
675    let dir = import_root.join(year);
676    fs::create_dir_all(&dir).map_err(|e| format!("Could not create {}: {e}", dir.display()))?;
677    let safe_name = sanitize_file_name(&media.file_name);
678    let preferred = dir.join(&safe_name);
679    if !preferred.exists() {
680        return Ok(preferred);
681    }
682    if calculate_hash(&preferred).ok().as_deref() == Some(hash) {
683        return Ok(preferred);
684    }
685    let source = Path::new(&safe_name);
686    let stem = source
687        .file_stem()
688        .and_then(|s| s.to_str())
689        .unwrap_or("photo");
690    let ext = source.extension().and_then(|s| s.to_str());
691    let mut candidate = match ext {
692        Some(ext) => dir.join(format!("{stem}_{}.{}", &hash[..10], ext)),
693        None => dir.join(format!("{stem}_{}", &hash[..10])),
694    };
695    let mut suffix = 2u32;
696    while candidate.exists() && calculate_hash(&candidate).ok().as_deref() != Some(hash) {
697        candidate = match ext {
698            Some(ext) => dir.join(format!("{stem}_{}_{}.{}", &hash[..10], suffix, ext)),
699            None => dir.join(format!("{stem}_{}_{}", &hash[..10], suffix)),
700        };
701        suffix += 1;
702    }
703    Ok(candidate)
704}
705
706fn atomic_move(source: &Path, destination: &Path) -> Result<(), String> {
707    if let Some(parent) = destination.parent() {
708        fs::create_dir_all(parent)
709            .map_err(|e| format!("Could not create {}: {e}", parent.display()))?;
710    }
711    fs::rename(source, destination).map_err(|e| {
712        format!(
713            "Could not place imported file at {}: {e}",
714            destination.display()
715        )
716    })
717}
718
719fn prepare_staging(staging: &Path) -> Result<(), String> {
720    if staging.exists() {
721        fs::remove_dir_all(staging)
722            .map_err(|e| format!("Could not clear stale Takeout staging data: {e}"))?;
723    }
724    fs::create_dir_all(staging)
725        .map_err(|e| format!("Could not create Takeout staging directory: {e}"))
726}
727
728fn check_available_space(root: &Path, required: u64) -> Result<(), String> {
729    let available = fs2::available_space(root)
730        .map_err(|e| format!("Could not determine available library space: {e}"))?;
731    let needed = required.saturating_add(SPACE_MARGIN_BYTES);
732    if available < needed {
733        return Err(format!(
734            "Not enough free space. The import needs about {}, but only {} is available",
735            human_bytes(needed),
736            human_bytes(available)
737        ));
738    }
739    Ok(())
740}
741
742fn validate_library_root(root: &Path) -> Result<(), String> {
743    if !root.is_dir() {
744        return Err(format!("Library folder does not exist: {}", root.display()));
745    }
746    let metadata = crate::db::connection::library_metadata_dir(root);
747    if !metadata.is_dir() {
748        return Err("Open the destination as a Smriti library before importing".into());
749    }
750    Ok(())
751}
752
753fn google_photos_relative(path: &str) -> Option<String> {
754    let components: Vec<&str> = path.split('/').filter(|s| !s.is_empty()).collect();
755    let idx = components
756        .iter()
757        .position(|component| component.eq_ignore_ascii_case("Google Photos"))?;
758    let rest = components.get(idx + 1..)?.join("/");
759    (!rest.is_empty()).then_some(rest)
760}
761
762fn normalize_zip_path(path: &Path) -> String {
763    path.components()
764        .map(|component| component.as_os_str().to_string_lossy())
765        .collect::<Vec<_>>()
766        .join("/")
767}
768
769fn folder_of(path: &str) -> String {
770    path.rsplit_once('/')
771        .map(|(folder, _)| folder.to_string())
772        .unwrap_or_default()
773}
774
775fn file_name_of(path: &str) -> &str {
776    path.rsplit('/').next().unwrap_or(path)
777}
778
779fn last_folder_component(path: &str) -> &str {
780    path.rsplit('/').next().unwrap_or(path)
781}
782
783fn sidecar_key(folder: &str, title: &str) -> String {
784    format!("{}\0{}", folder.to_lowercase(), title.to_lowercase())
785}
786
787fn sidecar_media_name(name: &str) -> String {
788    let lower = name.to_ascii_lowercase();
789    if lower.ends_with(".supplemental-metadata.json") {
790        name[..name.len() - ".supplemental-metadata.json".len()].to_string()
791    } else if lower.ends_with(".json") {
792        name[..name.len() - 5].to_string()
793    } else {
794        name.to_string()
795    }
796}
797
798fn is_year_bucket(folder: &str) -> bool {
799    let name = last_folder_component(folder).trim();
800    let lower = name.to_ascii_lowercase();
801    if let Some(year) = lower.strip_prefix("photos from ") {
802        return year.len() == 4 && year.chars().all(|c| c.is_ascii_digit());
803    }
804    false
805}
806
807fn sanitize_file_name(name: &str) -> String {
808    let mut safe: String = name
809        .chars()
810        .map(|c| {
811            if c.is_control() || matches!(c, '<' | '>' | ':' | '"' | '/' | '\\' | '|' | '?' | '*') {
812                '_'
813            } else {
814                c
815            }
816        })
817        .collect();
818    safe = safe.trim_matches([' ', '.']).to_string();
819    if safe.is_empty() {
820        safe = "photo".into();
821    }
822    let stem = Path::new(&safe)
823        .file_stem()
824        .and_then(|s| s.to_str())
825        .unwrap_or("")
826        .to_ascii_uppercase();
827    const RESERVED: &[&str] = &[
828        "CON", "PRN", "AUX", "NUL", "COM1", "COM2", "COM3", "COM4", "COM5", "COM6", "COM7", "COM8",
829        "COM9", "LPT1", "LPT2", "LPT3", "LPT4", "LPT5", "LPT6", "LPT7", "LPT8", "LPT9",
830    ];
831    if RESERVED.contains(&stem.as_str()) {
832        safe.insert(0, '_');
833    }
834    truncate_file_name(&safe, 180)
835}
836
837fn truncate_file_name(name: &str, max_bytes: usize) -> String {
838    if name.len() <= max_bytes {
839        return name.to_string();
840    }
841    let path = Path::new(name);
842    let extension = path
843        .extension()
844        .and_then(|s| s.to_str())
845        .filter(|ext| ext.len() <= 20)
846        .map(|ext| format!(".{ext}"))
847        .unwrap_or_default();
848    let stem = path.file_stem().and_then(|s| s.to_str()).unwrap_or("photo");
849    let budget = max_bytes.saturating_sub(extension.len());
850    let mut shortened = String::new();
851    for ch in stem.chars() {
852        if shortened.len() + ch.len_utf8() > budget {
853            break;
854        }
855        shortened.push(ch);
856    }
857    if shortened.is_empty() {
858        shortened.push_str("photo");
859    }
860    shortened.push_str(&extension);
861    shortened
862}
863
864fn human_bytes(bytes: u64) -> String {
865    if bytes >= 1024 * 1024 * 1024 {
866        format!("{:.1} GB", bytes as f64 / (1024.0 * 1024.0 * 1024.0))
867    } else {
868        format!("{:.1} MB", bytes as f64 / (1024.0 * 1024.0))
869    }
870}
871
872#[cfg(test)]
873mod tests {
874    use super::*;
875    use std::io::Write;
876    use tempfile::tempdir;
877    use zip::write::SimpleFileOptions;
878
879    fn write_zip(path: &Path, entries: &[(&str, &[u8])]) {
880        let file = File::create(path).unwrap();
881        let mut zip = zip::ZipWriter::new(file);
882        for (name, bytes) in entries {
883            zip.start_file(*name, SimpleFileOptions::default()).unwrap();
884            zip.write_all(bytes).unwrap();
885        }
886        zip.finish().unwrap();
887    }
888
889    fn jpeg_bytes(seed: u8) -> Vec<u8> {
890        let mut bytes = vec![seed; MIN_MEDIA_BYTES as usize + 1];
891        bytes[0..3].copy_from_slice(&[0xff, 0xd8, 0xff]);
892        bytes
893    }
894
895    #[test]
896    fn parses_google_sidecar_with_exif_geo_priority() {
897        let json = r#"{
898          "title":"IMG_0001.JPG",
899          "description":"At the beach",
900          "photoTakenTime":{"timestamp":"1704067200"},
901          "geoData":{"latitude":1.0,"longitude":2.0,"altitude":3.0},
902          "geoDataExif":{"latitude":10.0,"longitude":20.0,"altitude":30.0},
903          "favorited":true
904        }"#;
905        let (_, meta) = parse_sidecar("Photos from 2024/IMG_0001.JPG.json", json).unwrap();
906        assert_eq!(meta.taken_at_unix, Some(1_704_067_200));
907        assert_eq!(meta.latitude, Some(10.0));
908        assert_eq!(meta.longitude, Some(20.0));
909        assert!(meta.favorited);
910    }
911
912    #[test]
913    fn multi_part_import_collapses_album_copy_and_restores_metadata() {
914        let root = tempdir().unwrap();
915        let db = crate::db::Database::open_for_drive(root.path()).unwrap();
916        crate::db::create_schema(&db.conn).unwrap();
917        let first = root.path().join("takeout-001.zip");
918        let second = root.path().join("takeout-002.zip");
919        let photo = jpeg_bytes(7);
920        let sidecar = br#"{
921          "title":"IMG_0001.JPG",
922          "photoTakenTime":{"timestamp":"1704067200"},
923          "geoDataExif":{"latitude":10.0,"longitude":20.0,"altitude":30.0},
924          "favorited":true
925        }"#;
926        write_zip(
927            &first,
928            &[
929                (
930                    "Takeout/Google Photos/Photos from 2024/IMG_0001.JPG",
931                    &photo,
932                ),
933                (
934                    "Takeout/Google Photos/Photos from 2024/IMG_0001.JPG.json",
935                    sidecar,
936                ),
937            ],
938        );
939        write_zip(
940            &second,
941            &[
942                (
943                    "Takeout/Google Photos/Beach/metadata.json",
944                    br#"{"title":"Beach"}"#,
945                ),
946                ("Takeout/Google Photos/Beach/IMG_0001.JPG", &photo),
947            ],
948        );
949
950        let report = import_google_takeout(
951            &[first.clone(), second.clone()],
952            root.path(),
953            &db.conn,
954            &AtomicBool::new(false),
955            |_| {},
956        )
957        .unwrap();
958        assert_eq!(report.imported, 1);
959        assert_eq!(report.duplicates_collapsed, 1);
960
961        let cancel = std::sync::Arc::new(AtomicBool::new(false));
962        let db_arc = std::sync::Arc::new(tokio::sync::Mutex::new(db));
963        let runtime = tokio::runtime::Runtime::new().unwrap();
964        let scan_report = runtime.block_on(async {
965            let (rx, task) = crate::services::scanner::start_scan(
966                root.path().to_path_buf(),
967                db_arc.clone(),
968                cancel,
969                false,
970            );
971            let drain = tokio::spawn(async move { while rx.recv().await.is_ok() {} });
972            let report = task.await.unwrap();
973            drain.await.unwrap();
974            report
975        });
976        assert!(
977            scan_report.errors.is_empty(),
978            "scan errors: {:?}",
979            scan_report.errors
980        );
981        let guard = runtime.block_on(db_arc.lock());
982        let photo_count: i64 = guard
983            .conn
984            .query_row("SELECT COUNT(*) FROM photos", [], |r| r.get(0))
985            .unwrap();
986        let ledger: Vec<(String, Option<String>)> = guard
987            .conn
988            .prepare("SELECT file_path, metadata_json FROM google_takeout_items")
989            .unwrap()
990            .query_map([], |r| Ok((r.get(0)?, r.get(1)?)))
991            .unwrap()
992            .collect::<Result<_, _>>()
993            .unwrap();
994        assert_eq!(photo_count, 1, "ledger={ledger:?}");
995        let (metadata, albums) = apply_takeout_metadata_and_albums(&guard.conn).unwrap();
996        assert_eq!(metadata, 1);
997        assert_eq!(albums, 1);
998        let restored: (String, f64, i64) = guard
999            .conn
1000            .query_row(
1001                "SELECT date_taken_source, gps_latitude, is_favorite FROM photos",
1002                [],
1003                |r| Ok((r.get(0)?, r.get(1)?, r.get(2)?)),
1004            )
1005            .unwrap();
1006        assert_eq!(restored.0, "google_takeout");
1007        assert_eq!(restored.1, 10.0);
1008        assert_eq!(restored.2, 1);
1009        let album_count: i64 = guard
1010            .conn
1011            .query_row("SELECT COUNT(*) FROM album_photos", [], |r| r.get(0))
1012            .unwrap();
1013        assert_eq!(album_count, 1);
1014
1015        let repeated = import_google_takeout(
1016            &[first, second],
1017            root.path(),
1018            &guard.conn,
1019            &AtomicBool::new(false),
1020            |_| {},
1021        )
1022        .unwrap();
1023        assert_eq!(repeated.imported, 0);
1024        assert_eq!(repeated.reused_existing, 1);
1025        assert_eq!(repeated.duplicates_collapsed, 1);
1026    }
1027
1028    #[test]
1029    fn sidecar_can_live_in_a_different_takeout_part() {
1030        let root = tempdir().unwrap();
1031        let first = root.path().join("takeout-001.zip");
1032        let second = root.path().join("takeout-002.zip");
1033        write_zip(
1034            &first,
1035            &[(
1036                "Takeout/Google Photos/Photos from 2020/a.jpg",
1037                &jpeg_bytes(2),
1038            )],
1039        );
1040        write_zip(
1041            &second,
1042            &[(
1043                "Takeout/Google Photos/Photos from 2020/a.jpg.json",
1044                br#"{"title":"a.jpg","photoTakenTime":{"timestamp":"1577836800"}}"#,
1045            )],
1046        );
1047        let plan =
1048            inspect_archives(&[first, second], &AtomicBool::new(false), |_, _, _| {}).unwrap();
1049        assert_eq!(plan.media.len(), 1);
1050        assert_eq!(
1051            plan.sidecars
1052                .get(&plan.media[0].sidecar_key)
1053                .and_then(|metadata| metadata.taken_at_unix),
1054            Some(1_577_836_800)
1055        );
1056    }
1057
1058    #[test]
1059    fn cancellation_before_inspection_is_resumable() {
1060        let root = tempdir().unwrap();
1061        let db = crate::db::Database::open_for_drive(root.path()).unwrap();
1062        crate::db::create_schema(&db.conn).unwrap();
1063        let archive = root.path().join("takeout.zip");
1064        write_zip(
1065            &archive,
1066            &[(
1067                "Takeout/Google Photos/Photos from 2020/a.jpg",
1068                &jpeg_bytes(3),
1069            )],
1070        );
1071        let cancel = AtomicBool::new(true);
1072        let report =
1073            import_google_takeout(&[archive], root.path(), &db.conn, &cancel, |_| {}).unwrap();
1074        assert!(report.cancelled);
1075        assert_eq!(report.imported, 0);
1076    }
1077
1078    #[test]
1079    fn rejects_zip_slip_paths() {
1080        let dir = tempdir().unwrap();
1081        let path = dir.path().join("bad.zip");
1082        write_zip(
1083            &path,
1084            &[("../Google Photos/Photos from 2024/a.jpg", &jpeg_bytes(1))],
1085        );
1086        let err = inspect_archives(&[path], &AtomicBool::new(false), |_, _, _| {}).unwrap_err();
1087        assert!(err.contains("Unsafe path"));
1088    }
1089
1090    #[test]
1091    fn sanitizes_windows_names_and_collisions() {
1092        assert_eq!(sanitize_file_name("CON.jpg"), "_CON.jpg");
1093        assert_eq!(sanitize_file_name("bad:name?.jpg"), "bad_name_.jpg");
1094        assert_eq!(sanitize_file_name("  .  "), "photo");
1095        let long = format!("{}.jpeg", "🙂".repeat(100));
1096        let shortened = sanitize_file_name(&long);
1097        assert!(shortened.len() <= 180);
1098        assert!(shortened.ends_with(".jpeg"));
1099    }
1100}