Skip to main content

smugmug_cli/uploader/
mod.rs

1use anyhow::{Context, Result};
2use colored::*;
3use indicatif::{MultiProgress, ProgressBar, ProgressStyle};
4use std::collections::HashMap;
5use std::path::PathBuf;
6use std::sync::Arc;
7use tokio::sync::Mutex;
8
9pub mod album_series;
10pub mod collect;
11pub mod queue;
12pub mod worker;
13
14use crate::api::SmugMugClient;
15use crate::api::albums::Album;
16use crate::api::images::AlbumImage;
17use crate::cache::hash_store::HashStore;
18use crate::config::{RawHandling, RawMode};
19use album_series::{AlbumSeries, ClientAlbumSeries};
20use queue::UploadQueue;
21use worker::{UploadStatus, UploadWorkerContext, upload_worker};
22// Re-exported for use in tests and examples
23#[allow(unused_imports)]
24pub use worker::calculate_file_hash;
25
26/// Fetch the images already present in an album and build the lookup maps
27/// used to decide, per local file, whether to skip it (content already
28/// present under the same filename), replace it in place (same filename,
29/// different content), or create it fresh (new filename). This runs by
30/// default so re-running an upload after locally editing photos updates the
31/// existing images instead of failing with a 409 Conflict.
32///
33/// The MD5 index (used by `--check-remote` to detect the same content
34/// uploaded anywhere in the album, under any filename) is built from the
35/// same listing when `include_md5_index` is set, avoiding a second API call.
36async fn fetch_remote_image_maps(
37    client: &SmugMugClient,
38    album_key: &str,
39    include_md5_index: bool,
40) -> Result<(
41    Option<Arc<HashMap<String, AlbumImage>>>,
42    Option<Arc<HashMap<String, String>>>,
43)> {
44    let images = client.list_album_images(album_key).await?;
45
46    let by_filename: HashMap<String, AlbumImage> = images
47        .iter()
48        .cloned()
49        .map(|img| (img.file_name.clone(), img))
50        .collect();
51
52    let by_md5 = if include_md5_index {
53        let md5_map: HashMap<String, String> = images
54            .iter()
55            .filter_map(|img| {
56                img.archived_md5
57                    .as_ref()
58                    .map(|md5| (md5.to_lowercase(), img.image_key.clone()))
59            })
60            .collect();
61        Some(Arc::new(md5_map))
62    } else {
63        None
64    };
65
66    Ok((Some(Arc::new(by_filename)), by_md5))
67}
68
69pub struct UploadOptions {
70    /// Files to upload (RAW files already selected; see `select_raw_files`).
71    pub files: Vec<PathBuf>,
72    /// How the RAW files in `files` are uploaded.
73    pub raw_handling: RawHandling,
74    /// Albums the files go into; new images claim room album by album.
75    pub series: Arc<AlbumSeries<ClientAlbumSeries>>,
76    pub client: Arc<SmugMugClient>,
77    pub threads: usize,
78    pub dry_run: bool,
79    pub check_remote: bool,
80    pub no_cache: bool,
81    pub cache_path: PathBuf,
82    pub retry_attempts: u32,
83}
84
85pub struct UploadStats {
86    pub total_files: usize,
87    pub uploaded: usize,
88    /// Files that replaced an existing, differently-content image with the
89    /// same filename already in the album.
90    pub replaced: usize,
91    /// Files already on SmugMug that were added to the album instead of
92    /// uploaded again (see `collect`).
93    pub collected: usize,
94    pub skipped: usize,
95    pub failed: usize,
96    pub total_bytes: u64,
97    pub folders_created: usize,
98    pub albums_created: usize,
99    pub duration_secs: u64,
100}
101
102impl UploadStats {
103    fn empty(start_time: std::time::Instant) -> Self {
104        UploadStats {
105            total_files: 0,
106            uploaded: 0,
107            replaced: 0,
108            collected: 0,
109            skipped: 0,
110            failed: 0,
111            total_bytes: 0,
112            folders_created: 0,
113            albums_created: 0,
114            duration_secs: start_time.elapsed().as_secs(),
115        }
116    }
117}
118
119/// What `select_raw_files` did with the RAW files it was given.
120#[derive(Debug, Default, Clone, Copy, PartialEq, Eq)]
121pub struct RawSelection {
122    /// RAW files kept, to be uploaded as JPEGs rendered from them.
123    pub rendered: usize,
124    /// RAW files left out because RAW handling is `Skip`.
125    pub skipped: usize,
126    /// RAW files left out because a JPEG or HEIC with the same name sits in
127    /// the same directory (as when shooting RAW+JPEG), or because an
128    /// earlier RAW there renders to the same name.
129    pub skipped_with_sibling: usize,
130}
131
132/// Decide which RAW files to upload, per `handling`. Non-RAW files are
133/// always kept.
134///
135/// When rendering, a RAW file `IMG_1234.CR2` becomes `IMG_1234.jpg`, which
136/// would collide with the camera's own `IMG_1234.JPG` next to it; the
137/// camera's JPEG is kept and the RAW dropped.
138pub fn select_raw_files(
139    files: Vec<PathBuf>,
140    handling: RawHandling,
141) -> (Vec<PathBuf>, RawSelection) {
142    use std::collections::HashSet;
143
144    let mut selection = RawSelection::default();
145    let key = |f: &PathBuf| {
146        let stem = f.file_stem().map(|s| s.to_string_lossy().to_lowercase());
147        (f.parent().map(|p| p.to_path_buf()), stem)
148    };
149
150    // Names already taken by a camera JPEG/HEIC in each directory.
151    let mut taken: HashSet<_> = HashSet::new();
152    if handling == RawHandling::Render {
153        for f in &files {
154            let ext = f
155                .extension()
156                .map(|e| e.to_string_lossy().to_lowercase())
157                .unwrap_or_default();
158            if matches!(ext.as_str(), "jpg" | "jpeg" | "heic" | "heif") {
159                taken.insert(key(f));
160            }
161        }
162    }
163
164    let kept = files
165        .into_iter()
166        .filter(|f| {
167            if !crate::scanner::is_raw_file(f) {
168                return true;
169            }
170            match handling {
171                RawHandling::Upload => true,
172                RawHandling::Skip => {
173                    selection.skipped += 1;
174                    false
175                }
176                RawHandling::Render => {
177                    if taken.insert(key(f)) {
178                        selection.rendered += 1;
179                        true
180                    } else {
181                        selection.skipped_with_sibling += 1;
182                        false
183                    }
184                }
185            }
186        })
187        .collect();
188    (kept, selection)
189}
190
191/// Tell the user what's happening to their RAW files, if anything unusual.
192pub fn print_raw_selection(selection: &RawSelection, mode: RawMode, has_smugmug_source: bool) {
193    if selection.rendered > 0 {
194        println!(
195            "{} Uploading {} RAW files as JPEGs (the preview the camera embedded in each)",
196            "•".cyan(),
197            selection.rendered
198        );
199    }
200    if selection.skipped_with_sibling > 0 {
201        println!(
202            "{} Skipping {} RAW files that have a JPEG or HEIC of the same name next to them",
203            "•".cyan(),
204            selection.skipped_with_sibling
205        );
206    }
207    if selection.skipped > 0 {
208        let reason = if mode == RawMode::Original && !has_smugmug_source {
209            "RAW originals need a SmugMug Source subscription"
210        } else {
211            "raw_mode is \"skip\""
212        };
213        println!(
214            "{}",
215            format!("⚠ Skipping {} RAW files: {}", selection.skipped, reason).yellow()
216        );
217        if mode == RawMode::Original && !has_smugmug_source {
218            println!(
219                "  {}",
220                "(Set raw_mode = \"auto\" in the config to upload JPEGs rendered from them, or run 'smugmug-cli init' if you have Source)"
221                    .bright_black()
222            );
223        }
224    }
225    if selection.rendered + selection.skipped + selection.skipped_with_sibling > 0 {
226        println!();
227    }
228}
229
230pub async fn upload_files(options: UploadOptions) -> Result<UploadStats> {
231    let start_time = std::time::Instant::now();
232
233    // Initialize hash store for deduplication
234    let hash_store = Arc::new(Mutex::new(
235        HashStore::new(&options.cache_path.to_string_lossy())
236            .context("Failed to initialize hash store")?,
237    ));
238
239    let total_files = options.files.len();
240
241    if total_files == 0 {
242        println!("No files found to upload");
243        return Ok(UploadStats::empty(start_time));
244    }
245
246    // Images already in the series' albums, so unchanged files are skipped
247    // and locally-edited files are replaced in place instead of hitting a 409
248    // (or being uploaded again into a later album of the series).
249    let mut by_filename: HashMap<String, AlbumImage> = HashMap::new();
250    let mut by_md5: HashMap<String, String> = HashMap::new();
251    let mut remote_ok = true;
252    for existing in options.series.existing_albums().await {
253        println!(
254            "Fetching existing images from album '{}'...",
255            existing.album.name
256        );
257        match fetch_remote_image_maps(
258            &options.client,
259            &existing.album.album_key,
260            options.check_remote,
261        )
262        .await
263        {
264            Ok((images, md5s)) => {
265                if let Some(images) = images {
266                    println!("Found {} existing images", images.len());
267                    by_filename.extend(images.iter().map(|(k, v)| (k.clone(), v.clone())));
268                }
269                if let Some(md5s) = md5s {
270                    by_md5.extend(md5s.iter().map(|(k, v)| (k.clone(), v.clone())));
271                }
272            }
273            Err(e) => {
274                println!("Warning: Failed to fetch existing album images: {}", e);
275                remote_ok = false;
276            }
277        }
278    }
279    if !remote_ok {
280        println!("Continuing without remote duplicate/replace detection");
281    }
282    let (remote_images, remote_md5s) = if remote_ok {
283        (
284            Some(Arc::new(by_filename)),
285            options.check_remote.then(|| Arc::new(by_md5)),
286        )
287    } else {
288        (None, None)
289    };
290
291    println!("\nFound {} files to process", total_files);
292
293    // Files the cache knows are in other albums are collected into the series
294    // (the workers would skip them as duplicates). Needs the listing above to
295    // know what the albums already hold.
296    let mut files = options.files;
297    let mut collected = 0;
298    let mut collect_failed = 0;
299    let mut known_hashes = None;
300    if remote_ok && !options.no_cache {
301        let store = hash_store.lock().await.clone();
302        let target = collect::TargetContents {
303            album_keys: options
304                .series
305                .existing_albums()
306                .await
307                .into_iter()
308                .map(|a| a.album.album_key)
309                .collect(),
310            file_names: remote_images
311                .iter()
312                .flat_map(|images| images.keys().map(String::as_str))
313                .collect(),
314        };
315        let outcome = collect::plan_collects(
316            files,
317            &target,
318            &store,
319            options.dry_run,
320            options.client.as_ref(),
321            options.series.as_ref(),
322        )
323        .await?;
324        files = outcome.to_upload;
325        collected = outcome.collected;
326        collect_failed = outcome.failed;
327        known_hashes = Some(Arc::new(outcome.hashes));
328    }
329
330    // Setup progress bars
331    let multi_progress = MultiProgress::new();
332    let overall_progress = multi_progress.add(ProgressBar::new(files.len() as u64));
333    overall_progress.set_style(
334        ProgressStyle::default_bar()
335            .template(
336                "{spinner:.green} [{elapsed_precise}] [{bar:40.cyan/blue}] {pos}/{len} ({eta})",
337            )
338            .unwrap()
339            .progress_chars("#>-"),
340    );
341
342    // Track statistics
343    let stats = Arc::new(Mutex::new(UploadStats {
344        total_files,
345        ..UploadStats::empty(start_time)
346    }));
347
348    let context = Arc::new(UploadWorkerContext {
349        client: options.client.clone(),
350        album_uri: String::new(),
351        album_key: String::new(),
352        series: Some(options.series.clone()),
353        hash_store: hash_store.clone(),
354        remote_md5s,
355        remote_images,
356        dry_run: options.dry_run,
357        no_cache: options.no_cache,
358        retry_attempts: options.retry_attempts,
359        render_raw: options.raw_handling == RawHandling::Render,
360        skip_raw_files: Arc::new(std::sync::atomic::AtomicBool::new(false)),
361        known_hashes,
362    });
363
364    let mut queue = UploadQueue::new();
365    for file in files {
366        queue.add(file);
367    }
368
369    run_upload_workers(
370        queue,
371        context,
372        options.threads,
373        stats.clone(),
374        overall_progress.clone(),
375    )
376    .await?;
377
378    overall_progress.finish_with_message("Upload complete");
379
380    let albums_created = options
381        .series
382        .albums()
383        .await
384        .iter()
385        .filter(|a| a.created)
386        .count();
387
388    // Return final statistics
389    let final_stats = stats.lock().await;
390    Ok(UploadStats {
391        total_files: final_stats.total_files,
392        uploaded: final_stats.uploaded,
393        replaced: final_stats.replaced,
394        collected,
395        skipped: final_stats.skipped,
396        failed: final_stats.failed + collect_failed,
397        total_bytes: final_stats.total_bytes,
398        folders_created: 0,
399        albums_created,
400        duration_secs: start_time.elapsed().as_secs(),
401    })
402}
403
404/// Drain `queue` with `threads` concurrent workers uploading into the album
405/// described by `context`, recording results in `stats`.
406async fn run_upload_workers(
407    queue: UploadQueue,
408    context: Arc<UploadWorkerContext>,
409    threads: usize,
410    stats: Arc<Mutex<UploadStats>>,
411    progress: ProgressBar,
412) -> Result<()> {
413    let queue = Arc::new(Mutex::new(queue));
414    let mut handles = vec![];
415
416    for _ in 0..threads.max(1) {
417        let queue = queue.clone();
418        let context = context.clone();
419        let stats = stats.clone();
420        let progress = progress.clone();
421
422        let handle = tokio::spawn(async move {
423            loop {
424                // Get next file from queue
425                let file_path = {
426                    let mut q = queue.lock().await;
427                    q.next()
428                };
429
430                let Some(file_path) = file_path else {
431                    break;
432                };
433
434                // Process the file
435                match upload_worker(&file_path, context.clone()).await {
436                    Ok(status) => {
437                        let mut stats = stats.lock().await;
438                        match status {
439                            UploadStatus::Uploaded { file_size, .. } => {
440                                stats.uploaded += 1;
441                                stats.total_bytes += file_size;
442                                progress.set_message(format!("Uploaded: {}", file_path.display()));
443                            }
444                            UploadStatus::Replaced { file_size, .. } => {
445                                stats.replaced += 1;
446                                stats.total_bytes += file_size;
447                                progress.set_message(format!("Replaced: {}", file_path.display()));
448                            }
449                            UploadStatus::Skipped { .. } => {
450                                stats.skipped += 1;
451                                progress.set_message(format!("Skipped: {}", file_path.display()));
452                            }
453                            UploadStatus::DryRun { file_size, .. } => {
454                                stats.uploaded += 1;
455                                stats.total_bytes += file_size;
456                                progress
457                                    .set_message(format!("Would upload: {}", file_path.display()));
458                            }
459                        }
460                        progress.inc(1);
461                    }
462                    Err(e) => {
463                        let mut stats = stats.lock().await;
464                        stats.failed += 1;
465                        progress.set_message(format!("Failed: {} - {}", file_path.display(), e));
466                        progress.inc(1);
467                    }
468                }
469            }
470        });
471
472        handles.push(handle);
473    }
474
475    // Wait for all workers to complete
476    for handle in handles {
477        handle.await?;
478    }
479    Ok(())
480}
481
482pub struct UploadStructureOptions {
483    pub path: PathBuf,
484    pub client: Arc<SmugMugClient>,
485    pub dry_run: bool,
486    pub check_remote: bool,
487    pub no_cache: bool,
488    pub cache_path: PathBuf,
489    pub retry_attempts: u32,
490    pub has_smugmug_source: bool,
491    pub raw_mode: RawMode,
492}
493
494pub async fn upload_with_structure(options: UploadStructureOptions) -> Result<UploadStats> {
495    use std::collections::HashMap;
496
497    let start_time = std::time::Instant::now();
498
499    // Initialize hash store for deduplication
500    let hash_store = Arc::new(Mutex::new(
501        HashStore::new(&options.cache_path.to_string_lossy())
502            .context("Failed to initialize hash store")?,
503    ));
504
505    // Walk the directory structure and build a map of folders to files
506    println!("Scanning directory structure...");
507    let mut all_files: Vec<PathBuf> = Vec::new();
508
509    fn scan_directory_recursive(dir: &std::path::Path, all_files: &mut Vec<PathBuf>) -> Result<()> {
510        for entry in std::fs::read_dir(dir)? {
511            let entry = entry?;
512            let path = entry.path();
513
514            if path.is_dir() {
515                // Recursively scan subdirectories
516                scan_directory_recursive(&path, all_files)?;
517            } else if path.is_file() {
518                // Check if this is a supported file using the scanner module
519                if crate::scanner::is_supported_file(&path) {
520                    all_files.push(path);
521                }
522            }
523        }
524        Ok(())
525    }
526
527    scan_directory_recursive(&options.path, &mut all_files)?;
528
529    let raw_handling = options.raw_mode.handling(options.has_smugmug_source);
530    let (all_files, raw_selection) = select_raw_files(all_files, raw_handling);
531    print_raw_selection(&raw_selection, options.raw_mode, options.has_smugmug_source);
532
533    let mut folder_map: HashMap<PathBuf, Vec<PathBuf>> = HashMap::new();
534    for path in all_files {
535        let parent = path.parent().unwrap().to_path_buf();
536        folder_map.entry(parent).or_default().push(path);
537    }
538
539    if folder_map.is_empty() {
540        println!("No image files found to upload");
541        return Ok(UploadStats {
542            total_files: 0,
543            uploaded: 0,
544            replaced: 0,
545            collected: 0,
546            skipped: 0,
547            failed: 0,
548            total_bytes: 0,
549            folders_created: 0,
550            albums_created: 0,
551            duration_secs: start_time.elapsed().as_secs(),
552        });
553    }
554
555    // Count total files
556    let total_files: usize = folder_map.values().map(|v| v.len()).sum();
557    println!(
558        "Found {} files in {} folders\n",
559        total_files,
560        folder_map.len()
561    );
562
563    // Track statistics
564    let stats = Arc::new(Mutex::new(UploadStats {
565        total_files,
566        uploaded: 0,
567        replaced: 0,
568        collected: 0,
569        skipped: 0,
570        failed: 0,
571        total_bytes: 0,
572        folders_created: 0,
573        albums_created: 0,
574        duration_secs: 0,
575    }));
576
577    // Get the authenticated user's node URI
578    let auth_user_url = "https://api.smugmug.com/api/v2!authuser";
579    let response = options.client.get_with_auth(auth_user_url).await?;
580    let body_text = response.text().await?;
581
582    #[derive(serde::Deserialize)]
583    struct UserResponse {
584        #[serde(rename = "Response")]
585        response: UserResponseData,
586    }
587
588    #[derive(serde::Deserialize)]
589    struct UserResponseData {
590        #[serde(rename = "User")]
591        user: UserInfo,
592    }
593
594    #[derive(serde::Deserialize)]
595    struct UserInfo {
596        #[serde(rename = "Uris")]
597        uris: UserUris,
598    }
599
600    #[derive(serde::Deserialize)]
601    struct UserUris {
602        #[serde(rename = "Node")]
603        node: NodeUriInfo,
604    }
605
606    #[derive(serde::Deserialize)]
607    struct NodeUriInfo {
608        #[serde(rename = "Uri")]
609        uri: String,
610    }
611
612    let user_data: UserResponse = serde_json::from_str(&body_text)?;
613    let root_node_uri = user_data.response.user.uris.node.uri;
614
615    // Create a map to cache node URIs for each folder path
616    let node_cache: Arc<Mutex<HashMap<PathBuf, String>>> = Arc::new(Mutex::new(HashMap::new()));
617    node_cache
618        .lock()
619        .await
620        .insert(options.path.clone(), root_node_uri);
621
622    // Setup progress bar
623    let multi_progress = MultiProgress::new();
624    let overall_progress = multi_progress.add(ProgressBar::new(total_files as u64));
625    overall_progress.set_style(
626        ProgressStyle::default_bar()
627            .template(
628                "{spinner:.green} [{elapsed_precise}] [{bar:40.cyan/blue}] {pos}/{len} ({eta})",
629            )
630            .unwrap()
631            .progress_chars("#>-"),
632    );
633
634    // Process each folder
635    for (folder_path, files) in folder_map {
636        // Get or create the folder/album structure
637        let album = match get_or_create_folder_structure(
638            &options.client,
639            &options.path,
640            &folder_path,
641            &node_cache,
642            &stats,
643        )
644        .await
645        {
646            Ok(album) => album,
647            Err(e) => {
648                println!(
649                    "✗ Failed to create folder structure for {}: {}",
650                    folder_path.display(),
651                    e
652                );
653                let mut s = stats.lock().await;
654                s.failed += files.len();
655                continue;
656            }
657        };
658
659        // Fetch existing images in this album so unchanged files are skipped
660        // and locally-edited files are replaced in place instead of 409ing.
661        let (remote_images, remote_md5s) =
662            match fetch_remote_image_maps(&options.client, &album.album_key, options.check_remote)
663                .await
664            {
665                Ok((images, md5s)) => (images, md5s),
666                Err(e) => {
667                    println!(
668                        "Warning: Failed to fetch existing images for {}: {}",
669                        album.name, e
670                    );
671                    (None, None)
672                }
673            };
674
675        // Create worker context for this album
676        let context = Arc::new(UploadWorkerContext {
677            client: options.client.clone(),
678            album_uri: format!("/api/v2/album/{}", album.album_key),
679            album_key: album.album_key.clone(),
680            series: None,
681            hash_store: hash_store.clone(),
682            remote_md5s,
683            remote_images,
684            dry_run: options.dry_run,
685            no_cache: options.no_cache,
686            retry_attempts: options.retry_attempts,
687            render_raw: raw_handling == RawHandling::Render,
688            skip_raw_files: Arc::new(std::sync::atomic::AtomicBool::new(false)),
689            known_hashes: None,
690        });
691
692        // Upload files in this folder
693        for file_path in files {
694            match upload_worker(&file_path, context.clone()).await {
695                Ok(status) => {
696                    let mut s = stats.lock().await;
697                    match status {
698                        UploadStatus::Uploaded { file_size, .. } => {
699                            s.uploaded += 1;
700                            s.total_bytes += file_size;
701                            overall_progress
702                                .set_message(format!("Uploaded: {}", file_path.display()));
703                        }
704                        UploadStatus::Replaced { file_size, .. } => {
705                            s.replaced += 1;
706                            s.total_bytes += file_size;
707                            overall_progress
708                                .set_message(format!("Replaced: {}", file_path.display()));
709                        }
710                        UploadStatus::Skipped { .. } => {
711                            s.skipped += 1;
712                            overall_progress
713                                .set_message(format!("Skipped: {}", file_path.display()));
714                        }
715                        UploadStatus::DryRun { file_size, .. } => {
716                            s.uploaded += 1;
717                            s.total_bytes += file_size;
718                            overall_progress
719                                .set_message(format!("Would upload: {}", file_path.display()));
720                        }
721                    }
722                    overall_progress.inc(1);
723                }
724                Err(e) => {
725                    let mut s = stats.lock().await;
726                    s.failed += 1;
727                    overall_progress.set_message(format!(
728                        "Failed: {} - {}",
729                        file_path.display(),
730                        e
731                    ));
732                    overall_progress.inc(1);
733                }
734            }
735        }
736    }
737
738    overall_progress.finish_with_message("Upload complete");
739
740    // Return final statistics
741    let final_stats = stats.lock().await;
742    Ok(UploadStats {
743        total_files: final_stats.total_files,
744        uploaded: final_stats.uploaded,
745        replaced: final_stats.replaced,
746        collected: 0,
747        skipped: final_stats.skipped,
748        failed: final_stats.failed,
749        total_bytes: final_stats.total_bytes,
750        folders_created: final_stats.folders_created,
751        albums_created: final_stats.albums_created,
752        duration_secs: start_time.elapsed().as_secs(),
753    })
754}
755
756async fn get_or_create_folder_structure(
757    client: &Arc<SmugMugClient>,
758    base_path: &std::path::Path,
759    target_path: &std::path::Path,
760    node_cache: &Arc<Mutex<HashMap<PathBuf, String>>>,
761    stats: &Arc<Mutex<UploadStats>>,
762) -> Result<Album> {
763    // Get the relative path from base to target
764    let rel_path = target_path.strip_prefix(base_path)?;
765
766    // If this is the base directory, create an album at the root
767    if rel_path.as_os_str().is_empty() {
768        let album_name = base_path
769            .file_name()
770            .and_then(|n| n.to_str())
771            .unwrap_or("Uploads")
772            .to_string();
773
774        let cache = node_cache.lock().await;
775        let parent_node_uri = cache.get(base_path).unwrap().clone();
776        drop(cache);
777
778        let album = match client
779            .create_album(&album_name, Some(&parent_node_uri), "Private")
780            .await
781        {
782            Ok(album) => {
783                let mut s = stats.lock().await;
784                s.albums_created += 1;
785                album
786            }
787            Err(e) => {
788                // If we get a conflict (409), try to find the existing album
789                let error_msg = e.to_string();
790                if error_msg.contains("409") || error_msg.contains("Conflict") {
791                    // Album already exists, try to find it
792                    match client
793                        .find_album_in_folder(&parent_node_uri, &album_name)
794                        .await?
795                    {
796                        Some(album) => album,
797                        None => return Err(e), // Album should exist but we can't find it
798                    }
799                } else {
800                    return Err(e);
801                }
802            }
803        };
804        return Ok(album);
805    }
806
807    // Walk through each component of the path, creating folders/albums as needed
808    let mut current_path = base_path.to_path_buf();
809    let components: Vec<_> = rel_path.components().collect();
810
811    for (i, component) in components.iter().enumerate() {
812        let component_name = component.as_os_str().to_str().unwrap();
813        current_path.push(component_name);
814
815        // Check if we already have a node for this path
816        let existing_node = {
817            let cache = node_cache.lock().await;
818            cache.get(&current_path).cloned()
819        };
820
821        if existing_node.is_none() {
822            // Get parent node URI
823            let parent_path = current_path.parent().unwrap().to_path_buf();
824            let parent_node_uri = {
825                let cache = node_cache.lock().await;
826                cache.get(&parent_path).unwrap().clone()
827            };
828
829            // Determine if this should be a folder or album
830            let is_last = i == components.len() - 1;
831
832            if is_last {
833                // Last component - create an album (private by default)
834                let album = match client
835                    .create_album(component_name, Some(&parent_node_uri), "Private")
836                    .await
837                {
838                    Ok(album) => {
839                        let mut s = stats.lock().await;
840                        s.albums_created += 1;
841                        album
842                    }
843                    Err(e) => {
844                        // If we get a conflict (409), try to find the existing album
845                        let error_msg = e.to_string();
846                        if error_msg.contains("409") || error_msg.contains("Conflict") {
847                            // Album already exists, try to find it
848                            match client
849                                .find_album_in_folder(&parent_node_uri, component_name)
850                                .await?
851                            {
852                                Some(album) => album,
853                                None => return Err(e), // Album should exist but we can't find it
854                            }
855                        } else {
856                            return Err(e);
857                        }
858                    }
859                };
860
861                // Cache the album's node URI
862                let mut cache = node_cache.lock().await;
863                cache.insert(
864                    current_path.clone(),
865                    format!("/api/v2/node/{}", album.node_id),
866                );
867
868                return Ok(album);
869            } else {
870                // Intermediate component - create a folder
871                let folder_node_uri = match create_folder(client, component_name, &parent_node_uri)
872                    .await
873                {
874                    Ok(uri) => {
875                        let mut s = stats.lock().await;
876                        s.folders_created += 1;
877                        uri
878                    }
879                    Err(e) => {
880                        // If we get a conflict (409), try to find the existing folder
881                        let error_msg = e.to_string();
882                        if error_msg.contains("409") || error_msg.contains("Conflict") {
883                            // Folder already exists, find it
884                            find_existing_folder(client, component_name, &parent_node_uri).await?
885                        } else {
886                            return Err(e);
887                        }
888                    }
889                };
890
891                let mut cache = node_cache.lock().await;
892                cache.insert(current_path.clone(), folder_node_uri);
893            }
894        }
895    }
896
897    anyhow::bail!("Failed to create folder structure")
898}
899
900async fn create_folder(
901    client: &Arc<SmugMugClient>,
902    name: &str,
903    parent_node_uri: &str,
904) -> Result<String> {
905    let create_url = format!("https://api.smugmug.com{}!children", parent_node_uri);
906
907    let body = serde_json::json!({
908        "Type": "Folder",
909        "Name": name,
910    });
911
912    let response = client.post_with_auth(&create_url, body).await?;
913
914    let status = response.status();
915    let body_text = response.text().await?;
916
917    if !status.is_success() {
918        anyhow::bail!("Failed to create folder: {} - {}", status, body_text);
919    }
920
921    #[derive(serde::Deserialize)]
922    struct CreateNodeResponse {
923        #[serde(rename = "Response")]
924        response: CreateNodeResponseData,
925    }
926
927    #[derive(serde::Deserialize)]
928    struct CreateNodeResponseData {
929        #[serde(rename = "Node")]
930        node: NodeInfo,
931    }
932
933    #[derive(serde::Deserialize)]
934    struct NodeInfo {
935        #[serde(rename = "Uri")]
936        uri: String,
937    }
938
939    let node_response: CreateNodeResponse = serde_json::from_str(&body_text)?;
940    Ok(node_response.response.node.uri)
941}
942
943async fn find_existing_folder(
944    client: &Arc<SmugMugClient>,
945    name: &str,
946    parent_node_uri: &str,
947) -> Result<String> {
948    #[derive(serde::Deserialize)]
949    struct NodeInfo {
950        #[serde(rename = "Name")]
951        name: String,
952        #[serde(rename = "Type")]
953        node_type: String,
954        #[serde(rename = "Uri")]
955        uri: String,
956    }
957
958    // Get every child of the parent node (all pages)
959    let children_url = format!("https://api.smugmug.com{}!children", parent_node_uri);
960    let nodes: Vec<NodeInfo> = client.get_all_pages(&children_url, "Node").await?;
961
962    // Find the folder with the matching name
963    for node in nodes {
964        if node.node_type == "Folder" && node.name == name {
965            return Ok(node.uri);
966        }
967    }
968
969    anyhow::bail!("Folder '{}' not found in parent node", name)
970}
971
972#[cfg(test)]
973mod tests {
974    use super::*;
975
976    // Tests for upload orchestration, options, statistics, and concurrent operations
977    // Focuses on structure initialization, stats tracking, and thread safety
978
979    #[test]
980    fn test_select_raw_files_upload_and_skip() {
981        let files = vec![
982            PathBuf::from("a.jpg"),
983            PathBuf::from("b.CR2"),
984            PathBuf::from("c.png"),
985            PathBuf::from("d.nef"),
986        ];
987
988        let (kept, selection) = select_raw_files(files.clone(), RawHandling::Skip);
989        assert_eq!(kept, vec![PathBuf::from("a.jpg"), PathBuf::from("c.png")]);
990        assert_eq!(selection.skipped, 2);
991
992        let (kept, selection) = select_raw_files(files.clone(), RawHandling::Upload);
993        assert_eq!(kept, files);
994        assert_eq!(selection, RawSelection::default());
995    }
996
997    #[test]
998    fn test_select_raw_files_render_skips_siblings() {
999        let files = vec![
1000            PathBuf::from("/p/IMG_1.CR2"),
1001            PathBuf::from("/p/img_1.JPG"), // camera JPEG for IMG_1
1002            PathBuf::from("/p/IMG_2.CR2"),
1003            PathBuf::from("/p/IMG_2.dng"), // renders to the same name as IMG_2.CR2
1004            PathBuf::from("/q/IMG_1.CR2"), // other directory, no JPEG there
1005            PathBuf::from("/q/IMG_3.nef"),
1006            PathBuf::from("/q/IMG_3.heic"),
1007        ];
1008
1009        let (kept, selection) = select_raw_files(files, RawHandling::Render);
1010        assert_eq!(
1011            kept,
1012            vec![
1013                PathBuf::from("/p/img_1.JPG"),
1014                PathBuf::from("/p/IMG_2.CR2"),
1015                PathBuf::from("/q/IMG_1.CR2"),
1016                PathBuf::from("/q/IMG_3.heic"),
1017            ]
1018        );
1019        assert_eq!(
1020            selection,
1021            RawSelection {
1022                rendered: 2,
1023                skipped: 0,
1024                skipped_with_sibling: 3,
1025            }
1026        );
1027    }
1028
1029    #[test]
1030    fn test_upload_stats_initial() {
1031        let stats = UploadStats {
1032            total_files: 10,
1033            uploaded: 0,
1034            replaced: 0,
1035            collected: 0,
1036            skipped: 0,
1037            failed: 0,
1038            total_bytes: 0,
1039            folders_created: 0,
1040            albums_created: 0,
1041            duration_secs: 0,
1042        };
1043
1044        assert_eq!(stats.total_files, 10);
1045        assert_eq!(stats.uploaded, 0);
1046        assert_eq!(stats.skipped, 0);
1047        assert_eq!(stats.failed, 0);
1048        assert_eq!(stats.total_bytes, 0);
1049        assert_eq!(stats.folders_created, 0);
1050        assert_eq!(stats.albums_created, 0);
1051    }
1052
1053    #[test]
1054    fn test_upload_stats_progress() {
1055        let mut stats = UploadStats {
1056            total_files: 10,
1057            uploaded: 0,
1058            replaced: 0,
1059            collected: 0,
1060            skipped: 0,
1061            failed: 0,
1062            total_bytes: 0,
1063            folders_created: 0,
1064            albums_created: 0,
1065            duration_secs: 0,
1066        };
1067
1068        // Simulate progress
1069        stats.uploaded += 5;
1070        stats.skipped += 2;
1071        stats.failed += 1;
1072        stats.total_bytes += 5000000; // 5MB
1073
1074        assert_eq!(stats.uploaded, 5);
1075        assert_eq!(stats.skipped, 2);
1076        assert_eq!(stats.failed, 1);
1077        assert_eq!(stats.total_bytes, 5000000);
1078        assert_eq!(stats.uploaded + stats.skipped + stats.failed, 8);
1079    }
1080
1081    #[test]
1082    fn test_upload_structure_options_creation() {
1083        let client = Arc::new(SmugMugClient::new(
1084            "key".to_string(),
1085            "secret".to_string(),
1086            "token".to_string(),
1087            "token_secret".to_string(),
1088        ));
1089
1090        let options = UploadStructureOptions {
1091            path: PathBuf::from("/test/structure"),
1092            client: client.clone(),
1093            dry_run: false,
1094            check_remote: true,
1095            no_cache: false,
1096            cache_path: PathBuf::from("/cache/path"),
1097            retry_attempts: 3,
1098            has_smugmug_source: false,
1099            raw_mode: RawMode::Auto,
1100        };
1101
1102        assert_eq!(options.path, PathBuf::from("/test/structure"));
1103        assert!(!options.dry_run);
1104        assert!(options.check_remote);
1105        assert!(!options.no_cache);
1106        assert_eq!(options.cache_path, PathBuf::from("/cache/path"));
1107    }
1108
1109    #[tokio::test]
1110    async fn test_queue_thread_safety() {
1111        use std::sync::Arc;
1112        use tokio::sync::Mutex;
1113
1114        let queue = Arc::new(Mutex::new(UploadQueue::new()));
1115
1116        // Add files to queue
1117        for i in 0..10 {
1118            let mut q = queue.lock().await;
1119            q.add(PathBuf::from(format!("/test/file{}.jpg", i)));
1120        }
1121
1122        // Spawn multiple workers to consume from queue
1123        let mut handles = vec![];
1124        for _ in 0..3 {
1125            let q = queue.clone();
1126            let handle = tokio::spawn(async move {
1127                let mut count = 0;
1128                loop {
1129                    let item = {
1130                        let mut queue = q.lock().await;
1131                        queue.next()
1132                    };
1133                    if item.is_none() {
1134                        break;
1135                    }
1136                    count += 1;
1137                }
1138                count
1139            });
1140            handles.push(handle);
1141        }
1142
1143        // Collect results
1144        let mut total_processed = 0;
1145        for handle in handles {
1146            total_processed += handle.await.unwrap();
1147        }
1148
1149        assert_eq!(total_processed, 10);
1150
1151        // Queue should be empty
1152        let mut q = queue.lock().await;
1153        assert!(q.next().is_none());
1154    }
1155
1156    #[tokio::test]
1157    async fn test_upload_stats_concurrent_updates() {
1158        let stats = Arc::new(Mutex::new(UploadStats {
1159            total_files: 100,
1160            uploaded: 0,
1161            replaced: 0,
1162            collected: 0,
1163            skipped: 0,
1164            failed: 0,
1165            total_bytes: 0,
1166            folders_created: 0,
1167            albums_created: 0,
1168            duration_secs: 0,
1169        }));
1170
1171        let mut handles = vec![];
1172
1173        // Spawn multiple tasks updating stats
1174        for _ in 0..10 {
1175            let s = stats.clone();
1176            let handle = tokio::spawn(async move {
1177                for _ in 0..10 {
1178                    let mut stats = s.lock().await;
1179                    stats.uploaded += 1;
1180                    stats.total_bytes += 1024;
1181                }
1182            });
1183            handles.push(handle);
1184        }
1185
1186        // Wait for all tasks
1187        for handle in handles {
1188            handle.await.unwrap();
1189        }
1190
1191        let final_stats = stats.lock().await;
1192        assert_eq!(final_stats.uploaded, 100);
1193        assert_eq!(final_stats.total_bytes, 102400);
1194    }
1195
1196    #[test]
1197    fn test_upload_stats_zero_state() {
1198        let stats = UploadStats {
1199            total_files: 0,
1200            uploaded: 0,
1201            replaced: 0,
1202            collected: 0,
1203            skipped: 0,
1204            failed: 0,
1205            total_bytes: 0,
1206            folders_created: 0,
1207            albums_created: 0,
1208            duration_secs: 0,
1209        };
1210
1211        assert_eq!(stats.total_files, 0);
1212        assert_eq!(stats.uploaded + stats.skipped + stats.failed, 0);
1213    }
1214
1215    #[test]
1216    fn test_upload_stats_large_numbers() {
1217        let stats = UploadStats {
1218            total_files: 10000,
1219            uploaded: 8500,
1220            replaced: 0,
1221            collected: 0,
1222            skipped: 1200,
1223            failed: 300,
1224            total_bytes: 50_000_000_000, // 50GB
1225            folders_created: 100,
1226            albums_created: 50,
1227            duration_secs: 0,
1228        };
1229
1230        assert_eq!(stats.total_files, 10000);
1231        assert_eq!(stats.uploaded, 8500);
1232        assert_eq!(stats.skipped, 1200);
1233        assert_eq!(stats.failed, 300);
1234        assert_eq!(stats.total_bytes, 50_000_000_000);
1235        assert_eq!(stats.folders_created, 100);
1236        assert_eq!(stats.albums_created, 50);
1237    }
1238
1239    #[tokio::test]
1240    async fn test_multiple_workers_processing_queue() {
1241        let queue = Arc::new(Mutex::new(UploadQueue::new()));
1242
1243        // Add 100 files
1244        for i in 0..100 {
1245            let mut q = queue.lock().await;
1246            q.add(PathBuf::from(format!("/test/image{}.jpg", i)));
1247        }
1248
1249        let processed = Arc::new(Mutex::new(Vec::new()));
1250        let mut handles = vec![];
1251
1252        // Create 5 workers
1253        for worker_id in 0..5 {
1254            let q = queue.clone();
1255            let p = processed.clone();
1256
1257            let handle = tokio::spawn(async move {
1258                loop {
1259                    let file = {
1260                        let mut queue = q.lock().await;
1261                        queue.next()
1262                    };
1263
1264                    match file {
1265                        Some(path) => {
1266                            let mut proc = p.lock().await;
1267                            proc.push((worker_id, path));
1268                        }
1269                        None => break,
1270                    }
1271                }
1272            });
1273
1274            handles.push(handle);
1275        }
1276
1277        // Wait for all workers
1278        for handle in handles {
1279            handle.await.unwrap();
1280        }
1281
1282        let proc = processed.lock().await;
1283        assert_eq!(proc.len(), 100);
1284
1285        // Verify all files were processed
1286        let mut file_indices: Vec<usize> = proc
1287            .iter()
1288            .map(|(_, path)| {
1289                let name = path.file_name().unwrap().to_str().unwrap();
1290                let num_str = name
1291                    .strip_prefix("image")
1292                    .unwrap()
1293                    .strip_suffix(".jpg")
1294                    .unwrap();
1295                num_str.parse().unwrap()
1296            })
1297            .collect();
1298
1299        file_indices.sort();
1300        assert_eq!(file_indices, (0..100).collect::<Vec<_>>());
1301    }
1302
1303    #[test]
1304    fn test_upload_stats_calculation() {
1305        let stats = UploadStats {
1306            total_files: 100,
1307            uploaded: 70,
1308            replaced: 0,
1309            collected: 0,
1310            skipped: 20,
1311            failed: 10,
1312            total_bytes: 1_073_741_824, // 1GB
1313            folders_created: 5,
1314            albums_created: 3,
1315            duration_secs: 0,
1316        };
1317
1318        // Verify all files accounted for
1319        assert_eq!(
1320            stats.uploaded + stats.skipped + stats.failed,
1321            stats.total_files
1322        );
1323
1324        // Check individual values
1325        assert_eq!(stats.uploaded, 70);
1326        assert_eq!(stats.skipped, 20);
1327        assert_eq!(stats.failed, 10);
1328        assert_eq!(stats.total_bytes, 1_073_741_824);
1329    }
1330}