Skip to main content

smugmug_cli/uploader/
collect.rs

1//! Adding photos that are already on SmugMug to the target album
2//! ("collecting" them) instead of uploading a second copy.
3//!
4//! The cache maps each uploaded file to its SmugMug image, and the upload
5//! workers skip any file it knows. So a file uploaded to album A and then
6//! uploaded to album B would be skipped and never appear in B. SmugMug
7//! doesn't deduplicate uploads (uploading it again would store a second,
8//! separate image), but an image can be collected into any number of
9//! albums. Before the workers start, `plan_collects` collects files the
10//! cache knows are in another album into the target album, from the cached
11//! image URI, in batches of `COLLECT_BATCH` per album.
12//!
13//! Collected images take room in the album series like uploads do. Files
14//! SmugMug refuses (the image was deleted) go on to the upload workers.
15
16use anyhow::Result;
17use chrono::Utc;
18use colored::Colorize;
19use rayon::prelude::*;
20use std::collections::{HashMap, HashSet};
21use std::path::PathBuf;
22
23use crate::api::SmugMugClient;
24use crate::api::albums::Album;
25use crate::api::images::CollectResult;
26use crate::cache::hash_store::{HashStore, UploadedFile};
27use crate::uploader::album_series::{AlbumSeries, AlbumSeriesBackend};
28use crate::uploader::worker::calculate_file_hash;
29
30/// Image URIs per `!collectimages` request.
31pub(crate) const COLLECT_BATCH: usize = 100;
32
33/// How collects reach SmugMug.
34// Only implemented and used inside this crate with concrete types, so the
35// Send-bound caveat of async fns in public traits doesn't matter here.
36#[allow(async_fn_in_trait)]
37pub trait CollectBackend {
38    async fn collect(&self, album_key: &str, uris: &[String]) -> Result<CollectResult>;
39}
40
41impl CollectBackend for SmugMugClient {
42    async fn collect(&self, album_key: &str, uris: &[String]) -> Result<CollectResult> {
43        self.collect_images(album_key, uris).await
44    }
45}
46
47/// What the target albums already hold, as listed before the upload.
48pub struct TargetContents<'a> {
49    /// Keys of the series' albums.
50    pub album_keys: HashSet<String>,
51    /// File names of images in them.
52    pub file_names: HashSet<&'a str>,
53}
54
55#[derive(Debug, Default)]
56pub struct CollectOutcome {
57    /// Files left for the upload workers.
58    pub to_upload: Vec<PathBuf>,
59    /// SHA-256 of files the planner hashed, so the workers needn't again.
60    pub hashes: HashMap<PathBuf, String>,
61    /// Files collected (on a dry run: that would be).
62    pub collected: usize,
63    /// Files that couldn't be collected because a request failed. They're
64    /// neither collected nor uploaded (that would duplicate them); their
65    /// cache entries stay, so the next run tries again.
66    pub failed: usize,
67}
68
69/// A file to collect and the image it's already on SmugMug as.
70pub(crate) struct Pending {
71    pub(crate) path: PathBuf,
72    pub(crate) hash: String,
73    pub(crate) file_size: u64,
74    pub(crate) image_uri: String,
75    pub(crate) image_key: String,
76}
77
78/// Collect the files in `files` that the cache knows are in another album
79/// into the series' albums, and return the rest for uploading.
80pub async fn plan_collects<C: CollectBackend, S: AlbumSeriesBackend>(
81    files: Vec<PathBuf>,
82    target: &TargetContents<'_>,
83    store: &HashStore,
84    dry_run: bool,
85    backend: &C,
86    series: &AlbumSeries<S>,
87) -> Result<CollectOutcome> {
88    let mut outcome = CollectOutcome::default();
89
90    // Hash and look up every file, in parallel and without requests.
91    let classified: Vec<(PathBuf, Option<String>, Option<Pending>)> = files
92        .into_par_iter()
93        .map(|path| {
94            let (hash, pending) = classify(&path, target, store);
95            (path, hash, pending)
96        })
97        .collect();
98
99    let mut pending = Vec::new();
100    for (path, hash, collect) in classified {
101        if let Some(hash) = hash {
102            outcome.hashes.insert(path.clone(), hash);
103        }
104        match collect {
105            Some(p) => pending.push(p),
106            None => outcome.to_upload.push(path),
107        }
108    }
109
110    if !pending.is_empty() {
111        let collect = collect_pending(pending, store, dry_run, backend, series).await?;
112        outcome.collected = collect.collected;
113        outcome.failed = collect.failed;
114        outcome.to_upload.extend(collect.refused);
115    }
116
117    Ok(outcome)
118}
119
120/// The file's hash, and what to collect if the cache knows the file is in
121/// another album. Everything else is left to the worker.
122fn classify(
123    path: &PathBuf,
124    target: &TargetContents<'_>,
125    store: &HashStore,
126) -> (Option<String>, Option<Pending>) {
127    // Same name already in the album: the worker skips or replaces it.
128    let file_name = path.file_name().map(|n| n.to_string_lossy());
129    if file_name.is_some_and(|n| target.file_names.contains(n.as_ref())) {
130        return (None, None);
131    }
132    // Errors here resurface (and are reported) in the worker.
133    let Ok(hash) = calculate_file_hash(path) else {
134        return (None, None);
135    };
136    let Ok(Some(cached)) = store.get(&hash) else {
137        return (Some(hash), None);
138    };
139    // Already in the series: the worker skips it.
140    if target.album_keys.contains(&cached.album_key) {
141        return (Some(hash), None);
142    }
143    let pending = Pending {
144        path: path.clone(),
145        hash: hash.clone(),
146        file_size: std::fs::metadata(path).map(|m| m.len()).unwrap_or(0),
147        image_uri: cached.smugmug_uri,
148        image_key: cached.image_key,
149    };
150    (Some(hash), Some(pending))
151}
152
153/// How collecting went.
154#[derive(Default)]
155struct Collected {
156    collected: usize,
157    /// Files SmugMug refused to collect (the image is gone or can't be
158    /// collected), to upload instead.
159    refused: Vec<PathBuf>,
160    /// Files not collected because a request failed.
161    failed: usize,
162}
163
164/// Claim room for each file in the series and collect them album by album.
165async fn collect_pending<C: CollectBackend, S: AlbumSeriesBackend>(
166    pending: Vec<Pending>,
167    store: &HashStore,
168    dry_run: bool,
169    backend: &C,
170    series: &AlbumSeries<S>,
171) -> Result<Collected> {
172    // Group by album, keeping the series' order.
173    let mut groups: Vec<(Album, Vec<Pending>)> = Vec::new();
174    for p in pending {
175        let album = series.claim().await?;
176        match groups.iter_mut().find(|(a, _)| a.name == album.name) {
177            Some((_, members)) => members.push(p),
178            None => groups.push((album, vec![p])),
179        }
180    }
181
182    let total: usize = groups.iter().map(|(_, m)| m.len()).sum();
183    if dry_run {
184        return Ok(Collected {
185            collected: total,
186            ..Default::default()
187        });
188    }
189
190    println!("Adding {} files already on SmugMug to the album...", total);
191    let mut outcome = Collected::default();
192    for (album, members) in groups {
193        let mut members = members.into_iter().peekable();
194        while members.peek().is_some() {
195            let batch: Vec<Pending> = members.by_ref().take(COLLECT_BATCH).collect();
196            let batch = collect_batch(batch, &album, store, backend).await;
197            outcome.collected += batch.collected;
198            for _ in 0..batch.refused.len() + batch.failed.len() {
199                series.release(&album).await;
200            }
201            outcome
202                .refused
203                .extend(batch.refused.into_iter().map(|p| p.path));
204            outcome.failed += batch.failed.len();
205        }
206    }
207
208    if !outcome.refused.is_empty() {
209        println!(
210            "{}",
211            format!(
212                "⚠ SmugMug refused to add {} files (no longer there?); uploading them instead",
213                outcome.refused.len()
214            )
215            .yellow()
216        );
217    }
218    if outcome.failed > 0 {
219        println!(
220            "{}",
221            format!(
222                "⚠ {} files couldn't be added from SmugMug; run the upload again to retry",
223                outcome.failed
224            )
225            .yellow()
226        );
227    }
228    Ok(outcome)
229}
230
231/// How one batch went.
232pub(crate) struct BatchOutcome {
233    pub(crate) collected: usize,
234    pub(crate) refused: Vec<Pending>,
235    pub(crate) failed: Vec<Pending>,
236}
237
238/// Collect one batch into `album`, recording each collected file in the
239/// cache.
240pub(crate) async fn collect_batch<C: CollectBackend>(
241    batch: Vec<Pending>,
242    album: &Album,
243    store: &HashStore,
244    backend: &C,
245) -> BatchOutcome {
246    let uris: Vec<String> = batch.iter().map(|p| p.image_uri.clone()).collect();
247    let (mut good, refused) = match backend.collect(&album.album_key, &uris).await {
248        Ok(result) => batch
249            .into_iter()
250            .partition::<Vec<Pending>, _>(|p| !result.rejected.contains_key(&p.image_uri)),
251        Err(e) => {
252            println!("Warning: Couldn't add files to '{}': {}", album.name, e);
253            return BatchOutcome {
254                collected: 0,
255                refused: Vec::new(),
256                failed: batch,
257            };
258        }
259    };
260
261    // The cache points at an image that's gone or unusable: forget it, so
262    // the worker uploads the file instead of skipping it as a duplicate.
263    for p in &refused {
264        let _ = store.remove(&p.hash);
265    }
266
267    // A refused URI fails the whole request, though SmugMug may have
268    // collected the others; collecting them again settles it.
269    let mut failed = Vec::new();
270    if !refused.is_empty() && !good.is_empty() {
271        let uris: Vec<String> = good.iter().map(|p| p.image_uri.clone()).collect();
272        match backend.collect(&album.album_key, &uris).await {
273            Ok(result) if result.rejected.is_empty() => {}
274            outcome => {
275                if let Err(e) = outcome {
276                    println!("Warning: Couldn't add files to '{}': {}", album.name, e);
277                }
278                failed.append(&mut good);
279            }
280        }
281    }
282
283    for p in &good {
284        let _ = store.insert(
285            &p.hash,
286            UploadedFile {
287                smugmug_uri: p.image_uri.clone(),
288                album_key: album.album_key.clone(),
289                image_key: p.image_key.clone(),
290                uploaded_at: Utc::now(),
291                file_size: p.file_size,
292                original_path: p.path.to_string_lossy().to_string(),
293            },
294        );
295    }
296    BatchOutcome {
297        collected: good.len(),
298        refused,
299        failed,
300    }
301}
302
303#[cfg(test)]
304mod tests {
305    use super::*;
306    use crate::uploader::album_series::MAX_ALBUM_IMAGES;
307    use std::sync::Mutex as StdMutex;
308
309    /// In-memory SmugMug collects.
310    #[derive(Default)]
311    struct FakeSmugMug {
312        /// URIs `collect` refuses ("Does not exist").
313        missing: HashSet<String>,
314        fail_collect: bool,
315        /// Fail every collect after the first.
316        fail_on_retry: bool,
317        /// (album key, URIs) per collect call.
318        collects: StdMutex<Vec<(String, Vec<String>)>>,
319    }
320
321    impl CollectBackend for FakeSmugMug {
322        async fn collect(&self, album_key: &str, uris: &[String]) -> Result<CollectResult> {
323            self.collects
324                .lock()
325                .unwrap()
326                .push((album_key.to_string(), uris.to_vec()));
327            let calls = self.collects.lock().unwrap().len();
328            if self.fail_collect || (self.fail_on_retry && calls > 1) {
329                anyhow::bail!("503 Service Unavailable");
330            }
331            let rejected = uris
332                .iter()
333                .filter(|u| self.missing.contains(*u))
334                .map(|u| (u.clone(), vec!["Does not exist".to_string()]))
335                .collect();
336            Ok(CollectResult { rejected })
337        }
338    }
339
340    /// Albums of a series: name -> image count.
341    struct FakeSeries(StdMutex<HashMap<String, u64>>);
342
343    impl AlbumSeriesBackend for FakeSeries {
344        async fn find_album(&self, name: &str) -> Result<Option<Album>> {
345            Ok(self
346                .0
347                .lock()
348                .unwrap()
349                .contains_key(name)
350                .then(|| album(name)))
351        }
352        async fn image_count(&self, album: &Album) -> Result<u64> {
353            Ok(self.0.lock().unwrap()[&album.name])
354        }
355        async fn create_album(&self, name: &str) -> Result<Album> {
356            self.0.lock().unwrap().insert(name.to_string(), 0);
357            Ok(album(name))
358        }
359    }
360
361    fn album(name: &str) -> Album {
362        Album {
363            album_key: format!("key-{}", name),
364            name: name.to_string(),
365            url_name: String::new(),
366            node_id: String::new(),
367            uri: format!("/api/v2/album/key-{}", name),
368            web_uri: None,
369            uris: None,
370            image_count: None,
371        }
372    }
373
374    async fn series(albums: &[(&str, u64)], capacity: u64) -> AlbumSeries<FakeSeries> {
375        let backend = FakeSeries(StdMutex::new(
376            albums.iter().map(|(n, c)| (n.to_string(), *c)).collect(),
377        ));
378        AlbumSeries::load(backend, "Trip", capacity, false)
379            .await
380            .unwrap()
381    }
382
383    fn photo(dir: &std::path::Path, name: &str) -> PathBuf {
384        let path = dir.join(name);
385        std::fs::write(&path, format!("photo {}", name)).unwrap();
386        path
387    }
388
389    fn cached(store: &HashStore, path: &std::path::Path, album_key: &str, image_key: &str) {
390        store
391            .insert(
392                &calculate_file_hash(path).unwrap(),
393                UploadedFile {
394                    smugmug_uri: format!("/api/v2/image/{}-0", image_key),
395                    album_key: album_key.to_string(),
396                    image_key: image_key.to_string(),
397                    uploaded_at: Utc::now(),
398                    file_size: 1,
399                    original_path: path.to_string_lossy().to_string(),
400                },
401            )
402            .unwrap();
403    }
404
405    fn empty_target() -> TargetContents<'static> {
406        TargetContents {
407            album_keys: HashSet::from(["key-Trip".to_string()]),
408            file_names: HashSet::new(),
409        }
410    }
411
412    struct Setup {
413        dir: tempfile::TempDir,
414        store: HashStore,
415    }
416
417    fn setup() -> Setup {
418        let dir = tempfile::tempdir().unwrap();
419        let store = HashStore::new(&dir.path().join("cache").to_string_lossy()).unwrap();
420        Setup { dir, store }
421    }
422
423    #[tokio::test]
424    async fn test_cache_hit_elsewhere_is_collected() {
425        let s = setup();
426        let a = photo(s.dir.path(), "a.jpg");
427        let b = photo(s.dir.path(), "b.jpg");
428        cached(&s.store, &a, "key-Other", "AAA");
429        let smugmug = FakeSmugMug::default();
430        let series = series(&[("Trip", 0)], MAX_ALBUM_IMAGES).await;
431
432        let outcome = plan_collects(
433            vec![a.clone(), b.clone()],
434            &empty_target(),
435            &s.store,
436            false,
437            &smugmug,
438            &series,
439        )
440        .await
441        .unwrap();
442
443        assert_eq!(outcome.collected, 1);
444        assert_eq!(outcome.to_upload, vec![b.clone()]);
445        assert_eq!(
446            *smugmug.collects.lock().unwrap(),
447            vec![(
448                "key-Trip".to_string(),
449                vec!["/api/v2/image/AAA-0".to_string()]
450            )]
451        );
452        // The cache now knows the file is in this album.
453        let entry = s.store.get(&outcome.hashes[&a]).unwrap().unwrap();
454        assert_eq!(entry.album_key, "key-Trip");
455        // Both files were hashed for the workers.
456        assert!(outcome.hashes.contains_key(&b));
457        assert_eq!(series.albums().await[0].claimed, 1);
458    }
459
460    #[tokio::test]
461    async fn test_cache_hit_in_target_or_same_name_is_left_to_worker() {
462        let s = setup();
463        let a = photo(s.dir.path(), "a.jpg");
464        let b = photo(s.dir.path(), "b.jpg");
465        cached(&s.store, &a, "key-Trip", "AAA");
466        cached(&s.store, &b, "key-Other", "BBB");
467        let smugmug = FakeSmugMug::default();
468        let series = series(&[("Trip", 0)], MAX_ALBUM_IMAGES).await;
469        let target = TargetContents {
470            file_names: HashSet::from(["b.jpg"]),
471            ..empty_target()
472        };
473
474        let outcome = plan_collects(
475            vec![a.clone(), b.clone()],
476            &target,
477            &s.store,
478            false,
479            &smugmug,
480            &series,
481        )
482        .await
483        .unwrap();
484
485        assert_eq!(outcome.collected, 0);
486        assert_eq!(outcome.to_upload.len(), 2);
487        assert!(smugmug.collects.lock().unwrap().is_empty());
488    }
489
490    #[tokio::test]
491    async fn test_stale_cache_entry_is_forgotten_and_uploaded() {
492        let s = setup();
493        let a = photo(s.dir.path(), "a.jpg");
494        let b = photo(s.dir.path(), "b.jpg");
495        cached(&s.store, &a, "key-Other", "GONE");
496        cached(&s.store, &b, "key-Other", "BBB");
497        let smugmug = FakeSmugMug {
498            missing: HashSet::from(["/api/v2/image/GONE-0".to_string()]),
499            ..Default::default()
500        };
501        let series = series(&[("Trip", 0)], MAX_ALBUM_IMAGES).await;
502
503        let outcome = plan_collects(
504            vec![a.clone(), b.clone()],
505            &empty_target(),
506            &s.store,
507            false,
508            &smugmug,
509            &series,
510        )
511        .await
512        .unwrap();
513
514        assert_eq!(outcome.collected, 1);
515        assert_eq!(outcome.failed, 0);
516        assert_eq!(outcome.to_upload, vec![a.clone()]);
517        assert!(s.store.get(&outcome.hashes[&a]).unwrap().is_none());
518        // The refused file's room went back.
519        assert_eq!(series.albums().await[0].claimed, 1);
520        // The good one is collected again after the refused batch.
521        let collects = smugmug.collects.lock().unwrap();
522        assert_eq!(collects.len(), 2);
523        assert_eq!(collects[1].1, vec!["/api/v2/image/BBB-0".to_string()]);
524    }
525
526    #[tokio::test]
527    async fn test_failed_collect_fails_without_uploading_or_forgetting() {
528        let s = setup();
529        let a = photo(s.dir.path(), "a.jpg");
530        cached(&s.store, &a, "key-Other", "AAA");
531        let smugmug = FakeSmugMug {
532            fail_collect: true,
533            ..Default::default()
534        };
535        let series = series(&[("Trip", 0)], MAX_ALBUM_IMAGES).await;
536
537        let outcome = plan_collects(
538            vec![a.clone()],
539            &empty_target(),
540            &s.store,
541            false,
542            &smugmug,
543            &series,
544        )
545        .await
546        .unwrap();
547
548        // Not uploaded (it's on SmugMug; that would duplicate it) and not
549        // forgotten, so the next run tries again; reported as failed.
550        assert_eq!(outcome.collected, 0);
551        assert_eq!(outcome.failed, 1);
552        assert!(outcome.to_upload.is_empty());
553        assert!(s.store.get(&outcome.hashes[&a]).unwrap().is_some());
554        assert_eq!(series.albums().await[0].claimed, 0);
555    }
556
557    #[tokio::test]
558    async fn test_failed_retry_after_refusal_counts_as_failed() {
559        let s = setup();
560        let a = photo(s.dir.path(), "a.jpg");
561        let b = photo(s.dir.path(), "b.jpg");
562        cached(&s.store, &a, "key-Other", "GONE");
563        cached(&s.store, &b, "key-Other", "BBB");
564        let smugmug = FakeSmugMug {
565            missing: HashSet::from(["/api/v2/image/GONE-0".to_string()]),
566            fail_on_retry: true,
567            ..Default::default()
568        };
569        let series = series(&[("Trip", 0)], MAX_ALBUM_IMAGES).await;
570
571        let outcome = plan_collects(
572            vec![a.clone(), b.clone()],
573            &empty_target(),
574            &s.store,
575            false,
576            &smugmug,
577            &series,
578        )
579        .await
580        .unwrap();
581
582        // a was refused: uploaded, and its stale entry forgotten. b's retry
583        // failed: kept for the next run.
584        assert_eq!(outcome.collected, 0);
585        assert_eq!(outcome.to_upload, vec![a.clone()]);
586        assert_eq!(outcome.failed, 1);
587        assert!(s.store.get(&outcome.hashes[&a]).unwrap().is_none());
588        assert!(s.store.get(&outcome.hashes[&b]).unwrap().is_some());
589        assert_eq!(series.albums().await[0].claimed, 0);
590    }
591
592    #[tokio::test]
593    async fn test_collects_roll_over_into_the_next_album() {
594        let s = setup();
595        let files: Vec<PathBuf> = (0..3)
596            .map(|i| {
597                let p = photo(s.dir.path(), &format!("{i}.jpg"));
598                cached(&s.store, &p, "key-Other", &format!("K{i}"));
599                p
600            })
601            .collect();
602        let smugmug = FakeSmugMug::default();
603        let series = series(&[("Trip", 4)], 5).await;
604
605        let outcome = plan_collects(files, &empty_target(), &s.store, false, &smugmug, &series)
606            .await
607            .unwrap();
608
609        assert_eq!(outcome.collected, 3);
610        let collects = smugmug.collects.lock().unwrap();
611        let per_album: Vec<(String, usize)> =
612            collects.iter().map(|(k, u)| (k.clone(), u.len())).collect();
613        assert_eq!(
614            per_album,
615            vec![("key-Trip".to_string(), 1), ("key-Trip (2)".to_string(), 2)]
616        );
617    }
618
619    #[tokio::test]
620    async fn test_dry_run_collects_nothing() {
621        let s = setup();
622        let a = photo(s.dir.path(), "a.jpg");
623        cached(&s.store, &a, "key-Other", "AAA");
624        let smugmug = FakeSmugMug::default();
625        let series = series(&[("Trip", 0)], MAX_ALBUM_IMAGES).await;
626
627        let outcome = plan_collects(vec![a], &empty_target(), &s.store, true, &smugmug, &series)
628            .await
629            .unwrap();
630
631        assert_eq!(outcome.collected, 1);
632        assert!(outcome.to_upload.is_empty());
633        assert!(smugmug.collects.lock().unwrap().is_empty());
634    }
635}