1use anyhow::{Context, Result};
17use chrono::{NaiveDate, Utc};
18use colored::Colorize;
19use futures_util::future::join_all;
20use indicatif::{ProgressBar, ProgressStyle};
21use md5::Context as Md5Context;
22use serde::Serialize;
23use sha2::{Digest, Sha256};
24use std::collections::{BTreeMap, HashMap, VecDeque};
25use std::io::{IsTerminal, Read};
26use std::path::{Path, PathBuf};
27use std::sync::Arc;
28use std::sync::atomic::{AtomicBool, Ordering};
29use std::time::{Duration, Instant};
30
31use crate::api::SmugMugClient;
32use crate::api::upload::{
33 FileUpload, UploadPayload, UploadRejected, UploadResult, UploadTarget, upload_bytes,
34 upload_file,
35};
36use crate::cache::hash_store::{HashStore, UploadedFile};
37use crate::config::{LivePhotoVideos, RawHandling};
38use crate::dated::date::{self, DateSource};
39use crate::dated::index::{FileIndex, FileRecord};
40use crate::dated::plan::{DayAlbums, NameCheck, Smug};
41use crate::dated::scan::{self, ScannedFile};
42use crate::uploader::collect::{COLLECT_BATCH, CollectBackend, Pending, collect_batch};
43
44pub type StopFlag = Arc<AtomicBool>;
47
48pub struct RunOptions {
49 pub sources: Vec<PathBuf>,
50 pub excludes: Vec<String>,
51 pub root_folder: String,
53 pub root_privacy: Option<&'static str>,
56 pub upload_threads: usize,
57 pub read_threads: usize,
58 pub dry_run: bool,
59 pub no_cache: bool,
62 pub raw_handling: RawHandling,
63 pub live_photo_videos: LivePhotoVideos,
64 pub retry_attempts: u32,
65 pub cache_path: PathBuf,
66}
67
68#[derive(Debug, Default, Clone, Serialize)]
70pub struct RunStats {
71 pub scanned: usize,
73 pub unsupported: usize,
74 pub scan_errors: Vec<String>,
75 pub raw_rendered: usize,
76 pub raw_skipped: usize,
77 pub raw_skipped_with_sibling: usize,
78 pub live_photo_videos_skipped: usize,
79 pub unchanged: usize,
81 pub previously_refused: usize,
83 pub to_process: usize,
85 pub bytes_to_process: u64,
86 pub touched: usize,
88 pub uploaded: usize,
89 pub replaced: usize,
90 pub linked: usize,
92 pub collected: usize,
94 pub failed: usize,
95 pub refused: usize,
98 pub bytes_uploaded: u64,
99 pub albums_created: usize,
100 pub date_sources: BTreeMap<String, usize>,
102 pub days: BTreeMap<NaiveDate, usize>,
104 pub errors: Vec<String>,
106 pub duration_secs: u64,
107 pub interrupted: bool,
109}
110
111const MAX_ERRORS: usize = 20;
112
113impl RunStats {
114 fn error(&mut self, path: &Path, e: impl std::fmt::Display) {
115 if self.errors.len() < MAX_ERRORS {
116 self.errors.push(format!("{}: {}", path.display(), e));
117 }
118 }
119
120 fn dated(&mut self, date: NaiveDate, source: DateSource) {
121 *self.days.entry(date).or_default() += 1;
122 *self
123 .date_sources
124 .entry(source.describe().to_string())
125 .or_default() += 1;
126 }
127}
128
129pub enum Content {
131 File,
133 Rendered(bytes::Bytes),
135}
136
137#[allow(async_fn_in_trait)]
141pub trait Uploads {
142 async fn upload(
143 &self,
144 target: UploadTarget<'_>,
145 path: &Path,
146 content: &Content,
147 size: u64,
148 md5: &str,
149 name: &str,
150 ) -> Result<UploadResult>;
151
152 async fn root_folder(
155 &self,
156 path: &str,
157 create: bool,
158 privacy: Option<&str>,
159 ) -> Result<Option<String>>;
160}
161
162impl Uploads for SmugMugClient {
163 async fn upload(
164 &self,
165 target: UploadTarget<'_>,
166 path: &Path,
167 content: &Content,
168 size: u64,
169 md5: &str,
170 name: &str,
171 ) -> Result<UploadResult> {
172 match content {
173 Content::File => {
174 upload_file(
175 self,
176 target,
177 &FileUpload {
178 path,
179 size,
180 md5_hex: md5,
181 file_name: name,
182 },
183 )
184 .await
185 }
186 Content::Rendered(data) => {
187 upload_bytes(
188 self,
189 target,
190 &UploadPayload {
191 data: data.clone(),
192 file_name: name.to_string(),
193 mime_type: "image/jpeg".to_string(),
194 },
195 )
196 .await
197 }
198 }
199 }
200
201 async fn root_folder(
202 &self,
203 path: &str,
204 create: bool,
205 privacy: Option<&str>,
206 ) -> Result<Option<String>> {
207 if create {
208 self.find_or_create_sorted_folder_path(path, privacy)
209 .await
210 .map(Some)
211 } else {
212 self.find_folder_path(path).await
213 }
214 }
215}
216
217struct PendingFile {
219 file: ScannedFile,
220 previous: Option<FileRecord>,
221}
222
223enum Kind {
224 Same,
226 SameRender { md5: String },
229 Known(UploadedFile),
231 Upload {
232 content: Content,
233 name: String,
234 md5: String,
235 size: u64,
236 },
237}
238
239struct Prepared {
241 file: ScannedFile,
242 previous: Option<FileRecord>,
243 sha256: String,
244 date: NaiveDate,
245 date_source: DateSource,
246 kind: Kind,
247}
248
249impl Prepared {
250 fn record(
251 &self,
252 uploaded_md5: Option<String>,
253 image_uri: Option<String>,
254 album_key: Option<String>,
255 failed: Option<String>,
256 ) -> FileRecord {
257 FileRecord {
258 size: self.file.size,
259 mtime_ns: self.file.mtime_ns,
260 sha256: self.sha256.clone(),
261 uploaded_md5,
262 capture_date: self.date,
263 date_source: self.date_source,
264 image_uri,
265 album_key,
266 failed,
267 }
268 }
269
270 fn uploaded_file(&self, image_uri: &str, album_key: &str, size: u64) -> UploadedFile {
271 UploadedFile {
272 smugmug_uri: image_uri.to_string(),
273 album_key: album_key.to_string(),
274 image_key: image_uri.rsplit('/').next().unwrap_or("").to_string(),
275 uploaded_at: Utc::now(),
276 file_size: size,
277 original_path: self.file.path.to_string_lossy().to_string(),
278 }
279 }
280}
281
282fn hash_file(path: &Path) -> std::io::Result<(String, String, u64)> {
284 let mut file = std::fs::File::open(path)?;
285 let mut sha = Sha256::new();
286 let mut md5 = Md5Context::new();
287 let mut buffer = vec![0u8; 1024 * 1024];
288 let mut len = 0u64;
289 loop {
290 let n = file.read(&mut buffer)?;
291 if n == 0 {
292 break;
293 }
294 sha.update(&buffer[..n]);
295 md5.consume(&buffer[..n]);
296 len += n as u64;
297 }
298 Ok((
299 hex::encode(sha.finalize()),
300 format!("{:x}", md5.finalize()),
301 len,
302 ))
303}
304
305fn date_for(file: &ScannedFile, data: Option<&[u8]>, index: &FileIndex) -> (NaiveDate, DateSource) {
307 if let Some(photo) = &file.live_photo_of {
308 if let Ok(Some(record)) = index.get(photo) {
309 return (record.capture_date, DateSource::LivePhoto);
310 }
311 if let Some(date) = date::metadata_date(photo) {
312 return (date, DateSource::LivePhoto);
313 }
314 }
315 date::capture_date(&file.path, data, file.mtime_ns)
316}
317
318#[derive(Clone)]
320struct ReadContext {
321 store: HashStore,
322 index: FileIndex,
323 render_raw: bool,
324 no_cache: bool,
325}
326
327fn prepare(pending: PendingFile, ctx: &ReadContext) -> Result<Prepared> {
329 let PendingFile { file, previous } = pending;
330 let same = |sha: &str| {
331 previous
332 .as_ref()
333 .is_some_and(|r| r.sha256 == sha && r.is_settled())
334 };
335 let known = |sha: &str| -> Result<Option<UploadedFile>> {
336 if ctx.no_cache {
337 Ok(None)
338 } else {
339 ctx.store.get(sha)
340 }
341 };
342
343 if ctx.render_raw && crate::scanner::is_raw_file(&file.path) {
344 let data = std::fs::read(&file.path).context("Failed to read file")?;
345 let sha256 = hex::encode(Sha256::digest(&data));
346 let (date, date_source) = date_for(&file, Some(&data), &ctx.index);
347 let kind = if same(&sha256) {
348 Kind::Same
349 } else if let Some(cached) = known(&sha256)? {
350 Kind::Known(cached)
351 } else {
352 let rendered = crate::raw::render_jpeg_bytes(&data)
353 .context("Failed to render RAW file to JPEG")?
354 .data;
355 let md5 = format!("{:x}", md5::compute(&rendered));
356 let same_render = previous.as_ref().is_some_and(|r| {
357 r.image_uri.is_some() && r.uploaded_md5.as_deref() == Some(md5.as_str())
358 });
359 if same_render {
360 Kind::SameRender { md5 }
361 } else {
362 Kind::Upload {
363 size: rendered.len() as u64,
364 content: Content::Rendered(rendered.into()),
365 name: crate::raw::rendered_file_name(&file.path),
366 md5,
367 }
368 }
369 };
370 return Ok(Prepared {
371 file,
372 previous,
373 sha256,
374 date,
375 date_source,
376 kind,
377 });
378 }
379
380 let (sha256, md5, size) = hash_file(&file.path).context("Failed to read file")?;
381 let (date, date_source) = date_for(&file, None, &ctx.index);
382 let kind = if same(&sha256) {
383 Kind::Same
384 } else if let Some(cached) = known(&sha256)? {
385 Kind::Known(cached)
386 } else {
387 Kind::Upload {
388 content: Content::File,
389 name: file
390 .path
391 .file_name()
392 .map(|n| n.to_string_lossy().into_owned())
393 .unwrap_or_else(|| "image".to_string()),
394 md5,
395 size,
396 }
397 };
398 Ok(Prepared {
399 file,
400 previous,
401 sha256,
402 date,
403 date_source,
404 kind,
405 })
406}
407
408struct Workers<'a, T> {
410 smug: &'a T,
411 days: DayAlbums<T>,
412 store: HashStore,
413 index: FileIndex,
414 options: &'a RunOptions,
415 stats: std::sync::Mutex<RunStats>,
416 in_flight: std::sync::Mutex<HashMap<String, Vec<Prepared>>>,
420 collects: tokio::sync::Mutex<HashMap<NaiveDate, Vec<(Prepared, UploadedFile)>>>,
422 progress: ProgressBar,
423}
424
425impl<T: Smug + Uploads + CollectBackend> Workers<'_, T> {
426 fn stats(&self) -> std::sync::MutexGuard<'_, RunStats> {
427 self.stats.lock().unwrap()
428 }
429
430 fn put(&self, prepared: &Prepared, record: &FileRecord) {
431 if let Err(e) = self.index.put(&prepared.file.path, record) {
432 self.stats().error(&prepared.file.path, e);
433 }
434 }
435
436 async fn handle(&self, prepared: Prepared) {
437 self.stats().dated(prepared.date, prepared.date_source);
438 match &prepared.kind {
439 Kind::Same => {
440 let mut record = prepared.previous.clone().expect("Same needs a record");
441 record.size = prepared.file.size;
442 record.mtime_ns = prepared.file.mtime_ns;
443 self.put(&prepared, &record);
444 self.stats().touched += 1;
445 }
446 Kind::SameRender { md5 } => {
447 let previous = prepared
448 .previous
449 .as_ref()
450 .expect("SameRender needs a record");
451 let uri = previous.image_uri.clone().unwrap_or_default();
452 let album_key = previous.album_key.clone().unwrap_or_default();
453 let size = prepared.file.size;
454 let _ = self.store.insert_deferred(
455 &prepared.sha256,
456 &prepared.uploaded_file(&uri, &album_key, size),
457 );
458 let record = prepared.record(
459 Some(md5.clone()),
460 previous.image_uri.clone(),
461 previous.album_key.clone(),
462 None,
463 );
464 self.put(&prepared, &record);
465 self.stats().touched += 1;
466 }
467 Kind::Known(cached) => {
468 let cached = cached.clone();
469 self.place_known(prepared, cached).await;
470 }
471 Kind::Upload { .. } => self.upload(prepared).await,
472 }
473 self.progress.inc(1);
474 }
475
476 async fn place_known(&self, prepared: Prepared, cached: UploadedFile) {
479 let day = match self.days.day(prepared.date).await {
480 Ok(day) => day,
481 Err(e) => return self.fail_with(&prepared, &e),
482 };
483 if day.album_keys().await.contains(&cached.album_key) {
484 let record = prepared.record(
485 None,
486 Some(cached.smugmug_uri.clone()),
487 Some(cached.album_key.clone()),
488 None,
489 );
490 self.put(&prepared, &record);
491 self.stats().linked += 1;
492 return;
493 }
494 let batch = {
495 let mut collects = self.collects.lock().await;
496 let pending = collects.entry(prepared.date).or_default();
497 pending.push((prepared, cached));
498 if pending.len() >= COLLECT_BATCH {
499 std::mem::take(pending)
500 } else {
501 Vec::new()
502 }
503 };
504 if !batch.is_empty() {
505 self.collect(batch).await;
506 }
507 }
508
509 async fn collect(&self, items: Vec<(Prepared, UploadedFile)>) {
511 let Some(date) = items.first().map(|(p, _)| p.date) else {
512 return;
513 };
514 let day = match self.days.day(date).await {
515 Ok(day) => day,
516 Err(e) => {
517 for (p, _) in &items {
518 self.fail_with(p, &e);
519 }
520 return;
521 }
522 };
523
524 let mut groups: Vec<(crate::api::albums::Album, Vec<(Prepared, UploadedFile)>)> =
526 Vec::new();
527 for (prepared, cached) in items {
528 let album = match day.series.claim().await {
529 Ok(album) => album,
530 Err(e) => {
531 self.fail_with(&prepared, &e);
532 continue;
533 }
534 };
535 match groups.iter_mut().find(|(a, _)| a.name == album.name) {
536 Some((_, members)) => members.push((prepared, cached)),
537 None => groups.push((album, vec![(prepared, cached)])),
538 }
539 }
540
541 for (album, members) in groups {
542 let pending: Vec<Pending> = members
543 .iter()
544 .map(|(p, cached)| Pending {
545 path: p.file.path.clone(),
546 hash: p.sha256.clone(),
547 file_size: p.file.size,
548 image_uri: cached.smugmug_uri.clone(),
549 image_key: cached.image_key.clone(),
550 })
551 .collect();
552 let outcome = collect_batch(pending, &album, &self.store, self.smug).await;
553 let not_done: HashMap<PathBuf, &'static str> = outcome
554 .refused
555 .iter()
556 .map(|p| {
557 (
558 p.path.clone(),
559 "no longer on SmugMug; it'll be uploaded on the next run",
560 )
561 })
562 .chain(
563 outcome
564 .failed
565 .iter()
566 .map(|p| (p.path.clone(), "couldn't be added to its album; will retry")),
567 )
568 .collect();
569 for (prepared, cached) in members {
570 match not_done.get(&prepared.file.path) {
571 Some(reason) => {
572 day.series.release(&album).await;
573 self.fail(&prepared, reason);
574 }
575 None => {
576 let record = prepared.record(
577 None,
578 Some(cached.smugmug_uri.clone()),
579 Some(album.album_key.clone()),
580 None,
581 );
582 self.put(&prepared, &record);
583 self.stats().collected += 1;
584 }
585 }
586 }
587 }
588 }
589
590 async fn flush_collects(&self) {
591 let batches: Vec<_> = self.collects.lock().await.drain().map(|(_, v)| v).collect();
592 for mut batch in batches {
593 while !batch.is_empty() {
594 let rest = batch.split_off(batch.len().min(COLLECT_BATCH));
595 self.collect(batch).await;
596 batch = rest;
597 }
598 }
599 }
600
601 async fn upload(&self, prepared: Prepared) {
602 if !self.options.no_cache
604 && let Ok(Some(cached)) = self.store.get(&prepared.sha256)
605 {
606 return self.place_known(prepared, cached).await;
607 }
608 let sha = prepared.sha256.clone();
610 let prepared = {
611 let mut in_flight = self.in_flight.lock().unwrap();
612 match in_flight.get_mut(&sha) {
613 Some(waiting) => {
614 waiting.push(prepared);
615 return;
616 }
617 None => {
618 in_flight.insert(sha.clone(), Vec::new());
619 prepared
620 }
621 }
622 };
623
624 let uploaded = self.upload_one(&prepared).await;
625
626 let waiting = self
627 .in_flight
628 .lock()
629 .unwrap()
630 .remove(&sha)
631 .unwrap_or_default();
632 for copy in waiting {
634 match &uploaded {
635 Some(cached) => self.place_known(copy, cached.clone()).await,
636 None => self.fail(©, "an identical file failed to upload"),
637 }
638 }
639 }
640
641 async fn upload_one(&self, prepared: &Prepared) -> Option<UploadedFile> {
643 let Kind::Upload {
644 content,
645 name,
646 md5,
647 size,
648 } = &prepared.kind
649 else {
650 return None;
651 };
652
653 if let Some(previous) = &prepared.previous
656 && let Some(old_uri) = &previous.image_uri
657 && self.index.refs(old_uri).unwrap_or(u32::MAX) <= 1
658 {
659 match self
660 .send(
661 UploadTarget::ReplaceImage(old_uri),
662 prepared,
663 content,
664 *size,
665 md5,
666 name,
667 )
668 .await
669 {
670 Ok(result) => {
671 let album_key = previous.album_key.clone().unwrap_or_default();
672 let cached = prepared.uploaded_file(&result.image_uri, &album_key, *size);
673 if self
674 .store
675 .get(&previous.sha256)
676 .ok()
677 .flatten()
678 .is_some_and(|c| &c.smugmug_uri == old_uri)
679 {
680 let _ = self.store.remove_deferred(&previous.sha256);
681 }
682 let _ = self.store.insert_deferred(&prepared.sha256, &cached);
683 let record = prepared.record(
684 Some(md5.clone()),
685 Some(result.image_uri),
686 Some(album_key),
687 None,
688 );
689 self.put(prepared, &record);
690 let mut stats = self.stats();
691 stats.replaced += 1;
692 stats.bytes_uploaded += size;
693 return Some(cached);
694 }
695 Err(e)
697 if e.downcast_ref::<UploadRejected>()
698 .is_some_and(|r| r.status == 404) => {}
699 Err(e) => {
700 self.fail_with(prepared, &e);
701 return None;
702 }
703 }
704 }
705
706 let day = match self.days.day(prepared.date).await {
707 Ok(day) => day,
708 Err(e) => {
709 self.fail_with(prepared, &e);
710 return None;
711 }
712 };
713 let name = match day.check_name(name, md5).await {
714 Ok(NameCheck::Present(image)) => {
715 let cached = prepared.uploaded_file(&image.uri, &image.album_key, *size);
716 let _ = self.store.insert_deferred(&prepared.sha256, &cached);
717 let record = prepared.record(
718 Some(md5.clone()),
719 Some(image.uri),
720 Some(image.album_key),
721 None,
722 );
723 self.put(prepared, &record);
724 self.stats().linked += 1;
725 return Some(cached);
726 }
727 Ok(NameCheck::Upload(name)) => name,
728 Err(e) => {
729 self.fail_with(prepared, &e);
730 return None;
731 }
732 };
733 let album = match day.series.claim().await {
734 Ok(album) => album,
735 Err(e) => {
736 day.release_name(&name).await;
737 self.fail_with(prepared, &e);
738 return None;
739 }
740 };
741 match self
742 .send(
743 UploadTarget::Album(&album.uri),
744 prepared,
745 content,
746 *size,
747 md5,
748 &name,
749 )
750 .await
751 {
752 Ok(result) => {
753 day.name_uploaded(&name, &result.image_uri, &album.album_key)
754 .await;
755 let cached = prepared.uploaded_file(&result.image_uri, &album.album_key, *size);
756 let _ = self.store.insert_deferred(&prepared.sha256, &cached);
757 let record = prepared.record(
758 Some(md5.clone()),
759 Some(result.image_uri),
760 Some(album.album_key.clone()),
761 None,
762 );
763 self.put(prepared, &record);
764 let mut stats = self.stats();
765 stats.uploaded += 1;
766 stats.bytes_uploaded += size;
767 Some(cached)
768 }
769 Err(e) => {
770 day.series.release(&album).await;
771 day.release_name(&name).await;
772 self.fail_with(prepared, &e);
773 None
774 }
775 }
776 }
777
778 async fn send(
780 &self,
781 target: UploadTarget<'_>,
782 prepared: &Prepared,
783 content: &Content,
784 size: u64,
785 md5: &str,
786 name: &str,
787 ) -> Result<UploadResult> {
788 let attempts = self.options.retry_attempts.max(1);
789 let mut attempt = 1;
790 loop {
791 match self
792 .smug
793 .upload(target, &prepared.file.path, content, size, md5, name)
794 .await
795 {
796 Ok(result) => return Ok(result),
797 Err(e) if attempt < attempts && is_transient(&e, target) => {
798 tokio::time::sleep(Duration::from_secs(2u64.pow(attempt))).await;
799 attempt += 1;
800 }
801 Err(e) => return Err(e),
802 }
803 }
804 }
805
806 fn fail(&self, prepared: &Prepared, reason: &str) {
808 let mut stats = self.stats();
809 stats.failed += 1;
810 stats.error(&prepared.file.path, reason);
811 }
812
813 fn fail_with(&self, prepared: &Prepared, error: &anyhow::Error) {
817 let refused = error
818 .downcast_ref::<UploadRejected>()
819 .is_some_and(UploadRejected::is_permanent);
820 if !refused {
821 return self.fail(prepared, &format!("{:#}", error));
822 }
823 let record = prepared.record(None, None, None, Some(error.to_string()));
824 self.put(prepared, &record);
825 let mut stats = self.stats();
826 stats.refused += 1;
827 stats.error(&prepared.file.path, error);
828 }
829}
830
831fn is_transient(error: &anyhow::Error, target: UploadTarget<'_>) -> bool {
837 let repeatable = matches!(target, UploadTarget::ReplaceImage(_));
838 if let Some(rejected) = error.downcast_ref::<UploadRejected>() {
839 return rejected.status == 429 || (repeatable && rejected.status >= 500);
840 }
841 match error.downcast_ref::<reqwest::Error>() {
842 Some(e) if e.is_connect() => true,
843 Some(e) => repeatable && (e.is_timeout() || e.is_request() || e.is_body()),
844 None => false,
846 }
847}
848
849fn progress_bar(len: u64) -> ProgressBar {
850 let bar = ProgressBar::new(len);
851 bar.set_style(
852 ProgressStyle::default_bar()
853 .template(
854 "{spinner:.green} [{elapsed_precise}] [{bar:40.cyan/blue}] {pos}/{len} ({eta})",
855 )
856 .unwrap()
857 .progress_chars("#>-"),
858 );
859 bar
860}
861
862pub async fn run<T: Smug + Uploads + CollectBackend>(
864 smug: Arc<T>,
865 options: &RunOptions,
866 stop: &StopFlag,
867) -> Result<RunStats> {
868 let started = Instant::now();
869 let store = HashStore::new(&options.cache_path.to_string_lossy())
870 .context("Failed to open the cache")?;
871 let index = FileIndex::open(&store)?;
872 let mut stats = RunStats::default();
873
874 let scanned = {
876 let sources = options.sources.clone();
877 let excludes = options.excludes.clone();
878 let raw = options.raw_handling;
879 let live = options.live_photo_videos;
880 tokio::task::spawn_blocking(move || scan::scan(&sources, &excludes, raw, live)).await??
881 };
882 stats.scanned = scanned.files.len();
883 stats.unsupported = scanned.unsupported;
884 stats.scan_errors = scanned.errors;
885 stats.raw_rendered = scanned.raw_selection.rendered;
886 stats.raw_skipped = scanned.raw_selection.skipped;
887 stats.raw_skipped_with_sibling = scanned.raw_selection.skipped_with_sibling;
888 stats.live_photo_videos_skipped = scanned.live_photo_videos_skipped;
889
890 let (pending, unchanged, previously_refused) = {
892 let index = index.clone();
893 let no_cache = options.no_cache;
894 let files = scanned.files;
895 tokio::task::spawn_blocking(move || -> Result<_> {
896 let mut pending = VecDeque::new();
897 let (mut unchanged, mut refused) = (0, 0);
898 for file in files {
899 let previous = if no_cache {
900 None
901 } else {
902 index.get(&file.path)?
903 };
904 match &previous {
905 Some(r) if r.matches(&file) && r.is_settled() => {
906 unchanged += 1;
907 if r.failed.is_some() {
908 refused += 1;
909 }
910 }
911 _ => pending.push_back(PendingFile { file, previous }),
912 }
913 }
914 Ok((pending, unchanged, refused))
915 })
916 .await??
917 };
918 stats.unchanged = unchanged;
919 stats.previously_refused = previously_refused;
920 stats.to_process = pending.len();
921 stats.bytes_to_process = pending.iter().map(|p| p.file.size).sum();
922
923 if pending.is_empty() {
924 stats.duration_secs = started.elapsed().as_secs();
925 return Ok(stats);
926 }
927
928 if options.dry_run {
929 let dated = dry_run_dates(pending, options.read_threads, &index, stop).await?;
930 for (date, source) in dated {
931 stats.dated(date, source);
932 }
933 stats.interrupted = stop.load(Ordering::Relaxed);
934 stats.duration_secs = started.elapsed().as_secs();
935 return Ok(stats);
936 }
937
938 let root = smug
939 .root_folder(&options.root_folder, true, options.root_privacy)
940 .await
941 .with_context(|| format!("Failed to find or create folder {}", options.root_folder))?;
942 let progress = progress_bar(pending.len() as u64);
943 let workers = Workers {
944 smug: smug.as_ref(),
945 days: DayAlbums::new(smug.clone(), root, false),
946 store: store.clone(),
947 index: index.clone(),
948 options,
949 stats: std::sync::Mutex::new(stats),
950 in_flight: std::sync::Mutex::new(HashMap::new()),
951 collects: tokio::sync::Mutex::new(HashMap::new()),
952 progress: progress.clone(),
953 };
954
955 let queue = Arc::new(std::sync::Mutex::new(pending));
957 let (tx, rx) = tokio::sync::mpsc::channel::<Result<Prepared, (PathBuf, anyhow::Error)>>(
958 options.upload_threads.max(1) * 2,
959 );
960 let read_ctx = ReadContext {
961 store: store.clone(),
962 index: index.clone(),
963 render_raw: options.raw_handling == RawHandling::Render,
964 no_cache: options.no_cache,
965 };
966 let readers: Vec<_> = (0..options.read_threads.max(1))
967 .map(|_| {
968 let queue = queue.clone();
969 let tx = tx.clone();
970 let ctx = read_ctx.clone();
971 let stop = stop.clone();
972 tokio::task::spawn_blocking(move || {
973 loop {
974 if stop.load(Ordering::Relaxed) {
975 break;
976 }
977 let Some(next) = queue.lock().unwrap().pop_front() else {
978 break;
979 };
980 let path = next.file.path.clone();
981 let result = prepare(next, &ctx).map_err(|e| (path, e));
982 if tx.blocking_send(result).is_err() {
983 break;
984 }
985 }
986 })
987 })
988 .collect();
989 drop(tx);
990
991 let rx = tokio::sync::Mutex::new(rx);
993 let uploaders = (0..options.upload_threads.max(1)).map(|_| async {
994 loop {
995 let next = rx.lock().await.recv().await;
996 match next {
997 Some(Ok(prepared)) => workers.handle(prepared).await,
998 Some(Err((path, e))) => {
999 let mut stats = workers.stats();
1000 stats.failed += 1;
1001 stats.error(&path, format!("{:#}", e));
1002 drop(stats);
1003 workers.progress.inc(1);
1004 }
1005 None => break,
1006 }
1007 }
1008 });
1009
1010 let done = tokio::sync::Notify::new();
1011 let ticker = async {
1012 if std::io::stderr().is_terminal() {
1013 return;
1014 }
1015 loop {
1016 tokio::select! {
1017 _ = done.notified() => return,
1018 _ = tokio::time::sleep(Duration::from_secs(60)) => {
1019 let s = workers.stats();
1020 println!(
1021 "… {}/{} processed: {} uploaded ({}), {} replaced, {} already there, {} added from other albums, {} failed",
1022 progress.position(),
1023 progress.length().unwrap_or(0),
1024 s.uploaded,
1025 human_bytes(s.bytes_uploaded),
1026 s.replaced,
1027 s.linked + s.touched,
1028 s.collected,
1029 s.failed + s.refused,
1030 );
1031 }
1032 }
1033 }
1034 };
1035
1036 let work = async {
1037 let (_, readers) = tokio::join!(join_all(uploaders), join_all(readers));
1038 for reader in readers {
1039 if let Err(e) = reader {
1040 eprintln!("{}", format!("A reader thread failed: {}", e).red());
1041 }
1042 }
1043 workers.flush_collects().await;
1044 done.notify_one();
1045 };
1046 tokio::join!(work, ticker);
1047 progress.finish_and_clear();
1048
1049 store.flush()?;
1050 let mut stats = workers.stats.into_inner().unwrap();
1051 for day in workers.days.loaded_days() {
1052 stats.albums_created += day
1053 .series
1054 .albums()
1055 .await
1056 .iter()
1057 .filter(|a| a.created)
1058 .count();
1059 }
1060 stats.interrupted = stop.load(Ordering::Relaxed);
1061 stats.duration_secs = started.elapsed().as_secs();
1062 Ok(stats)
1063}
1064
1065async fn dry_run_dates(
1068 pending: VecDeque<PendingFile>,
1069 threads: usize,
1070 index: &FileIndex,
1071 stop: &StopFlag,
1072) -> Result<Vec<(NaiveDate, DateSource)>> {
1073 let queue = Arc::new(std::sync::Mutex::new(pending));
1074 let handles: Vec<_> = (0..threads.max(1))
1075 .map(|_| {
1076 let queue = queue.clone();
1077 let index = index.clone();
1078 let stop = stop.clone();
1079 tokio::task::spawn_blocking(move || {
1080 let mut out = Vec::new();
1081 while !stop.load(Ordering::Relaxed) {
1082 let Some(next) = queue.lock().unwrap().pop_front() else {
1083 break;
1084 };
1085 out.push(date_for(&next.file, None, &index));
1086 }
1087 out
1088 })
1089 })
1090 .collect();
1091 let mut dated = Vec::new();
1092 for handle in handles {
1093 dated.extend(handle.await?);
1094 }
1095 Ok(dated)
1096}
1097
1098pub fn human_bytes(bytes: u64) -> String {
1099 const UNITS: [&str; 5] = ["B", "KB", "MB", "GB", "TB"];
1100 let mut value = bytes as f64;
1101 let mut unit = 0;
1102 while value >= 1024.0 && unit < UNITS.len() - 1 {
1103 value /= 1024.0;
1104 unit += 1;
1105 }
1106 if unit == 0 {
1107 format!("{} B", bytes)
1108 } else {
1109 format!("{:.1} {}", value, UNITS[unit])
1110 }
1111}
1112
1113#[cfg(test)]
1114mod tests {
1115 use super::*;
1116 use crate::api::images::{AlbumImage, CollectResult};
1117 use crate::dated::plan::tests::FakeSmug;
1118 use std::fs;
1119 use std::time::SystemTime;
1120 use tempfile::TempDir;
1121
1122 impl Uploads for FakeSmug {
1123 async fn upload(
1124 &self,
1125 target: UploadTarget<'_>,
1126 _path: &Path,
1127 _content: &Content,
1128 size: u64,
1129 md5: &str,
1130 name: &str,
1131 ) -> Result<UploadResult> {
1132 if self.refuse.lock().unwrap().contains(name) {
1133 return Err(UploadRejected {
1134 status: 413,
1135 body: "too big".into(),
1136 }
1137 .into());
1138 }
1139 match target {
1140 UploadTarget::Album(album_uri) => {
1141 self.calls
1142 .lock()
1143 .unwrap()
1144 .push(format!("upload {} {}", album_uri, name));
1145 let key = album_uri.rsplit('/').next().unwrap().to_string();
1146 let uri = format!("{}/image/{}-{}", album_uri, self.id(), name);
1147 self.images
1148 .lock()
1149 .unwrap()
1150 .entry(key.clone())
1151 .or_default()
1152 .push(AlbumImage {
1153 image_key: uri.rsplit('/').next().unwrap().into(),
1154 file_name: name.into(),
1155 archived_uri: String::new(),
1156 file_size: size,
1157 format: "JPG".into(),
1158 uri: uri.clone(),
1159 title: None,
1160 archived_md5: Some(md5.into()),
1161 });
1162 Ok(UploadResult {
1163 image_key: uri.rsplit('/').next().unwrap().into(),
1164 image_uri: uri,
1165 status_code: 200,
1166 })
1167 }
1168 UploadTarget::ReplaceImage(uri) => {
1169 self.calls.lock().unwrap().push(format!("replace {}", uri));
1170 Ok(UploadResult {
1171 image_key: uri.rsplit('/').next().unwrap().into(),
1172 image_uri: uri.into(),
1173 status_code: 200,
1174 })
1175 }
1176 }
1177 }
1178
1179 async fn root_folder(
1180 &self,
1181 _path: &str,
1182 _create: bool,
1183 _privacy: Option<&str>,
1184 ) -> Result<Option<String>> {
1185 Ok(Some("/node/root".into()))
1186 }
1187 }
1188
1189 impl CollectBackend for FakeSmug {
1190 async fn collect(&self, album_key: &str, uris: &[String]) -> Result<CollectResult> {
1191 self.calls
1192 .lock()
1193 .unwrap()
1194 .push(format!("collect {} {}", album_key, uris.len()));
1195 Ok(CollectResult::default())
1196 }
1197 }
1198
1199 struct Setup {
1200 photos: TempDir,
1201 cache: TempDir,
1202 smug: Arc<FakeSmug>,
1203 }
1204
1205 impl Setup {
1206 fn new() -> Self {
1207 let (smug, _) = FakeSmug::with_root();
1208 Setup {
1209 photos: TempDir::new().unwrap(),
1210 cache: TempDir::new().unwrap(),
1211 smug,
1212 }
1213 }
1214
1215 fn write(&self, rel: &str, data: &[u8]) -> PathBuf {
1216 let path = self.photos.path().join(rel);
1217 fs::create_dir_all(path.parent().unwrap()).unwrap();
1218 fs::write(&path, data).unwrap();
1219 path
1220 }
1221
1222 fn options(&self) -> RunOptions {
1223 RunOptions {
1224 sources: vec![self.photos.path().to_path_buf()],
1225 excludes: Vec::new(),
1226 root_folder: "Backup".into(),
1227 root_privacy: Some("Private"),
1228 upload_threads: 4,
1229 read_threads: 2,
1230 dry_run: false,
1231 no_cache: false,
1232 raw_handling: RawHandling::Render,
1233 live_photo_videos: LivePhotoVideos::Upload,
1234 retry_attempts: 1,
1235 cache_path: self.cache.path().join("db"),
1236 }
1237 }
1238
1239 async fn run(&self) -> RunStats {
1240 self.run_with(self.options()).await
1241 }
1242
1243 async fn run_with(&self, options: RunOptions) -> RunStats {
1244 let stop: StopFlag = Arc::new(AtomicBool::new(false));
1245 let stats = run(self.smug.clone(), &options, &stop).await.unwrap();
1246 assert!(
1247 stats.errors.is_empty() || stats.failed + stats.refused > 0,
1248 "{:?}",
1249 stats.errors
1250 );
1251 stats
1252 }
1253
1254 fn calls(&self, prefix: &str) -> usize {
1255 self.smug.calls(prefix)
1256 }
1257
1258 fn album_names(&self) -> Vec<String> {
1259 let nodes = self.smug.nodes.lock().unwrap();
1260 let mut names: Vec<String> = nodes
1261 .values()
1262 .flatten()
1263 .filter(|c| c.node_type == "Album")
1264 .map(|c| c.name.clone())
1265 .collect();
1266 names.sort();
1267 names
1268 }
1269 }
1270
1271 fn set_mtime(path: &Path, secs_ago: u64) {
1272 let file = fs::File::options().write(true).open(path).unwrap();
1273 file.set_modified(SystemTime::now() - Duration::from_secs(secs_ago))
1274 .unwrap();
1275 }
1276
1277 #[tokio::test]
1278 async fn uploads_into_day_albums_then_skips_unchanged_files_unread() {
1279 let s = Setup::new();
1280 s.write("a/IMG_20140712_0001.jpg", b"one");
1281 s.write("b/IMG_20140712_0002.jpg", b"two");
1282 s.write("b/IMG_20150101_0003.jpg", b"three");
1283
1284 let stats = s.run().await;
1285 assert_eq!(stats.uploaded, 3);
1286 assert_eq!(stats.albums_created, 2);
1287 assert_eq!(s.album_names(), vec!["2014-07-12", "2015-01-01"]);
1288
1289 let calls_before = s.smug.calls.lock().unwrap().len();
1291 let stats = s.run().await;
1292 assert_eq!(stats.unchanged, 3);
1293 assert_eq!(stats.to_process, 0);
1294 assert_eq!(s.smug.calls.lock().unwrap().len(), calls_before);
1295 }
1296
1297 #[tokio::test]
1298 async fn touched_files_are_rerecorded_not_uploaded() {
1299 let s = Setup::new();
1300 let path = s.write("IMG_20140712_0001.jpg", b"one");
1301 s.run().await;
1302 set_mtime(&path, 3600);
1303
1304 let stats = s.run().await;
1305 assert_eq!(stats.touched, 1);
1306 assert_eq!(s.calls("upload"), 1);
1307
1308 assert_eq!(s.run().await.unchanged, 1);
1310 }
1311
1312 #[tokio::test]
1313 async fn edited_files_replace_their_image() {
1314 let s = Setup::new();
1315 let path = s.write("IMG_20140712_0001.jpg", b"one");
1316 s.run().await;
1317 fs::write(&path, b"one, edited").unwrap();
1318 set_mtime(&path, 10);
1319
1320 let stats = s.run().await;
1321 assert_eq!(stats.replaced, 1);
1322 assert_eq!(s.calls("replace"), 1);
1323 assert_eq!(s.calls("upload"), 1);
1324 }
1325
1326 #[tokio::test]
1327 async fn identical_copies_upload_once() {
1328 let s = Setup::new();
1329 for i in 0..6 {
1330 s.write(&format!("copy{}/IMG_20140712_0001.jpg", i), b"same bytes");
1331 }
1332 let stats = s.run().await;
1333 assert_eq!(stats.uploaded, 1);
1334 assert_eq!(stats.linked, 5);
1335 assert_eq!(s.calls("upload"), 1);
1336 }
1337
1338 #[tokio::test]
1339 async fn editing_one_of_two_copies_uploads_it_anew() {
1340 let s = Setup::new();
1341 let a = s.write("a/IMG_20140712_0001.jpg", b"same");
1342 s.write("b/IMG_20140712_0001.jpg", b"same");
1343 s.run().await;
1344 fs::write(&a, b"different now").unwrap();
1345 set_mtime(&a, 10);
1346
1347 let stats = s.run().await;
1348 assert_eq!(stats.replaced, 0);
1351 assert_eq!(stats.uploaded, 1);
1352 assert_eq!(s.calls("replace"), 0);
1353 assert!(
1354 s.smug
1355 .calls
1356 .lock()
1357 .unwrap()
1358 .iter()
1359 .any(|c| c.contains("IMG_20140712_0001~"))
1360 );
1361 }
1362
1363 #[tokio::test]
1364 async fn same_name_different_content_is_renamed() {
1365 let s = Setup::new();
1366 s.write("cam1/IMG_20140712_0001.jpg", b"camera one");
1367 s.write("cam2/IMG_20140712_0001.jpg", b"camera two");
1368 let stats = s.run().await;
1369 assert_eq!(stats.uploaded, 2);
1370 let names: Vec<String> = s
1371 .smug
1372 .images
1373 .lock()
1374 .unwrap()
1375 .values()
1376 .flatten()
1377 .map(|i| i.file_name.clone())
1378 .collect();
1379 assert!(names.contains(&"IMG_20140712_0001.jpg".to_string()));
1380 assert!(
1381 names
1382 .iter()
1383 .any(|n| n.starts_with("IMG_20140712_0001~") && n.ends_with(".jpg"))
1384 );
1385 }
1386
1387 #[tokio::test]
1388 async fn content_already_in_the_day_album_is_linked_without_cache() {
1389 let s = Setup::new();
1390 s.write("IMG_20140712_0001.jpg", b"one");
1391 s.run().await;
1392 let mut options = s.options();
1394 options.cache_path = s.cache.path().join("fresh");
1395 let stats = s.run_with(options).await;
1396 assert_eq!(stats.linked, 1);
1397 assert_eq!(s.calls("upload"), 1);
1398 }
1399
1400 #[tokio::test]
1401 async fn files_in_other_albums_are_collected() {
1402 let s = Setup::new();
1403 let path = s.write("IMG_20140712_0001.jpg", b"one");
1404 {
1406 let store = HashStore::new(&s.cache.path().join("db").to_string_lossy()).unwrap();
1407 let (sha, _, size) = hash_file(&path).unwrap();
1408 store
1409 .insert(
1410 &sha,
1411 UploadedFile {
1412 smugmug_uri: "/api/v2/album/ELSEWHERE/image/X-0".into(),
1413 album_key: "ELSEWHERE".into(),
1414 image_key: "X-0".into(),
1415 uploaded_at: Utc::now(),
1416 file_size: size,
1417 original_path: "elsewhere".into(),
1418 },
1419 )
1420 .unwrap();
1421 }
1422 let stats = s.run().await;
1423 assert_eq!(stats.collected, 1);
1424 assert_eq!(s.calls("collect"), 1);
1425 assert_eq!(s.calls("upload"), 0);
1426 assert_eq!(s.run().await.unchanged, 1);
1427 }
1428
1429 #[tokio::test]
1430 async fn live_photo_videos_go_to_their_photos_day() {
1431 let s = Setup::new();
1432 s.write("IMG_20140712_0001.HEIC", b"photo");
1433 let video = s.write("IMG_20140712_0001.MOV", b"video");
1434 set_mtime(&video, 0);
1436 let stats = s.run().await;
1437 assert_eq!(stats.uploaded, 2);
1438 assert_eq!(s.album_names(), vec!["2014-07-12"]);
1439 }
1440
1441 fn fake_raw(extra: &[u8]) -> Vec<u8> {
1442 let mut raw = crate::raw::exif::tests::fake_tiff_raw_be();
1443 raw.extend(crate::raw::jpeg::tests::fake_jpeg(0xC0, 6000, 4000, None));
1444 raw.extend_from_slice(extra);
1445 raw
1446 }
1447
1448 #[tokio::test]
1449 async fn raw_files_are_rendered_once() {
1450 let s = Setup::new();
1451 let path = s.write("IMG_0001.CR2", &fake_raw(b""));
1452 let stats = s.run().await;
1453 assert_eq!(stats.uploaded, 1);
1454 assert_eq!(s.album_names(), vec!["2024-06-01"]);
1456 let uploaded = s
1457 .smug
1458 .images
1459 .lock()
1460 .unwrap()
1461 .values()
1462 .flatten()
1463 .next()
1464 .unwrap()
1465 .clone();
1466 assert_eq!(uploaded.file_name, "IMG_0001.jpg");
1467
1468 set_mtime(&path, 3600);
1470 assert_eq!(s.run().await.touched, 1);
1471
1472 fs::write(&path, fake_raw(b"\0\0\0\0")).unwrap();
1475 set_mtime(&path, 10);
1476 let stats = s.run().await;
1477 assert_eq!(stats.touched, 1);
1478 assert_eq!(s.calls("upload") + s.calls("replace"), 1);
1479
1480 let moved = s.write("moved/IMG_0001.CR2", &fs::read(&path).unwrap());
1482 fs::remove_file(&path).unwrap();
1483 let _ = moved;
1484 let stats = s.run().await;
1485 assert_eq!(stats.linked, 1);
1486 assert_eq!(s.calls("upload") + s.calls("replace"), 1);
1487 }
1488
1489 #[tokio::test]
1490 async fn refused_files_are_not_retried_until_they_change() {
1491 let s = Setup::new();
1492 let path = s.write("IMG_20140712_0001.jpg", b"huge");
1493 s.smug
1494 .refuse
1495 .lock()
1496 .unwrap()
1497 .insert("IMG_20140712_0001.jpg".into());
1498 let stats = s.run().await;
1499 assert_eq!(stats.refused, 1);
1500
1501 let stats = s.run().await;
1502 assert_eq!(stats.previously_refused, 1);
1503 assert_eq!(stats.to_process, 0);
1504
1505 s.smug.refuse.lock().unwrap().clear();
1506 fs::write(&path, b"smaller").unwrap();
1507 set_mtime(&path, 10);
1508 assert_eq!(s.run().await.uploaded, 1);
1509 }
1510
1511 #[tokio::test]
1512 async fn dry_run_dates_files_and_touches_nothing() {
1513 let s = Setup::new();
1514 s.write("IMG_20140712_0001.jpg", b"one");
1515 s.write("IMG_20140713_0001.jpg", b"two");
1516 let mut options = s.options();
1517 options.dry_run = true;
1518 let stats = s.run_with(options).await;
1519 assert_eq!(stats.to_process, 2);
1520 assert_eq!(stats.days.len(), 2);
1521 assert_eq!(stats.date_sources.get("file name"), Some(&2));
1522 assert!(s.smug.calls.lock().unwrap().is_empty());
1523 assert_eq!(s.run().await.uploaded, 2);
1525 }
1526
1527 #[tokio::test]
1528 async fn excluded_paths_are_left_out() {
1529 let s = Setup::new();
1530 s.write("keep/IMG_20140712_0001.jpg", b"one");
1531 s.write("old_bak/IMG_20140712_0002.jpg", b"two");
1532 let mut options = s.options();
1533 options.excludes = vec!["old_bak/".into()];
1534 let stats = s.run_with(options).await;
1535 assert_eq!(stats.scanned, 1);
1536 assert_eq!(stats.uploaded, 1);
1537 }
1538
1539 #[test]
1540 fn new_uploads_are_retried_only_when_nothing_can_have_been_stored() {
1541 let rejected = |status| {
1542 anyhow::Error::from(UploadRejected {
1543 status,
1544 body: String::new(),
1545 })
1546 };
1547 let album = UploadTarget::Album("/api/v2/album/A");
1548 let replace = UploadTarget::ReplaceImage("/api/v2/album/A/image/B-0");
1549 assert!(is_transient(&rejected(429), album));
1550 assert!(!is_transient(&rejected(500), album));
1551 assert!(!is_transient(&rejected(503), album));
1552 assert!(is_transient(&rejected(503), replace));
1553 assert!(!is_transient(&rejected(413), replace));
1554 assert!(!is_transient(&anyhow::anyhow!("file vanished"), replace));
1555 }
1556
1557 #[test]
1558 fn human_sizes() {
1559 assert_eq!(human_bytes(512), "512 B");
1560 assert_eq!(human_bytes(1536), "1.5 KB");
1561 assert_eq!(human_bytes(5 * 1024 * 1024 * 1024), "5.0 GB");
1562 }
1563}