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/// Run/node status (design.md §1.2).
363///
364/// `Done`, `Failed`, and `Cancelled` are **terminal**: once a run or node
365/// reaches one of them its `status` must never change again. The reducer
366/// enforces this — `apply_run_status`, `apply_node_status`, and
367/// `apply_node_report` are all no-ops once [`Status::is_terminal`] holds — so
368/// a late-arriving event (e.g. an agent success report racing a `run cancel`)
369/// cannot resurrect a settled state. Only the `status` field is frozen;
370/// other projection fields may still be mutated by non-status events.
371#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
372#[serde(rename_all = "kebab-case")]
373pub enum Status {
374 /// Created but not yet started.
375 Pending,
376 /// Actively executing.
377 Running,
378 /// Stalled awaiting input (e.g. an open discussion).
379 Blocked,
380 /// Completed successfully (terminal).
381 Done,
382 /// Completed with failure (terminal).
383 Failed,
384 /// Terminated before completion by an operator or parent (terminal).
385 Cancelled,
386}
387
388impl Status {
389 /// True for the terminal states `Done | Failed | Cancelled`. A run or
390 /// node in a terminal state is settled: the reducer treats any further
391 /// *status* transition as a no-op. "Settled" applies to `status` only —
392 /// non-status projection fields (e.g. `Node::children` via
393 /// `child.spawned`, or manifest counters) can still change.
394 pub fn is_terminal(self) -> bool {
395 matches!(self, Status::Done | Status::Failed | Status::Cancelled)
396 }
397}
398
399/// Aggregate a set of node statuses into the run's rolled-up terminal status,
400/// or `None` when the run is not yet complete.
401///
402/// The single, shared roll-up rule — used both by the supervisor's per-tick
403/// `rollup_status` and by [`cancel_node`](crate::cancel_node)'s in-lock
404/// self-roll-up (so the two can never diverge). A **three-way** classification
405/// (design §2.5, "rollup terminalizes the run cancelled/done/failed once every
406/// node is terminal"):
407///
408/// - `None` if the set is empty (a freshly-created run must not vacuously
409/// complete) or if ANY node is still live (`Pending`/`Running`/`Blocked`);
410/// - `Some(Status::Failed)` if any node genuinely `Failed` (a real failure
411/// dominates the batch outcome);
412/// - `Some(Status::Cancelled)` if no node failed but at least one was
413/// `Cancelled` (a deliberate per-node/whole-run cancel — nothing failed, but
414/// the batch did not fully complete; branch-preserving work is untouched);
415/// - `Some(Status::Done)` when every node is `Done`.
416pub fn aggregate_terminal_status<I>(statuses: I) -> Option<Status>
417where
418 I: IntoIterator<Item = Status>,
419{
420 let mut any = false;
421 let mut any_failed = false;
422 let mut any_cancelled = false;
423 for s in statuses {
424 any = true;
425 match s {
426 Status::Done => {}
427 Status::Failed => any_failed = true,
428 Status::Cancelled => any_cancelled = true,
429 // Any live node means the run is not done yet.
430 Status::Pending | Status::Running | Status::Blocked => return None,
431 }
432 }
433 if !any {
434 return None;
435 }
436 Some(if any_failed {
437 Status::Failed
438 } else if any_cancelled {
439 Status::Cancelled
440 } else {
441 Status::Done
442 })
443}
444
445/// Compact profile resolution recorded at create time.
446#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
447pub struct AgentSelection {
448 /// Compact selection schema version (currently 1).
449 pub schema_version: u32,
450 /// Requested and selected user profile name.
451 pub profile: String,
452 /// Precedence layer that supplied the request.
453 pub selection_source: String,
454 /// Explicit create-time interaction mode.
455 pub interaction: String,
456 /// Declared profile capability tier.
457 pub capability: String,
458 /// Declared profile residency class.
459 pub residency: String,
460 /// Legacy harness alias requested at the winning layer, when applicable.
461 #[serde(default, skip_serializing_if = "Option::is_none")]
462 pub requested_harness: Option<String>,
463 /// First statically eligible candidate.
464 pub selected: SelectedAgentCandidate,
465 /// Earlier candidates and their single deterministic skip reasons.
466 #[serde(default)]
467 pub fallback: Vec<SkippedAgentCandidate>,
468}
469
470/// Exact selected candidate pinned for launch and retry.
471#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
472pub struct SelectedAgentCandidate {
473 /// Zero-based position in the profile's ordered candidate list.
474 pub candidate_index: u8,
475 /// Selected harness (`pi` or `claude`).
476 pub harness: String,
477 /// Exact user-owned argv; never a shell string.
478 pub command: Vec<String>,
479 /// Declared telemetry adapter protocol, when configured.
480 #[serde(default, skip_serializing_if = "Option::is_none")]
481 pub telemetry: Option<String>,
482}
483
484impl SelectedAgentCandidate {
485 /// Whether recorded policy configures the public worker telemetry v1
486 /// adapter. This is configuration support, not runtime attestation or a
487 /// claim that any sample has arrived.
488 #[must_use]
489 pub fn supports_worker_telemetry_v1(&self) -> bool {
490 self.harness == "pi" && self.telemetry.as_deref() == Some("worker-v1")
491 }
492}
493
494/// One rejected candidate and its first applicable reason.
495#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
496pub struct SkippedAgentCandidate {
497 /// Zero-based position in the profile's ordered candidate list.
498 pub candidate_index: u8,
499 /// Candidate harness.
500 pub harness: String,
501 /// Stable skip reason code.
502 pub reason: String,
503}
504
505impl AgentSelection {
506 /// Validate semantic bounds and closed vocabularies at the durable event
507 /// boundary, independently of the user-config parser that constructed it.
508 pub fn validate(&self) -> Result<(), String> {
509 if self.schema_version != 1 {
510 return Err(format!(
511 "unsupported selection schema_version {}",
512 self.schema_version
513 ));
514 }
515 let name = self.profile.as_bytes();
516 if name.is_empty()
517 || name.len() > 63
518 || !name[0].is_ascii_lowercase()
519 || !name
520 .iter()
521 .all(|b| b.is_ascii_lowercase() || b.is_ascii_digit() || *b == b'-')
522 || self.profile.ends_with('-')
523 || self.profile.contains("--")
524 {
525 return Err("invalid profile name".into());
526 }
527 if !matches!(
528 self.selection_source.as_str(),
529 "cli"
530 | "environment"
531 | "repository-per-kind"
532 | "user-per-kind"
533 | "repository-default"
534 | "user-default"
535 | "builtin-harness"
536 ) {
537 return Err("invalid selection_source".into());
538 }
539 if !matches!(
540 self.interaction.as_str(),
541 "autonomous" | "explicit-interactive"
542 ) || !matches!(
543 self.capability.as_str(),
544 "fast" | "capable" | "ultra-capable"
545 ) || !matches!(self.residency.as_str(), "local" | "remote")
546 {
547 return Err("invalid interaction, capability, or residency".into());
548 }
549 validate_selected_candidate(&self.selected)?;
550 if self.interaction == "autonomous"
551 && (self.selected.harness != "pi"
552 || self.selected.telemetry.as_deref() != Some("worker-v1"))
553 {
554 return Err("autonomous selection requires pi with worker-v1 telemetry".into());
555 }
556 if self.selected.candidate_index >= 8
557 || self.fallback.len() != usize::from(self.selected.candidate_index)
558 {
559 return Err("candidate index/count exceeds profile bound".into());
560 }
561 let mut prior = None;
562 for (expected_index, skipped) in self.fallback.iter().enumerate() {
563 if usize::from(skipped.candidate_index) != expected_index
564 || prior.is_some_and(|value| skipped.candidate_index <= value)
565 || !matches!(skipped.harness.as_str(), "pi" | "claude")
566 || !matches!(
567 skipped.reason.as_str(),
568 "executable_missing"
569 | "autonomous_harness_unsupported"
570 | "telemetry_unsupported"
571 )
572 {
573 return Err("invalid fallback candidate index, harness, or reason".into());
574 }
575 if skipped.reason == "autonomous_harness_unsupported"
576 && (self.interaction != "autonomous" || skipped.harness == "pi")
577 {
578 return Err("inconsistent autonomous harness skip reason".into());
579 }
580 if skipped.reason == "telemetry_unsupported"
581 && (self.interaction != "autonomous" || skipped.harness != "pi")
582 {
583 return Err("inconsistent telemetry skip reason".into());
584 }
585 prior = Some(skipped.candidate_index);
586 }
587 Ok(())
588 }
589}
590
591fn validate_selected_candidate(candidate: &SelectedAgentCandidate) -> Result<(), String> {
592 if !matches!(candidate.harness.as_str(), "pi" | "claude")
593 || candidate.command.is_empty()
594 || candidate.command.len() > 32
595 || candidate
596 .command
597 .iter()
598 .any(|arg| arg.is_empty() || arg.len() > 4096 || arg.contains('\0'))
599 || candidate.command.iter().map(String::len).sum::<usize>() > 16_384
600 {
601 return Err("invalid selected harness or command".into());
602 }
603 match (candidate.harness.as_str(), candidate.telemetry.as_deref()) {
604 (_, None) | ("pi", Some("worker-v1")) => Ok(()),
605 _ => Err("invalid selected telemetry declaration".into()),
606 }
607}
608
609/// `manifest.json` (design.md §1.2).
610#[derive(Debug, Clone, Serialize, Deserialize)]
611pub struct Manifest {
612 /// State-schema version this file was written with.
613 pub schema_version: u32,
614 /// Watermark: the highest event `seq` whose projection fold is durably
615 /// committed. Events in `events.jsonl` with `seq > applied_seq` are
616 /// *unapplied tail* events — replayed into the projections on the next
617 /// lock acquisition before any new append (see
618 /// [`crate::events::append_and_apply_event`]). This is what makes
619 /// append-then-apply atomic across a reducer crash: the event log can run
620 /// ahead of the projections, but the gap is always healed before the next
621 /// writer observes stale state.
622 ///
623 /// `#[serde(default)]` so a legacy `manifest.json` written before this
624 /// field existed deserializes with `applied_seq = 0`. Such a manifest
625 /// self-migrates on its next write: the catch-up replay re-folds the whole
626 /// log — every event a no-op, because legacy state was already projected
627 /// synchronously under the old append-then-apply path — and advances the
628 /// watermark to `last_seq`. No separate migration pass or schema bump is
629 /// required (the field is purely additive to a derived-cache file).
630 #[serde(default)]
631 pub applied_seq: u64,
632 /// Unique run identifier (ULID). Validated on read.
633 pub run_id: RunId,
634 /// Kind of work this run performs.
635 pub kind: Kind,
636 /// How-run state (autonomous vs interactive), set once at `run create` from
637 /// the explicit `--interactive` flag — never transitioned. See [`Lifecycle`].
638 pub lifecycle: Lifecycle,
639 /// Human-readable run title.
640 pub title: String,
641 /// Current aggregate run status.
642 pub status: Status,
643 /// When the run was created.
644 pub created_at: DateTime<Utc>,
645 /// When the manifest was last modified.
646 pub updated_at: DateTime<Utc>,
647 /// Source repository the run operates on, if any.
648 pub source_repo: Option<String>,
649 /// Branch the run was started from, if any.
650 pub source_branch: Option<String>,
651 /// Root directory under which this run's worktrees live, if any.
652 pub worktree_root: Option<String>,
653 /// tmux session taskfleet created to host this run's headless windows
654 /// (`--headless` / `--tmux-session <name>`), if any. `None` for a foreground
655 /// run whose window lives in the user's own session — that session is never
656 /// a teardown target. When set, the supervisor kills this session once its
657 /// last taskfleet-owned window is torn down and only the synthetic
658 /// bootstrap shell window remains, so an empty headless session is not left
659 /// behind (issue `headless-tmux-session-not-torn-down`). `#[serde(default)]`
660 /// keeps a manifest written before this field existed readable.
661 #[serde(default)]
662 pub managed_tmux_session: Option<String>,
663 /// Completion-notification command registered at `run create --notify`,
664 /// if any. When the run reaches a terminal state (`done | failed |
665 /// cancelled`) the supervisor runs this command (at-least-once, deduped on a
666 /// durable `run.notified` marker event — the healthy path fires once, a
667 /// crash between firing and recording may re-fire) with `TASKFLEET_RUN_ID` /
668 /// `TASKFLEET_STATUS` / `TASKFLEET_SUMMARY` (and `TASKFLEET_RUN_KIND` / `TASKFLEET_RUN_TITLE`)
669 /// in its environment, BEFORE teardown removes the worktree/window. This is
670 /// how a spawning session learns of completion without polling (issue
671 /// `no-completion-notification-to-parent`). `None` for a run created without
672 /// `--notify`; `#[serde(default)]` keeps a manifest written before this
673 /// field existed readable.
674 #[serde(default)]
675 pub notify_cmd: Option<String>,
676 /// The agent runtime selected for this run's worker
677 /// (`claude` | `pi`), resolved at `run create`
678 /// via the flag > env > config > default precedence and recorded here as
679 /// provenance. This is the *selected* harness — recorded before the worker is
680 /// spawned, so it reflects intent even if the spawn later fails. `None` for a
681 /// manifest written before this field existed
682 /// (`#[serde(default)]`) — such legacy runs predate harness selection and
683 /// were all `claude`. Surfaced on `run show` / `run list --json`.
684 #[serde(default)]
685 pub harness: Option<String>,
686 /// Compact create-time profile resolution. `None` keeps manifests written
687 /// before profile selection readable without inventing requested/selected
688 /// history. Retry and read paths consume this recorded value; they never
689 /// re-resolve current configuration.
690 #[serde(default)]
691 pub agent_selection: Option<AgentSelection>,
692 /// Number of nodes created in this run (denormalized counter).
693 pub node_count: u32,
694 /// Run that spawned this run, if it is itself a child.
695 pub parent_run_id: Option<RunId>,
696 /// Node in the parent run that spawned this run, if any.
697 pub parent_node_id: Option<NodeId>,
698}
699
700/// `(child_run_id, child_node_id)` pointer recorded in `Node::children`.
701#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
702pub struct ChildRef {
703 /// Run id of the spawned child. Validated on read.
704 pub run_id: RunId,
705 /// Node id within the child run. Validated on read.
706 pub node_id: NodeId,
707}
708
709/// State of durable native worker evidence capture.
710#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq)]
711#[serde(rename_all = "snake_case")]
712pub enum EvidenceStatus {
713 /// Capture has not completed yet.
714 Pending,
715 /// The latest capture attempt failed and cleanup remains vetoed.
716 Failed,
717 /// Every durable artifact was synced and recorded.
718 Complete,
719}
720
721impl std::fmt::Display for EvidenceStatus {
722 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
723 f.write_str(match self {
724 Self::Pending => "pending",
725 Self::Failed => "failed",
726 Self::Complete => "complete",
727 })
728 }
729}
730
731/// Durable native Pi session and terminal evidence for one worker attempt.
732///
733/// This is folded onto the node projection from append-only events. The live
734/// source is state-root-relative; completed artifact paths are run-relative.
735#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
736pub struct WorkerEvidence {
737 /// Worker attempt this evidence belongs to.
738 pub attempt: u32,
739 /// Exact Pi session identifier assigned before the candidate starts.
740 pub session_id: String,
741 /// Original cwd stored in the native transcript header.
742 pub original_cwd: String,
743 /// State-root-relative live native transcript path.
744 pub live_session_path: String,
745 /// Typed capture state.
746 pub status: EvidenceStatus,
747 /// Run-relative byte-identical archived transcript, once complete.
748 #[serde(default, skip_serializing_if = "Option::is_none")]
749 pub transcript_path: Option<String>,
750 /// Run-relative resume copy whose header names a surviving cwd.
751 #[serde(default, skip_serializing_if = "Option::is_none")]
752 pub resume_path: Option<String>,
753 /// Run-relative final tmux pane snapshot.
754 #[serde(default, skip_serializing_if = "Option::is_none")]
755 pub pane_path: Option<String>,
756 /// Run-relative exact terminal report JSON.
757 #[serde(default, skip_serializing_if = "Option::is_none")]
758 pub report_path: Option<String>,
759 /// SHA-256 of the archived original transcript.
760 #[serde(default, skip_serializing_if = "Option::is_none")]
761 pub transcript_sha256: Option<String>,
762 /// Explicit capture failure detail. Never present on complete evidence.
763 #[serde(default, skip_serializing_if = "Option::is_none")]
764 pub error: Option<String>,
765}
766
767/// `nodes/<node-id>.json` (design.md §1.3).
768#[derive(Debug, Clone, Serialize, Deserialize)]
769pub struct Node {
770 /// State-schema version this file was written with.
771 pub schema_version: u32,
772 /// Unique node identifier within its run (e.g. `n-0001`). Validated on
773 /// read; this is the projection's filename key, so it can never name a
774 /// path outside `nodes/`.
775 pub node_id: NodeId,
776 /// Run this node belongs to. Validated on read.
777 pub run_id: RunId,
778 /// Parent node within the same run, if this is a sub-node.
779 pub parent_node_id: Option<NodeId>,
780 /// Kind of work this node performs.
781 pub kind: Kind,
782 /// Current node status.
783 pub status: Status,
784 /// Task description / prompt driving the node, if recorded.
785 pub task: Option<String>,
786 /// Filesystem path of the node's git worktree, if created.
787 pub worktree_path: Option<String>,
788 /// Git branch the node works on, if any.
789 pub branch: Option<String>,
790 /// The commit SHA the node's branch/worktree was forked from at spawn
791 /// (the branch tip the moment `create.sh` materialized the worktree). It
792 /// is the fixed reference point that lets the supervisor tell "this branch
793 /// produced work that merged into source" from "this branch never diverged
794 /// from its fork point": a branch still at `base_sha` is trivially an
795 /// ancestor of its source branch but has merged nothing, so it must NOT be
796 /// reconciled to success or torn down (that would drop a live agent's
797 /// uncommitted work). Only a branch whose tip has moved past `base_sha`
798 /// *and* is now an ancestor of the run's `source_branch` is a confirmed
799 /// merge (issues `false-failed-after-merge` /
800 /// `supervisor-stuck-pending-after-self-merge`). `#[serde(default)]` keeps a
801 /// node written before this field existed readable (`None` → the
802 /// git-reconcile fallback simply does not fire for it).
803 #[serde(default)]
804 pub base_sha: Option<String>,
805 /// tmux window hosting the node's agent, if interactive. This is the
806 /// human-readable window *name* — not unique across sessions and blind to
807 /// non-default sockets. Kept for display and as the legacy liveness key;
808 /// prefer [`Node::tmux_identity`] when present.
809 pub tmux_window: Option<String>,
810 /// Fully-qualified tmux identity (`session:window_id` + socket path)
811 /// captured at spawn time. `None` for nodes registered before create.sh
812 /// emitted the qualified fields — those fall back to bare-name matching on
813 /// [`Node::tmux_window`]. New spawns always populate this when create.sh
814 /// returns it.
815 #[serde(default)]
816 pub tmux_identity: Option<TmuxIdentity>,
817 /// Native worker evidence, when this attempt uses Pi. Initialized before
818 /// publication and advanced only by locked evidence events.
819 #[serde(default, skip_serializing_if = "Option::is_none")]
820 pub evidence: Option<WorkerEvidence>,
821 /// PID of the running agent process, if live.
822 pub agent_pid: Option<i32>,
823 /// Start time of the agent process, used to detect PID reuse.
824 pub agent_pid_start_time: Option<DateTime<Utc>>,
825 /// PID of the supervisor watching this node, if live.
826 pub supervisor_pid: Option<i32>,
827 /// Children this node has spawned.
828 #[serde(default)]
829 pub children: Vec<ChildRef>,
830 /// When the node started executing, if it has.
831 pub started_at: Option<DateTime<Utc>>,
832 /// When the node file was last modified.
833 pub updated_at: DateTime<Utc>,
834 /// The `node.report` payload that drove this node to its terminal status.
835 /// Set only by the report that actually transitions the node (Done /
836 /// Failed / Cancelled). Once the node is terminal it is frozen: a late
837 /// report against an already-settled node is dropped without overwriting
838 /// this field (see `reducer::apply_node_report`). So for a node cancelled
839 /// by `run cancel`, this holds the synthesized cancel report, not a
840 /// later-arriving agent report — that payload remains only in
841 /// `events.jsonl`.
842 pub last_report: Option<Value>,
843 /// Highest report `seq` consumed per child run id, for idempotent
844 /// report processing across supervisor restarts.
845 #[serde(default)]
846 pub last_processed_report_seq_by_child: Map<String, Value>,
847 /// Number of times the supervisor has auto-retried this node after an
848 /// empty-handed `agent-died` (issue `autoretry-agent-died-worker`). The
849 /// DURABLE, restart-safe bound on the bounded-retry loop: each `node.retry`
850 /// event increments it, and the watchdog terminalizes the run `failed` once
851 /// it reaches `RETRY_MAX_ATTEMPTS`. `#[serde(default)]` keeps a node written
852 /// before this field existed readable (`0` — never retried).
853 #[serde(default)]
854 pub retry_attempts: u32,
855 /// The **told** exit status of the node's worker process, recorded durably by
856 /// the `run-worker` launcher shim (`crates/taskfleet/src/run_worker.rs`) when
857 /// it `wait()`s on the agent it wrapped. This is a *fact*, not an inference:
858 /// the supervisor consumes it via the typed outcome table instead of guessing
859 /// completion from pid/pane/activity proxies (design.md §2.1, issue
860 /// `thin-exit-status-launcher`). A non-zero code or a terminating signal is a
861 /// `failed` worker; `code == 0` with no `explicit-merge` transition is the
862 /// *finished-but-unmerged* case that must stay non-terminal (attention-
863 /// required), NOT be auto-failed. `None` until the shim records an exit — or
864 /// forever, for a worker never launched through the shim (the crash backstop
865 /// still covers that path). `#[serde(default)]` keeps a node written before
866 /// this field existed readable.
867 #[serde(default)]
868 pub worker_exit: Option<WorkerExit>,
869 /// The in-flight `run merge` transaction for this node, if one has been
870 /// STARTED but not yet completed. `run merge` records a `merge.started`
871 /// event (setting this field) BEFORE it mutates git, because the merge spans
872 /// two durability domains — git refs and the event log — and is not atomic
873 /// across them (design.md §2.1b / A2, issue `merge-transaction-recovery`). A
874 /// crash after the git merge but before the terminal `explicit-merge`
875 /// `node.report` would otherwise strand the work *merged in source* with *no
876 /// merge event* → a false `failed`.
877 ///
878 /// This field is the durable op-log record that lets recovery finish or
879 /// reject that ONE known transaction deterministically, by OID — never a
880 /// general branch-content heuristic. It is set by [`crate::MergeTxn`]-carrying
881 /// `merge.started`, and cleared when the transaction resolves: a terminal
882 /// `node.report` (the merge completed) or a `merge.aborted` (recovery found
883 /// the git mutation never landed). `#[serde(default)]` keeps a node written
884 /// before this field existed readable (`None` — no in-flight merge).
885 ///
886 /// Boxed so the (rare) in-flight transaction does not inflate every `Node` /
887 /// `ProjectionOp` by the full [`MergeTxn`] footprint.
888 #[serde(default)]
889 pub pending_merge: Option<Box<MergeTxn>>,
890 /// The durable, monotonic timestamp of the FIRST tick on which the supervisor
891 /// observed this node's worker process confirmed-dead with no told
892 /// `worker.exited` and no merge — the anchor for the residual crash
893 /// backstop's fixed post-death grace (design.md §2.1a, issue
894 /// `typed-supervisor-outcomes`).
895 ///
896 /// The backstop is the ONLY place pid liveness still governs an outcome
897 /// (pid liveness is a pure crash backstop now, never a primary signal). When
898 /// the launcher shim's exit fact is lost — a hard kill of the shim, host
899 /// death — the supervisor never sees a `worker.exited` event, so it falls
900 /// back to "process confirmed gone → `failed`". The grace exists only to let
901 /// an in-flight `worker.exited` / merge append land before that fires: on the
902 /// first confirmed death the supervisor records this timestamp (via a
903 /// `node.death_observed` event) and DEFERS; it terminalizes `failed` only on a
904 /// later tick once a fixed short window has elapsed AND an exclusive-lock
905 /// re-read confirms no exit/merge landed in the race window.
906 ///
907 /// Persisted (not in-memory) so the grace survives a supervisor restart in the
908 /// window — a restart re-reads it rather than restarting the clock. First-write
909 /// -wins in the reducer, so the anchor is monotonic. `None` until the first
910 /// confirmed-death observation, or forever for a worker that exits cleanly
911 /// (the shim records `worker.exited` and the backstop never engages).
912 /// `#[serde(default)]` keeps a node written before this field existed readable.
913 #[serde(default)]
914 pub first_death_at: Option<DateTime<Utc>>,
915 /// An open, agent-authored request for a human decision. The worker records
916 /// this through `node.awaiting_input` instead of blocking on interactive
917 /// stdin. It remains non-terminal and is cleared by `node.input_resolved`, a
918 /// terminal `node.report`, or `node.retry`.
919 ///
920 /// `opened_at` is stamped from the event envelope and is therefore a durable,
921 /// monotonic grace-window anchor that survives supervisor restarts. The
922 /// original discussion objects are retained verbatim so read surfaces and
923 /// notification hooks can carry the question, options, and recommended
924 /// default without inventing a second advisory schema.
925 #[serde(default)]
926 pub awaiting_input: Option<Box<AwaitingInput>>,
927}
928
929/// Durable open-discussion state projected from `node.awaiting_input`.
930#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
931pub struct AwaitingInput {
932 /// Timestamp of the first open signal in the current unresolved generation.
933 pub opened_at: DateTime<Utc>,
934 /// Event sequence that opened this generation, used to deduplicate its
935 /// delayed parent notification independently from later generations.
936 pub event_seq: u64,
937 /// Validated report-shaped discussion objects. Each carries `topic`,
938 /// `options`, and `recommended_default`.
939 pub discussion_items: Vec<Value>,
940}
941
942/// A durable, in-flight `run merge` transaction recorded by `merge.started`
943/// BEFORE the git mutation, and the sole input to deterministic merge-crash
944/// recovery (design.md §2.1b / A2, issue `merge-transaction-recovery`).
945///
946/// `run merge` spans git refs and the event log and is not atomic across them.
947/// Recording the transaction — the exact source ref it will move, the OID it
948/// expects that ref to be at (`expected_source_oid`, the compare half of the
949/// compare-and-swap), and the worker's tip — lets the supervisor (or a retried
950/// `run merge`) resolve the ONE recorded transaction by OID after a crash:
951///
952/// - source ref still at `expected_source_oid` → the mutation never landed →
953/// REJECT (`merge.aborted`), preserving the worker's branch + work.
954/// - source ref moved off `expected_source_oid` AND the worker's content is
955/// integrated (rebase-robust content verification) → COMPLETE (append the
956/// `explicit-merge` `node.report` the crash prevented).
957/// - source ref moved unexpectedly but the worker's content is not integrated →
958/// fail closed (REJECT), preserving the work.
959#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
960pub struct MergeTxn {
961 /// Opaque unique id for this merge attempt. A fresh id per `run merge`
962 /// invocation (each attempt re-reads `expected_source_oid`), so recovery can
963 /// name exactly which transaction it resolved in the `merge.aborted` audit.
964 pub op_id: String,
965 /// The source/target ref this merge moves — `manifest.source_branch`
966 /// (`main`, or an integration branch). Recovery reads this ref's current OID
967 /// to decide the transaction's fate.
968 pub source_branch: String,
969 /// The worker branch whose commits are being merged (`node.branch`). Its
970 /// content is what recovery verifies is integrated into `source_branch`.
971 pub worker_branch: String,
972 /// The OID `source_branch` was at when the transaction was recorded — the
973 /// compare half of the compare-and-swap. If the ref is still here at recovery
974 /// time, the git mutation never landed.
975 pub expected_source_oid: String,
976 /// The worker branch tip at record time. Retained for the audit trail and as
977 /// a secondary landing signal; the authoritative completion check is
978 /// content-based (rebase-robust) against `source_branch`.
979 pub worker_oid: String,
980 /// The worker branch's fork point (`node.base_sha`), used to bound the
981 /// content check to the worker's own commits. `None` when unrecorded.
982 #[serde(default)]
983 pub base_sha: Option<String>,
984 /// PID of the `run merge` process driving the transaction, so recovery can
985 /// tell a still-in-progress merge (driver alive — leave it) from a crashed
986 /// one (driver gone — resolve it), never racing a live merge. `None` when
987 /// unrecorded.
988 #[serde(default)]
989 pub driver_pid: Option<i32>,
990 /// Start time of `driver_pid` in Unix seconds (the same representation the
991 /// pid-file liveness check records), guarding against PID reuse the way the
992 /// agent/supervisor liveness checks do — a recycled PID must not look alive.
993 /// `None` when the platform could not read it.
994 #[serde(default)]
995 pub driver_pid_start_secs: Option<u64>,
996 /// When the transaction was recorded.
997 pub started_at: DateTime<Utc>,
998}
999
1000/// The observed exit status of a node's worker process, recorded by the
1001/// `run-worker` launcher shim under the run lock (design.md §2.1 / A1).
1002///
1003/// Exactly one of `code` / `signal` is meaningful: a worker that returned
1004/// normally carries `code = Some(n)` (and `signal = None`); a worker killed by a
1005/// signal carries `signal = Some(s)` (and, on Unix, `code = None`). A recorded
1006/// exit is a durable *told fact* — the supervisor reads it rather than inferring
1007/// completion from liveness proxies.
1008#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
1009pub struct WorkerExit {
1010 /// Normal-exit status code, if the worker was not killed by a signal.
1011 #[serde(default)]
1012 pub code: Option<i32>,
1013 /// Terminating signal number, if the worker was killed by a signal.
1014 #[serde(default)]
1015 pub signal: Option<i32>,
1016 /// When the shim observed the worker's exit.
1017 pub at: DateTime<Utc>,
1018}
1019
1020impl WorkerExit {
1021 /// A clean exit: not signalled, and a zero return code. This is the *only*
1022 /// success-shaped worker exit — but a clean exit alone is NOT a completed
1023 /// unit (the worker may have finished-but-skipped `run merge`); merge is the
1024 /// only success truth (design.md §2.6). Callers pair this with a merge check.
1025 pub fn is_clean(self) -> bool {
1026 self.signal.is_none() && self.code == Some(0)
1027 }
1028
1029 /// A failed worker: killed by a signal, or a non-zero return code. Mutually
1030 /// exclusive with [`WorkerExit::is_clean`].
1031 pub fn is_failure(self) -> bool {
1032 !self.is_clean()
1033 }
1034}
1035
1036/// A fully-qualified tmux window identity recorded at spawn time.
1037///
1038/// `tmux_window` (the human name) is not unique across sessions, and a bare
1039/// `tmux list-windows -a` cannot see windows on a non-default socket. This
1040/// triple pins the exact window the agent runs in — `session:window_id` is
1041/// unique per server, `window_id` (the `@NNNN` form) survives renames, and
1042/// `socket` disambiguates multiple tmux servers. The watchdog matches on this
1043/// when present (design.md §8.1).
1044///
1045/// `pane_id` (the `%NN` form) pins the agent's *specific* pane within that
1046/// window, recorded at spawn. Window-owning operations (`kill-window` teardown —
1047/// the supervisor owns the whole window per the cleanup invariants) key off
1048/// `window_id`; only per-pane operations that must not follow the window's
1049/// *active* pane — chiefly `pipe-pane` agent-log capture — use `pane_id`. It is
1050/// `None` for a run spawned before create.sh emitted the field; capture then
1051/// falls back to `window_id` (issue `capture-agent-pane-by-pane-id`).
1052///
1053/// The watchdog's liveness probe still keys off `window_id` (correct for the
1054/// single-pane autonomous path). A pane-aware liveness probe — needed so a split
1055/// interactive window whose agent pane dies while a user shell pane survives is
1056/// still seen as dead — is a follow-up (`watchdog-pane-aware-liveness`), not this
1057/// change.
1058#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
1059pub struct TmuxIdentity {
1060 /// Server socket path (`#{socket_path}`). `None` if create.sh could not
1061 /// read it; the watchdog then queries tmux on its default socket.
1062 #[serde(default)]
1063 pub socket: Option<String>,
1064 /// Session that owns the window (`#{session_name}`).
1065 pub session: String,
1066 /// Stable window id in `@NNNN` form (`#{window_id}`). Survives renames and
1067 /// is unique within the server.
1068 pub window_id: String,
1069 /// Stable pane id in `%NN` form (`#{pane_id}`), recorded at spawn — the
1070 /// agent's own pane. `None` for a run whose create.sh predates the field
1071 /// (back-compat: old state deserializes with `pane_id: None`). Prefer
1072 /// [`TmuxIdentity::capture_target`] over reading this directly.
1073 #[serde(default)]
1074 pub pane_id: Option<String>,
1075}
1076
1077impl TmuxIdentity {
1078 /// The tmux target for a per-pane operation that must hit the agent's own
1079 /// pane, not the window's *active* pane: the recorded `pane_id` when
1080 /// present, else the `window_id` (which resolves to the active pane).
1081 ///
1082 /// Used by agent-log capture (`pipe-pane`). Window-level operations
1083 /// (`kill-window`, liveness) must NOT use this — they key off `window_id`
1084 /// directly so they act on the whole window.
1085 ///
1086 /// A recorded `pane_id` is preferred only when non-empty; an empty string
1087 /// (a directly-deserialized/corrupt state that the reducer/spawn normalizers
1088 /// never produce) is treated as absent so capture never targets `-t ""`.
1089 pub fn capture_target(&self) -> &str {
1090 self.pane_id
1091 .as_deref()
1092 .filter(|id| !id.is_empty())
1093 .unwrap_or(&self.window_id)
1094 }
1095}
1096
1097/// One event-log line (design.md §1.4).
1098///
1099/// `run_id` / `node_id` are the typed id newtypes, so deserializing an
1100/// `events.jsonl` line validates the whole envelope on read: a malformed
1101/// `run_id` or `node_id` fails the `serde` parse at the read boundary (the
1102/// id newtypes' validating `Deserialize`) rather than being carried as an
1103/// unvalidated `String` until some later path helper. The parse failure
1104/// surfaces as whatever error the reader maps a bad line to — e.g. a
1105/// newline-terminated bad line is [`Error::CorruptEventLog`] from both
1106/// [`read_all_events`] and [`find_prior_with_key`], which share one physical
1107/// reader and torn-tail policy. The reducer still performs its own per-event
1108/// checks (envelope `run_id` matches the run it is folded into; `data`-borne
1109/// ids parse), but the envelope ids can no longer be the unvalidated party.
1110///
1111/// [`read_all_events`]: crate::events::read_all_events
1112/// [`find_prior_with_key`]: crate::events
1113/// [`Error::CorruptEventLog`]: crate::Error::CorruptEventLog
1114#[derive(Debug, Clone, Serialize, Deserialize)]
1115pub struct Event {
1116 /// Wall-clock timestamp the event was appended.
1117 pub ts: DateTime<Utc>,
1118 /// Monotonic per-run sequence number (recovered on append).
1119 pub seq: u64,
1120 /// Event kind discriminator (e.g. `node.created`, `discussion.opened`).
1121 pub kind: String,
1122 /// Run the event belongs to. Validated on read.
1123 pub run_id: RunId,
1124 /// Node the event concerns, when applicable. Validated on read.
1125 #[serde(skip_serializing_if = "Option::is_none", default)]
1126 pub node_id: Option<NodeId>,
1127 /// Caller-supplied key used to dedupe retried appends.
1128 #[serde(skip_serializing_if = "Option::is_none", default)]
1129 pub idempotency_key: Option<String>,
1130 /// Kind-specific payload applied by the reducer.
1131 #[serde(default)]
1132 pub data: Value,
1133}
1134
1135#[cfg(test)]
1136mod tests {
1137 use super::*;
1138
1139 #[test]
1140 fn aggregate_terminal_status_is_the_three_way_rule() {
1141 use Status::{Blocked, Cancelled, Done, Failed, Pending, Running};
1142 // Empty set → not complete.
1143 assert_eq!(aggregate_terminal_status([]), None);
1144 // Any live node → not complete.
1145 for live in [Pending, Running, Blocked] {
1146 assert_eq!(aggregate_terminal_status([Done, live]), None);
1147 }
1148 // All done → Done.
1149 assert_eq!(aggregate_terminal_status([Done, Done]), Some(Done));
1150 // Any failure dominates.
1151 assert_eq!(aggregate_terminal_status([Done, Failed]), Some(Failed));
1152 assert_eq!(aggregate_terminal_status([Failed, Cancelled]), Some(Failed));
1153 // Cancelled (no failure) — pure or mixed with done.
1154 assert_eq!(
1155 aggregate_terminal_status([Cancelled, Cancelled]),
1156 Some(Cancelled)
1157 );
1158 assert_eq!(
1159 aggregate_terminal_status([Done, Cancelled]),
1160 Some(Cancelled)
1161 );
1162 }
1163
1164 /// `Kind::wire_name` (and thus `Kind::WIRE_NAMES`) must stay identical
1165 /// to what serde actually (de)serializes. If the `rename_all` routing
1166 /// or a variant name ever diverges from `wire_name`, this fails — which
1167 /// is what keeps the report validator's `expected` hint honest.
1168 #[test]
1169 fn wire_names_match_serde_round_trip() {
1170 for &name in Kind::WIRE_NAMES {
1171 let kind: Kind = serde_json::from_value(Value::String(name.to_string()))
1172 .unwrap_or_else(|_| panic!("WIRE_NAMES entry {name:?} is not a valid Kind"));
1173 assert_eq!(
1174 serde_json::to_value(kind).unwrap(),
1175 Value::String(name.to_string()),
1176 "serde round-trip diverged from wire_name for {name:?}",
1177 );
1178 }
1179 }
1180
1181 /// The bounded auto-retry eligibility gate (issue `autoretry-agent-died-worker`)
1182 /// must include exactly the autonomous single-node worker kinds and exclude
1183 /// the fan-out driver (and the read-only `Unknown` catch-all).
1184 #[test]
1185 fn autonomous_single_node_worker_set_is_exact() {
1186 for k in [Kind::Spinoff, Kind::Research, Kind::TechnicalDecision] {
1187 assert!(
1188 k.is_autonomous_single_node_worker(),
1189 "{k:?} should be retry-eligible"
1190 );
1191 assert_eq!(k.lifecycle(), Lifecycle::Autonomous);
1192 }
1193 for k in [
1194 Kind::FanOut, // multi-unit driver
1195 Kind::Unknown, // legacy on-disk run — never freshly supervised
1196 ] {
1197 assert!(
1198 !k.is_autonomous_single_node_worker(),
1199 "{k:?} must NOT be retry-eligible"
1200 );
1201 }
1202 }
1203
1204 /// A legacy run recorded under a since-removed kind must still deserialize
1205 /// to the read-only [`Kind::Unknown`] catch-all rather than faulting the
1206 /// read — the ADR §D7 "report, never delete" contract for the on-disk
1207 /// evidence corpus. Every creatable kind still round-trips to itself.
1208 #[test]
1209 fn removed_kinds_deserialize_to_unknown() {
1210 for removed in [
1211 "code",
1212 "orchestrate",
1213 "orchestrated",
1214 "bugfix",
1215 "make-skill",
1216 ] {
1217 let kind: Kind = serde_json::from_value(Value::String(removed.to_string()))
1218 .expect("a removed kind must still deserialize, not fault");
1219 assert_eq!(kind, Kind::Unknown, "{removed:?} should map to Unknown");
1220 }
1221 // A wholly unknown value maps there too (forward-compat).
1222 assert_eq!(
1223 serde_json::from_value::<Kind>(Value::String("future-kind".into())).unwrap(),
1224 Kind::Unknown
1225 );
1226 // The surviving kinds are unaffected.
1227 for &name in Kind::WIRE_NAMES {
1228 let kind: Kind = serde_json::from_value(Value::String(name.to_string())).unwrap();
1229 assert_ne!(kind, Kind::Unknown, "{name:?} must not fold to Unknown");
1230 }
1231 }
1232
1233 /// Back-compat acceptance criterion (issue `capture-agent-pane-by-pane-id`):
1234 /// a `TmuxIdentity` persisted before `pane_id` existed — with the field
1235 /// entirely absent, or written as an explicit `null` — must still
1236 /// deserialize, yielding `pane_id: None` and a `window_id` capture target.
1237 #[test]
1238 fn tmux_identity_deserializes_legacy_state_without_pane_id() {
1239 // Field entirely absent (a state file written by an older binary).
1240 let absent: TmuxIdentity = serde_json::from_value(serde_json::json!({
1241 "socket": null,
1242 "session": "taskfleet",
1243 "window_id": "@42",
1244 }))
1245 .expect("legacy identity without pane_id must deserialize");
1246 assert_eq!(absent.pane_id, None);
1247 assert_eq!(absent.capture_target(), "@42");
1248
1249 // Field present but explicitly null.
1250 let null: TmuxIdentity = serde_json::from_value(serde_json::json!({
1251 "socket": null,
1252 "session": "taskfleet",
1253 "window_id": "@42",
1254 "pane_id": null,
1255 }))
1256 .expect("identity with explicit null pane_id must deserialize");
1257 assert_eq!(null.pane_id, None);
1258 assert_eq!(null.capture_target(), "@42");
1259 }
1260
1261 /// `capture_target` prefers a recorded `pane_id` (`%NN`) over the window id,
1262 /// but treats an empty `pane_id` as absent (never targets `-t ""`).
1263 #[test]
1264 fn capture_target_prefers_nonempty_pane_id() {
1265 let with_pane = TmuxIdentity {
1266 socket: None,
1267 session: "taskfleet".into(),
1268 window_id: "@42".into(),
1269 pane_id: Some("%7".into()),
1270 };
1271 assert_eq!(with_pane.capture_target(), "%7");
1272
1273 let empty_pane = TmuxIdentity {
1274 pane_id: Some(String::new()),
1275 ..with_pane.clone()
1276 };
1277 assert_eq!(empty_pane.capture_target(), "@42");
1278 }
1279}
1280
1281#[cfg(test)]
1282mod id_tests {
1283 use super::*;
1284
1285 /// Inputs every id type must reject — the path-traversal vectors plus the
1286 /// generic malformed cases called out in the issue's success criteria.
1287 const TRAVERSAL_VECTORS: &[&str] = &[
1288 "..",
1289 "../etc",
1290 "a/b",
1291 "a/../b",
1292 ".hidden",
1293 "./x",
1294 "foo/bar.json",
1295 "n-0001/../../etc",
1296 "",
1297 ];
1298
1299 #[test]
1300 fn run_id_accepts_generator_output_and_rejects_malformed() {
1301 let id = crate::new_run_id();
1302 assert!(
1303 RunId::parse_str(&id).is_ok(),
1304 "generator must validate: {id}"
1305 );
1306 for bad in [
1307 "tooshort",
1308 "01jxsnap0000000000000000000", // 27 chars
1309 "01JXSNAP000000000000000000", // uppercase
1310 "01jxiiiiiiiiiiiiiiiiiiiiii", // `i` not in Crockford
1311 "80000000000000000000000000", // first char exceeds ULID range
1312 "n-0001", // wrong shape entirely
1313 ] {
1314 assert!(RunId::parse_str(bad).is_err(), "expected reject: {bad:?}");
1315 }
1316 for bad in TRAVERSAL_VECTORS {
1317 assert!(
1318 RunId::parse_str(bad).is_err(),
1319 "traversal not rejected: {bad:?}"
1320 );
1321 }
1322 }
1323
1324 #[test]
1325 fn node_id_accepts_canonical_and_rejects_malformed() {
1326 for ok in ["n-0001", "n-0010", "n-123456"] {
1327 assert!(NodeId::parse_str(ok).is_ok(), "expected accept: {ok}");
1328 }
1329 // Wrong prefix is its own error variant.
1330 assert!(matches!(
1331 NodeId::parse_str("d-0001"),
1332 Err(IdValidationError::WrongPrefix { .. })
1333 ));
1334 assert!(matches!(
1335 NodeId::parse_str("0001"),
1336 Err(IdValidationError::WrongPrefix { .. })
1337 ));
1338 for bad in [
1339 "n-1", // too few digits
1340 "n-abcd", // non-digit body
1341 "n-", // empty body
1342 "n-00a1", // mixed
1343 "n-00000000000", // 11 digits — over the 10-digit ceiling
1344 ] {
1345 assert!(
1346 matches!(
1347 NodeId::parse_str(bad),
1348 Err(IdValidationError::InvalidFormat { .. })
1349 ),
1350 "expected InvalidFormat: {bad:?}",
1351 );
1352 }
1353 for bad in TRAVERSAL_VECTORS {
1354 assert!(
1355 NodeId::parse_str(bad).is_err(),
1356 "traversal not rejected: {bad:?}"
1357 );
1358 }
1359 }
1360
1361 #[test]
1362 fn deserialize_rejects_malformed_ids() {
1363 // The validating Deserialize impl is the on-read guard: a tampered
1364 // projection file whose key no longer validates must fail to parse.
1365 assert!(serde_json::from_str::<NodeId>("\"n-0001\"").is_ok());
1366 assert!(serde_json::from_str::<NodeId>("\"../../etc\"").is_err());
1367 assert!(serde_json::from_str::<NodeId>("\"n-../escape\"").is_err());
1368 }
1369
1370 #[test]
1371 fn serialize_round_trips_as_bare_string() {
1372 let nid = NodeId::parse_str("n-0042").unwrap();
1373 let json = serde_json::to_string(&nid).unwrap();
1374 assert_eq!(json, "\"n-0042\"");
1375 let back: NodeId = serde_json::from_str(&json).unwrap();
1376 assert_eq!(back, nid);
1377 assert_eq!(nid.as_str(), "n-0042");
1378 assert_eq!(nid.to_string(), "n-0042");
1379 }
1380
1381 #[test]
1382 fn error_exposes_kind_and_expected() {
1383 let err = NodeId::parse_str("n-x").unwrap_err();
1384 assert_eq!(err.kind(), "node");
1385 assert_eq!(err.expected(), "n-NNNN (n- followed by 4-10 ASCII digits)");
1386 }
1387
1388 #[test]
1389 fn event_deserialize_validates_envelope_ids() {
1390 // The whole `events.jsonl` envelope is now validated on read: the
1391 // typed `run_id` / `node_id` fields parse through the id newtypes, so
1392 // a malformed envelope id fails the deserialize rather than being
1393 // carried downstream as an unchecked string.
1394 let ok = r#"{"ts":"2026-06-12T00:00:00Z","seq":1,"kind":"node.created","run_id":"01jxsnap000000000000000000","node_id":"n-0001","data":{}}"#;
1395 assert!(serde_json::from_str::<Event>(ok).is_ok());
1396
1397 // Invalid `run_id` (not a 26-char ULID) fails the parse.
1398 let bad_run = r#"{"ts":"2026-06-12T00:00:00Z","seq":1,"kind":"run.status","run_id":"not-a-ulid","data":{}}"#;
1399 assert!(serde_json::from_str::<Event>(bad_run).is_err());
1400
1401 // Invalid top-level `node_id` (too few digits) also fails the parse.
1402 let bad_node = r#"{"ts":"2026-06-12T00:00:00Z","seq":1,"kind":"node.status","run_id":"01jxsnap000000000000000000","node_id":"n-1","data":{}}"#;
1403 assert!(serde_json::from_str::<Event>(bad_node).is_err());
1404 }
1405
1406 #[test]
1407 fn from_str_and_ord_delegate_to_inner() {
1408 use std::str::FromStr;
1409 // `FromStr` mirrors `parse_str`, so the `str::parse` ecosystem works.
1410 assert!(RunId::from_str("01jxsnap000000000000000000").is_ok());
1411 assert!("n-0001".parse::<NodeId>().is_ok());
1412 assert!("n-x".parse::<NodeId>().is_err());
1413
1414 // `Ord` is lexicographic over the inner string; for ULIDs that is the
1415 // natural time-encoded order.
1416 let a = RunId::parse_str("01jxsnap000000000000000000").unwrap();
1417 let b = RunId::parse_str("02jxsnap000000000000000000").unwrap();
1418 assert!(a < b);
1419 let mut v = vec![b.clone(), a.clone()];
1420 v.sort();
1421 assert_eq!(v, vec![a, b]);
1422 }
1423}
1424
1425#[cfg(test)]
1426mod agent_selection_validation_tests {
1427 use super::*;
1428
1429 fn valid() -> AgentSelection {
1430 AgentSelection {
1431 schema_version: 1,
1432 profile: "capable".into(),
1433 selection_source: "cli".into(),
1434 interaction: "autonomous".into(),
1435 capability: "capable".into(),
1436 residency: "remote".into(),
1437 requested_harness: None,
1438 selected: SelectedAgentCandidate {
1439 candidate_index: 1,
1440 harness: "pi".into(),
1441 command: vec!["pi".into()],
1442 telemetry: Some("worker-v1".into()),
1443 },
1444 fallback: vec![SkippedAgentCandidate {
1445 candidate_index: 0,
1446 harness: "claude".into(),
1447 reason: "autonomous_harness_unsupported".into(),
1448 }],
1449 }
1450 }
1451
1452 #[test]
1453 fn durable_selection_rejects_impossible_state() {
1454 assert!(valid().validate().is_ok());
1455 let mut invalid = valid();
1456 invalid.selected.candidate_index = 200;
1457 assert!(invalid.validate().is_err());
1458 let mut invalid = valid();
1459 invalid.fallback[0].reason = "made_up".into();
1460 assert!(invalid.validate().is_err());
1461 let mut invalid = valid();
1462 invalid.selected.harness = "claude".into();
1463 assert!(invalid.validate().is_err());
1464 }
1465}