1use 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
30pub(crate) const COLLECT_BATCH: usize = 100;
32
33#[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
47pub struct TargetContents<'a> {
49 pub album_keys: HashSet<String>,
51 pub file_names: HashSet<&'a str>,
53}
54
55#[derive(Debug, Default)]
56pub struct CollectOutcome {
57 pub to_upload: Vec<PathBuf>,
59 pub hashes: HashMap<PathBuf, String>,
61 pub collected: usize,
63 pub failed: usize,
67}
68
69pub(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
78pub 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 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
120fn classify(
123 path: &PathBuf,
124 target: &TargetContents<'_>,
125 store: &HashStore,
126) -> (Option<String>, Option<Pending>) {
127 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 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 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#[derive(Default)]
155struct Collected {
156 collected: usize,
157 refused: Vec<PathBuf>,
160 failed: usize,
162}
163
164async 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 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
231pub(crate) struct BatchOutcome {
233 pub(crate) collected: usize,
234 pub(crate) refused: Vec<Pending>,
235 pub(crate) failed: Vec<Pending>,
236}
237
238pub(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 for p in &refused {
264 let _ = store.remove(&p.hash);
265 }
266
267 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 #[derive(Default)]
311 struct FakeSmugMug {
312 missing: HashSet<String>,
314 fail_collect: bool,
315 fail_on_retry: bool,
317 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 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 let entry = s.store.get(&outcome.hashes[&a]).unwrap().unwrap();
454 assert_eq!(entry.album_key, "key-Trip");
455 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 assert_eq!(series.albums().await[0].claimed, 1);
520 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 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 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}