pigeon-cli 0.4.0

Pigeon: authenticate, sink, and transform personal data from external services.
Documentation
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
530
531
532
533
534
535
536
537
538
539
540
541
542
543
544
545
546
547
548
549
550
551
552
553
554
555
556
557
558
559
560
561
562
563
564
565
566
567
568
569
570
571
572
573
574
575
576
577
578
579
580
581
582
583
584
585
586
587
588
589
590
591
592
593
594
595
596
597
598
599
600
601
602
603
604
605
606
607
608
609
610
611
612
613
614
615
616
617
618
619
620
621
622
623
624
625
626
627
628
629
630
631
632
633
634
635
636
637
638
639
640
641
642
643
644
645
646
647
648
649
650
651
652
653
654
655
656
657
658
659
660
661
662
663
664
665
666
667
668
//! Concurrent download + zip-expansion pipeline (ADR-0082 §4): every item is
//! either a zip (expanded, never itself hashed/placed -- members requeued)
//! or a plain file (hashed via `download::sha256_file` and staged for the
//! sequential dedup+placement pass, `dedup.rs`). No classification beyond
//! "is this a zip" -- no media recoding, no date extraction, unlike
//! `pull-transform`; every file is always processed and every zip is always
//! expanded (ADR-0082 §1).

use std::collections::{HashMap, HashSet, VecDeque};
use std::fs;
use std::path::{Path, PathBuf};
use std::sync::atomic::{AtomicU64, AtomicUsize, Ordering};
use std::sync::{Arc, Mutex};
use std::time::Duration;

use indicatif::MultiProgress;
use tracing::Instrument;

use crate::commands::job::download;
use crate::commands::job::email_sync::sink;
use crate::commands::job::pull_transform::archive;
use crate::commands::job::upload::{self, UploadedIndex};
use crate::commands::keyring::bucket::store::BucketConfig;
use crate::core::data::ContentIndex;

use super::dedup::{self, DeduplicateDedup, HashedFile};
use super::manifest::{self, DeduplicateTask, extension_of};

/// A fresh, not-yet-existing path under `dir` named by `counter`
/// (monotonically increasing, shared across concurrent workers) plus
/// `extension` -- own small copy of `pull_transform::worker`'s helper
/// (that module is private, unreachable).
fn next_scratch_path(dir: &Path, counter: &AtomicU64, extension: &str) -> Result<PathBuf, String> {
    fs::create_dir_all(dir).map_err(|err| format!("failed to create {}: {err}", dir.display()))?;
    let name = counter.fetch_add(1, Ordering::SeqCst);
    Ok(dir.join(format!("{name:012}.{extension}")))
}

fn is_zip_key(key: &str) -> bool {
    extension_of(key) == "zip"
}

/// One item on the shared work queue -- a top-level bucket object not yet
/// downloaded (`source_key: Some`, `path: None`), or a file already on disk
/// (`path: Some`) -- true both for a completed top-level download and for a
/// zip member streamed straight to disk during a parent's expansion. Only a
/// `depth == 0` item is ever checkpointed. `root_key` is the top-level
/// task's own `display_key`, unchanged through every descendant -- it never
/// participates in path computation, only in tracking whether *any*
/// descendant of a given root failed or lost data, so that root can be
/// excluded from the checkpoint (ADR-0098).
struct QueueItem {
    source_key: Option<String>,
    display_key: String,
    path: Option<PathBuf>,
    depth: u32,
    size: u64,
    root_key: String,
}

#[derive(Debug, Default)]
pub(crate) struct FailureBreakdown {
    pub download: usize,
    pub archive: usize,
    pub hash: usize,
    pub placement: usize,
}

impl FailureBreakdown {
    fn merge(&mut self, other: &FailureBreakdown) {
        self.download += other.download;
        self.archive += other.archive;
        self.hash += other.hash;
        self.placement += other.placement;
    }
}

enum FailureCategory {
    Download,
    Archive,
    Hash,
}

enum ItemOutcome {
    /// `display_key`/`depth` identify the zip that was expanded (checkpoint
    /// candidate iff `depth == 0`); `members` are queued for the next pass.
    /// `dropped` is how many members this zip lost to the per-archive
    /// extraction-ratio cap (ADR-0098) -- nonzero here means real data was
    /// discarded, so the caller both counts it as a failure and taints this
    /// item's root out of the checkpoint.
    ZipExpanded {
        display_key: String,
        depth: u32,
        members: Vec<QueueItem>,
        dropped: usize,
    },
    Hashed {
        depth: u32,
        file: HashedFile,
    },
    Failed {
        category: FailureCategory,
    },
}

/// Downloads (if not already on disk) then either expands `item` (a zip --
/// its own bytes are discarded, never hashed or placed, per ADR-0082 §4) or
/// hashes it for the placement pass. The expansion/hash itself runs via
/// `tokio::task::spawn_blocking` (ADR-0088) -- both are CPU-bound, so they're
/// handed to tokio's blocking-thread pool rather than occupying one of the
/// runtime's own async worker threads for the whole call. Returns the
/// item's `root_key` alongside the outcome: neither `spawn_blocking` nor the
/// caller's own `tokio::spawn` propagate the ambient `tracing` span on
/// their own, so each blocking closure below explicitly re-enters the span
/// captured just before it was spawned (ADR-0098) -- without this, every
/// warning logged here would be missing the `command`/`instance` fields
/// the rest of `pigeon.jsonl` relies on.
async fn process_item(
    bucket_config: &BucketConfig,
    secret: &str,
    item: QueueItem,
    raw_dir: &Path,
    counter: &Arc<AtomicU64>,
    announce: &download::DownloadAnnounce,
) -> (String, ItemOutcome) {
    let root_key = item.root_key.clone();
    let depth = item.depth;
    let extension = extension_of(&item.display_key);
    let is_zip = extension == "zip";

    let path = match item.path {
        Some(path) => path,
        None => {
            let key = item.source_key.as_deref().unwrap_or(&item.display_key);
            if let Err(err) = download::check_disk_space(raw_dir, item.size) {
                tracing::warn!(key = %item.display_key, step = "download", error = %err, "not enough disk space");
                return (
                    root_key,
                    ItemOutcome::Failed {
                        category: FailureCategory::Download,
                    },
                );
            }
            let raw_path = match next_scratch_path(raw_dir, counter, &extension) {
                Ok(path) => path,
                Err(err) => {
                    tracing::warn!(key = %item.display_key, step = "download", error = %err, "failed to allocate a raw path");
                    return (
                        root_key,
                        ItemOutcome::Failed {
                            category: FailureCategory::Download,
                        },
                    );
                }
            };
            if let Err(err) = download::download_with_retry(
                bucket_config,
                secret,
                key,
                item.size,
                &raw_path,
                announce,
            )
            .await
            {
                tracing::warn!(key = %item.display_key, step = "download", error = %err, "download failed");
                let _ = fs::remove_file(&raw_path);
                return (
                    root_key,
                    ItemOutcome::Failed {
                        category: FailureCategory::Download,
                    },
                );
            }
            crate::observability::metrics::record_phase(
                "deduplicate",
                "download",
                "ok",
                Some(bucket_config.alias.as_str()),
            );
            raw_path
        }
    };

    if is_zip {
        if depth >= archive::MAX_ZIP_DEPTH {
            tracing::warn!(key = %item.display_key, step = "archive", depth, "zip nesting depth cap reached, not expanding further");
            let _ = fs::remove_file(&path);
            return (
                root_key,
                ItemOutcome::Failed {
                    category: FailureCategory::Archive,
                },
            );
        }
        if let Err(err) = download::check_disk_space(raw_dir, 0) {
            tracing::warn!(key = %item.display_key, step = "archive", error = %err, "not enough disk space to expand");
            let _ = fs::remove_file(&path);
            return (
                root_key,
                ItemOutcome::Failed {
                    category: FailureCategory::Archive,
                },
            );
        }
        let expand_path = path.clone();
        let expand_raw_dir = raw_dir.to_path_buf();
        let expand_counter = Arc::clone(counter);
        let span = tracing::Span::current();
        let expand_result = tokio::task::spawn_blocking(move || {
            span.in_scope(|| archive::expand_to_dir(&expand_path, &expand_raw_dir, &expand_counter))
        })
        .await;
        let outcome = match expand_result {
            Ok(Ok((raw_members, dropped))) => {
                // The zip container itself is never hashed or placed -- only
                // its extracted members are (ADR-0082 §4).
                let _ = fs::remove_file(&path);
                let members = raw_members
                    .into_iter()
                    .map(|member| QueueItem {
                        source_key: None,
                        display_key: format!("{}!{}", item.display_key, member.name),
                        path: Some(member.path),
                        depth: depth + 1,
                        size: member.size,
                        root_key: root_key.clone(),
                    })
                    .collect();
                ItemOutcome::ZipExpanded {
                    display_key: item.display_key,
                    depth,
                    members,
                    dropped,
                }
            }
            Ok(Err(err)) => {
                tracing::warn!(key = %item.display_key, step = "archive", error = %err, "failed to open zip archive");
                let _ = fs::remove_file(&path);
                ItemOutcome::Failed {
                    category: FailureCategory::Archive,
                }
            }
            Err(err) => {
                tracing::error!(key = %item.display_key, step = "archive", error = %err, "zip expansion task panicked");
                let _ = fs::remove_file(&path);
                ItemOutcome::Failed {
                    category: FailureCategory::Archive,
                }
            }
        };
        return (root_key, outcome);
    }

    let hash_path = path.clone();
    let span = tracing::Span::current();
    let hash_result =
        tokio::task::spawn_blocking(move || span.in_scope(|| download::sha256_file(&hash_path)))
            .await;
    let outcome = match hash_result {
        Ok(Ok(content_hash)) => ItemOutcome::Hashed {
            depth,
            file: HashedFile {
                original_key: item.display_key,
                scratch_path: path,
                extension,
                content_hash,
            },
        },
        Ok(Err(err)) => {
            tracing::warn!(key = %item.display_key, step = "hash", error = %err, "failed to hash file");
            let _ = fs::remove_file(&path);
            ItemOutcome::Failed {
                category: FailureCategory::Hash,
            }
        }
        Err(err) => {
            tracing::error!(key = %item.display_key, step = "hash", error = %err, "hash task panicked");
            let _ = fs::remove_file(&path);
            ItemOutcome::Failed {
                category: FailureCategory::Hash,
            }
        }
    };
    (root_key, outcome)
}

#[derive(Debug, Default)]
pub(crate) struct DeduplicateSummary {
    pub processed: usize,
    pub failed: usize,
    pub failure_breakdown: FailureBreakdown,
    pub duplicates_skipped: usize,
    /// Zip members dropped by the per-archive extraction-ratio cap
    /// (ADR-0098) -- already folded into `failed`/`failure_breakdown.archive`
    /// too, since dropped data is a real failure, but broken out here so
    /// the wizard can print it as its own distinct, named count.
    pub dropped_members: usize,
    pub uploaded: usize,
    pub unchanged: usize,
    pub upload_failed: usize,
}

/// Runs the full deduplicate pipeline: `tasks` (from `Job::gather`, already
/// filtered against the `.processed` checkpoint) are downloaded/expanded/
/// hashed concurrently at `concurrency`, placed sequentially (dedup +
/// report), then uploaded (if `remote` is given) -- always unencrypted.
pub(crate) async fn run_deduplicate_job(
    bucket_config: &BucketConfig,
    secret: &str,
    local_output: &Path,
    tasks: Vec<DeduplicateTask>,
    concurrency: usize,
    upload_concurrency: usize,
    remote: Option<(&BucketConfig, &str)>,
) -> Result<DeduplicateSummary, String> {
    crate::observability::metrics::set_macro_phase("deduplicate", false);
    let staging_dir = local_output.join(".staging");
    let raw_dir = staging_dir.join("raw");
    let result_dir = local_output.join("result");
    fs::create_dir_all(&raw_dir)
        .map_err(|err| format!("failed to create {}: {err}", raw_dir.display()))?;
    let counter = Arc::new(AtomicU64::new(0));

    let multi_progress = MultiProgress::new();
    let total = tasks.len() as u64;
    let _ = multi_progress.println(format!("Downloading and processing {total} object(s)..."));
    tracing::info!(total, "deduplicate: download/expand/hash phase starting");
    let bar = sink::new_progress_bar("deduplicate".to_string(), total, &multi_progress);
    let announce = download::DownloadAnnounce::new(bar.clone());

    let queue: Arc<Mutex<VecDeque<QueueItem>>> = Arc::new(Mutex::new(
        tasks
            .into_iter()
            .map(|task| QueueItem {
                source_key: Some(task.key.clone()),
                display_key: task.key.clone(),
                path: None,
                depth: 0,
                size: task.size,
                root_key: task.key,
            })
            .collect(),
    ));
    let in_flight = Arc::new(AtomicUsize::new(0));

    let hashed_files = Arc::new(Mutex::new(Vec::<HashedFile>::new()));
    let finished_root_keys = Arc::new(Mutex::new(Vec::<String>::new()));
    let failure_breakdown = Arc::new(Mutex::new(FailureBreakdown::default()));
    let dropped_members = Arc::new(AtomicUsize::new(0));
    // Any root whose descendant failed outright, or lost a member to the
    // extraction-ratio cap, is excluded from the checkpoint below -- a root
    // zip was previously checkpointed unconditionally the moment it
    // expanded, regardless of what happened to its members, so a rerun
    // could never retry silently-dropped data (ADR-0098).
    let tainted_roots = Arc::new(Mutex::new(HashSet::<String>::new()));
    let command_span = tracing::Span::current();

    let worker_count = concurrency.max(1);
    let mut handles = Vec::with_capacity(worker_count);
    for _ in 0..worker_count {
        let queue = Arc::clone(&queue);
        let in_flight = Arc::clone(&in_flight);
        let hashed_files = Arc::clone(&hashed_files);
        let finished_root_keys = Arc::clone(&finished_root_keys);
        let failure_breakdown = Arc::clone(&failure_breakdown);
        let dropped_members = Arc::clone(&dropped_members);
        let tainted_roots = Arc::clone(&tainted_roots);
        let counter = Arc::clone(&counter);
        let bucket_config = bucket_config.clone();
        let secret = secret.to_string();
        let raw_dir = raw_dir.clone();
        let announce = announce.clone();
        let bar = bar.clone();

        handles.push(tokio::spawn(
            async move {
                loop {
                    let item = { queue.lock().unwrap().pop_front() };
                    let Some(item) = item else {
                        if in_flight.load(Ordering::SeqCst) == 0 {
                            break;
                        }
                        tokio::time::sleep(Duration::from_millis(10)).await;
                        continue;
                    };
                    in_flight.fetch_add(1, Ordering::SeqCst);

                    let (root_key, outcome) =
                        process_item(&bucket_config, &secret, item, &raw_dir, &counter, &announce)
                            .await;

                    match outcome {
                        ItemOutcome::ZipExpanded {
                            display_key,
                            depth,
                            members,
                            dropped,
                        } => {
                            crate::observability::metrics::record_phase(
                                "deduplicate",
                                "archive",
                                "ok",
                                Some(bucket_config.alias.as_str()),
                            );
                            bar.inc_length(members.len() as u64);
                            queue.lock().unwrap().extend(members);
                            if depth == 0 {
                                finished_root_keys.lock().unwrap().push(display_key);
                            }
                            if dropped > 0 {
                                failure_breakdown.lock().unwrap().archive += dropped;
                                dropped_members.fetch_add(dropped, Ordering::SeqCst);
                                tainted_roots.lock().unwrap().insert(root_key);
                                crate::observability::metrics::record_phase(
                                    "deduplicate",
                                    "archive",
                                    "failed",
                                    Some(bucket_config.alias.as_str()),
                                );
                            }
                        }
                        ItemOutcome::Hashed { depth, file } => {
                            crate::observability::metrics::record_phase(
                                "deduplicate",
                                "hash",
                                "ok",
                                Some(bucket_config.alias.as_str()),
                            );
                            if depth == 0 {
                                finished_root_keys
                                    .lock()
                                    .unwrap()
                                    .push(file.original_key.clone());
                            }
                            hashed_files.lock().unwrap().push(file);
                        }
                        ItemOutcome::Failed { category } => {
                            tainted_roots.lock().unwrap().insert(root_key);
                            let mut breakdown = failure_breakdown.lock().unwrap();
                            let phase = match category {
                                FailureCategory::Download => {
                                    breakdown.download += 1;
                                    "download"
                                }
                                FailureCategory::Archive => {
                                    breakdown.archive += 1;
                                    "archive"
                                }
                                FailureCategory::Hash => {
                                    breakdown.hash += 1;
                                    "hash"
                                }
                            };
                            crate::observability::metrics::record_phase(
                                "deduplicate",
                                phase,
                                "failed",
                                Some(bucket_config.alias.as_str()),
                            );
                        }
                    }

                    bar.inc(1);
                    in_flight.fetch_sub(1, Ordering::SeqCst);
                }
            }
            .instrument(command_span.clone()),
        ));
    }

    let mut first_panic = None;
    for handle in handles {
        if let Err(err) = handle.await {
            tracing::error!(error = %err, "deduplicate worker task panicked");
            if first_panic.is_none() {
                first_panic = Some(format!("worker task panicked: {err}"));
            }
        }
    }
    if let Some(err) = first_panic {
        return Err(err);
    }
    bar.finish();

    let files = Arc::try_unwrap(hashed_files)
        .map_err(|_| "internal error: hashed file list still shared".to_string())?
        .into_inner()
        .map_err(|_| "internal error: hashed file list lock poisoned".to_string())?;
    let mut finished_root_keys = Arc::try_unwrap(finished_root_keys)
        .map_err(|_| "internal error: finished-key list still shared".to_string())?
        .into_inner()
        .map_err(|_| "internal error: finished-key list lock poisoned".to_string())?;
    let failure_breakdown = Arc::try_unwrap(failure_breakdown)
        .map_err(|_| "internal error: failure breakdown still shared".to_string())?
        .into_inner()
        .map_err(|_| "internal error: failure breakdown lock poisoned".to_string())?;
    let dropped_members = Arc::try_unwrap(dropped_members)
        .map_err(|_| "internal error: dropped-member counter still shared".to_string())?
        .into_inner();
    let tainted_roots = Arc::try_unwrap(tainted_roots)
        .map_err(|_| "internal error: tainted-root set still shared".to_string())?
        .into_inner()
        .map_err(|_| "internal error: tainted-root set lock poisoned".to_string())?;

    tracing::info!(
        hashed = files.len(),
        download_failed = failure_breakdown.download,
        archive_failed = failure_breakdown.archive,
        hash_failed = failure_breakdown.hash,
        dropped_members,
        "deduplicate: download/expand/hash phase complete"
    );

    // `dedup_index` (the full hash->path map), `merge_records` (the
    // human-readable report -- can run well into the GB range for a large,
    // duplicate-heavy bucket), and `placed_keys` all live only inside this
    // block, so they're dropped here, before the upload phase runs, instead
    // of surviving in `run_deduplicate_job`'s own scope through the whole upload
    // phase afterward (ADR-0089 -- this is what let a real run's RSS keep
    // climbing well past the fetch+hash phase and eventually get OOM-killed
    // partway through upload).
    let placement_summary = {
        let mut dedup_index = DeduplicateDedup(ContentIndex::load(
            &staging_dir,
            dedup::CONTENT_HASHES_FILE,
        )?);
        let (placement_summary, merge_records, placed_keys) =
            dedup::place_and_report(&result_dir, files, &mut dedup_index, &multi_progress);
        dedup::write_report(local_output, &merge_records)?;

        let placed_keys: std::collections::HashSet<String> = placed_keys.into_iter().collect();
        finished_root_keys.retain(|key| {
            !tainted_roots.contains(key) && (placed_keys.contains(key) || is_zip_key(key))
        });
        for key in &finished_root_keys {
            manifest::append_checkpoint(&staging_dir, key)?;
        }
        placement_summary
    };

    tracing::info!(
        placed = placement_summary.placed,
        duplicates_skipped = placement_summary.duplicates_skipped,
        placement_failed = placement_summary.failed,
        "deduplicate: placement phase complete"
    );

    let mut summary = DeduplicateSummary {
        processed: placement_summary.placed,
        failed: failure_breakdown.download
            + failure_breakdown.archive
            + failure_breakdown.hash
            + placement_summary.failed,
        duplicates_skipped: placement_summary.duplicates_skipped,
        dropped_members,
        ..Default::default()
    };
    summary.failure_breakdown.merge(&failure_breakdown);
    // `place_and_report`'s placement failures were previously logged but
    // never counted anywhere at all (ADR-0093 fixes this as a side effect
    // of adding the live per-phase metric at that same call site) --
    // `pull_transform`'s equivalent placement pass already folds its own
    // `PlacementSummary.failed` into `failure_breakdown.placement` the same
    // way.
    summary.failure_breakdown.placement = placement_summary.failed;

    if let Some((remote_bucket, remote_secret)) = remote {
        tracing::info!(bucket = %remote_bucket.alias, "deduplicate: upload phase starting");
        let upload_summary = upload_result(
            &remote_bucket.alias,
            local_output,
            (remote_bucket, remote_secret),
            upload_concurrency,
            &multi_progress,
        )
        .await?;
        summary.uploaded = upload_summary.uploaded;
        summary.unchanged = upload_summary.unchanged;
        summary.upload_failed = upload_summary.upload_failed;
    }

    Ok(summary)
}

/// Uploads `local_output/result/` to `remote`, resuming via the existing
/// `.staging/.uploaded` index (ADR-0019/ADR-0024) -- the shared upload tail
/// both `run_deduplicate_job` and `run_upload_only` call, so there's one code
/// path and one resume mechanism between a fresh run and a resumed
/// `--upload-only` one (ADR-0089). `label` is purely descriptive (tracing/
/// log context): both callers pass `remote`'s own alias, since `remote` is
/// always the actual upload destination (ADR-0098 -- `run_deduplicate_job`
/// previously passed the source bucket's alias here, mislabeling every
/// upload-phase log line with the wrong bucket).
async fn upload_result(
    label: &str,
    local_output: &Path,
    remote: (&BucketConfig, &str),
    upload_concurrency: usize,
    multi_progress: &MultiProgress,
) -> Result<upload::UploadSummary, String> {
    let staging_dir = local_output.join(".staging");
    let result_dir = local_output.join("result");
    let (remote_bucket, remote_secret) = remote;

    let (upload_tasks, uploaded_index) = upload::pending_upload_tasks(
        "deduplicate",
        label,
        &staging_dir,
        &result_dir,
        &result_dir,
        false,
    )?;
    let mut uploaded_indexes: HashMap<PathBuf, Arc<Mutex<UploadedIndex>>> = HashMap::new();
    uploaded_indexes.insert(staging_dir.clone(), Arc::new(Mutex::new(uploaded_index)));
    Ok(upload::run_upload_phase(
        upload_tasks,
        &uploaded_indexes,
        remote_bucket,
        remote_secret,
        None,
        upload_concurrency,
        multi_progress,
    )
    .await)
}

/// Resumes uploading an already-completed local deduplicate run, skipping the
/// bucket listing/download/hash/placement phases entirely (ADR-0089's
/// `--upload-only`). Reuses `upload_result`, the same helper
/// `run_deduplicate_job`'s own upload tail calls.
pub(crate) async fn run_upload_only(
    local_output: &Path,
    remote: (&BucketConfig, &str),
    upload_concurrency: usize,
) -> Result<upload::UploadSummary, String> {
    let multi_progress = MultiProgress::new();
    let label = remote.0.alias.clone();
    upload_result(
        &label,
        local_output,
        remote,
        upload_concurrency,
        &multi_progress,
    )
    .await
}

#[cfg(test)]
mod tests {
    use super::*;

    #[test]
    fn is_zip_key_detects_zip_extension_case_insensitively() {
        assert!(is_zip_key("archive.zip"));
        assert!(is_zip_key("nested/ARCHIVE.ZIP"));
        assert!(!is_zip_key("photo.jpg"));
    }

    #[test]
    fn next_scratch_path_produces_unique_paths() {
        let dir = tempfile::tempdir().unwrap();
        let counter = AtomicU64::new(0);
        let a = next_scratch_path(dir.path(), &counter, "jpg").unwrap();
        let b = next_scratch_path(dir.path(), &counter, "jpg").unwrap();
        assert_ne!(a, b);
    }
}