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