Skip to main content

taskfleet_core/
schema.rs

1//! On-disk state schema types per `design.md` §1.
2
3use chrono::{DateTime, Utc};
4use serde::{Deserialize, Serialize};
5use serde_json::{Map, Value};
6
7/// The current state-on-disk schema version this crate writes.
8pub const STATE_SCHEMA_VERSION: u32 = 1;
9
10/// All state-schema versions this crate can read.
11pub const SUPPORTED_STATE_SCHEMAS: &[u32] = &[1];
12
13/// Crockford base32 alphabet in lowercase (excludes `i`, `l`, `o`, `u`). The
14/// charset for the bare ULID of a [`RunId`].
15const CROCKFORD_LOWER: &[u8] = b"0123456789abcdefghjkmnpqrstvwxyz";
16
17/// True iff every byte of `s` is a lowercase Crockford base32 character.
18fn all_crockford_lower(s: &str) -> bool {
19    s.bytes().all(|b| CROCKFORD_LOWER.contains(&b))
20}
21
22/// True iff `s` is a syntactically valid (possibly partial) prefix of a
23/// [`RunId`]: non-empty, no longer than a full ULID, every character a lowercase
24/// Crockford base32 digit, and a first character within ULID's `0..=7`
25/// timestamp bound. Used by the CLI to resolve an unambiguous run-id prefix
26/// (like `git`) — a value failing this is a malformed argument (`invalid_run_id`),
27/// not a legitimate-but-unknown prefix. The first-char bound is enforced because
28/// no valid `RunId` can begin outside `0..=7`, so an `8…`/`9…` prefix is
29/// impossible rather than merely absent — reporting it as malformed keeps the
30/// error class honest and consistent with how [`RunId::parse_str`] rejects a
31/// full-length id with the same leading digit.
32pub fn is_run_id_prefix(s: &str) -> bool {
33    !s.is_empty()
34        && s.len() <= RunId::LEN
35        && all_crockford_lower(s)
36        && matches!(s.as_bytes().first(), Some(b'0'..=b'7'))
37}
38
39/// Error returned when a typed identifier fails parse-time validation.
40///
41/// Every [`RunId`] and [`NodeId`] is
42/// constructed only through its `parse_str` constructor (or the equivalent
43/// validating `Deserialize`), so any value that reaches a path helper has
44/// already been checked for prefix, charset, and length. This is the
45/// path-traversal guard: a raw id containing `/`, `..`, or a leading dot can
46/// never be turned into one of these newtypes, so it can never name a file
47/// outside the run directory.
48#[derive(Debug, Clone, PartialEq, Eq, thiserror::Error)]
49pub enum IdValidationError {
50    /// The value carried the right prefix (or needs none) but its body had the
51    /// wrong length or used characters outside the permitted charset.
52    #[error("invalid {kind} id {value:?}: expected {expected}")]
53    InvalidFormat {
54        /// Which id type rejected the value (`run`, `node`).
55        kind: &'static str,
56        /// The offending raw value.
57        value: String,
58        /// Human-readable description of the accepted shape (e.g. `n-NNNN`).
59        expected: &'static str,
60    },
61    /// The value did not start with the id type's required prefix (`n-`).
62    #[error("invalid {kind} id: wrong prefix, expected {expected}")]
63    WrongPrefix {
64        /// Which id type rejected the value.
65        kind: &'static str,
66        /// Human-readable description of the accepted shape.
67        expected: &'static str,
68    },
69}
70
71impl IdValidationError {
72    /// The id type that rejected the value (`run`, `node`).
73    pub fn kind(&self) -> &'static str {
74        match self {
75            Self::InvalidFormat { kind, .. } | Self::WrongPrefix { kind, .. } => kind,
76        }
77    }
78
79    /// The accepted-shape hint, suitable for the `expected` field of a CLI
80    /// error envelope.
81    pub fn expected(&self) -> &'static str {
82        match self {
83            Self::InvalidFormat { expected, .. } | Self::WrongPrefix { expected, .. } => expected,
84        }
85    }
86}
87
88/// Generate the shared trait surface for a validated id newtype: `as_str`,
89/// `FromStr`, `Display`, `Debug`, `Ord` / `PartialOrd` (lexicographic over the
90/// inner string), `Serialize` (as the bare string), and a validating
91/// `Deserialize` (delegates to `parse_str`, so reading an old file with a
92/// malformed id fails loudly rather than silently widening the type). Each
93/// newtype supplies its own `parse_str` in a separate `impl` block.
94///
95/// `Ord` / `PartialOrd` are derived, so they forward to the inner `String`'s
96/// ordering — i.e. plain `&str` byte comparison. For the fixed-width ULID form
97/// ([`RunId`]) this preserves the natural time ordering ULIDs encode in their
98/// lexical sort.
99///
100/// CAVEAT — this ordering is lexical, *not* numeric or semantic: [`NodeId`] is
101/// `n-` + a variable-width number, so once the counter grows a digit the byte
102/// order diverges from the numeric order: `n-10000 < n-9999`. Do not sort
103/// `NodeId`s expecting ascending node number; parse the body if you need that.
104///
105/// The trait is provided for `BTreeMap`/`BTreeSet` keys and stable sorts.
106macro_rules! id_newtype {
107    ($(#[$m:meta])* $name:ident) => {
108        $(#[$m])*
109        #[derive(Clone, PartialEq, Eq, Hash, PartialOrd, Ord)]
110        pub struct $name(String);
111
112        impl $name {
113            /// The validated id as a string slice. There is no mutable or
114            /// owned-`String` accessor by design: the inner value can never be
115            /// mutated into an unvalidated state after construction.
116            pub fn as_str(&self) -> &str {
117                &self.0
118            }
119        }
120
121        impl std::str::FromStr for $name {
122            type Err = IdValidationError;
123
124            /// Parse via the newtype's own `parse_str`; lets callers use the
125            /// `str::parse` / `FromStr` ecosystem (`s.parse::<RunId>()?`).
126            fn from_str(s: &str) -> Result<Self, Self::Err> {
127                Self::parse_str(s)
128            }
129        }
130
131        impl std::fmt::Display for $name {
132            fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
133                f.write_str(&self.0)
134            }
135        }
136
137        impl std::fmt::Debug for $name {
138            fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
139                write!(f, "{}({:?})", stringify!($name), self.0)
140            }
141        }
142
143        impl serde::Serialize for $name {
144            fn serialize<S: serde::Serializer>(&self, s: S) -> Result<S::Ok, S::Error> {
145                s.serialize_str(&self.0)
146            }
147        }
148
149        impl<'de> serde::Deserialize<'de> for $name {
150            fn deserialize<D: serde::Deserializer<'de>>(d: D) -> Result<Self, D::Error> {
151                let s = String::deserialize(d)?;
152                Self::parse_str(&s).map_err(serde::de::Error::custom)
153            }
154        }
155    };
156}
157
158id_newtype! {
159    /// A validated run identifier: a lowercase ULID (26 Crockford base32
160    /// characters whose first character keeps the encoded timestamp within
161    /// ULID's 48-bit range). Mirrors what [`crate::new_run_id`] emits.
162    RunId
163}
164
165impl RunId {
166    /// Accepted-shape hint shared by every rejection.
167    const EXPECTED: &'static str = "26-char lowercase Crockford base32 ULID";
168    /// Canonical length of a ULID in Crockford base32. Public so CLI-side prefix
169    /// resolution can branch on "full id vs. prefix" without mirroring the
170    /// constant (which would silently drift if the id shape ever changed).
171    pub const LEN: usize = 26;
172
173    /// Parse and validate a `run_id`. Accepts only the 26-character lowercase
174    /// ULID shape; rejects wrong length, non-Crockford characters, and a first
175    /// character outside `0..=7` (which would overflow ULID's 48-bit timestamp).
176    pub fn parse_str(s: &str) -> Result<Self, IdValidationError> {
177        let reject = || IdValidationError::InvalidFormat {
178            kind: "run",
179            value: s.to_string(),
180            expected: Self::EXPECTED,
181        };
182        if s.len() != Self::LEN || !all_crockford_lower(s) {
183            return Err(reject());
184        }
185        // The first base32 char carries the top 5 bits of the 128-bit ULID;
186        // the 48-bit timestamp cannot overflow only if it is in `0..=7`.
187        if !(b'0'..=b'7').contains(&s.as_bytes()[0]) {
188            return Err(reject());
189        }
190        Ok(Self(s.to_string()))
191    }
192}
193
194id_newtype! {
195    /// A validated node identifier: `n-` followed by 4 or more ASCII digits
196    /// (e.g. `n-0001`). Mirrors what [`crate::format_node_id`] emits.
197    NodeId
198}
199
200impl NodeId {
201    /// Accepted-shape hint shared by every rejection.
202    const EXPECTED: &'static str = "n-NNNN (n- followed by 4-10 ASCII digits)";
203
204    /// Parse and validate a `node_id`. Requires the `n-` prefix followed by
205    /// 4 to 10 ASCII digits; rejects anything else (wrong prefix, too few or
206    /// too many digits, non-digit body). The 10-digit ceiling covers the full
207    /// `u32` counter range [`crate::format_node_id`] draws from while bounding
208    /// the filename length (a defense against `ENAMETOOLONG` from a forged id).
209    pub fn parse_str(s: &str) -> Result<Self, IdValidationError> {
210        let body = s.strip_prefix("n-").ok_or(IdValidationError::WrongPrefix {
211            kind: "node",
212            expected: Self::EXPECTED,
213        })?;
214        if (4..=10).contains(&body.len()) && body.bytes().all(|b| b.is_ascii_digit()) {
215            Ok(Self(s.to_string()))
216        } else {
217            Err(IdValidationError::InvalidFormat {
218                kind: "node",
219                value: s.to_string(),
220                expected: Self::EXPECTED,
221            })
222        }
223    }
224}
225
226/// The run/node kind enum (design.md §1.2).
227///
228/// The 0.2 subtractive cut removed the `code`, `orchestrate`, `orchestrated`,
229/// `bugfix`, and `make-skill` kinds (the interactive + DAG-driver topologies and
230/// the two phantom variants that were behaviourally `Spinoff`). The surviving
231/// kinds are all autonomous. [`Kind::Unknown`] is a read-only catch-all so a
232/// legacy on-disk run recorded under a since-removed kind still deserializes —
233/// `doctor` / `run list` report it, never delete it (ADR §D7).
234#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize)]
235#[serde(rename_all = "kebab-case")]
236pub enum Kind {
237    /// Autonomous fire-and-forget task that merges itself back (`/worktree-spinoff`).
238    Spinoff,
239    /// Autonomous multi-source research worktree (`/worktree-research`).
240    Research,
241    /// Drives one architectural decision to an ADR (`/worktree-technical-decision`).
242    TechnicalDecision,
243    /// Parallel fan-out of many identical units (`/fan-out`).
244    FanOut,
245    /// A kind this build no longer models — a legacy run recorded on disk under a
246    /// kind removed in the 0.2 cut (`code` / `orchestrate` / `orchestrated` /
247    /// `bugfix` / `make-skill`), or any future/unknown wire value. Read-only:
248    /// `#[serde(other)]` maps every unrecognized kind here so `doctor` / `run
249    /// list` can still surface such a run rather than faulting on it (ADR §D7).
250    /// It is NEVER a creatable kind — it is absent from [`Kind::WIRE_NAMES`], so
251    /// no CLI surface or report validator accepts it as input.
252    #[serde(other)]
253    Unknown,
254}
255
256impl Kind {
257    /// The kebab-case wire name for this kind — the same string serde
258    /// (de)serializes via `rename_all = "kebab-case"`.
259    ///
260    /// The exhaustive `match` is deliberate: adding a `Kind` variant fails
261    /// to compile until its wire name is listed here, so [`Kind::WIRE_NAMES`]
262    /// and any caller that advertises the accepted kinds (e.g. the report
263    /// validator's `expected` hint) cannot silently drift from the enum.
264    #[must_use]
265    pub const fn wire_name(self) -> &'static str {
266        match self {
267            Kind::Spinoff => "spinoff",
268            Kind::Research => "research",
269            Kind::TechnicalDecision => "technical-decision",
270            Kind::FanOut => "fan-out",
271            Kind::Unknown => "unknown",
272        }
273    }
274
275    /// Every *creatable* kind's kebab-case wire name, in declaration order.
276    /// Single source of truth for "the set of accepted kinds" — see
277    /// [`Kind::wire_name`]. Excludes [`Kind::Unknown`], which is a read-only
278    /// catch-all, never a valid input.
279    pub const WIRE_NAMES: &'static [&'static str] = &[
280        Kind::Spinoff.wire_name(),
281        Kind::Research.wire_name(),
282        Kind::TechnicalDecision.wire_name(),
283        Kind::FanOut.wire_name(),
284    ];
285
286    /// Default how-run [`Lifecycle`] for a kind — the value a run gets when
287    /// created WITHOUT `--interactive`. Every kind defaults to autonomous; the 0.2
288    /// cut removed the `code` kind that used to imply interactivity, so
289    /// interactivity is no longer kind-derived — it is the explicit `--interactive`
290    /// flag ([`Lifecycle`] docs, design.md §2/§6). This method only seeds the
291    /// default; it must NOT be read as "this kind is (non-)interactive".
292    /// [`Kind::Unknown`] (a legacy on-disk run) reads as autonomous too; it is
293    /// never freshly supervised, so the value only ever feeds read-only display.
294    pub fn lifecycle(self) -> Lifecycle {
295        match self {
296            Kind::Spinoff
297            | Kind::Research
298            | Kind::TechnicalDecision
299            | Kind::FanOut
300            | Kind::Unknown => Lifecycle::Autonomous,
301        }
302    }
303
304    /// Whether this kind is a **top-level, single-node, autonomous worker** —
305    /// one detached agent that materializes its own worktree and self-merges,
306    /// with no children and no parent DAG driving it. These are exactly the
307    /// kinds eligible for the supervisor's bounded auto-retry on an empty-handed
308    /// `agent-died` (issue `autoretry-agent-died-worker`).
309    ///
310    /// Excludes `FanOut` (a multi-unit driver — its driver node has no agent of
311    /// its own) and [`Kind::Unknown`] (a legacy on-disk run, never freshly
312    /// supervised).
313    ///
314    /// The exhaustive `match` fails to compile when a new `Kind` is added, forcing
315    /// a deliberate eligibility decision rather than a silent default.
316    #[must_use]
317    pub fn is_autonomous_single_node_worker(self) -> bool {
318        match self {
319            Kind::Spinoff | Kind::Research | Kind::TechnicalDecision => true,
320            Kind::FanOut | Kind::Unknown => false,
321        }
322    }
323}
324
325/// How a run is driven — its **how-run** state (design.md §2, §6).
326///
327/// This is an **explicit told fact**, set once at `run create` from the
328/// `--interactive` flag and never transitioned. It is deliberately NOT derived
329/// from [`Kind`]: the 0.2 cut removed the `code` kind that used to carry
330/// interactivity accidentally, and interactivity is now orthogonal to topology —
331/// *any* run can be marked interactive (`told, not guessed`, `target-state-0.2.md
332/// §2`/§4). Do not reintroduce a `Kind`-derived inference; `Kind::lifecycle`
333/// exists only to seed the default for a run created without the flag.
334///
335/// `Lifecycle` is a *category*, not a progress signal — an agent tracking
336/// completion polls `manifest.status` (`Pending | Running | Done | Failed |
337/// Cancelled`), NEVER `lifecycle`, whose value never changes (state-integrity
338/// invariant 4).
339#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
340#[serde(rename_all = "kebab-case")]
341pub enum Lifecycle {
342    /// Agent runs to completion unattended; the supervisor adjudicates exit
343    /// (the told `worker.exited` fact, then the residual crash backstop).
344    Autonomous,
345    /// Human-driven: the supervisor **never** auto-terminalizes or auto-tears-down
346    /// from a dead pid or a worker exit — it waits for an explicit `run merge`
347    /// (→ teardown) or `run cancel`. The human owns the whole lifecycle
348    /// (design.md §6).
349    Interactive,
350}
351
352impl Lifecycle {
353    /// True for [`Lifecycle::Interactive`] — the human-driven, supervisor-hands-off
354    /// how-run state. The single predicate the supervisor consults to suppress its
355    /// automatic terminalization/teardown machinery (design.md §6).
356    #[must_use]
357    pub fn is_interactive(self) -> bool {
358        matches!(self, Lifecycle::Interactive)
359    }
360}
361
362/// Who launches and owns the agent. This is independent of lifecycle: an
363/// interactive Taskfleet worker is still a Taskfleet-owned worker.
364#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Serialize, Deserialize)]
365#[serde(rename_all = "kebab-case")]
366pub enum AgentOwner {
367    #[default]
368    /// Taskfleet launches and supervises its own worker.
369    Taskfleet,
370    /// An external process owns the agent; Taskfleet owns only the Git run.
371    Caller,
372}
373
374impl AgentOwner {
375    /// Keep the historical projection shape byte-compatible for normal runs.
376    /// Serde's `skip_serializing_if` passes a reference even for `Copy` types.
377    #[allow(clippy::trivially_copy_pass_by_ref)]
378    pub fn is_taskfleet(&self) -> bool {
379        *self == Self::Taskfleet
380    }
381}
382
383/// Run/node status (design.md §1.2).
384///
385/// `Done`, `Failed`, and `Cancelled` are **terminal**: once a run or node
386/// reaches one of them its `status` must never change again. The reducer
387/// enforces this — `apply_run_status`, `apply_node_status`, and
388/// `apply_node_report` are all no-ops once [`Status::is_terminal`] holds — so
389/// a late-arriving event (e.g. an agent success report racing a `run cancel`)
390/// cannot resurrect a settled state. Only the `status` field is frozen;
391/// other projection fields may still be mutated by non-status events.
392#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
393#[serde(rename_all = "kebab-case")]
394pub enum Status {
395    /// Created but not yet started.
396    Pending,
397    /// Actively executing.
398    Running,
399    /// Stalled awaiting input (e.g. an open discussion).
400    Blocked,
401    /// Completed successfully (terminal).
402    Done,
403    /// Completed with failure (terminal).
404    Failed,
405    /// Terminated before completion by an operator or parent (terminal).
406    Cancelled,
407}
408
409impl Status {
410    /// True for the terminal states `Done | Failed | Cancelled`. A run or
411    /// node in a terminal state is settled: the reducer treats any further
412    /// *status* transition as a no-op. "Settled" applies to `status` only —
413    /// non-status projection fields (e.g. `Node::children` via
414    /// `child.spawned`, or manifest counters) can still change.
415    pub fn is_terminal(self) -> bool {
416        matches!(self, Status::Done | Status::Failed | Status::Cancelled)
417    }
418}
419
420/// Aggregate a set of node statuses into the run's rolled-up terminal status,
421/// or `None` when the run is not yet complete.
422///
423/// The single, shared roll-up rule — used both by the supervisor's per-tick
424/// `rollup_status` and by [`cancel_node`](crate::cancel_node)'s in-lock
425/// self-roll-up (so the two can never diverge). A **three-way** classification
426/// (design §2.5, "rollup terminalizes the run cancelled/done/failed once every
427/// node is terminal"):
428///
429/// - `None` if the set is empty (a freshly-created run must not vacuously
430///   complete) or if ANY node is still live (`Pending`/`Running`/`Blocked`);
431/// - `Some(Status::Failed)` if any node genuinely `Failed` (a real failure
432///   dominates the batch outcome);
433/// - `Some(Status::Cancelled)` if no node failed but at least one was
434///   `Cancelled` (a deliberate per-node/whole-run cancel — nothing failed, but
435///   the batch did not fully complete; branch-preserving work is untouched);
436/// - `Some(Status::Done)` when every node is `Done`.
437pub fn aggregate_terminal_status<I>(statuses: I) -> Option<Status>
438where
439    I: IntoIterator<Item = Status>,
440{
441    let mut any = false;
442    let mut any_failed = false;
443    let mut any_cancelled = false;
444    for s in statuses {
445        any = true;
446        match s {
447            Status::Done => {}
448            Status::Failed => any_failed = true,
449            Status::Cancelled => any_cancelled = true,
450            // Any live node means the run is not done yet.
451            Status::Pending | Status::Running | Status::Blocked => return None,
452        }
453    }
454    if !any {
455        return None;
456    }
457    Some(if any_failed {
458        Status::Failed
459    } else if any_cancelled {
460        Status::Cancelled
461    } else {
462        Status::Done
463    })
464}
465
466/// Compact profile resolution recorded at create time.
467#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
468pub struct AgentSelection {
469    /// Compact selection schema version (currently 1).
470    pub schema_version: u32,
471    /// Requested and selected user profile name.
472    pub profile: String,
473    /// Precedence layer that supplied the request.
474    pub selection_source: String,
475    /// Explicit create-time interaction mode.
476    pub interaction: String,
477    /// Declared profile capability tier.
478    pub capability: String,
479    /// Declared profile residency class.
480    pub residency: String,
481    /// Legacy harness alias requested at the winning layer, when applicable.
482    #[serde(default, skip_serializing_if = "Option::is_none")]
483    pub requested_harness: Option<String>,
484    /// First statically eligible candidate.
485    pub selected: SelectedAgentCandidate,
486    /// Earlier candidates and their single deterministic skip reasons.
487    #[serde(default)]
488    pub fallback: Vec<SkippedAgentCandidate>,
489}
490
491/// Exact selected candidate pinned for launch and retry.
492#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
493pub struct SelectedAgentCandidate {
494    /// Zero-based position in the profile's ordered candidate list.
495    pub candidate_index: u8,
496    /// Selected harness (`pi` or `claude`).
497    pub harness: String,
498    /// Exact user-owned argv; never a shell string.
499    pub command: Vec<String>,
500    /// Declared telemetry adapter protocol, when configured.
501    #[serde(default, skip_serializing_if = "Option::is_none")]
502    pub telemetry: Option<String>,
503}
504
505impl SelectedAgentCandidate {
506    /// Whether recorded policy configures the public worker telemetry v1
507    /// adapter. This is configuration support, not runtime attestation or a
508    /// claim that any sample has arrived.
509    #[must_use]
510    pub fn supports_worker_telemetry_v1(&self) -> bool {
511        self.harness == "pi" && self.telemetry.as_deref() == Some("worker-v1")
512    }
513}
514
515/// One rejected candidate and its first applicable reason.
516#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
517pub struct SkippedAgentCandidate {
518    /// Zero-based position in the profile's ordered candidate list.
519    pub candidate_index: u8,
520    /// Candidate harness.
521    pub harness: String,
522    /// Stable skip reason code.
523    pub reason: String,
524}
525
526impl AgentSelection {
527    /// Validate semantic bounds and closed vocabularies at the durable event
528    /// boundary, independently of the user-config parser that constructed it.
529    pub fn validate(&self) -> Result<(), String> {
530        if self.schema_version != 1 {
531            return Err(format!(
532                "unsupported selection schema_version {}",
533                self.schema_version
534            ));
535        }
536        let name = self.profile.as_bytes();
537        if name.is_empty()
538            || name.len() > 63
539            || !name[0].is_ascii_lowercase()
540            || !name
541                .iter()
542                .all(|b| b.is_ascii_lowercase() || b.is_ascii_digit() || *b == b'-')
543            || self.profile.ends_with('-')
544            || self.profile.contains("--")
545        {
546            return Err("invalid profile name".into());
547        }
548        if !matches!(
549            self.selection_source.as_str(),
550            "cli"
551                | "environment"
552                | "repository-per-kind"
553                | "user-per-kind"
554                | "repository-default"
555                | "user-default"
556                | "builtin-harness"
557        ) {
558            return Err("invalid selection_source".into());
559        }
560        if !matches!(
561            self.interaction.as_str(),
562            "autonomous" | "explicit-interactive"
563        ) || !matches!(
564            self.capability.as_str(),
565            "fast" | "capable" | "ultra-capable"
566        ) || !matches!(self.residency.as_str(), "local" | "remote")
567        {
568            return Err("invalid interaction, capability, or residency".into());
569        }
570        validate_selected_candidate(&self.selected)?;
571        if self.interaction == "autonomous"
572            && (self.selected.harness != "pi"
573                || self.selected.telemetry.as_deref() != Some("worker-v1"))
574        {
575            return Err("autonomous selection requires pi with worker-v1 telemetry".into());
576        }
577        if self.selected.candidate_index >= 8
578            || self.fallback.len() != usize::from(self.selected.candidate_index)
579        {
580            return Err("candidate index/count exceeds profile bound".into());
581        }
582        let mut prior = None;
583        for (expected_index, skipped) in self.fallback.iter().enumerate() {
584            if usize::from(skipped.candidate_index) != expected_index
585                || prior.is_some_and(|value| skipped.candidate_index <= value)
586                || !matches!(skipped.harness.as_str(), "pi" | "claude")
587                || !matches!(
588                    skipped.reason.as_str(),
589                    "executable_missing"
590                        | "autonomous_harness_unsupported"
591                        | "telemetry_unsupported"
592                )
593            {
594                return Err("invalid fallback candidate index, harness, or reason".into());
595            }
596            if skipped.reason == "autonomous_harness_unsupported"
597                && (self.interaction != "autonomous" || skipped.harness == "pi")
598            {
599                return Err("inconsistent autonomous harness skip reason".into());
600            }
601            if skipped.reason == "telemetry_unsupported"
602                && (self.interaction != "autonomous" || skipped.harness != "pi")
603            {
604                return Err("inconsistent telemetry skip reason".into());
605            }
606            prior = Some(skipped.candidate_index);
607        }
608        Ok(())
609    }
610}
611
612fn validate_selected_candidate(candidate: &SelectedAgentCandidate) -> Result<(), String> {
613    if !matches!(candidate.harness.as_str(), "pi" | "claude")
614        || candidate.command.is_empty()
615        || candidate.command.len() > 32
616        || candidate
617            .command
618            .iter()
619            .any(|arg| arg.is_empty() || arg.len() > 4096 || arg.contains('\0'))
620        || candidate.command.iter().map(String::len).sum::<usize>() > 16_384
621    {
622        return Err("invalid selected harness or command".into());
623    }
624    match (candidate.harness.as_str(), candidate.telemetry.as_deref()) {
625        (_, None) | ("pi", Some("worker-v1")) => Ok(()),
626        _ => Err("invalid selected telemetry declaration".into()),
627    }
628}
629
630/// `manifest.json` (design.md §1.2).
631#[derive(Debug, Clone, Serialize, Deserialize)]
632pub struct Manifest {
633    /// State-schema version this file was written with.
634    pub schema_version: u32,
635    /// Watermark: the highest event `seq` whose projection fold is durably
636    /// committed. Events in `events.jsonl` with `seq > applied_seq` are
637    /// *unapplied tail* events — replayed into the projections on the next
638    /// lock acquisition before any new append (see
639    /// [`crate::events::append_and_apply_event`]). This is what makes
640    /// append-then-apply atomic across a reducer crash: the event log can run
641    /// ahead of the projections, but the gap is always healed before the next
642    /// writer observes stale state.
643    ///
644    /// `#[serde(default)]` so a legacy `manifest.json` written before this
645    /// field existed deserializes with `applied_seq = 0`. Such a manifest
646    /// self-migrates on its next write: the catch-up replay re-folds the whole
647    /// log — every event a no-op, because legacy state was already projected
648    /// synchronously under the old append-then-apply path — and advances the
649    /// watermark to `last_seq`. No separate migration pass or schema bump is
650    /// required (the field is purely additive to a derived-cache file).
651    #[serde(default)]
652    pub applied_seq: u64,
653    /// Unique run identifier (ULID). Validated on read.
654    pub run_id: RunId,
655    /// Kind of work this run performs.
656    pub kind: Kind,
657    /// How-run state (autonomous vs interactive), set once at `run create` from
658    /// the explicit `--interactive` flag — never transitioned. See [`Lifecycle`].
659    pub lifecycle: Lifecycle,
660    /// Absent on historic projections, which always launched Taskfleet workers.
661    #[serde(default, skip_serializing_if = "AgentOwner::is_taskfleet")]
662    pub agent_owner: AgentOwner,
663    /// Caller settlement admission fence; absent for historic/ordinary runs.
664    #[serde(default, skip_serializing_if = "Option::is_none")]
665    pub caller_settlement_intent: Option<CallerSettlementIntent>,
666    /// Human-readable run title.
667    pub title: String,
668    /// Current aggregate run status.
669    pub status: Status,
670    /// When the run was created.
671    pub created_at: DateTime<Utc>,
672    /// When the manifest was last modified.
673    pub updated_at: DateTime<Utc>,
674    /// Source repository the run operates on, if any.
675    pub source_repo: Option<String>,
676    /// Branch the run was started from, if any.
677    pub source_branch: Option<String>,
678    /// Root directory under which this run's worktrees live, if any.
679    pub worktree_root: Option<String>,
680    /// tmux session taskfleet created to host this run's headless windows
681    /// (`--headless` / `--tmux-session <name>`), if any. `None` for a foreground
682    /// run whose window lives in the user's own session — that session is never
683    /// a teardown target. When set, the supervisor kills this session once its
684    /// last taskfleet-owned window is torn down and only the synthetic
685    /// bootstrap shell window remains, so an empty headless session is not left
686    /// behind (issue `headless-tmux-session-not-torn-down`). `#[serde(default)]`
687    /// keeps a manifest written before this field existed readable.
688    #[serde(default)]
689    pub managed_tmux_session: Option<String>,
690    /// Create-time terminal-window retention policy. `None` preserves the
691    /// historical immediate-window/session cleanup behavior. The policy is
692    /// recorded so supervisors and timer maintenance never reinterpret a run
693    /// through later config drift.
694    #[serde(default, skip_serializing_if = "Option::is_none")]
695    pub tmux_retention: Option<Box<TmuxRetentionPolicy>>,
696    /// Completion-notification command registered at `run create --notify`,
697    /// if any. When the run reaches a terminal state (`done | failed |
698    /// cancelled`) the supervisor runs this command (at-least-once, deduped on a
699    /// durable `run.notified` marker event — the healthy path fires once, a
700    /// crash between firing and recording may re-fire) with `TASKFLEET_RUN_ID` /
701    /// `TASKFLEET_STATUS` / `TASKFLEET_SUMMARY` (and `TASKFLEET_RUN_KIND` / `TASKFLEET_RUN_TITLE`)
702    /// in its environment, BEFORE teardown removes the worktree/window. This is
703    /// how a spawning session learns of completion without polling (issue
704    /// `no-completion-notification-to-parent`). `None` for a run created without
705    /// `--notify`; `#[serde(default)]` keeps a manifest written before this
706    /// field existed readable.
707    #[serde(default)]
708    pub notify_cmd: Option<String>,
709    /// The agent runtime selected for this run's worker
710    /// (`claude` | `pi`), resolved at `run create`
711    /// via the flag > env > config > default precedence and recorded here as
712    /// provenance. This is the *selected* harness — recorded before the worker is
713    /// spawned, so it reflects intent even if the spawn later fails. `None` for a
714    /// manifest written before this field existed
715    /// (`#[serde(default)]`) — such legacy runs predate harness selection and
716    /// were all `claude`. Surfaced on `run show` / `run list --json`.
717    #[serde(default)]
718    pub harness: Option<String>,
719    /// Compact create-time profile resolution. `None` keeps manifests written
720    /// before profile selection readable without inventing requested/selected
721    /// history. Retry and read paths consume this recorded value; they never
722    /// re-resolve current configuration.
723    #[serde(default)]
724    pub agent_selection: Option<AgentSelection>,
725    /// Number of nodes created in this run (denormalized counter).
726    pub node_count: u32,
727    /// Run that spawned this run, if it is itself a child.
728    pub parent_run_id: Option<RunId>,
729    /// Node in the parent run that spawned this run, if any.
730    pub parent_node_id: Option<NodeId>,
731}
732
733/// `(child_run_id, child_node_id)` pointer recorded in `Node::children`.
734#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
735pub struct ChildRef {
736    /// Run id of the spawned child. Validated on read.
737    pub run_id: RunId,
738    /// Node id within the child run. Validated on read.
739    pub node_id: NodeId,
740}
741
742/// State of durable native worker evidence capture.
743#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq)]
744#[serde(rename_all = "snake_case")]
745pub enum EvidenceStatus {
746    /// Capture has not completed yet.
747    Pending,
748    /// The latest capture attempt failed and cleanup remains vetoed.
749    Failed,
750    /// Every durable artifact was synced and recorded.
751    Complete,
752}
753
754impl std::fmt::Display for EvidenceStatus {
755    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
756        f.write_str(match self {
757            Self::Pending => "pending",
758            Self::Failed => "failed",
759            Self::Complete => "complete",
760        })
761    }
762}
763
764/// Durable native Pi session and terminal evidence for one worker attempt.
765///
766/// This is folded onto the node projection from append-only events. The live
767/// source is state-root-relative; completed artifact paths are run-relative.
768#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
769pub struct WorkerEvidence {
770    /// Worker attempt this evidence belongs to.
771    pub attempt: u32,
772    /// Exact Pi session identifier assigned before the candidate starts.
773    pub session_id: String,
774    /// Original cwd stored in the native transcript header.
775    pub original_cwd: String,
776    /// State-root-relative live native transcript path.
777    pub live_session_path: String,
778    /// Typed capture state.
779    pub status: EvidenceStatus,
780    /// Run-relative byte-identical archived transcript, once complete.
781    #[serde(default, skip_serializing_if = "Option::is_none")]
782    pub transcript_path: Option<String>,
783    /// Run-relative resume copy whose header names a surviving cwd.
784    #[serde(default, skip_serializing_if = "Option::is_none")]
785    pub resume_path: Option<String>,
786    /// Run-relative final tmux pane snapshot.
787    #[serde(default, skip_serializing_if = "Option::is_none")]
788    pub pane_path: Option<String>,
789    /// Run-relative exact terminal report JSON.
790    #[serde(default, skip_serializing_if = "Option::is_none")]
791    pub report_path: Option<String>,
792    /// SHA-256 of the archived original transcript.
793    #[serde(default, skip_serializing_if = "Option::is_none")]
794    pub transcript_sha256: Option<String>,
795    /// Explicit capture failure detail. Never present on complete evidence.
796    #[serde(default, skip_serializing_if = "Option::is_none")]
797    pub error: Option<String>,
798}
799
800/// Opt-in policy recorded on a run placed in a persistent tmux session.
801#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
802pub struct TmuxRetentionPolicy {
803    /// Keep the selected session after the last worker completes.
804    pub persistent: bool,
805    /// Maximum age of a completed inert display.
806    pub completed_window_ttl_secs: u64,
807    /// Maximum completed inert displays retained in this session.
808    pub completed_window_max: u32,
809}
810
811/// Durable identity of one inert completed-worker display.
812#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
813pub struct RetainedDisplay {
814    /// Worker attempt represented by this display.
815    pub attempt: u32,
816    /// Exact socket path observed at worker spawn.
817    pub socket: String,
818    /// Exact session name observed at worker spawn.
819    pub session: String,
820    /// Window created for the inert display.
821    pub window_id: String,
822    /// Its sole dead pane.
823    pub pane_id: String,
824    /// tmux server PID observed at worker spawn and retention creation.
825    pub server_pid: u32,
826    /// OS process start identity for `server_pid`, guarding PID/socket reuse.
827    pub server_pid_start_secs: u64,
828    /// Runtime-owned server marker observed at spawn (empty when unset).
829    pub server_marker: String,
830    /// Exact Taskfleet ownership marker stored as a tmux window option.
831    pub ownership_marker: String,
832    /// Time the terminal display became durable.
833    pub retained_at: DateTime<Utc>,
834    /// Policy-derived expiry time.
835    pub expires_at: DateTime<Utc>,
836    /// Set after maintenance removed the exact owned display. Durable archives
837    /// and run events are intentionally untouched.
838    #[serde(default, skip_serializing_if = "Option::is_none")]
839    pub expired_at: Option<DateTime<Utc>>,
840}
841
842/// Immutable caller-owned native Pi session identity. This is an association,
843/// not evidence of writer quiescence, history retention, or settlement authority.
844#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
845pub struct CallerPiSession {
846    /// Native Pi UUID.
847    pub pi_session_id: String,
848    /// Absolute native session JSONL path (a locator, not a history guarantee).
849    pub session_path: String,
850    /// Checkout recorded when binding, immutable after teardown.
851    pub original_cwd: String,
852    /// Native file device at registration; catches path replacement on retry.
853    pub file_dev: u64,
854    /// Native file inode at registration; append growth is permitted.
855    pub file_ino: u64,
856    /// Reserved launch generation, absent for legacy direct bindings.
857    #[serde(default, skip_serializing_if = "Option::is_none")]
858    pub generation: Option<u64>,
859}
860
861/// A caller attests a Pi lifecycle transition. This is not a worker exit or
862/// proof of quiescence: only an external writer-fence protocol can authorize
863/// settlement. The event log retains every distinct transition.
864#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
865#[serde(rename_all = "kebab-case")]
866pub enum CallerPiState {
867    /// Durable launch attempt, before Pi creates any native file. Never liveness proof.
868    Reserved,
869    /// Caller confirmed Start returned successfully; not continuing liveness proof.
870    Started,
871    /// Pi launch failed before a writer existed.
872    LaunchFailed,
873    /// Caller reaped or otherwise proved the exact Pi stopped.
874    Exited,
875    /// Caller cannot prove whether it still controls the writer.
876    ControlUncertain,
877}
878
879impl CallerPiState {
880    /// A told nonterminal condition requiring intervention.
881    pub fn needs_attention(&self) -> bool {
882        !matches!(self, Self::Reserved | Self::Started)
883    }
884}
885
886/// Current generation projection of append-only caller Pi lifecycle history.
887#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
888pub struct CallerPiLifecycle {
889    /// Monotonically increasing launch-attempt number.
890    pub generation: u64,
891    /// Verified bound native Pi UUID.
892    pub pi_session_id: String,
893    /// Native history path, unknown until Pi creates and the caller binds it.
894    #[serde(default, skip_serializing_if = "Option::is_none")]
895    pub session_path: Option<String>,
896    /// Last attested state, not a writer fence.
897    pub state: CallerPiState,
898    /// Explicit attestation, not an inferred exit status.
899    pub reason: Option<String>,
900}
901
902/// Sticky caller-owned settlement operation. An intent is not a terminal report.
903#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
904#[serde(rename_all = "kebab-case")]
905pub enum SettlementOperation {
906    /// Integrate branch into source.
907    Merge,
908    /// Terminalize without discarding work.
909    Cancel,
910    /// Explicitly discard retained work.
911    Discard,
912}
913
914/// One durable, single-run settlement fence. Never cleared by failed attempts.
915#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
916pub struct CallerSettlementIntent {
917    /// Exact full run identity.
918    pub run_id: RunId,
919    /// Exact node identity.
920    pub node_id: NodeId,
921    /// Caller-supplied stable retry key.
922    pub key: String,
923    /// Requested action (not permission to perform it).
924    pub operation: SettlementOperation,
925    /// Audited initiator.
926    pub actor: String,
927    /// Required for discard and cancel; optional for merge.
928    pub reason: Option<String>,
929    /// Inode pair of the lifetime writer lock.
930    pub writer_dev: u64,
931    /// Writer inode.
932    pub writer_ino: u64,
933    /// Inode pair of the launch gate.
934    pub gate_dev: u64,
935    /// Gate inode.
936    pub gate_ino: u64,
937    /// Current reservation generation; zero means no reservation yet.
938    pub generation: u64,
939    /// Event sequence, filled by the reducer.
940    pub seq: u64,
941}
942
943/// Durable `nodes/<node-id>.json` projection (design.md §1.3).
944#[derive(Debug, Clone, Serialize, Deserialize)]
945pub struct Node {
946    /// State-schema version this file was written with.
947    pub schema_version: u32,
948    /// Unique node identifier within its run (e.g. `n-0001`). Validated on
949    /// read; this is the projection's filename key, so it can never name a
950    /// path outside `nodes/`.
951    pub node_id: NodeId,
952    /// Run this node belongs to. Validated on read.
953    pub run_id: RunId,
954    /// Parent node within the same run, if this is a sub-node.
955    pub parent_node_id: Option<NodeId>,
956    /// Kind of work this node performs.
957    pub kind: Kind,
958    /// Current node status.
959    pub status: Status,
960    /// Task description / prompt driving the node, if recorded.
961    pub task: Option<String>,
962    /// Filesystem path of the node's git worktree, if created.
963    pub worktree_path: Option<String>,
964    /// Git branch the node works on, if any.
965    pub branch: Option<String>,
966    /// The commit SHA the node's branch/worktree was forked from at spawn
967    /// (the branch tip the moment `create.sh` materialized the worktree). It
968    /// is the fixed reference point that lets the supervisor tell "this branch
969    /// produced work that merged into source" from "this branch never diverged
970    /// from its fork point": a branch still at `base_sha` is trivially an
971    /// ancestor of its source branch but has merged nothing, so it must NOT be
972    /// reconciled to success or torn down (that would drop a live agent's
973    /// uncommitted work). Only a branch whose tip has moved past `base_sha`
974    /// *and* is now an ancestor of the run's `source_branch` is a confirmed
975    /// merge (issues `false-failed-after-merge` /
976    /// `supervisor-stuck-pending-after-self-merge`). `#[serde(default)]` keeps a
977    /// node written before this field existed readable (`None` → the
978    /// git-reconcile fallback simply does not fire for it).
979    #[serde(default)]
980    pub base_sha: Option<String>,
981    /// tmux window hosting the node's agent, if interactive. This is the
982    /// human-readable window *name* — not unique across sessions and blind to
983    /// non-default sockets. Kept for display and as the legacy liveness key;
984    /// prefer [`Node::tmux_identity`] when present.
985    pub tmux_window: Option<String>,
986    /// Fully-qualified tmux identity (`session:window_id` + socket path)
987    /// captured at spawn time. `None` for nodes registered before create.sh
988    /// emitted the qualified fields — those fall back to bare-name matching on
989    /// [`Node::tmux_window`]. New spawns always populate this when create.sh
990    /// returns it.
991    #[serde(default)]
992    pub tmux_identity: Option<Box<TmuxIdentity>>,
993    /// Native worker evidence, when this attempt uses Pi. Initialized before
994    /// publication and advanced only by locked evidence events.
995    #[serde(default, skip_serializing_if = "Option::is_none")]
996    pub evidence: Option<WorkerEvidence>,
997    /// Separate from Taskfleet-launched worker evidence; survives checkout removal.
998    #[serde(default, skip_serializing_if = "Option::is_none")]
999    pub caller_pi_session: Option<CallerPiSession>,
1000    #[serde(default, skip_serializing_if = "Option::is_none")]
1001    /// Current caller Pi lifecycle projection; absent for ordinary workers.
1002    pub caller_pi_lifecycle: Option<CallerPiLifecycle>,
1003    /// Inert completed-window identity, when this run opted into retention.
1004    #[serde(default, skip_serializing_if = "Option::is_none")]
1005    pub retained_display: Option<Box<RetainedDisplay>>,
1006    /// Permanent reason why a configured terminal display could not exist
1007    /// (for example the recorded tmux server generation was lost).
1008    #[serde(default, skip_serializing_if = "Option::is_none")]
1009    pub retention_unavailable: Option<String>,
1010    /// PID of the running agent process, if live.
1011    pub agent_pid: Option<i32>,
1012    /// Start time of the agent process, used to detect PID reuse.
1013    pub agent_pid_start_time: Option<DateTime<Utc>>,
1014    /// PID of the supervisor watching this node, if live.
1015    pub supervisor_pid: Option<i32>,
1016    /// Children this node has spawned.
1017    #[serde(default)]
1018    pub children: Vec<ChildRef>,
1019    /// When the node started executing, if it has.
1020    pub started_at: Option<DateTime<Utc>>,
1021    /// When the node file was last modified.
1022    pub updated_at: DateTime<Utc>,
1023    /// The `node.report` payload that drove this node to its terminal status.
1024    /// Set only by the report that actually transitions the node (Done /
1025    /// Failed / Cancelled). Once the node is terminal it is frozen: a late
1026    /// report against an already-settled node is dropped without overwriting
1027    /// this field (see `reducer::apply_node_report`). So for a node cancelled
1028    /// by `run cancel`, this holds the synthesized cancel report, not a
1029    /// later-arriving agent report — that payload remains only in
1030    /// `events.jsonl`.
1031    pub last_report: Option<Value>,
1032    /// Highest report `seq` consumed per child run id, for idempotent
1033    /// report processing across supervisor restarts.
1034    #[serde(default)]
1035    pub last_processed_report_seq_by_child: Map<String, Value>,
1036    /// Number of times the supervisor has auto-retried this node after an
1037    /// empty-handed `agent-died` (issue `autoretry-agent-died-worker`). The
1038    /// DURABLE, restart-safe bound on the bounded-retry loop: each `node.retry`
1039    /// event increments it, and the watchdog terminalizes the run `failed` once
1040    /// it reaches `RETRY_MAX_ATTEMPTS`. `#[serde(default)]` keeps a node written
1041    /// before this field existed readable (`0` — never retried).
1042    #[serde(default)]
1043    pub retry_attempts: u32,
1044    /// The **told** exit status of the node's worker process, recorded durably by
1045    /// the `run-worker` launcher shim (`crates/taskfleet/src/run_worker.rs`) when
1046    /// it `wait()`s on the agent it wrapped. This is a *fact*, not an inference:
1047    /// the supervisor consumes it via the typed outcome table instead of guessing
1048    /// completion from pid/pane/activity proxies (design.md §2.1, issue
1049    /// `thin-exit-status-launcher`). A non-zero code or a terminating signal is a
1050    /// `failed` worker; `code == 0` with no `explicit-merge` transition is the
1051    /// *finished-but-unmerged* case that must stay non-terminal (attention-
1052    /// required), NOT be auto-failed. `None` until the shim records an exit — or
1053    /// forever, for a worker never launched through the shim (the crash backstop
1054    /// still covers that path). `#[serde(default)]` keeps a node written before
1055    /// this field existed readable.
1056    #[serde(default)]
1057    pub worker_exit: Option<WorkerExit>,
1058    /// The in-flight `run merge` transaction for this node, if one has been
1059    /// STARTED but not yet completed. `run merge` records a `merge.started`
1060    /// event (setting this field) BEFORE it mutates git, because the merge spans
1061    /// two durability domains — git refs and the event log — and is not atomic
1062    /// across them (design.md §2.1b / A2, issue `merge-transaction-recovery`). A
1063    /// crash after the git merge but before the terminal `explicit-merge`
1064    /// `node.report` would otherwise strand the work *merged in source* with *no
1065    /// merge event* → a false `failed`.
1066    ///
1067    /// This field is the durable op-log record that lets recovery finish or
1068    /// reject that ONE known transaction deterministically, by OID — never a
1069    /// general branch-content heuristic. It is set by [`crate::MergeTxn`]-carrying
1070    /// `merge.started`, and cleared when the transaction resolves: a terminal
1071    /// `node.report` (the merge completed) or a `merge.aborted` (recovery found
1072    /// the git mutation never landed). `#[serde(default)]` keeps a node written
1073    /// before this field existed readable (`None` — no in-flight merge).
1074    ///
1075    /// Boxed so the (rare) in-flight transaction does not inflate every `Node` /
1076    /// `ProjectionOp` by the full [`MergeTxn`] footprint.
1077    #[serde(default)]
1078    pub pending_merge: Option<Box<MergeTxn>>,
1079    /// The durable, monotonic timestamp of the FIRST tick on which the supervisor
1080    /// observed this node's worker process confirmed-dead with no told
1081    /// `worker.exited` and no merge — the anchor for the residual crash
1082    /// backstop's fixed post-death grace (design.md §2.1a, issue
1083    /// `typed-supervisor-outcomes`).
1084    ///
1085    /// The backstop is the ONLY place pid liveness still governs an outcome
1086    /// (pid liveness is a pure crash backstop now, never a primary signal). When
1087    /// the launcher shim's exit fact is lost — a hard kill of the shim, host
1088    /// death — the supervisor never sees a `worker.exited` event, so it falls
1089    /// back to "process confirmed gone → `failed`". The grace exists only to let
1090    /// an in-flight `worker.exited` / merge append land before that fires: on the
1091    /// first confirmed death the supervisor records this timestamp (via a
1092    /// `node.death_observed` event) and DEFERS; it terminalizes `failed` only on a
1093    /// later tick once a fixed short window has elapsed AND an exclusive-lock
1094    /// re-read confirms no exit/merge landed in the race window.
1095    ///
1096    /// Persisted (not in-memory) so the grace survives a supervisor restart in the
1097    /// window — a restart re-reads it rather than restarting the clock. First-write
1098    /// -wins in the reducer, so the anchor is monotonic. `None` until the first
1099    /// confirmed-death observation, or forever for a worker that exits cleanly
1100    /// (the shim records `worker.exited` and the backstop never engages).
1101    /// `#[serde(default)]` keeps a node written before this field existed readable.
1102    #[serde(default)]
1103    pub first_death_at: Option<DateTime<Utc>>,
1104    /// An open, agent-authored request for a human decision. The worker records
1105    /// this through `node.awaiting_input` instead of blocking on interactive
1106    /// stdin. It remains non-terminal and is cleared by `node.input_resolved`, a
1107    /// terminal `node.report`, or `node.retry`.
1108    ///
1109    /// `opened_at` is stamped from the event envelope and is therefore a durable,
1110    /// monotonic grace-window anchor that survives supervisor restarts. The
1111    /// original discussion objects are retained verbatim so read surfaces and
1112    /// notification hooks can carry the question, options, and recommended
1113    /// default without inventing a second advisory schema.
1114    #[serde(default)]
1115    pub awaiting_input: Option<Box<AwaitingInput>>,
1116}
1117
1118/// Durable open-discussion state projected from `node.awaiting_input`.
1119#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
1120pub struct AwaitingInput {
1121    /// Timestamp of the first open signal in the current unresolved generation.
1122    pub opened_at: DateTime<Utc>,
1123    /// Event sequence that opened this generation, used to deduplicate its
1124    /// delayed parent notification independently from later generations.
1125    pub event_seq: u64,
1126    /// Validated report-shaped discussion objects. Each carries `topic`,
1127    /// `options`, and `recommended_default`.
1128    pub discussion_items: Vec<Value>,
1129}
1130
1131/// A durable, in-flight `run merge` transaction recorded by `merge.started`
1132/// BEFORE the git mutation, and the sole input to deterministic merge-crash
1133/// recovery (design.md §2.1b / A2, issue `merge-transaction-recovery`).
1134///
1135/// `run merge` spans git refs and the event log and is not atomic across them.
1136/// Recording the transaction — the exact source ref it will move, the OID it
1137/// expects that ref to be at (`expected_source_oid`, the compare half of the
1138/// compare-and-swap), and the worker's tip — lets the supervisor (or a retried
1139/// `run merge`) resolve the ONE recorded transaction by OID after a crash:
1140///
1141/// - source ref still at `expected_source_oid` → the mutation never landed →
1142///   REJECT (`merge.aborted`), preserving the worker's branch + work.
1143/// - source ref moved off `expected_source_oid` AND the worker's content is
1144///   integrated (rebase-robust content verification) → COMPLETE (append the
1145///   `explicit-merge` `node.report` the crash prevented).
1146/// - source ref moved unexpectedly but the worker's content is not integrated →
1147///   fail closed (REJECT), preserving the work.
1148#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
1149pub struct MergeTxn {
1150    /// Opaque unique id for this merge attempt. A fresh id per `run merge`
1151    /// invocation (each attempt re-reads `expected_source_oid`), so recovery can
1152    /// name exactly which transaction it resolved in the `merge.aborted` audit.
1153    pub op_id: String,
1154    /// The source/target ref this merge moves — `manifest.source_branch`
1155    /// (`main`, or an integration branch). Recovery reads this ref's current OID
1156    /// to decide the transaction's fate.
1157    pub source_branch: String,
1158    /// The worker branch whose commits are being merged (`node.branch`). Its
1159    /// content is what recovery verifies is integrated into `source_branch`.
1160    pub worker_branch: String,
1161    /// The OID `source_branch` was at when the transaction was recorded — the
1162    /// compare half of the compare-and-swap. If the ref is still here at recovery
1163    /// time, the git mutation never landed.
1164    pub expected_source_oid: String,
1165    /// The worker branch tip at record time. Retained for the audit trail and as
1166    /// a secondary landing signal; the authoritative completion check is
1167    /// content-based (rebase-robust) against `source_branch`.
1168    pub worker_oid: String,
1169    /// The worker branch's fork point (`node.base_sha`), used to bound the
1170    /// content check to the worker's own commits. `None` when unrecorded.
1171    #[serde(default)]
1172    pub base_sha: Option<String>,
1173    /// PID of the `run merge` process driving the transaction, so recovery can
1174    /// tell a still-in-progress merge (driver alive — leave it) from a crashed
1175    /// one (driver gone — resolve it), never racing a live merge. `None` when
1176    /// unrecorded.
1177    #[serde(default)]
1178    pub driver_pid: Option<i32>,
1179    /// Start time of `driver_pid` in Unix seconds (the same representation the
1180    /// pid-file liveness check records), guarding against PID reuse the way the
1181    /// agent/supervisor liveness checks do — a recycled PID must not look alive.
1182    /// `None` when the platform could not read it.
1183    #[serde(default)]
1184    pub driver_pid_start_secs: Option<u64>,
1185    /// Caller-only provenance tying the transaction to the sticky intent and
1186    /// the exact writer inode held by the merge driver. Absent on normal merges.
1187    #[serde(default, skip_serializing_if = "Option::is_none")]
1188    pub caller_authority: Option<CallerMergeLink>,
1189    /// When the transaction was recorded.
1190    pub started_at: DateTime<Utc>,
1191}
1192
1193/// Immutable provenance for a caller-owned merge. Checked against the intent
1194/// projection when the transaction is folded, not inferred from Git ancestry.
1195#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
1196pub struct CallerMergeLink {
1197    /// Durable sequence of `caller.settlement_intent`.
1198    pub intent_seq: u64,
1199    /// Stable operation key recorded by that intent.
1200    pub intent_key: String,
1201    /// Recorded writer device.
1202    pub writer_dev: u64,
1203    /// Recorded writer inode.
1204    pub writer_ino: u64,
1205}
1206
1207/// The observed exit status of a node's worker process, recorded by the
1208/// `run-worker` launcher shim under the run lock (design.md §2.1 / A1).
1209///
1210/// Exactly one of `code` / `signal` is meaningful: a worker that returned
1211/// normally carries `code = Some(n)` (and `signal = None`); a worker killed by a
1212/// signal carries `signal = Some(s)` (and, on Unix, `code = None`). A recorded
1213/// exit is a durable *told fact* — the supervisor reads it rather than inferring
1214/// completion from liveness proxies.
1215#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
1216pub struct WorkerExit {
1217    /// Normal-exit status code, if the worker was not killed by a signal.
1218    #[serde(default)]
1219    pub code: Option<i32>,
1220    /// Terminating signal number, if the worker was killed by a signal.
1221    #[serde(default)]
1222    pub signal: Option<i32>,
1223    /// When the shim observed the worker's exit.
1224    pub at: DateTime<Utc>,
1225}
1226
1227impl WorkerExit {
1228    /// A clean exit: not signalled, and a zero return code. This is the *only*
1229    /// success-shaped worker exit — but a clean exit alone is NOT a completed
1230    /// unit (the worker may have finished-but-skipped `run merge`); merge is the
1231    /// only success truth (design.md §2.6). Callers pair this with a merge check.
1232    pub fn is_clean(self) -> bool {
1233        self.signal.is_none() && self.code == Some(0)
1234    }
1235
1236    /// A failed worker: killed by a signal, or a non-zero return code. Mutually
1237    /// exclusive with [`WorkerExit::is_clean`].
1238    pub fn is_failure(self) -> bool {
1239        !self.is_clean()
1240    }
1241}
1242
1243/// A fully-qualified tmux window identity recorded at spawn time.
1244///
1245/// `tmux_window` (the human name) is not unique across sessions, and a bare
1246/// `tmux list-windows -a` cannot see windows on a non-default socket. This
1247/// triple pins the exact window the agent runs in — `session:window_id` is
1248/// unique per server, `window_id` (the `@NNNN` form) survives renames, and
1249/// `socket` disambiguates multiple tmux servers. The watchdog matches on this
1250/// when present (design.md §8.1).
1251///
1252/// `pane_id` (the `%NN` form) pins the agent's *specific* pane within that
1253/// window, recorded at spawn. Window-owning operations (`kill-window` teardown —
1254/// the supervisor owns the whole window per the cleanup invariants) key off
1255/// `window_id`; only per-pane operations that must not follow the window's
1256/// *active* pane — chiefly `pipe-pane` agent-log capture — use `pane_id`. It is
1257/// `None` for a run spawned before create.sh emitted the field; capture then
1258/// falls back to `window_id` (issue `capture-agent-pane-by-pane-id`).
1259///
1260/// The watchdog's liveness probe still keys off `window_id` (correct for the
1261/// single-pane autonomous path). A pane-aware liveness probe — needed so a split
1262/// interactive window whose agent pane dies while a user shell pane survives is
1263/// still seen as dead — is a follow-up (`watchdog-pane-aware-liveness`), not this
1264/// change.
1265#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
1266pub struct TmuxIdentity {
1267    /// Server socket path (`#{socket_path}`). `None` if create.sh could not
1268    /// read it; the watchdog then queries tmux on its default socket.
1269    #[serde(default)]
1270    pub socket: Option<String>,
1271    /// Session that owns the window (`#{session_name}`).
1272    pub session: String,
1273    /// Stable window id in `@NNNN` form (`#{window_id}`). Survives renames and
1274    /// is unique within the server.
1275    pub window_id: String,
1276    /// Stable pane id in `%NN` form (`#{pane_id}`), recorded at spawn — the
1277    /// agent's own pane. `None` for a run whose create.sh predates the field
1278    /// (back-compat: old state deserializes with `pane_id: None`). Prefer
1279    /// [`TmuxIdentity::capture_target`] over reading this directly.
1280    #[serde(default)]
1281    pub pane_id: Option<String>,
1282    /// tmux server PID and its OS start identity at spawn. New persistent
1283    /// sessions require both; legacy nodes deserialize them as absent.
1284    #[serde(default)]
1285    pub server_pid: Option<u32>,
1286    /// Unix-second process start identity for the recorded tmux server PID.
1287    #[serde(default)]
1288    pub server_pid_start_secs: Option<u64>,
1289    /// Runtime-owned server marker (for example Homebase's owner marker).
1290    #[serde(default)]
1291    pub server_marker: Option<String>,
1292}
1293
1294impl TmuxIdentity {
1295    /// The tmux target for a per-pane operation that must hit the agent's own
1296    /// pane, not the window's *active* pane: the recorded `pane_id` when
1297    /// present, else the `window_id` (which resolves to the active pane).
1298    ///
1299    /// Used by agent-log capture (`pipe-pane`). Window-level operations
1300    /// (`kill-window`, liveness) must NOT use this — they key off `window_id`
1301    /// directly so they act on the whole window.
1302    ///
1303    /// A recorded `pane_id` is preferred only when non-empty; an empty string
1304    /// (a directly-deserialized/corrupt state that the reducer/spawn normalizers
1305    /// never produce) is treated as absent so capture never targets `-t ""`.
1306    pub fn capture_target(&self) -> &str {
1307        self.pane_id
1308            .as_deref()
1309            .filter(|id| !id.is_empty())
1310            .unwrap_or(&self.window_id)
1311    }
1312}
1313
1314/// One event-log line (design.md §1.4).
1315///
1316/// `run_id` / `node_id` are the typed id newtypes, so deserializing an
1317/// `events.jsonl` line validates the whole envelope on read: a malformed
1318/// `run_id` or `node_id` fails the `serde` parse at the read boundary (the
1319/// id newtypes' validating `Deserialize`) rather than being carried as an
1320/// unvalidated `String` until some later path helper. The parse failure
1321/// surfaces as whatever error the reader maps a bad line to — e.g. a
1322/// newline-terminated bad line is [`Error::CorruptEventLog`] from both
1323/// [`read_all_events`] and [`find_prior_with_key`], which share one physical
1324/// reader and torn-tail policy. The reducer still performs its own per-event
1325/// checks (envelope `run_id` matches the run it is folded into; `data`-borne
1326/// ids parse), but the envelope ids can no longer be the unvalidated party.
1327///
1328/// [`read_all_events`]: crate::events::read_all_events
1329/// [`find_prior_with_key`]: crate::events
1330/// [`Error::CorruptEventLog`]: crate::Error::CorruptEventLog
1331#[derive(Debug, Clone, Serialize, Deserialize)]
1332pub struct Event {
1333    /// Wall-clock timestamp the event was appended.
1334    pub ts: DateTime<Utc>,
1335    /// Monotonic per-run sequence number (recovered on append).
1336    pub seq: u64,
1337    /// Event kind discriminator (e.g. `node.created`, `discussion.opened`).
1338    pub kind: String,
1339    /// Run the event belongs to. Validated on read.
1340    pub run_id: RunId,
1341    /// Node the event concerns, when applicable. Validated on read.
1342    #[serde(skip_serializing_if = "Option::is_none", default)]
1343    pub node_id: Option<NodeId>,
1344    /// Caller-supplied key used to dedupe retried appends.
1345    #[serde(skip_serializing_if = "Option::is_none", default)]
1346    pub idempotency_key: Option<String>,
1347    /// Kind-specific payload applied by the reducer.
1348    #[serde(default)]
1349    pub data: Value,
1350}
1351
1352#[cfg(test)]
1353mod tests {
1354    use super::*;
1355
1356    #[test]
1357    fn aggregate_terminal_status_is_the_three_way_rule() {
1358        use Status::{Blocked, Cancelled, Done, Failed, Pending, Running};
1359        // Empty set → not complete.
1360        assert_eq!(aggregate_terminal_status([]), None);
1361        // Any live node → not complete.
1362        for live in [Pending, Running, Blocked] {
1363            assert_eq!(aggregate_terminal_status([Done, live]), None);
1364        }
1365        // All done → Done.
1366        assert_eq!(aggregate_terminal_status([Done, Done]), Some(Done));
1367        // Any failure dominates.
1368        assert_eq!(aggregate_terminal_status([Done, Failed]), Some(Failed));
1369        assert_eq!(aggregate_terminal_status([Failed, Cancelled]), Some(Failed));
1370        // Cancelled (no failure) — pure or mixed with done.
1371        assert_eq!(
1372            aggregate_terminal_status([Cancelled, Cancelled]),
1373            Some(Cancelled)
1374        );
1375        assert_eq!(
1376            aggregate_terminal_status([Done, Cancelled]),
1377            Some(Cancelled)
1378        );
1379    }
1380
1381    /// `Kind::wire_name` (and thus `Kind::WIRE_NAMES`) must stay identical
1382    /// to what serde actually (de)serializes. If the `rename_all` routing
1383    /// or a variant name ever diverges from `wire_name`, this fails — which
1384    /// is what keeps the report validator's `expected` hint honest.
1385    #[test]
1386    fn wire_names_match_serde_round_trip() {
1387        for &name in Kind::WIRE_NAMES {
1388            let kind: Kind = serde_json::from_value(Value::String(name.to_string()))
1389                .unwrap_or_else(|_| panic!("WIRE_NAMES entry {name:?} is not a valid Kind"));
1390            assert_eq!(
1391                serde_json::to_value(kind).unwrap(),
1392                Value::String(name.to_string()),
1393                "serde round-trip diverged from wire_name for {name:?}",
1394            );
1395        }
1396    }
1397
1398    /// The bounded auto-retry eligibility gate (issue `autoretry-agent-died-worker`)
1399    /// must include exactly the autonomous single-node worker kinds and exclude
1400    /// the fan-out driver (and the read-only `Unknown` catch-all).
1401    #[test]
1402    fn autonomous_single_node_worker_set_is_exact() {
1403        for k in [Kind::Spinoff, Kind::Research, Kind::TechnicalDecision] {
1404            assert!(
1405                k.is_autonomous_single_node_worker(),
1406                "{k:?} should be retry-eligible"
1407            );
1408            assert_eq!(k.lifecycle(), Lifecycle::Autonomous);
1409        }
1410        for k in [
1411            Kind::FanOut,  // multi-unit driver
1412            Kind::Unknown, // legacy on-disk run — never freshly supervised
1413        ] {
1414            assert!(
1415                !k.is_autonomous_single_node_worker(),
1416                "{k:?} must NOT be retry-eligible"
1417            );
1418        }
1419    }
1420
1421    /// A legacy run recorded under a since-removed kind must still deserialize
1422    /// to the read-only [`Kind::Unknown`] catch-all rather than faulting the
1423    /// read — the ADR §D7 "report, never delete" contract for the on-disk
1424    /// evidence corpus. Every creatable kind still round-trips to itself.
1425    #[test]
1426    fn removed_kinds_deserialize_to_unknown() {
1427        for removed in [
1428            "code",
1429            "orchestrate",
1430            "orchestrated",
1431            "bugfix",
1432            "make-skill",
1433        ] {
1434            let kind: Kind = serde_json::from_value(Value::String(removed.to_string()))
1435                .expect("a removed kind must still deserialize, not fault");
1436            assert_eq!(kind, Kind::Unknown, "{removed:?} should map to Unknown");
1437        }
1438        // A wholly unknown value maps there too (forward-compat).
1439        assert_eq!(
1440            serde_json::from_value::<Kind>(Value::String("future-kind".into())).unwrap(),
1441            Kind::Unknown
1442        );
1443        // The surviving kinds are unaffected.
1444        for &name in Kind::WIRE_NAMES {
1445            let kind: Kind = serde_json::from_value(Value::String(name.to_string())).unwrap();
1446            assert_ne!(kind, Kind::Unknown, "{name:?} must not fold to Unknown");
1447        }
1448    }
1449
1450    /// Back-compat acceptance criterion (issue `capture-agent-pane-by-pane-id`):
1451    /// a `TmuxIdentity` persisted before `pane_id` existed — with the field
1452    /// entirely absent, or written as an explicit `null` — must still
1453    /// deserialize, yielding `pane_id: None` and a `window_id` capture target.
1454    #[test]
1455    fn tmux_identity_deserializes_legacy_state_without_pane_id() {
1456        // Field entirely absent (a state file written by an older binary).
1457        let absent: TmuxIdentity = serde_json::from_value(serde_json::json!({
1458            "socket": null,
1459            "session": "taskfleet",
1460            "window_id": "@42",
1461        }))
1462        .expect("legacy identity without pane_id must deserialize");
1463        assert_eq!(absent.pane_id, None);
1464        assert_eq!(absent.capture_target(), "@42");
1465
1466        // Field present but explicitly null.
1467        let null: TmuxIdentity = serde_json::from_value(serde_json::json!({
1468            "socket": null,
1469            "session": "taskfleet",
1470            "window_id": "@42",
1471            "pane_id": null,
1472        }))
1473        .expect("identity with explicit null pane_id must deserialize");
1474        assert_eq!(null.pane_id, None);
1475        assert_eq!(null.capture_target(), "@42");
1476    }
1477
1478    /// `capture_target` prefers a recorded `pane_id` (`%NN`) over the window id,
1479    /// but treats an empty `pane_id` as absent (never targets `-t ""`).
1480    #[test]
1481    fn capture_target_prefers_nonempty_pane_id() {
1482        let with_pane = TmuxIdentity {
1483            socket: None,
1484            session: "taskfleet".into(),
1485            window_id: "@42".into(),
1486            pane_id: Some("%7".into()),
1487            server_pid: None,
1488            server_pid_start_secs: None,
1489            server_marker: None,
1490        };
1491        assert_eq!(with_pane.capture_target(), "%7");
1492
1493        let empty_pane = TmuxIdentity {
1494            pane_id: Some(String::new()),
1495            ..with_pane.clone()
1496        };
1497        assert_eq!(empty_pane.capture_target(), "@42");
1498    }
1499}
1500
1501#[cfg(test)]
1502mod id_tests {
1503    use super::*;
1504
1505    /// Inputs every id type must reject — the path-traversal vectors plus the
1506    /// generic malformed cases called out in the issue's success criteria.
1507    const TRAVERSAL_VECTORS: &[&str] = &[
1508        "..",
1509        "../etc",
1510        "a/b",
1511        "a/../b",
1512        ".hidden",
1513        "./x",
1514        "foo/bar.json",
1515        "n-0001/../../etc",
1516        "",
1517    ];
1518
1519    #[test]
1520    fn run_id_accepts_generator_output_and_rejects_malformed() {
1521        let id = crate::new_run_id();
1522        assert!(
1523            RunId::parse_str(&id).is_ok(),
1524            "generator must validate: {id}"
1525        );
1526        for bad in [
1527            "tooshort",
1528            "01jxsnap0000000000000000000", // 27 chars
1529            "01JXSNAP000000000000000000",  // uppercase
1530            "01jxiiiiiiiiiiiiiiiiiiiiii",  // `i` not in Crockford
1531            "80000000000000000000000000",  // first char exceeds ULID range
1532            "n-0001",                      // wrong shape entirely
1533        ] {
1534            assert!(RunId::parse_str(bad).is_err(), "expected reject: {bad:?}");
1535        }
1536        for bad in TRAVERSAL_VECTORS {
1537            assert!(
1538                RunId::parse_str(bad).is_err(),
1539                "traversal not rejected: {bad:?}"
1540            );
1541        }
1542    }
1543
1544    #[test]
1545    fn node_id_accepts_canonical_and_rejects_malformed() {
1546        for ok in ["n-0001", "n-0010", "n-123456"] {
1547            assert!(NodeId::parse_str(ok).is_ok(), "expected accept: {ok}");
1548        }
1549        // Wrong prefix is its own error variant.
1550        assert!(matches!(
1551            NodeId::parse_str("d-0001"),
1552            Err(IdValidationError::WrongPrefix { .. })
1553        ));
1554        assert!(matches!(
1555            NodeId::parse_str("0001"),
1556            Err(IdValidationError::WrongPrefix { .. })
1557        ));
1558        for bad in [
1559            "n-1",           // too few digits
1560            "n-abcd",        // non-digit body
1561            "n-",            // empty body
1562            "n-00a1",        // mixed
1563            "n-00000000000", // 11 digits — over the 10-digit ceiling
1564        ] {
1565            assert!(
1566                matches!(
1567                    NodeId::parse_str(bad),
1568                    Err(IdValidationError::InvalidFormat { .. })
1569                ),
1570                "expected InvalidFormat: {bad:?}",
1571            );
1572        }
1573        for bad in TRAVERSAL_VECTORS {
1574            assert!(
1575                NodeId::parse_str(bad).is_err(),
1576                "traversal not rejected: {bad:?}"
1577            );
1578        }
1579    }
1580
1581    #[test]
1582    fn deserialize_rejects_malformed_ids() {
1583        // The validating Deserialize impl is the on-read guard: a tampered
1584        // projection file whose key no longer validates must fail to parse.
1585        assert!(serde_json::from_str::<NodeId>("\"n-0001\"").is_ok());
1586        assert!(serde_json::from_str::<NodeId>("\"../../etc\"").is_err());
1587        assert!(serde_json::from_str::<NodeId>("\"n-../escape\"").is_err());
1588    }
1589
1590    #[test]
1591    fn serialize_round_trips_as_bare_string() {
1592        let nid = NodeId::parse_str("n-0042").unwrap();
1593        let json = serde_json::to_string(&nid).unwrap();
1594        assert_eq!(json, "\"n-0042\"");
1595        let back: NodeId = serde_json::from_str(&json).unwrap();
1596        assert_eq!(back, nid);
1597        assert_eq!(nid.as_str(), "n-0042");
1598        assert_eq!(nid.to_string(), "n-0042");
1599    }
1600
1601    #[test]
1602    fn error_exposes_kind_and_expected() {
1603        let err = NodeId::parse_str("n-x").unwrap_err();
1604        assert_eq!(err.kind(), "node");
1605        assert_eq!(err.expected(), "n-NNNN (n- followed by 4-10 ASCII digits)");
1606    }
1607
1608    #[test]
1609    fn event_deserialize_validates_envelope_ids() {
1610        // The whole `events.jsonl` envelope is now validated on read: the
1611        // typed `run_id` / `node_id` fields parse through the id newtypes, so
1612        // a malformed envelope id fails the deserialize rather than being
1613        // carried downstream as an unchecked string.
1614        let ok = r#"{"ts":"2026-06-12T00:00:00Z","seq":1,"kind":"node.created","run_id":"01jxsnap000000000000000000","node_id":"n-0001","data":{}}"#;
1615        assert!(serde_json::from_str::<Event>(ok).is_ok());
1616
1617        // Invalid `run_id` (not a 26-char ULID) fails the parse.
1618        let bad_run = r#"{"ts":"2026-06-12T00:00:00Z","seq":1,"kind":"run.status","run_id":"not-a-ulid","data":{}}"#;
1619        assert!(serde_json::from_str::<Event>(bad_run).is_err());
1620
1621        // Invalid top-level `node_id` (too few digits) also fails the parse.
1622        let bad_node = r#"{"ts":"2026-06-12T00:00:00Z","seq":1,"kind":"node.status","run_id":"01jxsnap000000000000000000","node_id":"n-1","data":{}}"#;
1623        assert!(serde_json::from_str::<Event>(bad_node).is_err());
1624    }
1625
1626    #[test]
1627    fn from_str_and_ord_delegate_to_inner() {
1628        use std::str::FromStr;
1629        // `FromStr` mirrors `parse_str`, so the `str::parse` ecosystem works.
1630        assert!(RunId::from_str("01jxsnap000000000000000000").is_ok());
1631        assert!("n-0001".parse::<NodeId>().is_ok());
1632        assert!("n-x".parse::<NodeId>().is_err());
1633
1634        // `Ord` is lexicographic over the inner string; for ULIDs that is the
1635        // natural time-encoded order.
1636        let a = RunId::parse_str("01jxsnap000000000000000000").unwrap();
1637        let b = RunId::parse_str("02jxsnap000000000000000000").unwrap();
1638        assert!(a < b);
1639        let mut v = vec![b.clone(), a.clone()];
1640        v.sort();
1641        assert_eq!(v, vec![a, b]);
1642    }
1643}
1644
1645#[cfg(test)]
1646mod agent_selection_validation_tests {
1647    use super::*;
1648
1649    fn valid() -> AgentSelection {
1650        AgentSelection {
1651            schema_version: 1,
1652            profile: "capable".into(),
1653            selection_source: "cli".into(),
1654            interaction: "autonomous".into(),
1655            capability: "capable".into(),
1656            residency: "remote".into(),
1657            requested_harness: None,
1658            selected: SelectedAgentCandidate {
1659                candidate_index: 1,
1660                harness: "pi".into(),
1661                command: vec!["pi".into()],
1662                telemetry: Some("worker-v1".into()),
1663            },
1664            fallback: vec![SkippedAgentCandidate {
1665                candidate_index: 0,
1666                harness: "claude".into(),
1667                reason: "autonomous_harness_unsupported".into(),
1668            }],
1669        }
1670    }
1671
1672    #[test]
1673    fn durable_selection_rejects_impossible_state() {
1674        assert!(valid().validate().is_ok());
1675        let mut invalid = valid();
1676        invalid.selected.candidate_index = 200;
1677        assert!(invalid.validate().is_err());
1678        let mut invalid = valid();
1679        invalid.fallback[0].reason = "made_up".into();
1680        assert!(invalid.validate().is_err());
1681        let mut invalid = valid();
1682        invalid.selected.harness = "claude".into();
1683        assert!(invalid.validate().is_err());
1684    }
1685}