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#[derive(
209 Serialize, Deserialize, schemars::JsonSchema, Debug, Clone, Copy, PartialEq, Eq, Default,
210)]
211#[serde(rename_all = "snake_case")]
212pub enum RunAs {
213 /// LocalSystem privileges in Session 0. No GUI. Historical
214 /// default — every pre-v0.21 job ran this way.
215 #[default]
216 System,
217 /// The currently-logged-in console user's identity, in their
218 /// session. Can write HKCU / %APPDATA% / show GUI to the user.
219 /// Privileges are whatever the user has (admin users get the
220 /// UAC-filtered limited token, not the elevated one).
221 User,
222 /// LocalSystem privileges in the user's session — admin power
223 /// with GUI visibility. Niche but real (force-restart dialogs,
224 /// admin installers with progress UI).
225 SystemGui,
226}
227
228#[cfg(test)]
229mod tests {
230 use super::*;
231
232 fn sample_command() -> Command {
233 Command {
234 id: "echo-test".into(),
235 version: "1.0.0".into(),
236 request_id: "req-1".into(),
237 exec_id: Some("dep-1".into()),
238 shell: Shell::Powershell,
239 script: "echo hi".into(),
240 script_object: None,
241 script_object_sha256: None,
242 timeout_secs: 30,
243 bypass_local_limit: false,
244 jitter_secs: Some(5),
245 run_as: RunAs::System,
246 cwd: None,
247 deadline_at: None,
248 staleness: Staleness::Cached,
249 emit: None,
250 check: None,
251 collect: None,
252 retry: None,
253 finalize: None,
254 }
255 }
256
257 #[test]
258 fn shell_serialises_lowercase() {
259 let json = serde_json::to_string(&Shell::Powershell).unwrap();
260 assert_eq!(json, "\"powershell\"");
261 let json = serde_json::to_string(&Shell::Cmd).unwrap();
262 assert_eq!(json, "\"cmd\"");
263 let json = serde_json::to_string(&Shell::Sh).unwrap();
264 assert_eq!(json, "\"sh\"");
265 let json = serde_json::to_string(&Shell::Pwsh).unwrap();
266 assert_eq!(json, "\"pwsh\"");
267 // Round-trip the new variants from the wire form.
268 assert_eq!(serde_json::from_str::<Shell>("\"sh\"").unwrap(), Shell::Sh);
269 assert_eq!(
270 serde_json::from_str::<Shell>("\"pwsh\"").unwrap(),
271 Shell::Pwsh
272 );
273 }
274
275 #[test]
276 fn run_as_serialises_snake_case() {
277 for (mode, expected) in [
278 (RunAs::System, "\"system\""),
279 (RunAs::User, "\"user\""),
280 (RunAs::SystemGui, "\"system_gui\""),
281 ] {
282 let json = serde_json::to_string(&mode).unwrap();
283 assert_eq!(json, expected, "serialise {mode:?}");
284 let back: RunAs = serde_json::from_str(expected).unwrap();
285 assert_eq!(back, mode, "round-trip {expected}");
286 }
287 }
288
289 #[test]
290 fn run_as_defaults_to_system() {
291 assert_eq!(RunAs::default(), RunAs::System);
292 }
293
294 #[test]
295 fn command_round_trips_through_json() {
296 let orig = sample_command();
297 let json = serde_json::to_string(&orig).expect("encode");
298 let decoded: Command = serde_json::from_str(&json).expect("decode");
299 assert_eq!(decoded.id, orig.id);
300 assert_eq!(decoded.version, orig.version);
301 assert_eq!(decoded.request_id, orig.request_id);
302 assert_eq!(decoded.exec_id, orig.exec_id);
303 assert_eq!(decoded.shell, orig.shell);
304 assert_eq!(decoded.script, orig.script);
305 assert_eq!(decoded.timeout_secs, orig.timeout_secs);
306 assert_eq!(decoded.jitter_secs, orig.jitter_secs);
307 assert_eq!(decoded.run_as, orig.run_as);
308 }
309
310 #[test]
311 fn command_round_trips_each_run_as_variant() {
312 for mode in [RunAs::System, RunAs::User, RunAs::SystemGui] {
313 let cmd = Command {
314 run_as: mode,
315 ..sample_command()
316 };
317 let json = serde_json::to_string(&cmd).unwrap();
318 let back: Command = serde_json::from_str(&json).unwrap();
319 assert_eq!(back.run_as, mode);
320 }
321 }
322
323 #[test]
324 fn command_accepts_missing_optional_fields() {
325 let json = r#"{
326 "id": "x",
327 "version": "1.0.0",
328 "request_id": "r",
329 "shell": "cmd",
330 "script": "echo",
331 "timeout_secs": 5
332 }"#;
333 let cmd: Command = serde_json::from_str(json).expect("decode");
334 assert!(cmd.exec_id.is_none());
335 assert!(cmd.jitter_secs.is_none());
336 assert_eq!(cmd.shell, Shell::Cmd);
337 // Pre-v0.21 wire payloads omit run_as → falls back to System.
338 assert_eq!(cmd.run_as, RunAs::System);
339 // Pre-v0.21.1 omit cwd → None (= inherit agent cwd).
340 assert!(cmd.cwd.is_none());
341 // Pre-v0.22 omit deadline_at → None (= no deadline).
342 assert!(cmd.deadline_at.is_none());
343 // Pre-v0.43 wire omits both script_object fields — agent
344 // falls back to the inline `script` body.
345 assert!(cmd.script_object.is_none());
346 assert!(cmd.script_object_sha256.is_none());
347 }
348
349 #[test]
350 fn command_round_trips_script_object_fields() {
351 // kanadehq/kanade#210: backend builds Commands carrying an
352 // OBJECT_SCRIPTS reference + the operator-approved digest;
353 // agent resolves on fetch. Both fields must survive a JSON
354 // round-trip with the same shape.
355 let cmd = Command {
356 script: String::new(),
357 script_object: Some("cleanup-disk-temp/1.0.1".into()),
358 script_object_sha256: Some(
359 "deadbeefdeadbeefdeadbeefdeadbeefdeadbeefdeadbeefdeadbeefdeadbeef".into(),
360 ),
361 ..sample_command()
362 };
363 let json = serde_json::to_string(&cmd).expect("encode");
364 let back: Command = serde_json::from_str(&json).expect("decode");
365 assert_eq!(back.script, "");
366 assert_eq!(
367 back.script_object.as_deref(),
368 Some("cleanup-disk-temp/1.0.1")
369 );
370 assert_eq!(
371 back.script_object_sha256.as_deref(),
372 Some("deadbeefdeadbeefdeadbeefdeadbeefdeadbeefdeadbeefdeadbeefdeadbeef")
373 );
374 }
375
376 #[test]
377 fn command_decodes_legacy_job_id_field_as_exec_id() {
378 // v0.29 / Issue #19: Commands sitting in STREAM_EXEC published
379 // by a pre-v0.29 backend still carry the field named `job_id`.
380 // The `#[serde(alias = "job_id")]` on `exec_id` keeps them
381 // decodable through the upgrade window so the agent doesn't
382 // start dropping replays on first boot of a new binary.
383 let json = r#"{
384 "id": "x",
385 "version": "1.0.0",
386 "request_id": "r",
387 "job_id": "legacy-exec-uuid",
388 "shell": "powershell",
389 "script": "echo",
390 "timeout_secs": 5
391 }"#;
392 let cmd: Command = serde_json::from_str(json).expect("decode legacy");
393 assert_eq!(cmd.exec_id.as_deref(), Some("legacy-exec-uuid"));
394 }
395
396 #[test]
397 fn command_deadline_at_round_trips() {
398 use chrono::TimeZone;
399 let deadline = Utc.with_ymd_and_hms(2026, 5, 18, 9, 30, 0).unwrap();
400 let cmd = Command {
401 deadline_at: Some(deadline),
402 ..sample_command()
403 };
404 let json = serde_json::to_string(&cmd).unwrap();
405 let back: Command = serde_json::from_str(&json).unwrap();
406 assert_eq!(back.deadline_at, Some(deadline));
407 }
408
409 #[test]
410 fn command_retry_round_trips() {
411 // #418 Phase 4: a stamped retry policy must survive the wire
412 // so the agent can apply it on a live publish or a STREAM_EXEC
413 // replay.
414 let cmd = Command {
415 retry: Some(RetrySpec {
416 max: 3,
417 backoff_secs: 600,
418 }),
419 ..sample_command()
420 };
421 let json = serde_json::to_string(&cmd).unwrap();
422 let back: Command = serde_json::from_str(&json).unwrap();
423 assert_eq!(
424 back.retry,
425 Some(RetrySpec {
426 max: 3,
427 backoff_secs: 600
428 })
429 );
430 }
431
432 #[test]
433 fn command_omits_retry_when_absent() {
434 // skip_serializing_if keeps the field off the wire for the
435 // common (no-retry) case, and pre-Phase-4 payloads that never
436 // had it still decode (serde default → None).
437 let json = serde_json::to_string(&sample_command()).unwrap();
438 assert!(
439 !json.contains("retry"),
440 "retry must not appear when None: {json}"
441 );
442 }
443
444 #[test]
445 fn command_collect_round_trips_and_omits_when_absent() {
446 // #219: `collect` is off the wire when None (skip_serializing_if),
447 // so pre-#219 readers don't trip over it...
448 let json = serde_json::to_string(&sample_command()).unwrap();
449 assert!(
450 !json.contains("collect"),
451 "collect must be absent when None: {json}"
452 );
453 // ...and a forwarded CollectHint survives the round-trip.
454 let cmd = Command {
455 collect: Some(CollectHint {
456 name: "diag".into(),
457 description: Some("logs".into()),
458 max_size: Some("50MB".into()),
459 files_field: "files".into(),
460 }),
461 ..sample_command()
462 };
463 let back: Command = serde_json::from_str(&serde_json::to_string(&cmd).unwrap()).unwrap();
464 let c = back.collect.expect("collect survived round-trip");
465 assert_eq!(c.name, "diag");
466 assert_eq!(c.max_size.as_deref(), Some("50MB"));
467 assert_eq!(c.files_field, "files");
468 }
469
470 #[test]
471 fn command_finalize_round_trips_and_omits_when_absent() {
472 // Off the wire when None (skip_serializing_if), so pre-finalize
473 // readers don't trip over it...
474 let json = serde_json::to_string(&sample_command()).unwrap();
475 assert!(
476 !json.contains("finalize"),
477 "finalize must be absent when None: {json}"
478 );
479 // ...and a forwarded FinalizeCommand survives the round-trip.
480 let cmd = Command {
481 finalize: Some(FinalizeCommand {
482 shell: Shell::Powershell,
483 script: "Remove-Item $env:FILE".into(),
484 timeout_secs: 30,
485 run_as: RunAs::System,
486 cwd: None,
487 on_each_bundle: false,
488 }),
489 ..sample_command()
490 };
491 let back: Command = serde_json::from_str(&serde_json::to_string(&cmd).unwrap()).unwrap();
492 let f = back.finalize.expect("finalize survived round-trip");
493 assert_eq!(f.shell, Shell::Powershell);
494 assert_eq!(f.script, "Remove-Item $env:FILE");
495 assert_eq!(f.timeout_secs, 30);
496 assert_eq!(f.run_as, RunAs::System);
497 }
498}