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