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}