Skip to main content

kanade_shared/wire/
result.rs

1use serde::{Deserialize, Serialize};
2use uuid::Uuid;
3
4/// Prefix injected into the UUIDv5 name string for deriving legacy
5/// `result_id`s. Fixed marker so two backends (or one backend across
6/// restarts) projecting the same legacy payload arrive at the same
7/// id. Tied to the standard `Uuid::NAMESPACE_OID` namespace below.
8/// Bumping this prefix would break dedupe of legacy redeliveries
9/// crossing the upgrade — don't.
10const LEGACY_RESULT_ID_PREFIX: &str = "kanade-issue-19/legacy-result-id:";
11
12#[derive(Serialize, Deserialize, Debug, Clone)]
13pub struct ExecResult {
14    /// v0.29 / Issue #19: agent-minted UUID, unique per (Command, PC)
15    /// run. Replaces `request_id` as the projector's primary key so
16    /// broadcast Commands (commands.all / commands.group.X) — where N
17    /// PCs share one `request_id` — finally persist all N results
18    /// instead of silently dropping all but the first. Pre-v0.29
19    /// agents omit this field; it deserialises as the empty string,
20    /// and [`Self::stable_result_id`] derives a deterministic UUIDv5
21    /// from `(request_id, pc_id)` so legacy payloads (a) get distinct
22    /// ids across broadcast PCs (PC #2's row stops being dropped) and
23    /// (b) get the SAME id on JetStream redelivery (the new `ON
24    /// CONFLICT(result_id) DO NOTHING` path correctly dedupes, so
25    /// `executions.success_count` doesn't double-count across retries).
26    #[serde(default)]
27    pub result_id: String,
28    /// The NATS reply token. Still surfaced for joining back to the
29    /// `kanade run` request/reply path. No longer unique across rows
30    /// (broadcast Commands share it).
31    pub request_id: String,
32    /// v0.29 / Issue #19: back-link to `executions.exec_id`. Copied
33    /// from `Command.exec_id` by the agent. `None` for ad-hoc
34    /// `kanade run` (no deployment) and for results emitted by
35    /// pre-v0.29 agents (decoded via `serde(default)`).
36    #[serde(default, skip_serializing_if = "Option::is_none")]
37    pub exec_id: Option<String>,
38    /// #955: back-link to the parent run's `result_id` for a
39    /// `finalize:` hook's own result row. Set by the agent to the
40    /// triggering run's `result_id` so the SPA can link the
41    /// `<job>__finalize` row to (and from) the run whose collect it
42    /// cleaned up. `None` for every ordinary run and every pre-#955
43    /// payload (`serde(default)` keeps older results decodable).
44    #[serde(default, skip_serializing_if = "Option::is_none")]
45    pub parent_result_id: Option<String>,
46    pub pc_id: String,
47    pub exit_code: i32,
48    /// stdout. Empty string when [`Self::stdout_object`] is set — the
49    /// agent overflowed the bytes into [`crate::kv::OBJECT_RESULT_OUTPUT`]
50    /// because the inline payload would have exceeded NATS's default
51    /// `max_payload` (#227). The backend projector derefs the pointer
52    /// before inserting; SQLite still stores the full text inline so
53    /// the SPA Activity page reads unchanged.
54    pub stdout: String,
55    pub stderr: String,
56    pub started_at: chrono::DateTime<chrono::Utc>,
57    pub finished_at: chrono::DateTime<chrono::Utc>,
58    /// Object Store key under [`crate::kv::OBJECT_RESULT_OUTPUT`] when
59    /// `stdout` overflowed the agent's inline threshold (#227). Set to
60    /// `Some("<request_id>/stdout")` by the agent's outbox drain; the
61    /// backend projector fetches the bytes from that key and uses them
62    /// in place of the (empty) `stdout` field. `None` for the common
63    /// small-stdout case + every pre-#227 payload (`serde(default)`
64    /// keeps older results decodable).
65    #[serde(default, skip_serializing_if = "Option::is_none")]
66    pub stdout_object: Option<String>,
67    /// Sibling of `stdout_object` for the stderr stream. Same key
68    /// shape (`<request_id>/stderr`).
69    #[serde(default, skip_serializing_if = "Option::is_none")]
70    pub stderr_object: Option<String>,
71    /// v0.13: the manifest id that produced this result. Sourced
72    /// from `Command.id` (which is the YAML `manifest.id`, e.g.
73    /// `"inventory-hw"`). Distinct from the per-deploy UUID stored
74    /// in `Command.exec_id`. The results projector uses this to
75    /// look up the manifest's `inventory:` hint and upsert
76    /// `inventory_facts` rows for inventory-tagged jobs.
77    #[serde(default, skip_serializing_if = "Option::is_none")]
78    pub manifest_id: Option<String>,
79    /// #219: Object Store key under [`crate::kv::OBJECT_COLLECTIONS`] for
80    /// the bundle this run collected, when the job carried a `collect:`
81    /// hint and the run succeeded. Set by the agent to
82    /// `Some("<pc_id>/<job_id>/<rfc3339>.zip")` after it zips the
83    /// script's listed files and uploads the archive. `None` for every
84    /// non-collect job + every pre-#219 payload (`serde(default)` keeps
85    /// older results decodable). The SPA Collect page lists / downloads
86    /// these straight from the bucket.
87    #[serde(default, skip_serializing_if = "Option::is_none")]
88    pub collect_object: Option<String>,
89}
90
91/// Synthetic exit code for a run the agent skipped because the
92/// Command's `version` didn't match the `script_current` pin. One of
93/// the reserved synthetic codes (122–127): the agent publishes a
94/// normal [`ExecResult`] so the operator can see *why* nothing ran,
95/// but the script itself never executed — consumers that derive state
96/// from a run's output (e.g. the backend's `check_status` projection)
97/// must treat these as "no new evidence", not as a run (#909).
98pub const EXIT_SKIP_VERSION_PIN: i32 = 124;
99/// Synthetic exit code: `deadline_at` passed before the agent could
100/// fire. See [`EXIT_SKIP_VERSION_PIN`] for the shared contract.
101pub const EXIT_SKIP_DEADLINE: i32 = 125;
102/// Synthetic exit code: the script is revoked in
103/// `BUCKET_SCRIPT_STATUS`. See [`EXIT_SKIP_VERSION_PIN`].
104pub const EXIT_SKIP_REVOKED: i32 = 126;
105/// Synthetic exit code: the `staleness.mode: strict` policy suppressed
106/// the fire. See [`EXIT_SKIP_VERSION_PIN`].
107pub const EXIT_SKIP_STALENESS: i32 = 127;
108/// Synthetic exit code: the agent **refused** the command because its
109/// provenance signature did not check out (#1165 stage 3).
110///
111/// Extends the reserved block downwards, because 124–127 was full. It shares
112/// the block's contract — the script never ran, so this is not evidence about
113/// its outcome — but it is not a *skip*: the others mean "policy (or this
114/// OS) said not now", while this one means "this command was not authorised". The specific
115/// code is what carries that distinction; [`is_synthetic_skip`] deliberately
116/// does not, because every consumer of that predicate is asking the narrower
117/// question of whether the script ran.
118///
119/// Emitting a result at all is the point. `kanade run` waits on
120/// `results.<request_id>`, so a refusal that published nothing would be
121/// indistinguishable from an agent that is simply gone — and the most likely
122/// refusal in practice is an operator's own break-glass command going stale on
123/// a host whose clock is wrong, during an incident, with the backend down and
124/// the obs event stuck in the outbox.
125pub const EXIT_REJECTED_UNSIGNED: i32 = 123;
126/// Synthetic exit code: the schedule needs a feature this agent's OS
127/// **cannot evaluate or source** — a `constraints.require` gate with no
128/// sensing on this platform, or a `when.on` event this OS never emits.
129/// Emitted by the agent's local scheduler, which fails closed (the job is
130/// not run and no per-pc completion is recorded, so it still runs once the
131/// gate becomes supported). Grew the reserved block downwards again (123
132/// was the edge); shares the [`EXIT_SKIP_VERSION_PIN`] contract — the
133/// script never ran.
134pub const EXIT_SKIP_UNSUPPORTED: i32 = 122;
135
136/// True when `exit_code` is one of the reserved synthetic codes (122–127) —
137/// the agent published this result *instead of* running the script, so it
138/// carries no evidence about the script's outcome.
139pub fn is_synthetic_skip(exit_code: i32) -> bool {
140    (EXIT_SKIP_UNSUPPORTED..=EXIT_SKIP_STALENESS).contains(&exit_code)
141}
142
143impl ExecResult {
144    /// Return the `result_id` if the agent supplied one (v0.29+
145    /// payloads always do), otherwise derive a stable UUIDv5 from
146    /// `(request_id, pc_id)`. The projector calls this before INSERT
147    /// so legacy payloads still get a non-empty PK, AND so that
148    /// JetStream redeliveries of the same legacy payload hash to the
149    /// same id and dedupe via `ON CONFLICT`. Per-PC fan-out stays
150    /// distinct (different `pc_id` → different hash).
151    pub fn stable_result_id(&self) -> String {
152        if !self.result_id.is_empty() {
153            return self.result_id.clone();
154        }
155        let name = format!(
156            "{LEGACY_RESULT_ID_PREFIX}{}:{}",
157            self.request_id, self.pc_id
158        );
159        Uuid::new_v5(&Uuid::NAMESPACE_OID, name.as_bytes()).to_string()
160    }
161}
162
163#[cfg(test)]
164mod tests {
165    use super::*;
166    use chrono::TimeZone;
167
168    #[test]
169    fn synthetic_skip_covers_exactly_the_reserved_codes() {
170        for code in [
171            EXIT_SKIP_UNSUPPORTED,
172            EXIT_REJECTED_UNSIGNED,
173            EXIT_SKIP_VERSION_PIN,
174            EXIT_SKIP_DEADLINE,
175            EXIT_SKIP_REVOKED,
176            EXIT_SKIP_STALENESS,
177        ] {
178            assert!(is_synthetic_skip(code), "{code} is a reserved skip code");
179        }
180        // The block grew downwards twice: 123 for the signature rejection
181        // (#1165), then 122 for "unsupported on this OS". 121 is the new
182        // edge, and pinning both edges is what makes a future widening a
183        // deliberate act rather than an accident.
184        for code in [0, 1, -1, 121, 128, 255] {
185            assert!(!is_synthetic_skip(code), "{code} is a real exit code");
186        }
187    }
188
189    #[test]
190    fn exec_result_round_trips_through_json() {
191        let t0 = chrono::Utc.with_ymd_and_hms(2026, 5, 16, 0, 0, 0).unwrap();
192        let t1 = chrono::Utc.with_ymd_and_hms(2026, 5, 16, 0, 0, 5).unwrap();
193        let r = ExecResult {
194            result_id: "result-uuid-1".into(),
195            request_id: "req-1".into(),
196            exec_id: Some("exec-uuid-1".into()),
197            parent_result_id: None,
198            pc_id: "pc-01".into(),
199            exit_code: 0,
200            stdout: "hello\n".into(),
201            stderr: String::new(),
202            started_at: t0,
203            finished_at: t1,
204            stdout_object: None,
205            stderr_object: None,
206            manifest_id: Some("inventory-hw".into()),
207            collect_object: None,
208        };
209        let json = serde_json::to_string(&r).unwrap();
210        let back: ExecResult = serde_json::from_str(&json).unwrap();
211        assert_eq!(back.result_id, r.result_id);
212        assert_eq!(back.request_id, r.request_id);
213        assert_eq!(back.exec_id.as_deref(), Some("exec-uuid-1"));
214        assert_eq!(back.exit_code, r.exit_code);
215        assert_eq!(back.stdout, r.stdout);
216        assert_eq!(back.started_at, t0);
217        assert_eq!(back.finished_at, t1);
218        assert_eq!(back.manifest_id.as_deref(), Some("inventory-hw"));
219    }
220
221    #[test]
222    fn exec_result_without_manifest_id_decodes() {
223        // Older agents (pre-0.13) sent ExecResult with no manifest_id field.
224        let json = r#"{
225            "request_id":"r","pc_id":"x","exit_code":0,
226            "stdout":"","stderr":"",
227            "started_at":"2026-05-16T00:00:00Z",
228            "finished_at":"2026-05-16T00:00:00Z"
229        }"#;
230        let r: ExecResult = serde_json::from_str(json).unwrap();
231        assert_eq!(r.manifest_id, None);
232    }
233
234    #[test]
235    fn exec_result_without_result_id_decodes_empty() {
236        // v0.29 / Issue #19: pre-v0.29 agents don't send `result_id`.
237        // `#[serde(default)]` decodes it as the empty string so the
238        // projector can detect "legacy payload" and call
239        // `stable_result_id()` to derive a deterministic PK.
240        let json = r#"{
241            "request_id":"r","pc_id":"x","exit_code":0,
242            "stdout":"","stderr":"",
243            "started_at":"2026-05-16T00:00:00Z",
244            "finished_at":"2026-05-16T00:00:00Z"
245        }"#;
246        let r: ExecResult = serde_json::from_str(json).unwrap();
247        assert_eq!(r.result_id, "");
248        assert!(r.exec_id.is_none());
249    }
250
251    #[test]
252    fn stable_result_id_is_deterministic_for_legacy_payload() {
253        // Gemini #65 medium fix: legacy redeliveries (same request_id +
254        // pc_id) must hash to the SAME result_id so the projector's
255        // ON CONFLICT(result_id) DO NOTHING dedupes — otherwise
256        // `executions.success_count` double-counts on JetStream ack
257        // timeouts.
258        let json = r#"{
259            "request_id":"r","pc_id":"x","exit_code":0,
260            "stdout":"","stderr":"",
261            "started_at":"2026-05-16T00:00:00Z",
262            "finished_at":"2026-05-16T00:00:00Z"
263        }"#;
264        let a: ExecResult = serde_json::from_str(json).unwrap();
265        let b: ExecResult = serde_json::from_str(json).unwrap();
266        assert_eq!(
267            a.stable_result_id(),
268            b.stable_result_id(),
269            "same legacy payload must hash to the same result_id",
270        );
271    }
272
273    #[test]
274    fn stable_result_id_differs_across_pcs_for_broadcast() {
275        // The other half: a broadcast Command published to two PCs
276        // produces two legacy ExecResults sharing one request_id but
277        // with different pc_ids. Each must get its OWN result_id so
278        // both rows persist (the whole point of Issue #19).
279        let json_a = r#"{
280            "request_id":"shared","pc_id":"pc-1","exit_code":0,
281            "stdout":"","stderr":"",
282            "started_at":"2026-05-16T00:00:00Z",
283            "finished_at":"2026-05-16T00:00:00Z"
284        }"#;
285        let json_b = r#"{
286            "request_id":"shared","pc_id":"pc-2","exit_code":0,
287            "stdout":"","stderr":"",
288            "started_at":"2026-05-16T00:00:00Z",
289            "finished_at":"2026-05-16T00:00:00Z"
290        }"#;
291        let a: ExecResult = serde_json::from_str(json_a).unwrap();
292        let b: ExecResult = serde_json::from_str(json_b).unwrap();
293        assert_ne!(
294            a.stable_result_id(),
295            b.stable_result_id(),
296            "different pc_id must produce a different result_id",
297        );
298    }
299
300    #[test]
301    fn stable_result_id_passes_through_explicit_value() {
302        // v0.29 agents always supply result_id; the helper must
303        // return that as-is (no surprise re-hashing).
304        let r = ExecResult {
305            result_id: "agent-minted-uuid".into(),
306            request_id: "r".into(),
307            exec_id: None,
308            parent_result_id: None,
309            pc_id: "x".into(),
310            exit_code: 0,
311            stdout: String::new(),
312            stderr: String::new(),
313            started_at: chrono::Utc.with_ymd_and_hms(2026, 5, 16, 0, 0, 0).unwrap(),
314            finished_at: chrono::Utc.with_ymd_and_hms(2026, 5, 16, 0, 0, 0).unwrap(),
315            stdout_object: None,
316            stderr_object: None,
317            manifest_id: None,
318            collect_object: None,
319        };
320        assert_eq!(r.stable_result_id(), "agent-minted-uuid");
321    }
322
323    #[test]
324    fn exec_result_collect_object_round_trips_and_omits_when_absent() {
325        // #219: collect_object is off the wire when None
326        // (skip_serializing_if) so pre-#219 readers stay compatible...
327        let t0 = chrono::Utc.with_ymd_and_hms(2026, 6, 15, 0, 0, 0).unwrap();
328        let mut r = ExecResult {
329            result_id: "r1".into(),
330            request_id: "req".into(),
331            exec_id: None,
332            parent_result_id: None,
333            pc_id: "PC1".into(),
334            exit_code: 0,
335            stdout: String::new(),
336            stderr: String::new(),
337            started_at: t0,
338            finished_at: t0,
339            stdout_object: None,
340            stderr_object: None,
341            manifest_id: Some("collect-diagnostics".into()),
342            collect_object: None,
343        };
344        let json = serde_json::to_string(&r).unwrap();
345        assert!(
346            !json.contains("collect_object"),
347            "collect_object must be absent when None: {json}"
348        );
349        // ...and a set key survives the round-trip.
350        r.collect_object = Some("PC1/collect-diagnostics/20260615T000000Z.zip".into());
351        let back: ExecResult = serde_json::from_str(&serde_json::to_string(&r).unwrap()).unwrap();
352        assert_eq!(
353            back.collect_object.as_deref(),
354            Some("PC1/collect-diagnostics/20260615T000000Z.zip"),
355        );
356    }
357
358    #[test]
359    fn exec_result_parent_result_id_round_trips_and_omits_when_absent() {
360        // #955: a finalize row carries `parent_result_id`; ordinary runs
361        // leave it None, and it stays off the wire so pre-#955 readers
362        // are unaffected (skip_serializing_if).
363        let t0 = chrono::Utc.with_ymd_and_hms(2026, 7, 4, 0, 0, 0).unwrap();
364        let mut r = ExecResult {
365            result_id: "fin-1".into(),
366            request_id: "req__finalize".into(),
367            exec_id: None,
368            parent_result_id: None,
369            pc_id: "PC1".into(),
370            exit_code: 0,
371            stdout: String::new(),
372            stderr: String::new(),
373            started_at: t0,
374            finished_at: t0,
375            stdout_object: None,
376            stderr_object: None,
377            manifest_id: Some("screenshot-collect__finalize".into()),
378            collect_object: None,
379        };
380        let json = serde_json::to_string(&r).unwrap();
381        assert!(
382            !json.contains("parent_result_id"),
383            "parent_result_id must be absent when None: {json}"
384        );
385        r.parent_result_id = Some("parent-run-uuid".into());
386        let back: ExecResult = serde_json::from_str(&serde_json::to_string(&r).unwrap()).unwrap();
387        assert_eq!(back.parent_result_id.as_deref(), Some("parent-run-uuid"));
388    }
389}