Skip to main content

kanade_shared/wire/
command.rs

1use chrono::{DateTime, Utc};
2use serde::{Deserialize, Serialize};
3
4use super::Staleness;
5use crate::manifest::{CheckHint, CollectHint, EmitConfig};
6
7#[derive(Serialize, Deserialize, schemars::JsonSchema, Debug, Clone)]
8pub struct Command {
9    pub id: String,
10    pub version: String,
11    pub request_id: String,
12    /// v0.29 / Issue #19: the deployment / scheduler-fire UUID this
13    /// Command belongs to. Forwarded into `ExecResult.exec_id` by the
14    /// agent so the projector can attribute results back to the
15    /// originating `executions` row. `None` for ad-hoc `kanade run`
16    /// (no deployment row exists). Pre-v0.29 wire used the field name
17    /// `job_id` for this same value — `serde(alias)` keeps old
18    /// publishes in STREAM_EXEC decodable across the upgrade window.
19    #[serde(alias = "job_id")]
20    pub exec_id: Option<String>,
21    pub shell: Shell,
22    /// Inline script body, OR empty when [`script_object`] is set.
23    /// Mutually exclusive with `script_object` at the wire level —
24    /// backend builders fill one or the other (never both) and the
25    /// agent's resolver picks the populated one. Pre-v0.43 wire
26    /// always carries this populated.
27    ///
28    /// [`script_object`]: Self::script_object
29    pub script: String,
30    /// SPEC §2.4.1 / yukimemi/kanade#210: Object Store reference
31    /// (`<name>/<version>` key into `OBJECT_SCRIPTS`). When set,
32    /// the agent fetches the body via `script_cache` and verifies
33    /// its sha256 against [`script_object_sha256`] before launching.
34    /// `None` ⇒ inline `script` carries the body (legacy + the
35    /// majority of jobs).
36    ///
37    /// [`script_object_sha256`]: Self::script_object_sha256
38    #[serde(default, skip_serializing_if = "Option::is_none")]
39    pub script_object: Option<String>,
40    /// Hex-encoded sha256 of the bytes the operator approved at
41    /// Command-build time. Required when [`script_object`] is set;
42    /// the agent treats a mismatch on fetch as "operator
43    /// re-uploaded the script between exec submission and agent
44    /// fire" and aborts the run rather than silently executing the
45    /// new bytes. Pre-v0.43 wire omits this; the resolver path
46    /// requires both fields to be `Some`.
47    ///
48    /// [`script_object`]: Self::script_object
49    #[serde(default, skip_serializing_if = "Option::is_none")]
50    pub script_object_sha256: Option<String>,
51    pub timeout_secs: u64,
52    pub jitter_secs: Option<u64>,
53    /// Which (token, session) combination the agent should launch the
54    /// child process under (v0.21). Defaults to [`RunAs::System`] for
55    /// back-compat with pre-v0.21 backends that don't send this field.
56    #[serde(default)]
57    pub run_as: RunAs,
58    /// Working directory for the spawned child (v0.21.1). `None` ⇒
59    /// inherit the agent's cwd. Pre-v0.21.1 wire payloads omit this
60    /// field and parse fine via `#[serde(default)]`.
61    #[serde(default, skip_serializing_if = "Option::is_none")]
62    pub cwd: Option<String>,
63    /// Absolute time after which the agent should refuse to run
64    /// this Command (v0.22). Set by the scheduler from
65    /// `Schedule.starting_deadline` (humantime) measured against
66    /// the cron tick time. `None` ⇒ no deadline, run whenever
67    /// received (default for ad-hoc `kanade exec` + back-compat
68    /// for pre-v0.22 wire). The agent stamps a synthetic
69    /// `ExecResult { exit_code: 125, stderr: "skipped: deadline
70    /// expired ..." }` when it skips, so the operator sees the
71    /// outcome on the Results / Dashboard pages instead of silence.
72    #[serde(default, skip_serializing_if = "Option::is_none")]
73    pub deadline_at: Option<DateTime<Utc>>,
74    /// v0.26: Manifest-declared Layer 2 staleness policy
75    /// (see SPEC.md §2.6.2). Forwarded from `Manifest.staleness` so
76    /// the agent can evaluate it at fire time without re-fetching the
77    /// Manifest from `BUCKET_JOBS`. Pre-v0.26 wire omits this and
78    /// `#[serde(default)]` falls back to `Staleness::Cached`, matching
79    /// pre-v0.26 behaviour (silently use cached KV values).
80    #[serde(default)]
81    pub staleness: Staleness,
82    /// Issue #246: forwarded from `Manifest.emit` so the agent
83    /// doesn't have to re-fetch the manifest at fire time. When
84    /// `Some` and `EmitKind::Events`, the agent parses script
85    /// stdout as NDJSON `ObsEvent` and publishes each line on
86    /// `obs.<pc_id>`. Pre-#246 wire omits this; the `#[serde(default)]`
87    /// fallback to `None` preserves prior behaviour (stdout flows
88    /// to `ExecResult` unchanged).
89    #[serde(default, skip_serializing_if = "Option::is_none")]
90    pub emit: Option<EmitConfig>,
91    /// #290: forwarded from `Manifest.check` so the agent can build a
92    /// KLP Health-tab [`Check`](crate::ipc::state::Check) from the
93    /// job's stdout without re-fetching the Manifest. When `Some`, the
94    /// agent reads the `status_field` / `detail_field` values out of
95    /// the stdout JSON object after a successful run and caches the
96    /// result into `StateSnapshot.checks`. Pre-#290 wire omits this;
97    /// `#[serde(default)]` → `None` preserves prior behaviour.
98    #[serde(default, skip_serializing_if = "Option::is_none")]
99    pub check: Option<CheckHint>,
100    /// #219: forwarded from `Manifest.collect` so the agent can bundle
101    /// the script's listed files without re-fetching the Manifest. When
102    /// `Some`, the agent — after a successful run — reads the
103    /// `files_field` path array out of the stdout JSON object, zips those
104    /// files (capped at `max_size`), uploads the archive to
105    /// `OBJECT_COLLECTIONS`, and records the key in
106    /// [`ExecResult::collect_object`](super::ExecResult::collect_object).
107    /// Pre-#219 wire omits this; `#[serde(default)]` → `None` preserves
108    /// prior behaviour.
109    #[serde(default, skip_serializing_if = "Option::is_none")]
110    pub collect: Option<CollectHint>,
111    /// #418 Phase 4: lowered from `Schedule.on_failure.retry` by the
112    /// command builders (backend `exec_manifest` + the agent's local
113    /// scheduler). When `Some`, the agent re-runs the script
114    /// in-process on a non-zero exit / timeout, up to `max` extra
115    /// attempts with `backoff_secs` between them, before publishing
116    /// the final outcome. `None` (default) ⇒ no retry, the historical
117    /// behaviour and what ad-hoc `kanade run` / `kanade exec` use.
118    /// Pre-Phase-4 wire omits this; `#[serde(default)]` → `None`.
119    #[serde(default, skip_serializing_if = "Option::is_none")]
120    pub retry: Option<RetrySpec>,
121    /// Job-generic post-step hook lowered from `Manifest.finalize`. When
122    /// `Some` and the main script exits cleanly, the agent runs this hook
123    /// after the collect step (injecting `KANADE_COLLECT_RESULT` for a
124    /// `collect:` job) so the operator can delete / move / notify.
125    /// Best-effort — a finalize failure is logged, never published as the
126    /// run's outcome. Pre-finalize wire omits this; `#[serde(default)]` →
127    /// `None`.
128    #[serde(default, skip_serializing_if = "Option::is_none")]
129    pub finalize: Option<FinalizeCommand>,
130}
131
132/// Lowered, engine-vocabulary form of [`crate::manifest::FinalizeSpec`]
133/// — the post-step hook stamped onto a [`Command`]. The operator-facing
134/// humantime `timeout` is reduced to whole seconds at build time
135/// (mirrors `timeout_secs`), and the manifest `ExecuteShell` to the wire
136/// [`Shell`], so the agent's fire path does no parsing.
137#[derive(Serialize, Deserialize, schemars::JsonSchema, Debug, Clone)]
138pub struct FinalizeCommand {
139    pub shell: Shell,
140    /// Inline script body (inline-only in P1).
141    pub script: String,
142    pub timeout_secs: u64,
143    #[serde(default)]
144    pub run_as: RunAs,
145    #[serde(default, skip_serializing_if = "Option::is_none")]
146    pub cwd: Option<String>,
147    /// #965: for a `collect:` job, run this hook once per uploaded
148    /// bundle (single-bundle `KANADE_COLLECT_RESULT`) as each bundle
149    /// uploads, instead of once after the whole set — so an interrupted
150    /// collect still cleans up the days it managed to ship. `false`
151    /// (default, pre-#965 wire) keeps the one-call-after-all contract.
152    #[serde(default)]
153    pub on_each_bundle: bool,
154}
155
156/// Lowered, engine-vocabulary form of [`crate::manifest::Retry`] — a
157/// fixed-backoff retry policy stamped onto a [`Command`]. The
158/// operator-facing humantime `backoff` is reduced to whole seconds at
159/// build time (mirrors how `jitter_secs` / `timeout_secs` are
160/// pre-lowered) so the agent's fire path does no humantime parsing.
161#[derive(Serialize, Deserialize, schemars::JsonSchema, Debug, Clone, Copy, PartialEq, Eq)]
162pub struct RetrySpec {
163    /// Max additional attempts after the first failure (1..=10,
164    /// enforced by `Schedule::validate`).
165    pub max: u32,
166    /// Seconds slept between attempts.
167    pub backoff_secs: u64,
168}
169
170#[derive(Serialize, Deserialize, schemars::JsonSchema, Debug, Clone, Copy, PartialEq, Eq)]
171#[serde(rename_all = "lowercase")]
172pub enum Shell {
173    /// Windows PowerShell 5.1 — the literal `powershell` on PATH. The
174    /// agent stages the script to a temp `.ps1` and runs it via a
175    /// UTF-8-console launcher (see `process.rs`).
176    Powershell,
177    /// `cmd.exe` — `cmd /C <script>` inline. Windows only.
178    Cmd,
179    /// POSIX shell — `sh -c <script>` inline. Linux/macOS. The agent
180    /// spawns the literal `sh` on PATH; there is no per-OS gate, so a
181    /// misdirected `sh` job to a host without `sh` fails at spawn.
182    Sh,
183    /// PowerShell 7 (cross-platform) — the literal `pwsh` on PATH.
184    /// Distinct from [`Shell::Powershell`] (Windows PowerShell 5.1):
185    /// reuses the same temp-`.ps1` launcher, but skips
186    /// `-ExecutionPolicy Bypass` off Windows (no execution policy
187    /// there).
188    Pwsh,
189}
190
191/// **Token + session combination** the agent uses to spawn a job's
192/// child process. Two orthogonal axes — *whose privileges* and *which
193/// session* — collapse into three meaningful combinations:
194///
195/// | variant            | session                | privileges  | GUI |
196/// |--------------------|------------------------|-------------|-----|
197/// | `System` (default) | Session 0 (services)   | LocalSystem | ❌  |
198/// | `User`             | active console session | logged-in user (UAC-filtered when admin) | ✅ |
199/// | `SystemGui`        | active console session | LocalSystem | ✅  |
200///
201/// `SystemGui` is the "PsExec `-i -s`" pattern: the agent duplicates
202/// its own SYSTEM token and rewrites `TokenSessionId` to the user's
203/// console session, then launches with that hybrid token — useful
204/// when an installer needs admin power *and* needs the user to see
205/// its UI.
206#[derive(
207    Serialize, Deserialize, schemars::JsonSchema, Debug, Clone, Copy, PartialEq, Eq, Default,
208)]
209#[serde(rename_all = "snake_case")]
210pub enum RunAs {
211    /// LocalSystem privileges in Session 0. No GUI. Historical
212    /// default — every pre-v0.21 job ran this way.
213    #[default]
214    System,
215    /// The currently-logged-in console user's identity, in their
216    /// session. Can write HKCU / %APPDATA% / show GUI to the user.
217    /// Privileges are whatever the user has (admin users get the
218    /// UAC-filtered limited token, not the elevated one).
219    User,
220    /// LocalSystem privileges in the user's session — admin power
221    /// with GUI visibility. Niche but real (force-restart dialogs,
222    /// admin installers with progress UI).
223    SystemGui,
224}
225
226#[cfg(test)]
227mod tests {
228    use super::*;
229
230    fn sample_command() -> Command {
231        Command {
232            id: "echo-test".into(),
233            version: "1.0.0".into(),
234            request_id: "req-1".into(),
235            exec_id: Some("dep-1".into()),
236            shell: Shell::Powershell,
237            script: "echo hi".into(),
238            script_object: None,
239            script_object_sha256: None,
240            timeout_secs: 30,
241            jitter_secs: Some(5),
242            run_as: RunAs::System,
243            cwd: None,
244            deadline_at: None,
245            staleness: Staleness::Cached,
246            emit: None,
247            check: None,
248            collect: None,
249            retry: None,
250            finalize: None,
251        }
252    }
253
254    #[test]
255    fn shell_serialises_lowercase() {
256        let json = serde_json::to_string(&Shell::Powershell).unwrap();
257        assert_eq!(json, "\"powershell\"");
258        let json = serde_json::to_string(&Shell::Cmd).unwrap();
259        assert_eq!(json, "\"cmd\"");
260        let json = serde_json::to_string(&Shell::Sh).unwrap();
261        assert_eq!(json, "\"sh\"");
262        let json = serde_json::to_string(&Shell::Pwsh).unwrap();
263        assert_eq!(json, "\"pwsh\"");
264        // Round-trip the new variants from the wire form.
265        assert_eq!(serde_json::from_str::<Shell>("\"sh\"").unwrap(), Shell::Sh);
266        assert_eq!(
267            serde_json::from_str::<Shell>("\"pwsh\"").unwrap(),
268            Shell::Pwsh
269        );
270    }
271
272    #[test]
273    fn run_as_serialises_snake_case() {
274        for (mode, expected) in [
275            (RunAs::System, "\"system\""),
276            (RunAs::User, "\"user\""),
277            (RunAs::SystemGui, "\"system_gui\""),
278        ] {
279            let json = serde_json::to_string(&mode).unwrap();
280            assert_eq!(json, expected, "serialise {mode:?}");
281            let back: RunAs = serde_json::from_str(expected).unwrap();
282            assert_eq!(back, mode, "round-trip {expected}");
283        }
284    }
285
286    #[test]
287    fn run_as_defaults_to_system() {
288        assert_eq!(RunAs::default(), RunAs::System);
289    }
290
291    #[test]
292    fn command_round_trips_through_json() {
293        let orig = sample_command();
294        let json = serde_json::to_string(&orig).expect("encode");
295        let decoded: Command = serde_json::from_str(&json).expect("decode");
296        assert_eq!(decoded.id, orig.id);
297        assert_eq!(decoded.version, orig.version);
298        assert_eq!(decoded.request_id, orig.request_id);
299        assert_eq!(decoded.exec_id, orig.exec_id);
300        assert_eq!(decoded.shell, orig.shell);
301        assert_eq!(decoded.script, orig.script);
302        assert_eq!(decoded.timeout_secs, orig.timeout_secs);
303        assert_eq!(decoded.jitter_secs, orig.jitter_secs);
304        assert_eq!(decoded.run_as, orig.run_as);
305    }
306
307    #[test]
308    fn command_round_trips_each_run_as_variant() {
309        for mode in [RunAs::System, RunAs::User, RunAs::SystemGui] {
310            let cmd = Command {
311                run_as: mode,
312                ..sample_command()
313            };
314            let json = serde_json::to_string(&cmd).unwrap();
315            let back: Command = serde_json::from_str(&json).unwrap();
316            assert_eq!(back.run_as, mode);
317        }
318    }
319
320    #[test]
321    fn command_accepts_missing_optional_fields() {
322        let json = r#"{
323          "id": "x",
324          "version": "1.0.0",
325          "request_id": "r",
326          "shell": "cmd",
327          "script": "echo",
328          "timeout_secs": 5
329        }"#;
330        let cmd: Command = serde_json::from_str(json).expect("decode");
331        assert!(cmd.exec_id.is_none());
332        assert!(cmd.jitter_secs.is_none());
333        assert_eq!(cmd.shell, Shell::Cmd);
334        // Pre-v0.21 wire payloads omit run_as → falls back to System.
335        assert_eq!(cmd.run_as, RunAs::System);
336        // Pre-v0.21.1 omit cwd → None (= inherit agent cwd).
337        assert!(cmd.cwd.is_none());
338        // Pre-v0.22 omit deadline_at → None (= no deadline).
339        assert!(cmd.deadline_at.is_none());
340        // Pre-v0.43 wire omits both script_object fields — agent
341        // falls back to the inline `script` body.
342        assert!(cmd.script_object.is_none());
343        assert!(cmd.script_object_sha256.is_none());
344    }
345
346    #[test]
347    fn command_round_trips_script_object_fields() {
348        // yukimemi/kanade#210: backend builds Commands carrying an
349        // OBJECT_SCRIPTS reference + the operator-approved digest;
350        // agent resolves on fetch. Both fields must survive a JSON
351        // round-trip with the same shape.
352        let cmd = Command {
353            script: String::new(),
354            script_object: Some("cleanup-disk-temp/1.0.1".into()),
355            script_object_sha256: Some(
356                "deadbeefdeadbeefdeadbeefdeadbeefdeadbeefdeadbeefdeadbeefdeadbeef".into(),
357            ),
358            ..sample_command()
359        };
360        let json = serde_json::to_string(&cmd).expect("encode");
361        let back: Command = serde_json::from_str(&json).expect("decode");
362        assert_eq!(back.script, "");
363        assert_eq!(
364            back.script_object.as_deref(),
365            Some("cleanup-disk-temp/1.0.1")
366        );
367        assert_eq!(
368            back.script_object_sha256.as_deref(),
369            Some("deadbeefdeadbeefdeadbeefdeadbeefdeadbeefdeadbeefdeadbeefdeadbeef")
370        );
371    }
372
373    #[test]
374    fn command_decodes_legacy_job_id_field_as_exec_id() {
375        // v0.29 / Issue #19: Commands sitting in STREAM_EXEC published
376        // by a pre-v0.29 backend still carry the field named `job_id`.
377        // The `#[serde(alias = "job_id")]` on `exec_id` keeps them
378        // decodable through the upgrade window so the agent doesn't
379        // start dropping replays on first boot of a new binary.
380        let json = r#"{
381          "id": "x",
382          "version": "1.0.0",
383          "request_id": "r",
384          "job_id": "legacy-exec-uuid",
385          "shell": "powershell",
386          "script": "echo",
387          "timeout_secs": 5
388        }"#;
389        let cmd: Command = serde_json::from_str(json).expect("decode legacy");
390        assert_eq!(cmd.exec_id.as_deref(), Some("legacy-exec-uuid"));
391    }
392
393    #[test]
394    fn command_deadline_at_round_trips() {
395        use chrono::TimeZone;
396        let deadline = Utc.with_ymd_and_hms(2026, 5, 18, 9, 30, 0).unwrap();
397        let cmd = Command {
398            deadline_at: Some(deadline),
399            ..sample_command()
400        };
401        let json = serde_json::to_string(&cmd).unwrap();
402        let back: Command = serde_json::from_str(&json).unwrap();
403        assert_eq!(back.deadline_at, Some(deadline));
404    }
405
406    #[test]
407    fn command_retry_round_trips() {
408        // #418 Phase 4: a stamped retry policy must survive the wire
409        // so the agent can apply it on a live publish or a STREAM_EXEC
410        // replay.
411        let cmd = Command {
412            retry: Some(RetrySpec {
413                max: 3,
414                backoff_secs: 600,
415            }),
416            ..sample_command()
417        };
418        let json = serde_json::to_string(&cmd).unwrap();
419        let back: Command = serde_json::from_str(&json).unwrap();
420        assert_eq!(
421            back.retry,
422            Some(RetrySpec {
423                max: 3,
424                backoff_secs: 600
425            })
426        );
427    }
428
429    #[test]
430    fn command_omits_retry_when_absent() {
431        // skip_serializing_if keeps the field off the wire for the
432        // common (no-retry) case, and pre-Phase-4 payloads that never
433        // had it still decode (serde default → None).
434        let json = serde_json::to_string(&sample_command()).unwrap();
435        assert!(
436            !json.contains("retry"),
437            "retry must not appear when None: {json}"
438        );
439    }
440
441    #[test]
442    fn command_collect_round_trips_and_omits_when_absent() {
443        // #219: `collect` is off the wire when None (skip_serializing_if),
444        // so pre-#219 readers don't trip over it...
445        let json = serde_json::to_string(&sample_command()).unwrap();
446        assert!(
447            !json.contains("collect"),
448            "collect must be absent when None: {json}"
449        );
450        // ...and a forwarded CollectHint survives the round-trip.
451        let cmd = Command {
452            collect: Some(CollectHint {
453                name: "diag".into(),
454                description: Some("logs".into()),
455                max_size: Some("50MB".into()),
456                files_field: "files".into(),
457            }),
458            ..sample_command()
459        };
460        let back: Command = serde_json::from_str(&serde_json::to_string(&cmd).unwrap()).unwrap();
461        let c = back.collect.expect("collect survived round-trip");
462        assert_eq!(c.name, "diag");
463        assert_eq!(c.max_size.as_deref(), Some("50MB"));
464        assert_eq!(c.files_field, "files");
465    }
466
467    #[test]
468    fn command_finalize_round_trips_and_omits_when_absent() {
469        // Off the wire when None (skip_serializing_if), so pre-finalize
470        // readers don't trip over it...
471        let json = serde_json::to_string(&sample_command()).unwrap();
472        assert!(
473            !json.contains("finalize"),
474            "finalize must be absent when None: {json}"
475        );
476        // ...and a forwarded FinalizeCommand survives the round-trip.
477        let cmd = Command {
478            finalize: Some(FinalizeCommand {
479                shell: Shell::Powershell,
480                script: "Remove-Item $env:FILE".into(),
481                timeout_secs: 30,
482                run_as: RunAs::System,
483                cwd: None,
484                on_each_bundle: false,
485            }),
486            ..sample_command()
487        };
488        let back: Command = serde_json::from_str(&serde_json::to_string(&cmd).unwrap()).unwrap();
489        let f = back.finalize.expect("finalize survived round-trip");
490        assert_eq!(f.shell, Shell::Powershell);
491        assert_eq!(f.script, "Remove-Item $env:FILE");
492        assert_eq!(f.timeout_secs, 30);
493        assert_eq!(f.run_as, RunAs::System);
494    }
495}