Skip to main content

pigeon/commands/job/
cli.rs

1use std::path::PathBuf;
2
3use clap::{Args, Subcommand};
4
5use crate::core::observability::Observable;
6
7#[derive(Args, Debug)]
8pub struct JobArgs {
9    #[command(subcommand)]
10    pub command: JobCommands,
11}
12
13#[derive(Subcommand, Debug)]
14pub enum JobCommands {
15    /// Run a job (see subcommands for available job types)
16    Run(RunArgs),
17}
18
19#[derive(Args, Debug)]
20pub struct RunArgs {
21    #[command(subcommand)]
22    pub job_type: JobType,
23}
24
25#[derive(Subcommand, Debug)]
26pub enum JobType {
27    /// Fetch, transform, deduplicate, and optionally upload mail for one or
28    /// more authenticated email identities. Replaces `pigeon email sync`
29    /// (ADR-0021).
30    EmailSync {
31        /// Aliases of the identities to sync, comma-separated. Interactively
32        /// selected from the authenticated identities when omitted and
33        /// stdin is a terminal; required otherwise.
34        #[arg(long, value_delimiter = ',')]
35        identities: Option<Vec<String>>,
36
37        /// Local directory to stage and store output under, shared across
38        /// every selected identity (each gets its own subdirectory
39        /// underneath). Defaults to a directory under the OS temp directory
40        /// when omitted.
41        #[arg(long)]
42        local_output: Option<PathBuf>,
43
44        /// Alias of a configured bucket-config (see `pigeon dataops
45        /// bucket-config new`) to upload each identity's local result tree
46        /// to, once its local fetch/transform/dedupe phase is complete.
47        #[arg(long)]
48        remote_output: Option<String>,
49
50        /// Alias of a configured encryption key (see `pigeon keyring add
51        /// encryption-key`) to use, overriding the target bucket-config's
52        /// own default (if any). Interactively selected/confirmed when
53        /// omitted; falls back to the bucket's default non-interactively
54        /// (ADR-0027).
55        #[arg(long)]
56        encryption_key: Option<String>,
57
58        /// Maximum number of fetch/transform workers to run concurrently,
59        /// spanning every selected identity's every mailbox. Interactively
60        /// prompted (with a rough time estimate) when omitted and stdin is
61        /// a terminal; required otherwise.
62        #[arg(long)]
63        concurrency: Option<usize>,
64
65        /// Maximum number of files to upload concurrently, independent of
66        /// `--concurrency` (which sizes IMAP fetch/transform work) -- the
67        /// upload phase is network-round-trip-bound, not IMAP-bound, so it
68        /// benefits from its own, separately-tuned concurrency (ADR-0091).
69        /// Interactively prompted when omitted and stdin is a terminal;
70        /// defaults to 16 otherwise. With `--upload-only`, this is the only
71        /// concurrency flag that has any effect.
72        #[arg(long)]
73        upload_concurrency: Option<usize>,
74
75        /// Maximum number of simultaneous IMAP connections opened to any
76        /// one identity, regardless of `--concurrency` (ADR-0071) -- caps
77        /// worker concurrency per-account rather than only globally, so a
78        /// mailbox with enough pending batches can't cause more than this
79        /// many workers to log in to the same account at once and trip a
80        /// provider's simultaneous-connection limit. Defaults to 6 (well
81        /// under Gmail's documented 15-connection cap) when omitted; not
82        /// interactively prompted.
83        #[arg(long)]
84        max_connections_per_identity: Option<usize>,
85
86        /// Resumes uploading already-completed local runs for the selected
87        /// identities instead of starting new ones: skips the IMAP
88        /// connect/fetch/transform/dedup phases entirely (and the
89        /// per-identity IMAP credentials they'd otherwise need) and
90        /// uploads straight from each identity's existing local result
91        /// tree, picking up where a prior run's upload phase left off via
92        /// the same `.staging/.uploaded` index per identity (ADR-0090). An
93        /// identity with no completed local run is skipped with a warning
94        /// rather than failing the whole command. Requires a mandatory
95        /// `--remote-output`.
96        #[arg(long)]
97        upload_only: bool,
98
99        /// Alias of a configured bucket-config this run's report, the
100        /// shared observability log, and a transcript of its printed
101        /// output are uploaded to, always unencrypted, under a
102        /// `YYYY-MM-DD-job-name-{run-id}/` prefix (ADR-0100). Mandatory --
103        /// interactively selected when omitted and stdin is a terminal;
104        /// required otherwise.
105        #[arg(long)]
106        report_bucket: Option<String>,
107
108        /// Skip the final "proceed?" confirmation. Every other omitted
109        /// input (identities, concurrency) still follows its own
110        /// independent flag-or-prompt rule -- this only answers the last
111        /// prompt.
112        #[arg(long)]
113        yes: bool,
114    },
115
116    /// Decrypts every `*.enc` file under `--input-dir` into `--output-dir`
117    /// (`.enc` suffix stripped, relative structure preserved), using a
118    /// configured encryption key (ADR-0028).
119    DecryptFiles {
120        /// Directory containing `*.enc` files to decrypt. Interactively
121        /// prompted when omitted and stdin is a terminal; required
122        /// otherwise.
123        #[arg(long)]
124        input_dir: Option<PathBuf>,
125
126        /// Directory decrypted files are written under, mirroring
127        /// `--input-dir`'s relative structure. Must not be the same
128        /// directory as `--input-dir`. Interactively prompted when omitted
129        /// and stdin is a terminal; required otherwise.
130        #[arg(long)]
131        output_dir: Option<PathBuf>,
132
133        /// Alias of a configured encryption key (see `pigeon keyring add
134        /// encryption-key`) to decrypt with. Interactively selected when
135        /// omitted and stdin is a terminal; required otherwise.
136        #[arg(long)]
137        encryption_key: Option<String>,
138
139        /// Maximum number of files to decrypt concurrently.
140        #[arg(long)]
141        concurrency: Option<usize>,
142
143        /// Alias of a configured bucket-config this run's report, the
144        /// shared observability log, and a transcript of its printed
145        /// output are uploaded to, always unencrypted, under a
146        /// `YYYY-MM-DD-job-name-{run-id}/` prefix (ADR-0100). Mandatory --
147        /// interactively selected when omitted and stdin is a terminal;
148        /// required otherwise.
149        #[arg(long)]
150        report_bucket: Option<String>,
151
152        /// Skip the final "proceed?" confirmation.
153        #[arg(long)]
154        yes: bool,
155    },
156
157    /// Fetches raw `.eml` files and unpacked attachments (no Markdown/
158    /// frontmatter transform) for one or more authenticated email
159    /// identities, deduplicating attachments by content, and optionally
160    /// uploads the result unencrypted to a bucket-config (ADR-0081).
161    EmailPull {
162        /// Aliases of the identities to pull, comma-separated.
163        /// Interactively selected from the authenticated identities when
164        /// omitted and stdin is a terminal; required otherwise.
165        #[arg(long, value_delimiter = ',')]
166        identities: Option<Vec<String>>,
167
168        /// Local directory to stage and store output under, shared across
169        /// every selected identity. Defaults to a directory under the OS
170        /// temp directory when omitted.
171        #[arg(long)]
172        local_output: Option<PathBuf>,
173
174        /// Alias of a configured bucket-config to upload each identity's
175        /// local result tree to, once its local fetch/dedupe phase is
176        /// complete. Always uploaded unencrypted -- this job never offers
177        /// encryption (ADR-0081).
178        #[arg(long)]
179        remote_output: Option<String>,
180
181        /// Maximum number of fetch/extract workers to run concurrently,
182        /// spanning every selected identity's every mailbox. Interactively
183        /// prompted (with a rough time estimate) when omitted and stdin is
184        /// a terminal; required otherwise.
185        #[arg(long)]
186        concurrency: Option<usize>,
187
188        /// Maximum number of files to upload concurrently, independent of
189        /// `--concurrency` (which sizes IMAP fetch/extract work) -- the
190        /// upload phase is network-round-trip-bound, not IMAP-bound, so it
191        /// benefits from its own, separately-tuned concurrency (ADR-0091).
192        /// Interactively prompted when omitted and stdin is a terminal;
193        /// defaults to 16 otherwise. With `--upload-only`, this is the only
194        /// concurrency flag that has any effect.
195        #[arg(long)]
196        upload_concurrency: Option<usize>,
197
198        /// Maximum number of simultaneous IMAP connections opened to any
199        /// one identity, regardless of `--concurrency`. Defaults to 6 when
200        /// omitted; not interactively prompted.
201        #[arg(long)]
202        max_connections_per_identity: Option<usize>,
203
204        /// Resumes uploading already-completed local runs for the selected
205        /// identities instead of starting new ones: skips the IMAP
206        /// connect/fetch/dedup phases entirely (and the per-identity IMAP
207        /// credentials they'd otherwise need) and uploads straight from
208        /// each identity's existing local result tree, picking up where a
209        /// prior run's upload phase left off via the same
210        /// `.staging/.uploaded` index per identity (ADR-0090). An identity
211        /// with no completed local run is skipped with a warning rather
212        /// than failing the whole command. Requires a mandatory
213        /// `--remote-output`.
214        #[arg(long)]
215        upload_only: bool,
216
217        /// Alias of a configured bucket-config this run's report, the
218        /// shared observability log, and a transcript of its printed
219        /// output are uploaded to, always unencrypted, under a
220        /// `YYYY-MM-DD-job-name-{run-id}/` prefix (ADR-0100). Mandatory --
221        /// interactively selected when omitted and stdin is a terminal;
222        /// required otherwise.
223        #[arg(long)]
224        report_bucket: Option<String>,
225
226        /// Skip the final "proceed?" confirmation.
227        #[arg(long)]
228        yes: bool,
229    },
230
231    /// Recursively pulls every object from a bucket-config, expands zips,
232    /// recodes media into a size-optimized canonical format per category
233    /// (photo/screenshot -> jpg, video -> mp4, audio -> m4a), dates and
234    /// dedups everything by content, and organizes the result by extension
235    /// -- then optionally encrypts and uploads it to a (possibly
236    /// different) bucket-config (ADR-0074). Requires `ffmpeg`/`ffprobe` on
237    /// `PATH`.
238    PullTransform {
239        /// Alias of a configured bucket-config (see `pigeon keyring add
240        /// bucket`) to pull from. Interactively selected from the
241        /// configured bucket-configs when omitted and stdin is a terminal;
242        /// required otherwise.
243        #[arg(long)]
244        source_bucket: Option<String>,
245
246        /// Local directory to stage and store output under. Defaults to a
247        /// directory under the OS temp directory when omitted.
248        #[arg(long)]
249        local_output: Option<PathBuf>,
250
251        /// Alias of a configured bucket-config to upload the organized
252        /// result to, once local processing is complete.
253        #[arg(long)]
254        remote_output: Option<String>,
255
256        /// Alias of a configured encryption key, overriding the target
257        /// bucket-config's own default (if any).
258        #[arg(long)]
259        encryption_key: Option<String>,
260
261        /// File extensions to pull/transform/upload, comma-separated (e.g.
262        /// `jpg,mp4,pdf`; use the literal `none` for extensionless keys).
263        /// Everything else is left pending, untouched, for a future run --
264        /// never checkpointed as done (ADR-0077). Interactively selected
265        /// (all pre-checked) from the pending-summary table when omitted
266        /// and stdin is a terminal; defaults to everything otherwise.
267        #[arg(long, value_delimiter = ',')]
268        file_types: Option<Vec<String>>,
269
270        /// Keys of pending zip objects to expand and transform;
271        /// comma-separated. Every other pending zip is uploaded as-is,
272        /// untouched (ADR-0077). Interactively selected (all pre-checked)
273        /// when omitted and stdin is a terminal; defaults to expanding
274        /// every pending zip otherwise.
275        #[arg(long, value_delimiter = ',')]
276        expand_zips: Option<Vec<String>>,
277
278        /// Recode target for photos/screenshots: `jpg` (default) or `png`.
279        #[arg(long)]
280        image_format: Option<String>,
281
282        /// Recode target for video: `mp4` (default), `mkv`, or `webm`.
283        #[arg(long)]
284        video_format: Option<String>,
285
286        /// Recode target for audio: `m4a` (default), `mp3`, or `flac`.
287        #[arg(long)]
288        audio_format: Option<String>,
289
290        /// Maximum number of files to download/recode concurrently.
291        #[arg(long)]
292        concurrency: Option<usize>,
293
294        /// Maximum number of files to upload concurrently, independent of
295        /// `--concurrency` (which sizes download/recode work) -- the upload
296        /// phase is network-round-trip-bound, not CPU-bound, so it benefits
297        /// from its own, separately-tuned concurrency (ADR-0091).
298        /// Interactively prompted when omitted and stdin is a terminal;
299        /// defaults to 16 otherwise. With `--upload-only`, this is the only
300        /// concurrency flag that has any effect.
301        #[arg(long)]
302        upload_concurrency: Option<usize>,
303
304        /// Resumes uploading an already-completed local pull-transform run
305        /// instead of starting a new one: skips the bucket listing/
306        /// download/classify/recode/placement phases entirely (and the
307        /// source bucket credentials and `ffmpeg`/`ffprobe` check they'd
308        /// otherwise need) and uploads straight from an existing
309        /// `--local-output`, picking up where a prior run's upload phase
310        /// left off via the same `.staging/.uploaded` index (ADR-0090).
311        /// Requires a `--local-output` from a completed prior run (its
312        /// `.processed` checkpoint must exist and it must hold at least
313        /// one placed-content subdirectory) and a mandatory
314        /// `--remote-output`.
315        #[arg(long)]
316        upload_only: bool,
317
318        /// Alias of a configured bucket-config this run's report, the
319        /// shared observability log, and a transcript of its printed
320        /// output are uploaded to, always unencrypted, under a
321        /// `YYYY-MM-DD-job-name-{run-id}/` prefix (ADR-0100). Mandatory --
322        /// interactively selected when omitted and stdin is a terminal;
323        /// required otherwise.
324        #[arg(long)]
325        report_bucket: Option<String>,
326
327        /// Skip the final "proceed?" confirmation.
328        #[arg(long)]
329        yes: bool,
330    },
331
332    /// Recursively scans a bucket, always inflates every zip found (the
333    /// zip container itself is never uploaded, only its inflated
334    /// contents), content-hashes every file bucket-wide to keep one
335    /// byte-identical copy of each, writes a human-readable merge report,
336    /// and optionally uploads the result unencrypted to a (possibly
337    /// different) bucket-config (ADR-0082). Unlike `pull-transform`, every
338    /// file is always processed and every zip is always expanded -- there
339    /// is no file-type or zip-expansion selection, and this job never
340    /// offers encryption.
341    Deduplicate {
342        /// Alias of a configured bucket-config to pull from. Interactively
343        /// selected from the configured bucket-configs when omitted and
344        /// stdin is a terminal; required otherwise.
345        #[arg(long)]
346        source_bucket: Option<String>,
347
348        /// Local directory to stage and store output under. Defaults to a
349        /// directory under the OS temp directory when omitted.
350        #[arg(long)]
351        local_output: Option<PathBuf>,
352
353        /// Alias of a configured bucket-config to upload the deduped
354        /// result to, once local processing is complete. Always uploaded
355        /// unencrypted.
356        #[arg(long)]
357        remote_output: Option<String>,
358
359        /// Maximum number of files to download/hash concurrently.
360        #[arg(long)]
361        concurrency: Option<usize>,
362
363        /// Maximum number of files to upload concurrently, independent of
364        /// `--concurrency` (which sizes download/hash work) -- the upload
365        /// phase is network-round-trip-bound, not CPU-bound, so it benefits
366        /// from its own, separately-tuned concurrency (ADR-0091).
367        /// Interactively prompted when omitted and stdin is a terminal;
368        /// defaults to 16 otherwise. With `--upload-only`, this is the only
369        /// concurrency flag that has any effect.
370        #[arg(long)]
371        upload_concurrency: Option<usize>,
372
373        /// Resumes uploading an already-completed local deduplicate run instead
374        /// of starting a new one: skips the bucket listing/download/hash/
375        /// placement phases entirely (and the source bucket credentials
376        /// they'd otherwise need) and uploads straight from an existing
377        /// `--local-output`'s `result/` tree, picking up where a prior
378        /// run's upload phase left off via the same `.staging/.uploaded`
379        /// index (ADR-0089). Requires a `--local-output` from a completed
380        /// prior run (its `.staging/.processed` checkpoint must exist and
381        /// its `result/` must be non-empty) and a mandatory
382        /// `--remote-output`.
383        #[arg(long)]
384        upload_only: bool,
385
386        /// Alias of a configured bucket-config this run's report, the
387        /// shared observability log, and a transcript of its printed
388        /// output are uploaded to, always unencrypted, under a
389        /// `YYYY-MM-DD-job-name-{run-id}/` prefix (ADR-0100). Mandatory --
390        /// interactively selected when omitted and stdin is a terminal;
391        /// required otherwise.
392        #[arg(long)]
393        report_bucket: Option<String>,
394
395        /// Skip the final "proceed?" confirmation.
396        #[arg(long)]
397        yes: bool,
398    },
399
400    /// Runs after `deduplicate`: recursively scans a source bucket
401    /// already organized into top-level `<extension>/` folders, classifies
402    /// each extension as either genuinely valuable or an artifact/piece of
403    /// media (TV, movie, software installer, disk image) that's easily
404    /// reproduced from an external canonical source, and forwards only the
405    /// valuable extensions' objects to a mandatory destination bucket
406    /// (ADR-0096). Reproducible extensions are never even downloaded. Never
407    /// offers encryption; never expands zips (input is already flat).
408    Reduce {
409        /// Alias of a configured bucket-config to pull from. Interactively
410        /// selected from the configured bucket-configs when omitted and
411        /// stdin is a terminal; required otherwise.
412        #[arg(long)]
413        source_bucket: Option<String>,
414
415        /// Local directory to stage and store output under. Defaults to a
416        /// directory under the OS temp directory when omitted.
417        #[arg(long)]
418        local_output: Option<PathBuf>,
419
420        /// Alias of a configured bucket-config to upload the forwarded
421        /// result to. Always uploaded unencrypted. Interactively selected
422        /// from the configured bucket-configs when omitted and stdin is a
423        /// terminal; required otherwise -- uploading is mandatory for this
424        /// job.
425        #[arg(long)]
426        remote_output: Option<String>,
427
428        /// Maximum number of files to download concurrently.
429        #[arg(long)]
430        concurrency: Option<usize>,
431
432        /// Maximum number of files to upload concurrently, independent of
433        /// `--concurrency` (which sizes download work) -- the upload phase
434        /// is network-round-trip-bound, not CPU-bound, so it benefits from
435        /// its own, separately-tuned concurrency (ADR-0091). Interactively
436        /// prompted when omitted and stdin is a terminal; defaults to 16
437        /// otherwise. With `--upload-only`, this is the only concurrency
438        /// flag that has any effect.
439        #[arg(long)]
440        upload_concurrency: Option<usize>,
441
442        /// Resumes uploading an already-completed local reduce run instead
443        /// of starting a new one: skips the bucket listing/download/
444        /// placement phases entirely (and the source bucket credentials
445        /// they'd otherwise need) and uploads straight from an existing
446        /// `--local-output`'s `result/` tree, picking up where a prior
447        /// run's upload phase left off via the same `.staging/.uploaded`
448        /// index (ADR-0090). Requires a `--local-output` from a completed
449        /// prior run (its `.staging/.processed` checkpoint must exist and
450        /// its `result/` must be non-empty).
451        #[arg(long)]
452        upload_only: bool,
453
454        /// Extension (without the leading dot, e.g. `mp3`) to always treat
455        /// as valuable regardless of the built-in classification table.
456        /// Repeatable.
457        #[arg(long)]
458        force_valuable: Vec<String>,
459
460        /// Extension (without the leading dot, e.g. `pdf`) to always treat
461        /// as a reproducible artifact/media file regardless of the
462        /// built-in classification table. Repeatable.
463        #[arg(long)]
464        force_reproducible: Vec<String>,
465
466        /// Alias of a configured bucket-config this run's report, the
467        /// shared observability log, and a transcript of its printed
468        /// output are uploaded to, always unencrypted, under a
469        /// `YYYY-MM-DD-job-name-{run-id}/` prefix (ADR-0100). Mandatory --
470        /// interactively selected when omitted and stdin is a terminal;
471        /// required otherwise.
472        #[arg(long)]
473        report_bucket: Option<String>,
474
475        /// Skip the final "proceed?" confirmation.
476        #[arg(long)]
477        yes: bool,
478    },
479
480    /// Copies data from a configurable source to a configurable
481    /// destination by shelling out to the external `rclone` binary, with
482    /// its performance/retry flags fixed (not configurable here). Pigeon
483    /// manages no rclone credentials/config -- `--source`/`--destination`
484    /// are raw `remote:path` strings passed straight through to `rclone
485    /// copy`'s argv; `rclone.conf` is provisioned by an external process,
486    /// outside this crate's scope (ADR-0101). Requires `rclone` on `PATH`.
487    Import {
488        /// rclone source, e.g. `source:media/`. Interactively prompted
489        /// when omitted and stdin is a terminal; required otherwise.
490        #[arg(long)]
491        source: Option<String>,
492
493        /// rclone destination, e.g. `destination:`. Interactively
494        /// prompted when omitted and stdin is a terminal; required
495        /// otherwise.
496        #[arg(long)]
497        destination: Option<String>,
498
499        /// Local directory this run's rclone log (also serving as this
500        /// job's report) and transcript are written under. Defaults to a
501        /// directory under the OS temp directory when omitted. Unlike
502        /// every other job, this is not a staging area for transferred
503        /// data -- rclone transfers directly source -> destination with no
504        /// pigeon-side staging.
505        #[arg(long)]
506        local_output: Option<PathBuf>,
507
508        /// Alias of a configured bucket-config this run's report (the
509        /// rclone log itself, for this job), the shared observability log,
510        /// and a transcript of its printed output are uploaded to, always
511        /// unencrypted, under a `YYYY-MM-DD-job-name-{run-id}/` prefix
512        /// (ADR-0100). Mandatory -- interactively selected when omitted
513        /// and stdin is a terminal; required otherwise.
514        #[arg(long)]
515        report_bucket: Option<String>,
516
517        /// Skip the final "proceed?" confirmation.
518        #[arg(long)]
519        yes: bool,
520    },
521}
522
523impl Observable for JobType {
524    fn command_name(&self) -> &'static str {
525        match self {
526            JobType::EmailSync { .. } => "job.email-sync",
527            JobType::DecryptFiles { .. } => "job.decrypt-files",
528            JobType::EmailPull { .. } => "job.email-pull",
529            JobType::PullTransform { .. } => "job.pull-transform",
530            JobType::Deduplicate { .. } => "job.deduplicate",
531            JobType::Reduce { .. } => "job.reduce",
532            JobType::Import { .. } => "job.import",
533        }
534    }
535}
536
537impl JobType {
538    /// The bare dash-form job name (e.g. `"deduplicate"`) used in the
539    /// report-bucket upload prefix (ADR-0100) -- `command_name()`'s value
540    /// with its `"job."` prefix stripped, rather than a 6th place in this
541    /// codebase repeating the same job-name literals already duplicated
542    /// across every job's `worker.rs` metrics call sites.
543    pub(crate) fn job_name(&self) -> &'static str {
544        self.command_name().trim_start_matches("job.")
545    }
546}