Skip to main content

subc_control/
lib.rs

1//! Client-facing subc channel-0 control wire shapes.
2//!
3//! This crate is the client ↔ subc control-plane boundary. It depends only on
4//! [`subc-protocol`] for shared primitives such as `RouteTarget` and
5//! `BindIdentity`; clients can use it without depending on the
6//! daemon implementation.
7
8#![forbid(unsafe_code)]
9
10use std::path::PathBuf;
11
12use serde::{
13    de::{Error as _, MapAccess, SeqAccess, Visitor},
14    ser::SerializeMap,
15    Deserialize, Deserializer, Serialize, Serializer,
16};
17use subc_protocol::{
18    manifest::{CapabilityDeclarations, ManifestProvenance, ProviderRole, SelfSignalDeclaration},
19    session::HealthStatus,
20    BindIdentity, RouteTarget,
21};
22
23pub use subc_protocol::RouteCloseReason;
24
25macro_rules! open_string_enum {
26    (
27        $(#[$meta:meta])*
28        $name:ident {
29            $( $(#[$variant_meta:meta])* $variant:ident => $wire_name:literal ),+ $(,)?
30        }
31    ) => {
32        $(#[$meta])*
33        #[derive(Debug, Clone, PartialEq, Eq)]
34        pub enum $name {
35            $( $(#[$variant_meta])* $variant, )+
36            Unknown(String),
37        }
38
39        impl $name {
40            fn wire_name(&self) -> &str {
41                match self {
42                    $( Self::$variant => $wire_name, )+
43                    Self::Unknown(value) => value,
44                }
45            }
46        }
47
48        impl Serialize for $name {
49            fn serialize<S>(&self, serializer: S) -> Result<S::Ok, S::Error>
50            where
51                S: serde::Serializer,
52            {
53                serializer.serialize_str(self.wire_name())
54            }
55        }
56
57        impl<'de> Deserialize<'de> for $name {
58            fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
59            where
60                D: serde::Deserializer<'de>,
61            {
62                let value = String::deserialize(deserializer)?;
63                Ok(match value.as_str() {
64                    $( $wire_name => Self::$variant, )+
65                    _ => Self::Unknown(value),
66                })
67            }
68        }
69    };
70}
71
72/// Daemon-spawned consumer identity presented on route.open.
73#[derive(Clone, Serialize, Deserialize, PartialEq, Eq, Hash)]
74pub struct ConsumerIdentity {
75    pub module_id: String,
76    pub launch_nonce: String,
77}
78
79// Hand-written so the launch nonce is never printed. The nonce is the credential
80// that attributes a connection to a supervised module, and a derived Debug would
81// write it into any log line or panic message that formats this value. Same
82// reasoning as ConnectionInfo's Debug in subc-transport.
83impl std::fmt::Debug for ConsumerIdentity {
84    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
85        f.debug_struct("ConsumerIdentity")
86            .field("module_id", &self.module_id)
87            .field(
88                "launch_nonce",
89                &format_args!("<{} bytes redacted>", self.launch_nonce.len()),
90            )
91            .finish()
92    }
93}
94
95/// Reserved dotted operation prefixes for the v0.4 control vocabulary.
96///
97/// `scheduler.` and `watch.` were reserved here from v0.4 until 2026-08-10 and
98/// were removed deliberately rather than left as placeholders: neither was ever
99/// implemented, and both capabilities are now owned elsewhere by ruling --
100/// scheduled tasks belong to the session runtime (prefrontal) because the
101/// daemon is state-free routing, and external-event watching belongs to the
102/// connectors module (plexus). A reserved name for something that will never be
103/// built here reads as a roadmap commitment to anyone surveying the protocol,
104/// and it recruited exactly that misunderstanding from an outside contributor.
105pub mod ops {
106    pub const SERVER: &str = "server.";
107    pub const CATALOG: &str = "catalog.";
108    pub const ROUTE: &str = "route.";
109    pub const SUPERVISOR: &str = "supervisor.";
110    pub const CONFIG: &str = "config.";
111
112    pub const SERVER_DESCRIBE: &str = "server.describe";
113    pub const CATALOG_LIST: &str = "catalog.list";
114    pub const ROUTE_OPEN: &str = "route.open";
115    pub const ROUTE_POLL: &str = "route.poll";
116    pub const ROUTE_CLOSING: &str = "route.closing";
117    pub const ROUTE_CLOSED: &str = "route.closed";
118    pub const SUPERVISOR_LIST: &str = "supervisor.list";
119    pub const SUPERVISOR_RESTART: &str = "supervisor.restart";
120    pub const SUPERVISOR_SWAP: &str = "supervisor.swap";
121    pub const SUPERVISOR_RELOAD: &str = "supervisor.reload";
122    pub const SUPERVISOR_RESCAN: &str = "supervisor.rescan";
123    pub const SUPERVISOR_RELEASE_RESERVED: &str = "supervisor.release_reserved";
124    pub const SUPERVISOR_SET_ENABLED: &str = "supervisor.set_enabled";
125    pub const SUPERVISOR_HEALTH_PROBE: &str = "supervisor.health_probe";
126    pub const SUPERVISOR_HEALTH: &str = "supervisor.health";
127    pub const SUPERVISOR_STDERR_TAIL: &str = "supervisor.stderr_tail";
128    pub const SUPERVISOR_TERMINALS: &str = "supervisor.terminals";
129    pub const SUPERVISOR_ROUTES: &str = "supervisor.routes";
130    pub const SUPERVISOR_PROVENANCE: &str = "supervisor.provenance";
131    pub const SUPERVISOR_SPAWN_SNAPSHOT: &str = "supervisor.spawn_snapshot";
132    pub const SUPERVISOR_SPAWN_SUBSCRIBE: &str = "supervisor.spawn_subscribe";
133}
134
135/// Client-originated channel-0 control RPC body.
136#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
137#[serde(tag = "op")]
138// RouteOpen carries the complete route metadata, while several control operations
139// are markers; retain the direct public wire shape instead of boxing its fields.
140#[allow(clippy::large_enum_variant)]
141pub enum ClientControlRequest {
142    #[serde(rename = "server.describe")]
143    ServerDescribe {},
144    #[serde(rename = "catalog.list")]
145    CatalogList {
146        /// Absent lists every registered module; present narrows to one. A
147        /// narrowed list for an unregistered id is an empty list rather than an
148        /// error, so absent and unregistered are distinguishable only by which
149        /// question you asked.
150        #[serde(default)]
151        module_id: Option<String>,
152    },
153    #[serde(rename = "route.open")]
154    RouteOpen {
155        target: RouteTarget,
156        identity: BindIdentity,
157        /// The consumer's claim to a supervised launch, which the daemon verifies
158        /// against its live spawn nonces before stamping a principal.
159        ///
160        /// Absent is a legitimate shape, not an omission: a direct key-holder has
161        /// no launch nonce to present, and the daemon stamps `Direct`. So absence
162        /// means NO CLAIM WAS MADE, never that a claim was refused — a refused
163        /// claim is an error frame and the route never opens. A provider deciding
164        /// what to trust reads the stamped principal on the bind, not this.
165        #[serde(default, skip_serializing_if = "Option::is_none")]
166        consumer_identity: Option<ConsumerIdentity>,
167        /// Consumer-declared reverse-request capabilities for the route. This is
168        /// an unverified declaration, not a privilege grant; if a consumer
169        /// over-declares, providers may still send reverse requests that later
170        /// time out or deny. Providers must treat an absent field as no
171        /// reverse-request capability. The vocabulary is open strings; known MCP
172        /// method-family values today are "elicitation", "sampling", and
173        /// "roots".
174        #[serde(default, skip_serializing_if = "Option::is_none")]
175        consumer_capabilities: Option<Vec<String>>,
176        /// Opaque admission facts supplied by the configured carrier module.
177        #[serde(default, skip_serializing_if = "Option::is_none")]
178        admission_facts: Option<serde_json::Value>,
179    },
180    #[serde(rename = "route.poll")]
181    RoutePoll {
182        route_channel: u16,
183        route_epoch: u32,
184        kind: PollKind,
185    },
186    #[serde(rename = "supervisor.list")]
187    SupervisorList {},
188    /// Read the live supervised processes and the event cursor atomically.
189    #[serde(rename = "supervisor.spawn_snapshot")]
190    SupervisorSpawnSnapshot {},
191    /// Replay spawn events after `since`, then remain open for live events.
192    ///
193    /// The cursor is one value copied from a snapshot or event. It includes the
194    /// daemon incarnation so a restarted daemon rejects an earlier instance's
195    /// sequence instead of treating it as a position in the current stream.
196    ///
197    /// Refusals and terminal errors, as `Error` frames on the request's corr:
198    /// - `spawn_cursor_incarnation_mismatch` (detail `current_daemon_incarnation`):
199    ///   `since` names another daemon incarnation.
200    /// - `spawn_cursor_too_old` (detail `oldest_retained_cursor`): `since`
201    ///   predates the retained event ring.
202    /// - `spawn_subscriber_lagged` (detail `first_undelivered_cursor`): the open
203    ///   stream fell too far behind and the daemon dropped it. It arrives after
204    ///   every event that was already queued and ends the stream; resubscribe
205    ///   with `since` set to the last cursor received.
206    #[serde(rename = "supervisor.spawn_subscribe")]
207    SupervisorSpawnSubscribe {
208        #[serde(default, skip_serializing_if = "Option::is_none")]
209        since: Option<SpawnCursor>,
210    },
211    #[serde(rename = "supervisor.restart")]
212    SupervisorRestart {
213        module_id: String,
214        /// Optional per-restart override of the module's drain budget, in ms.
215        /// Absent: the module's configured `drain_timeout_ms` (or the daemon
216        /// default) applies. `0` tears down without waiting — the wedge-bounce
217        /// escape, where a stuck in-flight request would never settle anyway.
218        /// Additive; older daemons that predate this field reject unknown
219        /// fields on channel-0 requests, so senders must omit it unless asked
220        /// for (the CLI only sends it when a flag is passed).
221        #[serde(default, skip_serializing_if = "Option::is_none")]
222        drain_timeout_ms: Option<u64>,
223    },
224    /// Blue/green restart: start a replacement beside the running process, keep
225    /// routing to the running one until the replacement declares itself ready,
226    /// then move new routes over and drain the old process.
227    ///
228    /// Refused before anything is spawned unless the module's daemon config
229    /// declares `overlap: "safe"`: two processes on one single-writer store is
230    /// a data hazard, so the default is exclusive. Answered once the swap has
231    /// either cut over (the old process is still draining) or failed; a failure
232    /// leaves the old process serving and undrained, and names the arm in
233    /// `ErrorBody.detail`.
234    ///
235    /// A new op rather than a flag on `supervisor.restart`: a daemon that
236    /// predates swap rejects an unknown op, whereas an unknown field could be
237    /// dropped and turned into a plain restart.
238    #[serde(rename = "supervisor.swap")]
239    SupervisorSwap {
240        module_id: String,
241        /// How long the replacement may take to register and declare itself
242        /// ready before the swap is abandoned. Absent: the daemon default.
243        #[serde(default, skip_serializing_if = "Option::is_none")]
244        ready_timeout_ms: Option<u64>,
245    },
246    #[serde(rename = "supervisor.reload")]
247    SupervisorReload { module_id: String },
248    #[serde(rename = "supervisor.rescan")]
249    SupervisorRescan {
250        /// Compute the reconciliation and return it WITHOUT applying it.
251        ///
252        /// Rescan retires any supervised module absent from the config, which
253        /// stops live processes. Both halves of that decision are inspectable in
254        /// advance -- the config is a file, the running set is `supervisor.list`
255        /// -- but nothing reconstructs the diff for the operator, so it is read
256        /// from the result table AFTER the retires have happened.
257        ///
258        /// A preview must be computed daemon-side rather than by a client, because
259        /// a client would have to locate the daemon's config itself: two rules
260        /// selecting one subject, agreeing until someone runs a daemon with a
261        /// non-default config. A preview that can describe a different file than
262        /// the operation reads is worse than none, because it is believed.
263        ///
264        /// Defaults to false so an existing client sending `{}` still executes,
265        /// and is OMITTED when false so the bytes an existing client sends are
266        /// unchanged. Serialising `preview:false` would have altered the request's
267        /// wire form for every caller that never asked for a preview -- caught by
268        /// the golden fixture, which is the whole reason that pin exists.
269        #[serde(default, skip_serializing_if = "std::ops::Not::not")]
270        preview: bool,
271    },
272    /// Retire the retained exact-id reservation after its configuration entry has
273    /// been removed. This is intentionally separate from rescan so deleting
274    /// configuration never silently opens a protected module id to registration.
275    #[serde(rename = "supervisor.release_reserved")]
276    SupervisorReleaseReserved { module_id: String },
277    #[serde(rename = "supervisor.set_enabled")]
278    SupervisorSetEnabled { module_id: String, enabled: bool },
279    #[serde(rename = "supervisor.health_probe")]
280    SupervisorHealthProbe { module_id: String },
281    #[serde(rename = "supervisor.health")]
282    SupervisorHealth {},
283    /// Enumerate the routes currently served by one supervised module, or every
284    /// module when omitted.
285    ///
286    /// This privileged census is control-plane-only. It is deliberately not an
287    /// MCP facade or agent-tool operation: callers holding the daemon control
288    /// connection may inspect live route ownership, while agent-facing modules
289    /// must not be able to address that surface at all.
290    ///
291    /// The daemon answers from its forwarding table under a read lock and never
292    /// consults a module. That makes the read safe during a drain, when a module
293    /// cannot be queried without recreating the hang/restart hazard that route
294    /// status reads avoid.
295    #[serde(rename = "supervisor.routes")]
296    SupervisorRoutes {
297        #[serde(default, skip_serializing_if = "Option::is_none")]
298        module_id: Option<String>,
299    },
300    /// Report source-tagged provenance for supervised modules, optionally narrowed
301    /// to one module.
302    #[serde(rename = "supervisor.provenance")]
303    SupervisorProvenance {
304        #[serde(default, skip_serializing_if = "Option::is_none")]
305        module_id: Option<String>,
306    },
307    /// Retained stderr for one module.
308    ///
309    /// A separate op rather than a field on `supervisor.list`: the tail is
310    /// kilobytes per module and `list` renders every module, so carrying it in
311    /// the snapshot would charge every status read for a payload almost no
312    /// caller wants. Caps ride on the REQUEST so a caller wanting twenty lines
313    /// and one wanting the whole ring need no separate fields anywhere.
314    #[serde(rename = "supervisor.stderr_tail")]
315    SupervisorStderrTail {
316        module_id: String,
317        #[serde(default, skip_serializing_if = "Option::is_none")]
318        max_lines: Option<u32>,
319        #[serde(default, skip_serializing_if = "Option::is_none")]
320        max_bytes: Option<u32>,
321    },
322    /// Retained terminal exits for one module.
323    ///
324    /// This stays separate from `supervisor.list`: a history grows with every
325    /// incident, while the list is a current-state read most callers issue often.
326    ///
327    /// The read MUST stay off the supervisor command channel — it reads the
328    /// module's shared ring directly. This is a requirement, not an
329    /// optimisation: when the supervision task itself dies, every
330    /// command-channel op returns `CommandClosed`, and that is precisely the
331    /// moment an operator needs the exit history most. A history reachable only
332    /// through the machinery whose death you are diagnosing is unreachable when
333    /// it matters. Proven failure mode, not a hypothetical.
334    #[serde(rename = "supervisor.terminals")]
335    SupervisorTerminals { module_id: String },
336}
337
338/// subc's channel-0 response body for client control RPCs.
339#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
340#[serde(tag = "op")]
341pub enum ClientControlResponse {
342    #[serde(rename = "server.describe")]
343    ServerDescribe {
344        protocol_ver: u8,
345        subc_ops: Vec<String>,
346        capabilities: Vec<String>,
347        connected_clients: u64,
348        #[serde(default, skip_serializing_if = "Option::is_none")]
349        counters: Option<serde_json::Value>,
350        /// Git commit the daemon was built from, or "unavailable" when the
351        /// build could not read it. The crate version cannot discriminate a
352        /// skewed daemon/CLI pair (it moves per release, not per commit), so
353        /// this is the identity a consumer compares against its own embedded
354        /// commit to detect that it is talking to an older build than it was
355        /// compiled with. Absent from daemons predating the field.
356        #[serde(default, skip_serializing_if = "Option::is_none")]
357        build_git_sha: Option<String>,
358        /// sha256 of the workspace Cargo.lock at build time, or "unavailable".
359        /// Answers "which dependency set" where the commit answers "which
360        /// source"; a commit match with a digest mismatch means a rebuild
361        /// against edited dependencies. Absent from daemons predating the
362        /// field.
363        #[serde(default, skip_serializing_if = "Option::is_none")]
364        build_lock_digest: Option<String>,
365        /// Daemon-evaluated capability requirements. Present when the configured
366        /// fleet has declarations to evaluate, so operators can inspect an absent
367        /// required capability without parsing daemon logs.
368        #[serde(default, skip_serializing_if = "Vec::is_empty")]
369        capability_requirements: Vec<CapabilityRequirementStatus>,
370        /// The daemon's machine id (`subc_protocol::MachineId`), the same value
371        /// every module receives on `HELLO_ACK`. A name for this machine, never
372        /// an authority: nothing may grant trust because two parties report the
373        /// same value. Absent from daemons predating the field.
374        #[serde(default, skip_serializing_if = "Option::is_none")]
375        machine_id: Option<String>,
376    },
377    #[serde(rename = "catalog.list")]
378    CatalogList {
379        generation: u64,
380        modules: Vec<CatalogEntry>,
381        subc_ops: Vec<String>,
382    },
383    #[serde(rename = "route.open")]
384    RouteOpen {
385        route_channel: u16,
386        route_epoch: u32,
387    },
388    #[serde(rename = "route.poll")]
389    RoutePoll {
390        route_channel: u16,
391        route_epoch: u32,
392        status: Option<String>,
393        live: Option<bool>,
394    },
395    #[serde(rename = "supervisor.list")]
396    SupervisorList {
397        generation: u64,
398        modules: Vec<SupervisorEntry>,
399    },
400    #[serde(rename = "supervisor.spawn_snapshot")]
401    SupervisorSpawnSnapshot {
402        #[serde(flatten)]
403        snapshot: SpawnSnapshot,
404    },
405    #[serde(rename = "supervisor.ack")]
406    SupervisorAck { module_id: String, applied: bool },
407    #[serde(rename = "supervisor.rescan")]
408    SupervisorRescan {
409        #[serde(flatten)]
410        result: SupervisorRescanResult,
411    },
412    #[serde(rename = "supervisor.health_probe")]
413    SupervisorHealthProbe {
414        module_id: String,
415        status: HealthStatus,
416        #[serde(default, skip_serializing_if = "Option::is_none")]
417        detail: Option<String>,
418        #[serde(default, skip_serializing_if = "Option::is_none")]
419        metrics: Option<serde_json::Value>,
420    },
421    #[serde(rename = "supervisor.health")]
422    SupervisorHealth {
423        generation: u64,
424        modules: Vec<SupervisorHealthEntry>,
425    },
426    #[serde(rename = "supervisor.routes")]
427    SupervisorRoutes { modules: Vec<SupervisorRouteModule> },
428    #[serde(rename = "supervisor.provenance")]
429    SupervisorProvenance {
430        daemon: SupervisorDaemonProvenance,
431        modules: Vec<SupervisorModuleProvenance>,
432    },
433    #[serde(rename = "supervisor.stderr_tail")]
434    SupervisorStderrTail {
435        module_id: String,
436        #[serde(flatten)]
437        tail: StderrTail,
438    },
439    #[serde(rename = "supervisor.terminals")]
440    SupervisorTerminals {
441        module_id: String,
442        #[serde(flatten)]
443        terminals: TerminalHistory,
444    },
445}
446
447/// Daemon-originated channel-0 control push body.
448///
449/// A module cannot originate these pushes: subc creates them from its own
450/// forwarding state and enqueues them directly to client connection sinks.
451#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
452#[serde(tag = "op")]
453pub enum ClientControlPush {
454    #[serde(rename = "route.closing")]
455    RouteClosing {
456        module_id: String,
457        reason: RouteCloseReason,
458    },
459    #[serde(rename = "route.closed")]
460    RouteClosed {
461        module_id: String,
462        reason: RouteCloseReason,
463        /// The exact result of the forwarding-quiescence wait for live routes.
464        drained: bool,
465        /// Pending route.bind relays forced down before that wait. They are not
466        /// covered by `drained`, even when live routes quiesced.
467        abandoned: u32,
468        /// Subscription credits captured and excluded from this drain's wire predicate.
469        #[serde(default)]
470        excluded_subscriptions: u32,
471        /// Whether subc will leave this module down until operator action.
472        ///
473        /// The claim covers daemon-owned recovery only. `None` is accepted only
474        /// from daemons that predate this field; every current daemon emission is
475        /// `Some`.
476        #[serde(default, skip_serializing_if = "Option::is_none")]
477        terminal: Option<bool>,
478    },
479}
480
481/// A daemon-incarnation-scoped position in the supervised spawn event stream.
482#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
483pub struct SpawnCursor {
484    pub daemon_incarnation: String,
485    pub seq: u64,
486}
487
488/// One process present in an atomic supervisor spawn snapshot.
489#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
490pub struct LiveSpawn {
491    pub module_id: String,
492    pub spawn_generation: u64,
493    pub pid: u32,
494    pub spawned_at_ms: u64,
495}
496
497/// Atomic live-process census and the cursor at which it was observed.
498#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
499pub struct SpawnSnapshot {
500    pub cursor: SpawnCursor,
501    /// Maximum retained event count for this daemon.
502    pub ring_bound: u64,
503    pub live: Vec<LiveSpawn>,
504}
505
506/// Fact observed by the supervisor when a child process starts or exits.
507#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq)]
508#[serde(rename_all = "snake_case")]
509pub enum SpawnEventKind {
510    Spawned,
511    Exited,
512}
513
514/// One retained or live spawn event.
515///
516/// Exit events intentionally carry no disposition or reason because exit
517/// classification is recorded separately; credential consumers revoke on every
518/// exit regardless of the cause.
519#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
520pub struct SpawnEvent {
521    pub cursor: SpawnCursor,
522    pub kind: SpawnEventKind,
523    pub module_id: String,
524    pub spawn_generation: u64,
525    pub pid: u32,
526    #[serde(default, skip_serializing_if = "Option::is_none")]
527    pub exit_code: Option<i32>,
528    #[serde(default, skip_serializing_if = "Option::is_none")]
529    pub exit_signal: Option<i32>,
530}
531
532/// A module's retained stderr, oldest entry first.
533#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
534pub struct StderrTail {
535    pub capture: StderrCaptureState,
536    pub entries: Vec<StderrTailEntry>,
537    /// Lines not present above: evicted by the ring, or held back by this
538    /// request's own caps.
539    ///
540    /// Non-zero means the first entry is not the first line the module wrote. A
541    /// reader hunting a cause needs that, or an absent explanation reads as a
542    /// module that never gave one.
543    ///
544    /// Zero is skipped so the common complete-tail case stays compact.
545    #[serde(default, skip_serializing_if = "is_zero_u64")]
546    pub dropped_lines: u64,
547}
548
549/// Live routes served by one module.
550#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
551pub struct SupervisorRouteModule {
552    pub module_id: String,
553    pub routes: Vec<SupervisorRoute>,
554}
555
556/// One live consumer route in a [`SupervisorRouteModule`].
557#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
558pub struct SupervisorRoute {
559    pub consumer: SupervisorRouteConsumer,
560    /// Milliseconds since the daemon bound this route.
561    pub age_ms: u64,
562    /// True once the endpoint began draining. Draining routes remain visible so
563    /// a census does not misreport an already-closing route as live.
564    pub draining: bool,
565    /// WHY the endpoint is draining — the same reason vocabulary the
566    /// route.closing push carries — present exactly when `draining` is true.
567    /// Additive: older daemons omit it, and a census consumer must treat a
568    /// draining route without a reason as draining-for-an-unstated-reason,
569    /// never as not-draining.
570    #[serde(default, skip_serializing_if = "Option::is_none")]
571    pub drain_reason: Option<RouteCloseReason>,
572}
573
574/// Source-tagged provenance for one supervised module.
575#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
576pub struct SupervisorModuleProvenance {
577    pub module_id: String,
578    pub module_declared: ModuleDeclaredProvenance,
579    pub daemon_observed: SupervisorObservedProcess,
580}
581
582/// A module's declared build metadata, if its HELLO manifest carried it.
583#[derive(Debug, Clone, PartialEq)]
584pub enum ModuleDeclaredProvenance {
585    Reported {
586        build: ManifestProvenance,
587    },
588    Unverifiable,
589    /// Future discriminator. `body` retains the complete ordered object; `tag`
590    /// is its decoded discriminator projection.
591    Unknown {
592        tag: String,
593        body: OrderedJsonObject,
594    },
595}
596
597/// Process facts observed by the daemon for a supervised module.
598///
599/// Build claims remain under `module_declared`; mixing them here would imply the
600/// daemon independently observed module-provided metadata.
601#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
602pub struct SupervisorObservedProcess {
603    #[serde(default, skip_serializing_if = "Option::is_none")]
604    pub pid: Option<u32>,
605    #[serde(default, skip_serializing_if = "Option::is_none")]
606    pub spawned_at_ms: Option<u64>,
607    #[serde(default, skip_serializing_if = "Option::is_none")]
608    pub spawned_from: Option<PathBuf>,
609    pub running_image: RunningImageAgreement,
610}
611
612/// Independent comparisons of configured path and running image at list time.
613/// An absent verdict means the daemon predates this field, not agreement.
614#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
615pub struct PendingReloadVerdict {
616    pub path: ReloadPathAgreement,
617    pub image: RunningImageAgreement,
618}
619
620/// Whether the running process was spawned from the currently configured program.
621#[derive(Debug, Clone, PartialEq)]
622pub enum ReloadPathAgreement {
623    Match,
624    Mismatch {
625        configured: PathBuf,
626        spawned_from: PathBuf,
627    },
628    Unavailable {
629        reason: ReloadPathUnavailableReason,
630    },
631    Unknown {
632        tag: String,
633        body: OrderedJsonObject,
634    },
635}
636
637open_string_enum! {
638    /// Why configured and spawned paths cannot be compared.
639    ReloadPathUnavailableReason {
640        NotRunning => "not_running",
641        SpawnedPathUnavailable => "spawned_path_unavailable",
642    }
643}
644
645/// Daemon provenance paired with its runtime process observation.
646#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
647pub struct SupervisorDaemonProvenance {
648    pub daemon_build: DaemonBuildProvenance,
649    pub daemon_observed: DaemonObservedProcess,
650}
651
652/// Build metadata embedded in the daemon binary.
653#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
654pub struct DaemonBuildProvenance {
655    #[serde(default, skip_serializing_if = "Option::is_none")]
656    pub build_git_sha: Option<String>,
657    #[serde(default, skip_serializing_if = "Option::is_none")]
658    pub build_lock_digest: Option<String>,
659}
660
661/// Runtime process facts observed for the daemon itself.
662#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
663pub struct DaemonObservedProcess {
664    #[serde(default, skip_serializing_if = "Option::is_none")]
665    pub pid: Option<u32>,
666    /// Wall time derived from suspend-inclusive elapsed time at each read. Clock
667    /// correction can move it by the size of a clock step, and even without a
668    /// step it may vary by about a second between reads. Do not equality-compare
669    /// it. Use raw process start ticks for stable identity.
670    #[serde(default, skip_serializing_if = "Option::is_none")]
671    pub started_at_ms: Option<u64>,
672    pub running_image: RunningImageAgreement,
673}
674
675/// Whether the executable currently running agrees with the spawned image.
676#[derive(Debug, Clone, PartialEq)]
677pub enum RunningImageAgreement {
678    Match {
679        evidence: RunningImageEvidence,
680    },
681    Mismatch {
682        running: RunningImageEvidence,
683        disk: RunningImageEvidence,
684    },
685    Unavailable {
686        reason: RunningImageUnavailableReason,
687    },
688    /// Future discriminator. `body` retains the complete ordered object; `tag`
689    /// is its decoded discriminator projection.
690    Unknown {
691        tag: String,
692        body: OrderedJsonObject,
693    },
694}
695
696/// Platform-specific evidence used to compare a running image with its spawn path.
697#[derive(Debug, Clone, PartialEq)]
698pub enum RunningImageEvidence {
699    LinuxProcSha256 {
700        digest: String,
701    },
702    MacosSpawnInode {
703        device: u64,
704        inode: u64,
705    },
706    /// Future discriminator. `body` retains the complete ordered object; `tag`
707    /// is its decoded discriminator projection.
708    Unknown {
709        tag: String,
710        body: OrderedJsonObject,
711    },
712}
713
714open_string_enum! {
715    /// Reasons why an executable identity could not be observed.
716    RunningImageUnavailableReason {
717        NotRunning => "not_running",
718        UnsupportedPlatform => "unsupported_platform",
719        RunningExecutableUnreadable => "running_executable_unreadable",
720        SpawnedPathUnreadable => "spawned_path_unreadable",
721        HashFailed => "hash_failed",
722        ProcessIdentityUnconfirmed => "process_identity_unconfirmed",
723    }
724}
725
726/// The identity tier the daemon can honestly report for a route consumer.
727///
728/// A caller that proved a live daemon-issued launch nonce is named `reserved`.
729/// A direct key-holder has no such attestation, so it is reported as `direct`
730/// with its connection counter instead of an invented module name.
731#[derive(Debug, Clone, PartialEq)]
732pub enum SupervisorRouteConsumer {
733    Reserved {
734        module_id: String,
735    },
736    Direct {
737        connection_id: u64,
738    },
739    /// Future discriminator. `body` retains the complete ordered object; `tag`
740    /// is its decoded discriminator projection.
741    Unknown {
742        tag: String,
743        body: OrderedJsonObject,
744    },
745}
746
747/// Whether stderr is being captured for a module, and if not, why not.
748///
749/// A typed state rather than an empty-tail convention. "The module printed
750/// nothing before dying" and "nobody was capturing" send an operator in opposite
751/// directions, and rendering them alike is the defect this op exists to fix --
752/// the same shape as a `detail -` that means both no-detail and never-probed.
753#[derive(Debug, Clone, PartialEq)]
754pub enum StderrCaptureState {
755    /// A reader is attached, or was attached and saw clean EOF. An empty
756    /// `entries` under this state means the module genuinely wrote nothing.
757    Captured,
758    /// Retained entries are valid, but the stderr reader ended before clean EOF.
759    Incomplete { reason: String },
760    /// No reader was attached. `entries` says nothing about what the module wrote.
761    NotCaptured { reason: String },
762    /// Future discriminator. `body` retains the complete ordered object; `tag`
763    /// is its decoded discriminator projection.
764    Unknown {
765        tag: String,
766        body: OrderedJsonObject,
767    },
768}
769
770#[derive(Debug, Clone, PartialEq)]
771pub enum StderrTailEntry {
772    Line {
773        text: String,
774        /// The line was cut at the per-line cap and `text` is a prefix.
775        ///
776        /// Carried as a field rather than left to a marker in `text` so a
777        /// consumer can branch on it without string matching.
778        truncated: bool,
779        /// Wall-clock Unix milliseconds at which the daemon read this line off
780        /// the module's pipe. `None` from a daemon that predates the field; a
781        /// reader must then show no time rather than make one up.
782        at_ms: Option<u64>,
783    },
784    /// The supervisor spawned a new process. Entries after this came from it.
785    ///
786    /// In-band because position is the information: which side of the restart a
787    /// line falls on is unanswerable from a count.
788    ProcessStart,
789    /// Future discriminator. `body` retains the complete ordered object; `tag`
790    /// is its decoded discriminator projection.
791    Unknown {
792        tag: String,
793        body: OrderedJsonObject,
794    },
795}
796
797#[derive(Debug, Serialize, Deserialize)]
798#[serde(tag = "status", rename_all = "snake_case")]
799enum ModuleDeclaredProvenanceWire {
800    Reported { build: ManifestProvenance },
801    Unverifiable,
802}
803
804#[derive(Debug, Serialize, Deserialize)]
805#[serde(tag = "status", rename_all = "snake_case")]
806enum RunningImageAgreementWire {
807    Match {
808        evidence: RunningImageEvidence,
809    },
810    Mismatch {
811        running: RunningImageEvidence,
812        disk: RunningImageEvidence,
813    },
814    Unavailable {
815        reason: RunningImageUnavailableReason,
816    },
817}
818
819#[derive(Debug, Serialize, Deserialize)]
820#[serde(tag = "status", rename_all = "snake_case")]
821enum ReloadPathAgreementWire {
822    Match,
823    Mismatch {
824        configured: PathBuf,
825        spawned_from: PathBuf,
826    },
827    Unavailable {
828        reason: ReloadPathUnavailableReason,
829    },
830}
831
832#[derive(Debug, Serialize, Deserialize)]
833#[serde(tag = "method", rename_all = "snake_case")]
834enum RunningImageEvidenceWire {
835    LinuxProcSha256 { digest: String },
836    MacosSpawnInode { device: u64, inode: u64 },
837}
838
839#[derive(Debug, Serialize, Deserialize)]
840#[serde(tag = "kind", rename_all = "snake_case")]
841enum SupervisorRouteConsumerWire {
842    Reserved { module_id: String },
843    Direct { connection_id: u64 },
844}
845
846#[derive(Debug, Serialize, Deserialize)]
847#[serde(tag = "state", rename_all = "snake_case")]
848enum StderrCaptureStateWire {
849    Captured,
850    Incomplete { reason: String },
851    NotCaptured { reason: String },
852}
853
854#[derive(Debug, Serialize, Deserialize)]
855#[serde(tag = "status", rename_all = "snake_case")]
856enum ChildResourceUsageWire {
857    Measured(ChildResourceReading),
858    Unavailable {
859        reason: ChildResourceUnavailableReason,
860    },
861}
862
863#[derive(Debug, Serialize, Deserialize)]
864#[serde(tag = "kind", rename_all = "snake_case")]
865enum StderrTailEntryWire {
866    Line {
867        text: String,
868        #[serde(default, skip_serializing_if = "std::ops::Not::not")]
869        truncated: bool,
870        // Omitted when absent so a reply without it is byte-identical to what
871        // an older daemon sends. Older decoders ignore the member when present.
872        #[serde(default, skip_serializing_if = "Option::is_none")]
873        at_ms: Option<u64>,
874    },
875    ProcessStart,
876}
877
878/// JSON values whose object members retain wire order at every depth.
879#[derive(Debug, Clone, PartialEq)]
880pub enum OrderedJsonValue {
881    Null,
882    Bool(bool),
883    Number(serde_json::Number),
884    String(String),
885    Array(Vec<Self>),
886    Object(OrderedJsonObject),
887}
888
889/// Ordered JSON members retained for an unknown tagged value.
890#[derive(Debug, Clone, PartialEq)]
891pub struct OrderedJsonObject(Vec<(String, OrderedJsonValue)>);
892
893impl OrderedJsonObject {
894    /// Returns the members in the order they appeared on the wire.
895    pub fn as_entries(&self) -> &[(String, OrderedJsonValue)] {
896        &self.0
897    }
898
899    fn into_value(self) -> serde_json::Value {
900        serde_json::Value::Object(
901            self.0
902                .into_iter()
903                .map(|(key, value)| (key, value.into_value()))
904                .collect(),
905        )
906    }
907}
908
909impl OrderedJsonValue {
910    fn into_value(self) -> serde_json::Value {
911        match self {
912            Self::Null => serde_json::Value::Null,
913            Self::Bool(value) => serde_json::Value::Bool(value),
914            Self::Number(value) => serde_json::Value::Number(value),
915            Self::String(value) => serde_json::Value::String(value),
916            Self::Array(values) => {
917                serde_json::Value::Array(values.into_iter().map(Self::into_value).collect())
918            }
919            Self::Object(value) => value.into_value(),
920        }
921    }
922}
923
924impl Serialize for OrderedJsonValue {
925    fn serialize<S>(&self, serializer: S) -> Result<S::Ok, S::Error>
926    where
927        S: Serializer,
928    {
929        match self {
930            Self::Null => serializer.serialize_unit(),
931            Self::Bool(value) => serializer.serialize_bool(*value),
932            Self::Number(value) => value.serialize(serializer),
933            Self::String(value) => serializer.serialize_str(value),
934            Self::Array(values) => values.serialize(serializer),
935            Self::Object(value) => value.serialize(serializer),
936        }
937    }
938}
939
940impl<'de> Deserialize<'de> for OrderedJsonValue {
941    fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
942    where
943        D: Deserializer<'de>,
944    {
945        struct OrderedValueVisitor;
946
947        impl<'de> Visitor<'de> for OrderedValueVisitor {
948            type Value = OrderedJsonValue;
949
950            fn expecting(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
951                formatter.write_str("a JSON value with ordered object members")
952            }
953
954            fn visit_unit<E>(self) -> Result<Self::Value, E>
955            where
956                E: serde::de::Error,
957            {
958                Ok(OrderedJsonValue::Null)
959            }
960
961            fn visit_none<E>(self) -> Result<Self::Value, E>
962            where
963                E: serde::de::Error,
964            {
965                Ok(OrderedJsonValue::Null)
966            }
967
968            fn visit_some<D>(self, deserializer: D) -> Result<Self::Value, D::Error>
969            where
970                D: Deserializer<'de>,
971            {
972                OrderedJsonValue::deserialize(deserializer)
973            }
974
975            fn visit_bool<E>(self, value: bool) -> Result<Self::Value, E>
976            where
977                E: serde::de::Error,
978            {
979                Ok(OrderedJsonValue::Bool(value))
980            }
981
982            fn visit_i64<E>(self, value: i64) -> Result<Self::Value, E>
983            where
984                E: serde::de::Error,
985            {
986                Ok(OrderedJsonValue::Number(value.into()))
987            }
988
989            fn visit_u64<E>(self, value: u64) -> Result<Self::Value, E>
990            where
991                E: serde::de::Error,
992            {
993                Ok(OrderedJsonValue::Number(value.into()))
994            }
995
996            fn visit_f64<E>(self, value: f64) -> Result<Self::Value, E>
997            where
998                E: serde::de::Error,
999            {
1000                serde_json::Number::from_f64(value)
1001                    .map(OrderedJsonValue::Number)
1002                    .ok_or_else(|| E::custom("non-finite JSON number"))
1003            }
1004
1005            fn visit_str<E>(self, value: &str) -> Result<Self::Value, E>
1006            where
1007                E: serde::de::Error,
1008            {
1009                Ok(OrderedJsonValue::String(value.to_owned()))
1010            }
1011
1012            fn visit_string<E>(self, value: String) -> Result<Self::Value, E>
1013            where
1014                E: serde::de::Error,
1015            {
1016                Ok(OrderedJsonValue::String(value))
1017            }
1018
1019            fn visit_seq<A>(self, mut sequence: A) -> Result<Self::Value, A::Error>
1020            where
1021                A: SeqAccess<'de>,
1022            {
1023                let mut values = Vec::new();
1024                while let Some(value) = sequence.next_element()? {
1025                    values.push(value);
1026                }
1027                Ok(OrderedJsonValue::Array(values))
1028            }
1029
1030            fn visit_map<A>(self, mut map: A) -> Result<Self::Value, A::Error>
1031            where
1032                A: MapAccess<'de>,
1033            {
1034                let mut entries = Vec::new();
1035                while let Some((key, value)) = map.next_entry()? {
1036                    entries.push((key, value));
1037                }
1038                Ok(OrderedJsonValue::Object(OrderedJsonObject(entries)))
1039            }
1040        }
1041
1042        deserializer.deserialize_any(OrderedValueVisitor)
1043    }
1044}
1045
1046impl Serialize for OrderedJsonObject {
1047    fn serialize<S>(&self, serializer: S) -> Result<S::Ok, S::Error>
1048    where
1049        S: Serializer,
1050    {
1051        let mut map = serializer.serialize_map(Some(self.0.len()))?;
1052        for (key, value) in &self.0 {
1053            map.serialize_entry(key, value)?;
1054        }
1055        map.end()
1056    }
1057}
1058
1059impl<'de> Deserialize<'de> for OrderedJsonObject {
1060    fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
1061    where
1062        D: Deserializer<'de>,
1063    {
1064        struct OrderedObjectVisitor;
1065
1066        impl<'de> Visitor<'de> for OrderedObjectVisitor {
1067            type Value = OrderedJsonObject;
1068
1069            fn expecting(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
1070                formatter.write_str("an object with ordered JSON members")
1071            }
1072
1073            fn visit_map<A>(self, mut map: A) -> Result<Self::Value, A::Error>
1074            where
1075                A: MapAccess<'de>,
1076            {
1077                let mut entries = Vec::new();
1078                while let Some((key, value)) = map.next_entry()? {
1079                    entries.push((key, value));
1080                }
1081                Ok(OrderedJsonObject(entries))
1082            }
1083        }
1084
1085        deserializer.deserialize_map(OrderedObjectVisitor)
1086    }
1087}
1088
1089fn read_tagged<'de, D>(
1090    deserializer: D,
1091    field: &'static str,
1092) -> Result<(String, OrderedJsonObject), D::Error>
1093where
1094    D: Deserializer<'de>,
1095{
1096    let body = OrderedJsonObject::deserialize(deserializer)?;
1097    let mut tag = None;
1098    for (key, value) in body.as_entries() {
1099        if key != field {
1100            continue;
1101        }
1102        if tag.is_some() {
1103            return Err(D::Error::custom(format!(
1104                "tagged object has duplicate `{field}` field"
1105            )));
1106        }
1107        let OrderedJsonValue::String(value) = value else {
1108            return Err(D::Error::custom(format!(
1109                "tagged object has no string `{field}` field"
1110            )));
1111        };
1112        tag = Some(value);
1113    }
1114    let Some(tag) = tag else {
1115        return Err(D::Error::custom(format!(
1116            "tagged object has no string `{field}` field"
1117        )));
1118    };
1119    Ok((tag.to_string(), body))
1120}
1121
1122fn read_ordered_tagged(
1123    value: OrderedJsonValue,
1124    field: &'static str,
1125) -> Result<(String, OrderedJsonObject), String> {
1126    let OrderedJsonValue::Object(body) = value else {
1127        return Err(format!("expected tagged object with `{field}` field"));
1128    };
1129    let mut tag = None;
1130    for (key, value) in body.as_entries() {
1131        if key != field {
1132            continue;
1133        }
1134        if tag.is_some() {
1135            return Err(format!("tagged object has duplicate `{field}` field"));
1136        }
1137        let OrderedJsonValue::String(value) = value else {
1138            return Err(format!("tagged object has no string `{field}` field"));
1139        };
1140        tag = Some(value);
1141    }
1142    let Some(tag) = tag else {
1143        return Err(format!("tagged object has no string `{field}` field"));
1144    };
1145    Ok((tag.to_string(), body))
1146}
1147
1148fn ordered_field<'a>(body: &'a OrderedJsonObject, field: &str) -> Option<&'a OrderedJsonValue> {
1149    body.as_entries()
1150        .iter()
1151        .find_map(|(key, value)| (key == field).then_some(value))
1152}
1153
1154fn ordered_string(body: &OrderedJsonObject, field: &str) -> Result<String, String> {
1155    match ordered_field(body, field) {
1156        Some(OrderedJsonValue::String(value)) => Ok(value.clone()),
1157        Some(_) => Err(format!("tagged object field `{field}` is not a string")),
1158        None => Err(format!("tagged object has no `{field}` field")),
1159    }
1160}
1161
1162fn decode_running_image_evidence(value: OrderedJsonValue) -> Result<RunningImageEvidence, String> {
1163    let (tag, body) = read_ordered_tagged(value, "method")?;
1164    match tag.as_str() {
1165        "linux_proc_sha256" => Ok(RunningImageEvidence::LinuxProcSha256 {
1166            digest: ordered_string(&body, "digest")?,
1167        }),
1168        "macos_spawn_inode" => {
1169            let device = ordered_field(&body, "device")
1170                .and_then(|value| match value {
1171                    OrderedJsonValue::Number(number) => number.as_u64(),
1172                    _ => None,
1173                })
1174                .ok_or_else(|| "tagged object has no unsigned `device` field".to_string())?;
1175            let inode = ordered_field(&body, "inode")
1176                .and_then(|value| match value {
1177                    OrderedJsonValue::Number(number) => number.as_u64(),
1178                    _ => None,
1179                })
1180                .ok_or_else(|| "tagged object has no unsigned `inode` field".to_string())?;
1181            Ok(RunningImageEvidence::MacosSpawnInode { device, inode })
1182        }
1183        _ => Ok(RunningImageEvidence::Unknown { tag, body }),
1184    }
1185}
1186
1187impl Serialize for ModuleDeclaredProvenance {
1188    fn serialize<S>(&self, serializer: S) -> Result<S::Ok, S::Error>
1189    where
1190        S: Serializer,
1191    {
1192        match self {
1193            Self::Reported { build } => ModuleDeclaredProvenanceWire::Reported {
1194                build: build.clone(),
1195            }
1196            .serialize(serializer),
1197            Self::Unverifiable => ModuleDeclaredProvenanceWire::Unverifiable.serialize(serializer),
1198            Self::Unknown { body, .. } => body.serialize(serializer),
1199        }
1200    }
1201}
1202
1203impl<'de> Deserialize<'de> for ModuleDeclaredProvenance {
1204    fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
1205    where
1206        D: serde::Deserializer<'de>,
1207    {
1208        let (tag, value) = read_tagged(deserializer, "status")?;
1209        match tag.as_str() {
1210            "reported" => match serde_json::from_value(value.into_value())
1211                .map_err(D::Error::custom)?
1212            {
1213                ModuleDeclaredProvenanceWire::Reported { build } => Ok(Self::Reported { build }),
1214                ModuleDeclaredProvenanceWire::Unverifiable => unreachable!(),
1215            },
1216            "unverifiable" => {
1217                match serde_json::from_value(value.into_value()).map_err(D::Error::custom)? {
1218                    ModuleDeclaredProvenanceWire::Unverifiable => Ok(Self::Unverifiable),
1219                    ModuleDeclaredProvenanceWire::Reported { .. } => unreachable!(),
1220                }
1221            }
1222            _ => Ok(Self::Unknown { tag, body: value }),
1223        }
1224    }
1225}
1226
1227impl Serialize for RunningImageAgreement {
1228    fn serialize<S>(&self, serializer: S) -> Result<S::Ok, S::Error>
1229    where
1230        S: Serializer,
1231    {
1232        match self {
1233            Self::Match { evidence } => RunningImageAgreementWire::Match {
1234                evidence: evidence.clone(),
1235            }
1236            .serialize(serializer),
1237            Self::Mismatch { running, disk } => RunningImageAgreementWire::Mismatch {
1238                running: running.clone(),
1239                disk: disk.clone(),
1240            }
1241            .serialize(serializer),
1242            Self::Unavailable { reason } => RunningImageAgreementWire::Unavailable {
1243                reason: reason.clone(),
1244            }
1245            .serialize(serializer),
1246            Self::Unknown { body, .. } => body.serialize(serializer),
1247        }
1248    }
1249}
1250
1251impl Serialize for ReloadPathAgreement {
1252    fn serialize<S>(&self, serializer: S) -> Result<S::Ok, S::Error>
1253    where
1254        S: Serializer,
1255    {
1256        match self {
1257            Self::Match => ReloadPathAgreementWire::Match.serialize(serializer),
1258            Self::Mismatch {
1259                configured,
1260                spawned_from,
1261            } => ReloadPathAgreementWire::Mismatch {
1262                configured: configured.clone(),
1263                spawned_from: spawned_from.clone(),
1264            }
1265            .serialize(serializer),
1266            Self::Unavailable { reason } => ReloadPathAgreementWire::Unavailable {
1267                reason: reason.clone(),
1268            }
1269            .serialize(serializer),
1270            Self::Unknown { body, .. } => body.serialize(serializer),
1271        }
1272    }
1273}
1274
1275impl<'de> Deserialize<'de> for ReloadPathAgreement {
1276    fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
1277    where
1278        D: Deserializer<'de>,
1279    {
1280        let (tag, body) = read_tagged(deserializer, "status")?;
1281        match tag.as_str() {
1282            "match" => Ok(Self::Match),
1283            "mismatch" => {
1284                match serde_json::from_value(body.into_value()).map_err(D::Error::custom)? {
1285                    ReloadPathAgreementWire::Mismatch {
1286                        configured,
1287                        spawned_from,
1288                    } => Ok(Self::Mismatch {
1289                        configured,
1290                        spawned_from,
1291                    }),
1292                    _ => unreachable!(),
1293                }
1294            }
1295            "unavailable" => match serde_json::from_value(body.into_value())
1296                .map_err(D::Error::custom)?
1297            {
1298                ReloadPathAgreementWire::Unavailable { reason } => Ok(Self::Unavailable { reason }),
1299                _ => unreachable!(),
1300            },
1301            _ => Ok(Self::Unknown { tag, body }),
1302        }
1303    }
1304}
1305
1306impl<'de> Deserialize<'de> for RunningImageAgreement {
1307    fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
1308    where
1309        D: serde::Deserializer<'de>,
1310    {
1311        let (tag, value) = read_tagged(deserializer, "status")?;
1312        match tag.as_str() {
1313            "match" => Ok(Self::Match {
1314                evidence: decode_running_image_evidence(
1315                    ordered_field(&value, "evidence")
1316                        .cloned()
1317                        .ok_or_else(|| D::Error::custom("tagged object has no `evidence` field"))?,
1318                )
1319                .map_err(D::Error::custom)?,
1320            }),
1321            "mismatch" => Ok(Self::Mismatch {
1322                running: decode_running_image_evidence(
1323                    ordered_field(&value, "running")
1324                        .cloned()
1325                        .ok_or_else(|| D::Error::custom("tagged object has no `running` field"))?,
1326                )
1327                .map_err(D::Error::custom)?,
1328                disk: decode_running_image_evidence(
1329                    ordered_field(&value, "disk")
1330                        .cloned()
1331                        .ok_or_else(|| D::Error::custom("tagged object has no `disk` field"))?,
1332                )
1333                .map_err(D::Error::custom)?,
1334            }),
1335            "unavailable" => Ok(Self::Unavailable {
1336                reason: serde_json::from_value(
1337                    ordered_field(&value, "reason")
1338                        .cloned()
1339                        .ok_or_else(|| D::Error::custom("tagged object has no `reason` field"))?
1340                        .into_value(),
1341                )
1342                .map_err(D::Error::custom)?,
1343            }),
1344            _ => Ok(Self::Unknown { tag, body: value }),
1345        }
1346    }
1347}
1348
1349impl Serialize for ChildResourceUsage {
1350    fn serialize<S>(&self, serializer: S) -> Result<S::Ok, S::Error>
1351    where
1352        S: Serializer,
1353    {
1354        match self {
1355            Self::Measured(reading) => {
1356                ChildResourceUsageWire::Measured(reading.clone()).serialize(serializer)
1357            }
1358            Self::Unavailable { reason } => ChildResourceUsageWire::Unavailable {
1359                reason: reason.clone(),
1360            }
1361            .serialize(serializer),
1362            Self::Unknown { body, .. } => body.serialize(serializer),
1363        }
1364    }
1365}
1366
1367impl<'de> Deserialize<'de> for ChildResourceUsage {
1368    fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
1369    where
1370        D: Deserializer<'de>,
1371    {
1372        let (tag, body) = read_tagged(deserializer, "status")?;
1373        match tag.as_str() {
1374            "measured" | "unavailable" => {
1375                match serde_json::from_value(body.into_value()).map_err(D::Error::custom)? {
1376                    ChildResourceUsageWire::Measured(reading) => Ok(Self::Measured(reading)),
1377                    ChildResourceUsageWire::Unavailable { reason } => {
1378                        Ok(Self::Unavailable { reason })
1379                    }
1380                }
1381            }
1382            _ => Ok(Self::Unknown { tag, body }),
1383        }
1384    }
1385}
1386
1387impl Serialize for RunningImageEvidence {
1388    fn serialize<S>(&self, serializer: S) -> Result<S::Ok, S::Error>
1389    where
1390        S: Serializer,
1391    {
1392        match self {
1393            Self::LinuxProcSha256 { digest } => RunningImageEvidenceWire::LinuxProcSha256 {
1394                digest: digest.clone(),
1395            }
1396            .serialize(serializer),
1397            Self::MacosSpawnInode { device, inode } => RunningImageEvidenceWire::MacosSpawnInode {
1398                device: *device,
1399                inode: *inode,
1400            }
1401            .serialize(serializer),
1402            Self::Unknown { body, .. } => body.serialize(serializer),
1403        }
1404    }
1405}
1406
1407impl<'de> Deserialize<'de> for RunningImageEvidence {
1408    fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
1409    where
1410        D: serde::Deserializer<'de>,
1411    {
1412        let (tag, value) = read_tagged(deserializer, "method")?;
1413        match tag.as_str() {
1414            "linux_proc_sha256" => {
1415                match serde_json::from_value(value.into_value()).map_err(D::Error::custom)? {
1416                    RunningImageEvidenceWire::LinuxProcSha256 { digest } => {
1417                        Ok(Self::LinuxProcSha256 { digest })
1418                    }
1419                    _ => unreachable!(),
1420                }
1421            }
1422            "macos_spawn_inode" => {
1423                match serde_json::from_value(value.into_value()).map_err(D::Error::custom)? {
1424                    RunningImageEvidenceWire::MacosSpawnInode { device, inode } => {
1425                        Ok(Self::MacosSpawnInode { device, inode })
1426                    }
1427                    _ => unreachable!(),
1428                }
1429            }
1430            _ => Ok(Self::Unknown { tag, body: value }),
1431        }
1432    }
1433}
1434
1435impl Serialize for SupervisorRouteConsumer {
1436    fn serialize<S>(&self, serializer: S) -> Result<S::Ok, S::Error>
1437    where
1438        S: Serializer,
1439    {
1440        match self {
1441            Self::Reserved { module_id } => SupervisorRouteConsumerWire::Reserved {
1442                module_id: module_id.clone(),
1443            }
1444            .serialize(serializer),
1445            Self::Direct { connection_id } => SupervisorRouteConsumerWire::Direct {
1446                connection_id: *connection_id,
1447            }
1448            .serialize(serializer),
1449            Self::Unknown { body, .. } => body.serialize(serializer),
1450        }
1451    }
1452}
1453
1454impl<'de> Deserialize<'de> for SupervisorRouteConsumer {
1455    fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
1456    where
1457        D: serde::Deserializer<'de>,
1458    {
1459        let (tag, value) = read_tagged(deserializer, "kind")?;
1460        match tag.as_str() {
1461            "reserved" => {
1462                match serde_json::from_value(value.into_value()).map_err(D::Error::custom)? {
1463                    SupervisorRouteConsumerWire::Reserved { module_id } => {
1464                        Ok(Self::Reserved { module_id })
1465                    }
1466                    _ => unreachable!(),
1467                }
1468            }
1469            "direct" => {
1470                match serde_json::from_value(value.into_value()).map_err(D::Error::custom)? {
1471                    SupervisorRouteConsumerWire::Direct { connection_id } => {
1472                        Ok(Self::Direct { connection_id })
1473                    }
1474                    _ => unreachable!(),
1475                }
1476            }
1477            _ => Ok(Self::Unknown { tag, body: value }),
1478        }
1479    }
1480}
1481
1482impl Serialize for StderrCaptureState {
1483    fn serialize<S>(&self, serializer: S) -> Result<S::Ok, S::Error>
1484    where
1485        S: Serializer,
1486    {
1487        match self {
1488            Self::Captured => StderrCaptureStateWire::Captured.serialize(serializer),
1489            Self::Incomplete { reason } => StderrCaptureStateWire::Incomplete {
1490                reason: reason.clone(),
1491            }
1492            .serialize(serializer),
1493            Self::NotCaptured { reason } => StderrCaptureStateWire::NotCaptured {
1494                reason: reason.clone(),
1495            }
1496            .serialize(serializer),
1497            Self::Unknown { body, .. } => body.serialize(serializer),
1498        }
1499    }
1500}
1501
1502impl<'de> Deserialize<'de> for StderrCaptureState {
1503    fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
1504    where
1505        D: serde::Deserializer<'de>,
1506    {
1507        let (tag, value) = read_tagged(deserializer, "state")?;
1508        match tag.as_str() {
1509            "captured" => {
1510                match serde_json::from_value(value.into_value()).map_err(D::Error::custom)? {
1511                    StderrCaptureStateWire::Captured => Ok(Self::Captured),
1512                    _ => unreachable!(),
1513                }
1514            }
1515            "incomplete" => match serde_json::from_value(value.into_value())
1516                .map_err(D::Error::custom)?
1517            {
1518                StderrCaptureStateWire::Incomplete { reason } => Ok(Self::Incomplete { reason }),
1519                _ => unreachable!(),
1520            },
1521            "not_captured" => match serde_json::from_value(value.into_value())
1522                .map_err(D::Error::custom)?
1523            {
1524                StderrCaptureStateWire::NotCaptured { reason } => Ok(Self::NotCaptured { reason }),
1525                _ => unreachable!(),
1526            },
1527            _ => Ok(Self::Unknown { tag, body: value }),
1528        }
1529    }
1530}
1531
1532impl Serialize for StderrTailEntry {
1533    fn serialize<S>(&self, serializer: S) -> Result<S::Ok, S::Error>
1534    where
1535        S: Serializer,
1536    {
1537        match self {
1538            Self::Line {
1539                text,
1540                truncated,
1541                at_ms,
1542            } => StderrTailEntryWire::Line {
1543                text: text.clone(),
1544                truncated: *truncated,
1545                at_ms: *at_ms,
1546            }
1547            .serialize(serializer),
1548            Self::ProcessStart => StderrTailEntryWire::ProcessStart.serialize(serializer),
1549            Self::Unknown { body, .. } => body.serialize(serializer),
1550        }
1551    }
1552}
1553
1554impl<'de> Deserialize<'de> for StderrTailEntry {
1555    fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
1556    where
1557        D: serde::Deserializer<'de>,
1558    {
1559        let (tag, value) = read_tagged(deserializer, "kind")?;
1560        match tag.as_str() {
1561            "line" => match serde_json::from_value(value.into_value()).map_err(D::Error::custom)? {
1562                StderrTailEntryWire::Line {
1563                    text,
1564                    truncated,
1565                    at_ms,
1566                } => Ok(Self::Line {
1567                    text,
1568                    truncated,
1569                    at_ms,
1570                }),
1571                _ => unreachable!(),
1572            },
1573            "process_start" => {
1574                match serde_json::from_value(value.into_value()).map_err(D::Error::custom)? {
1575                    StderrTailEntryWire::ProcessStart => Ok(Self::ProcessStart),
1576                    _ => unreachable!(),
1577                }
1578            }
1579            _ => Ok(Self::Unknown { tag, body: value }),
1580        }
1581    }
1582}
1583
1584fn is_zero_u64(value: &u64) -> bool {
1585    *value == 0
1586}
1587
1588fn default_true() -> bool {
1589    true
1590}
1591
1592/// Bounded terminal history for one module, oldest retained record first.
1593#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
1594pub struct TerminalHistory {
1595    /// Unix milliseconds at the current daemon's start; entries may predate it.
1596    pub daemon_started_at_ms: u64,
1597    pub entries: Vec<TerminalEntry>,
1598    /// Exits evicted by the current daemon's ring, possibly recovered from its
1599    /// journal. Not a count of missing exits: expired journal totals are unknown.
1600    #[serde(default, skip_serializing_if = "is_zero_u64")]
1601    pub dropped: u64,
1602    /// Unparseable or incomplete lines across the shared journal, including
1603    /// lines whose module cannot be determined. Zero on older daemons.
1604    #[serde(default, skip_serializing_if = "is_zero_u64")]
1605    pub journal_skipped_lines: u64,
1606    /// Files that could not be read completely, excluding absent generations.
1607    #[serde(default, skip_serializing_if = "is_zero_u64")]
1608    pub journal_read_errors: u64,
1609    /// Failed journal appends across all modules in the current daemon.
1610    #[serde(default, skip_serializing_if = "is_zero_u64")]
1611    pub journal_write_failures: u64,
1612}
1613
1614/// One terminal child exit and the supervisor action it selected.
1615#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
1616pub struct TerminalEntry {
1617    /// Opaque identity of the daemon that observed this exit; absent on older
1618    /// daemons. Different tokens mean different lifetimes, not chronological order.
1619    #[serde(default, skip_serializing_if = "Option::is_none")]
1620    pub daemon_incarnation: Option<String>,
1621    #[serde(default, skip_serializing_if = "Option::is_none")]
1622    pub exit_code: Option<i32>,
1623    #[serde(default, skip_serializing_if = "Option::is_none")]
1624    pub exit_signal: Option<i32>,
1625    pub at_ms: u64,
1626    pub disposition: TerminalDisposition,
1627    /// Supervisor classification of this exit. Absent on daemons that predate
1628    /// the field; unknown future kinds remain readable instead of failing the
1629    /// enclosing terminal record.
1630    #[serde(default, skip_serializing_if = "Option::is_none")]
1631    pub exit_kind: Option<TerminalExitKind>,
1632    /// Why the supervisor chose this disposition, when the disposition alone
1633    /// does not say. A `failed` record carries the exhausted crash budget here
1634    /// (`crash budget exhausted: max_restarts=3 within window_secs=600`), which
1635    /// is the difference between an operator seeing "it failed" and seeing which
1636    /// limit stopped it. Prose for humans: render it, never parse it. Absent for
1637    /// ordinary dispositions and on daemons predating the field.
1638    #[serde(default, skip_serializing_if = "Option::is_none")]
1639    pub disposition_detail: Option<String>,
1640}
1641
1642/// Exit classification carried by supervisor history and census records.
1643///
1644/// This is an open string enum so future daemon variants degrade to a readable
1645/// unknown kind rather than making a consumer discard the enclosing record.
1646#[derive(Debug, Clone, PartialEq, Eq)]
1647pub enum TerminalExitKind {
1648    Clean,
1649    Crash,
1650    DeliberateSeverance,
1651    Unknown(String),
1652}
1653
1654impl TerminalExitKind {
1655    fn wire_name(&self) -> &str {
1656        match self {
1657            Self::Clean => "clean",
1658            Self::Crash => "crash",
1659            Self::DeliberateSeverance => "deliberate_severance",
1660            Self::Unknown(value) => value,
1661        }
1662    }
1663}
1664
1665impl Serialize for TerminalExitKind {
1666    fn serialize<S>(&self, serializer: S) -> Result<S::Ok, S::Error>
1667    where
1668        S: serde::Serializer,
1669    {
1670        serializer.serialize_str(self.wire_name())
1671    }
1672}
1673
1674impl<'de> Deserialize<'de> for TerminalExitKind {
1675    fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
1676    where
1677        D: serde::Deserializer<'de>,
1678    {
1679        let value = String::deserialize(deserializer)?;
1680        Ok(match value.as_str() {
1681            "clean" => Self::Clean,
1682            "crash" => Self::Crash,
1683            "deliberate_severance" => Self::DeliberateSeverance,
1684            _ => Self::Unknown(value),
1685        })
1686    }
1687}
1688
1689open_string_enum! {
1690    /// The supervisor disposition selected after observing a terminal exit.
1691    TerminalDisposition {
1692        Stopped => "stopped",
1693        Disabled => "disabled",
1694        Failed => "failed",
1695        Restarting => "restarting",
1696        /// The child exited after the daemon had begun its own announced
1697        /// shutdown, whatever its exit code or signal. Such an exit is not
1698        /// a crash and is never followed by a respawn. Readers that predate
1699        /// this value decode it as `Unknown("daemon_shutdown")`.
1700        DaemonShutdown => "daemon_shutdown",
1701    }
1702}
1703
1704#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq)]
1705#[serde(rename_all = "snake_case")]
1706pub enum PollKind {
1707    Status,
1708    Liveness,
1709}
1710
1711#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
1712pub struct CatalogEntry {
1713    pub module_id: String,
1714    /// Whether the registered module currently accepts new route binds.
1715    ///
1716    /// This is the module's EFFECTIVE readiness: its own declared readiness
1717    /// AND every one of its `need: required` capabilities having a registered
1718    /// provider. It is exactly the condition `route.open` checks, so a caller
1719    /// reading `false` here will be refused with `module_warming`; `not_ready`
1720    /// says why.
1721    ///
1722    /// Older daemons omit this field and are interpreted as ready. Daemons that
1723    /// predate `not_ready` report declared readiness only.
1724    #[serde(default = "default_true")]
1725    pub ready: bool,
1726    /// Why `ready` is false, in the same shape `route.open` puts in the
1727    /// `detail` of its `module_warming` refusal. Absent when the module is
1728    /// ready, and absent from daemons that predate the field.
1729    #[serde(default, skip_serializing_if = "Option::is_none")]
1730    pub not_ready: Option<NotReadyReason>,
1731    /// The registered module's self-declared build version, projected from its
1732    /// manifest so a consumer can tell WHICH BUILD of a module it is talking
1733    /// to at connect time.
1734    ///
1735    /// Without this, a client compiled against a module's current source reads
1736    /// a contract that is true of the repository and false of the running
1737    /// process -- the types match, the JSON decodes, and the meaning has
1738    /// changed. That failure carries no error to notice; the version in the
1739    /// catalog turns a semantic skew into a log line at connect instead of a
1740    /// wrong sentence on a user's screen.
1741    ///
1742    /// Optional on the wire only because entries serialized by older daemons
1743    /// lack it: absent means "daemon predates the field", never "module has
1744    /// no version" (the manifest field is required at registration).
1745    ///
1746    /// The reading is ARMED BY OBSERVATION, not by this documentation: until
1747    /// a consumer has seen at least one populated entry from the daemon it is
1748    /// connected to, an all-None catalog is indistinguishable from an old
1749    /// daemon, and a client shipping the documented reading against it would
1750    /// hold a guarantee it does not have.
1751    #[serde(default, skip_serializing_if = "Option::is_none")]
1752    pub module_version: Option<String>,
1753    pub roles: Vec<ProviderRole>,
1754    pub control_ops: Vec<String>,
1755    /// Static capability declarations from the registering module's manifest.
1756    ///
1757    /// Optional on the wire so consumers connected to a daemon that predates the
1758    /// capability grammar retain their existing catalog decoding behavior.
1759    #[serde(default, skip_serializing_if = "Option::is_none")]
1760    pub capabilities: Option<CapabilityDeclarations>,
1761    /// Self-signal declarations mirrored verbatim from the registering module's
1762    /// manifest. The daemon relays these declarations without interpreting them.
1763    #[serde(default, skip_serializing_if = "Option::is_none")]
1764    pub self_signals: Option<Vec<SelfSignalDeclaration>>,
1765}
1766
1767/// Why a registered module is not accepting new route binds.
1768#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
1769pub struct NotReadyReason {
1770    /// `declared_not_ready` when the module itself said it is not ready, or
1771    /// `required_capability_unprovided` when a capability it declares
1772    /// `need: required` has no registered provider. Open vocabulary: a newer
1773    /// daemon may add reasons.
1774    pub reason: String,
1775    /// For `required_capability_unprovided`, the lexicographically first
1776    /// required capability that has no registered provider.
1777    #[serde(default, skip_serializing_if = "Option::is_none")]
1778    pub capability: Option<String>,
1779}
1780
1781impl NotReadyReason {
1782    pub const DECLARED_NOT_READY: &'static str = "declared_not_ready";
1783    pub const REQUIRED_CAPABILITY_UNPROVIDED: &'static str = "required_capability_unprovided";
1784}
1785
1786#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
1787pub struct CapabilityRequirementStatus {
1788    pub consumer: String,
1789    pub capability: String,
1790    pub need: String,
1791    pub verdict: String,
1792    pub episode_seq: u64,
1793    pub config_satisfiable: bool,
1794    pub runtime_available: bool,
1795    pub detail: String,
1796}
1797
1798#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
1799pub struct SupervisorRescanResult {
1800    pub added: Vec<String>,
1801    pub removed: Vec<String>,
1802    pub changed_pending_reload: Vec<String>,
1803    /// Modules whose enabled flag differs between config and running state.
1804    ///
1805    /// Rescan calls `set_enabled` for these, so omitting them made the preview
1806    /// describe two of the three mutation classes it performs. A module changing
1807    /// only its enabled flag landed in no bucket at all -- not added, removed or
1808    /// changed, and deliberately not counted as unchanged either -- so the sole
1809    /// evidence was that the buckets no longer summed to the configured module
1810    /// count. A preview is consulted precisely when someone is being careful,
1811    /// which is the worst place to under-report.
1812    ///
1813    /// Empty is skipped so consumers written against the older shape keep
1814    /// parsing.
1815    #[serde(default, skip_serializing_if = "Vec::is_empty")]
1816    pub enabled_changes: Vec<String>,
1817    pub unchanged: u32,
1818    /// True when this reconciliation was computed but NOT applied.
1819    ///
1820    /// Carried on the result rather than left to the caller's memory of what it
1821    /// asked for. A preview and an execution are otherwise byte-identical, so a
1822    /// reader who meets this output later -- in a log, a transcript, a pasted
1823    /// snippet -- cannot tell which one happened. Absent when false, so existing
1824    /// consumers see the shape they already parse.
1825    #[serde(default, skip_serializing_if = "std::ops::Not::not")]
1826    pub preview: bool,
1827    /// Config sections that changed but which rescan CANNOT apply, so the
1828    /// operator learns a daemon restart is required from the command they just
1829    /// ran rather than from the journal.
1830    ///
1831    /// The daemon has always detected this and logged a warning. A warning in a
1832    /// log is addressed to whoever is reading the log, and the person who just
1833    /// edited the config is by construction looking at the CLI instead: reported
1834    /// by an outside contributor after a module crash-looped through four
1835    /// respawns because a new top-level `storage` section was silently not
1836    /// applied, diagnosable only by journal archaeology.
1837    ///
1838    /// Names the SECTIONS rather than a boolean, because "something else
1839    /// changed" sends the operator back to diffing their own file -- which is
1840    /// the work the message exists to save.
1841    ///
1842    /// Empty is skipped, so consumers written against the older shape keep
1843    /// parsing.
1844    #[serde(default, skip_serializing_if = "Vec::is_empty")]
1845    pub restart_required: Vec<String>,
1846    /// Required capabilities that a dry-run's resulting module set would leave
1847    /// unprovided. Rows are human-readable because the preview is an operator
1848    /// explanation, not a second manifest schema.
1849    #[serde(default, skip_serializing_if = "Vec::is_empty")]
1850    pub capability_warnings: Vec<String>,
1851}
1852
1853/// Which wire protocol a supervised module speaks to subc, as DECLARED in
1854/// daemon config. Never inferred from observed behaviour.
1855///
1856/// The distinction this exists to keep is between a module that should have
1857/// registered and has not yet, and one that never will. A `Subc` module that has
1858/// not registered is a subc module that is LATE -- it may be booting, it may be
1859/// wedged, and the supervisor's health probing and restart escalation are the
1860/// right response. A `None` module is a third-party process (the NATS server is
1861/// the first) that subc launches, supervises, and stops, and that is all: it
1862/// speaks no subc wire at all, so treating its silence as a fault would restart
1863/// a perfectly healthy process forever.
1864///
1865/// Inferring the difference from "has not registered within N seconds" would
1866/// collapse exactly the two cases that must stay apart, which is why this is a
1867/// declaration.
1868#[derive(Debug, Clone, Copy, Default, Serialize, Deserialize, PartialEq, Eq)]
1869#[serde(rename_all = "snake_case")]
1870pub enum ModuleProtocol {
1871    /// The module registers over channel 0, answers `health.check`, and can
1872    /// serve routes. Every module predating this field is one of these, which is
1873    /// why it is the default.
1874    #[default]
1875    Subc,
1876    /// The module speaks no subc wire. It is supervised as a process only.
1877    ///
1878    /// A clean exit (status 0) that the daemon did not request is restarted as
1879    /// a crash, counting against the restart budget, instead of being recorded
1880    /// as a stop. Such a module is usually a stock program that exits 0 on
1881    /// SIGTERM, so a stray outside signal would otherwise leave it down for
1882    /// good; a subc-wire module re-raises SIGTERM instead, so this rule is not
1883    /// needed for it.
1884    None,
1885}
1886
1887#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
1888pub struct SupervisorEntry {
1889    pub module_id: String,
1890    pub state: String,
1891    pub enabled: bool,
1892    /// Whether this module is serving.
1893    ///
1894    /// For a `Subc` module: enabled, running, process alive, AND registered.
1895    /// For a `None` module the registration term is dropped, because a module
1896    /// that speaks no subc wire never registers and the daemon cannot assert
1897    /// more than "the process it launched is alive". READ IT WITH `protocol`:
1898    /// `live: true` means something weaker for a `None` module, and a renderer
1899    /// that prints it as a bare boolean for one is claiming more than the daemon
1900    /// knows.
1901    pub live: bool,
1902    /// The module's declared wire protocol. Absent on daemons predating the
1903    /// field, where every module was a subc module, so the default is exactly
1904    /// what those daemons meant.
1905    #[serde(default)]
1906    pub protocol: ModuleProtocol,
1907    pub health: SupervisorHealthStatus,
1908    /// Computed from the stored launch spec and observed process at list time;
1909    /// None means an older daemon did not report this comparison.
1910    #[serde(default, skip_serializing_if = "Option::is_none")]
1911    pub pending_reload: Option<PendingReloadVerdict>,
1912    /// When the daemon last collected this module's health, as unix
1913    /// milliseconds. Absent means NEVER PROBED (a module inside its first probe
1914    /// window, whose `health` is therefore `Unknown` rather than good), not
1915    /// probed-long-ago. An old value and an absent one call for opposite
1916    /// readings, so do not render them alike.
1917    #[serde(default)]
1918    pub last_probe_ms: Option<u64>,
1919    /// Exit code of the module's most recent process exit, if the process has
1920    /// exited at least once. Survives respawn so a now-`running` module still
1921    /// reports what killed its previous incarnation.
1922    #[serde(default, skip_serializing_if = "Option::is_none")]
1923    pub last_exit_code: Option<i32>,
1924    /// Terminating signal of the module's most recent process exit (Unix), if
1925    /// any. `Some(9)` = SIGKILL (OOM/jetsam/kill-on-drop), `Some(6)` = SIGABRT
1926    /// (often a panic-abort). Survives respawn.
1927    #[serde(default, skip_serializing_if = "Option::is_none")]
1928    pub last_exit_signal: Option<i32>,
1929    /// Unix milliseconds when the most recent child exit was observed. Present
1930    /// even when the terminal ring is not queried, so existing list readers can
1931    /// order their latest observed exit against events they already received.
1932    #[serde(default, skip_serializing_if = "Option::is_none")]
1933    pub last_exit_ms: Option<u64>,
1934    /// Classification of the most recent child exit. Absent on daemons that
1935    /// predate exit-kind reporting.
1936    #[serde(default, skip_serializing_if = "Option::is_none")]
1937    pub last_exit_kind: Option<TerminalExitKind>,
1938    /// Replacement processes spawned for this module so far, against the budget
1939    /// that disables it.
1940    ///
1941    /// THIS IS THE COUNTER THAT ENDS A MODULE, and it is not the one beside it.
1942    /// `SupervisorHealthEntry::consecutive_failures` returns to zero on any
1943    /// successful probe, so a module can miss probes all day and read zero; this
1944    /// one only decreases when an operator restarts, reloads, or re-enables the
1945    /// module. Reaching the budget moves it to `Failed` and it stays there until
1946    /// somebody intervenes.
1947    ///
1948    /// So a module one restart from being disabled is indistinguishable from a
1949    /// freshly booted one unless this pair is read. Both are reported together
1950    /// because the count alone does not say how close it is.
1951    ///
1952    /// Absent from daemons predating the field, which is why it is optional
1953    /// rather than defaulted to zero: zero would assert a full budget.
1954    #[serde(default, skip_serializing_if = "Option::is_none")]
1955    pub restart_count: Option<u32>,
1956    /// Replacement processes this module is allowed before it is disabled. See
1957    /// `restart_count`; absent on daemons predating the field.
1958    #[serde(default, skip_serializing_if = "Option::is_none")]
1959    pub max_restarts: Option<u32>,
1960    /// Replacement processes spawned over this module's entire supervisor lifetime.
1961    /// Unlike `restart_count`, this value is never reset by an operator action.
1962    #[serde(default, skip_serializing_if = "Option::is_none")]
1963    pub lifetime_restarts: Option<u32>,
1964    /// Successful child spawns in this daemon incarnation. Zero means the
1965    /// module has not successfully spawned; every successful spawn increments
1966    /// the value exactly once.
1967    #[serde(default, skip_serializing_if = "Option::is_none")]
1968    pub spawn_generation: Option<u64>,
1969    /// The span `restart_count` is counted over, in seconds. The crash budget is
1970    /// a RATE, not a lifetime total: `restart_count` counts only the restarts
1971    /// inside the last `restart_window_secs`, and older ones no longer hold a
1972    /// slot. Without this field a reader cannot tell "2 of 3 crashes, ever" from
1973    /// "2 of 3 crashes in the last ten minutes", and those two call for opposite
1974    /// reactions.
1975    ///
1976    /// Absent on daemons predating the windowed budget, where the count really
1977    /// was a lifetime total.
1978    #[serde(default, skip_serializing_if = "Option::is_none")]
1979    pub restart_window_secs: Option<u64>,
1980    /// Effective drain budget for this module, in milliseconds. This is the
1981    /// resolved policy the running supervisor uses, not a config-file reread.
1982    /// Absent on older daemons.
1983    #[serde(default, skip_serializing_if = "Option::is_none")]
1984    pub drain_timeout_ms: Option<u64>,
1985    /// Effective base delay before a crash restart, in milliseconds. Absent on
1986    /// older daemons.
1987    #[serde(default, skip_serializing_if = "Option::is_none")]
1988    pub restart_backoff_ms: Option<u64>,
1989    /// Effective maximum delay before a crash restart, in milliseconds. Absent
1990    /// on older daemons.
1991    #[serde(default, skip_serializing_if = "Option::is_none")]
1992    pub restart_max_backoff_ms: Option<u64>,
1993    /// Memory and cumulative CPU time of the module's process, read when this
1994    /// list was answered. Report only: the daemon keeps no history and acts on
1995    /// none of it.
1996    ///
1997    /// It describes the one process the supervisor spawned (its `pid`), not
1998    /// processes that one has started in turn, so a module that forks workers
1999    /// reports only its own share.
2000    ///
2001    /// Absent means the daemon predates the field. A daemon that has the field
2002    /// but could not read the process (not running, unsupported platform, read
2003    /// failed) says so with `Unavailable` and a reason, so neither case can be
2004    /// mistaken for a process using nothing.
2005    #[serde(default, skip_serializing_if = "Option::is_none")]
2006    pub resources: Option<ChildResourceUsage>,
2007}
2008
2009/// A module process's memory and CPU time as read at list time, or why none
2010/// could be read.
2011#[derive(Debug, Clone, PartialEq)]
2012pub enum ChildResourceUsage {
2013    Measured(ChildResourceReading),
2014    Unavailable {
2015        reason: ChildResourceUnavailableReason,
2016    },
2017    /// Future discriminator. `body` retains the complete ordered object; `tag`
2018    /// is its decoded discriminator projection.
2019    Unknown {
2020        tag: String,
2021        body: OrderedJsonObject,
2022    },
2023}
2024
2025/// One reading of a module process's memory and CPU time.
2026#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
2027pub struct ChildResourceReading {
2028    /// Memory in bytes, measured as `memory_kind` says. The figures differ
2029    /// by platform and are not comparable across kinds.
2030    pub memory_bytes: u64,
2031    pub memory_kind: ChildMemoryKind,
2032    /// Bytes swapped out, where the platform reports it per process (Linux).
2033    /// Absent means not reported, not zero.
2034    #[serde(default, skip_serializing_if = "Option::is_none")]
2035    pub swap_bytes: Option<u64>,
2036    /// CPU time spent in user mode since the process started, in
2037    /// milliseconds. Cumulative, not a rate: a percentage needs two readings
2038    /// and the elapsed time between them.
2039    pub cpu_user_ms: u64,
2040    /// CPU time spent in the kernel on the process's behalf since it started,
2041    /// in milliseconds.
2042    pub cpu_system_ms: u64,
2043}
2044
2045open_string_enum! {
2046    /// What `ChildResourceReading::memory_bytes` measures.
2047    ChildMemoryKind {
2048        /// macOS `phys_footprint`: memory the kernel charges to the process,
2049        /// the figure jetsam acts on. Unlike resident size it does not count
2050        /// pages an allocator has already released with `MADV_FREE`.
2051        PhysFootprint => "phys_footprint",
2052        /// Linux `VmRSS`: pages resident in RAM, shared file-backed pages
2053        /// included and swapped-out pages excluded.
2054        ResidentSet => "resident_set",
2055    }
2056}
2057
2058open_string_enum! {
2059    /// Why a module process's resources could not be read.
2060    ChildResourceUnavailableReason {
2061        /// The module has no running process.
2062        NotRunning => "not_running",
2063        /// The daemon's platform has no per-process source.
2064        UnsupportedPlatform => "unsupported_platform",
2065        /// The process could not be read, typically because it exited while
2066        /// the list was being answered.
2067        Unreadable => "unreadable",
2068        /// The pid no longer names the process the supervisor spawned, so a
2069        /// reading would describe some other process.
2070        ProcessIdentityUnconfirmed => "process_identity_unconfirmed",
2071    }
2072}
2073
2074#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq)]
2075#[serde(rename_all = "snake_case")]
2076pub enum SupervisorHealthStatus {
2077    Ok,
2078    Degraded,
2079    Failing,
2080    Unresponsive,
2081    Unknown,
2082}
2083
2084#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
2085pub struct SupervisorHealthEntry {
2086    pub module_id: String,
2087    pub status: SupervisorHealthStatus,
2088    /// The module's own human-readable note on its state. Absent means the
2089    /// module said nothing, which is the ordinary shape for a healthy module and
2090    /// is NOT a claim that nothing is wrong. Never parse it: it is prose the
2091    /// module may reword freely, and `status` plus `metrics` are the machine
2092    /// surface.
2093    #[serde(default, skip_serializing_if = "Option::is_none")]
2094    pub detail: Option<String>,
2095    /// The module's own metrics object, relayed opaquely. Absent means the module
2096    /// published none on this probe — either it reports no metrics at all, or the
2097    /// probe did not reach it — so absence cannot distinguish "nothing to report"
2098    /// from "nobody asked". Read `last_probe_ms` to tell those apart.
2099    #[serde(default, skip_serializing_if = "Option::is_none")]
2100    pub metrics: Option<serde_json::Value>,
2101    pub consecutive_failures: u32,
2102    /// Number of recurring health replies received after their daemon deadline.
2103    /// Each increment is evidence that the module remained alive despite a miss.
2104    #[serde(default)]
2105    pub late_answer_count: u64,
2106    /// End-to-end latency of the newest late reply, measured from probe start.
2107    #[serde(default, skip_serializing_if = "Option::is_none")]
2108    pub last_late_answer_latency_ms: Option<u64>,
2109    /// The escalation the supervisor last took for this module (report, restart,
2110    /// alert). Absent means NO ACTION HAS EVER BEEN TAKEN, not that the last one
2111    /// succeeded — a module that has never misbehaved and one whose action record
2112    /// predates a daemon restart both present as absent.
2113    #[serde(default)]
2114    pub last_action: Option<String>,
2115    /// When `last_action` was taken, as unix milliseconds. Absent exactly when
2116    /// `last_action` is absent; the pair moves together.
2117    #[serde(default)]
2118    pub last_action_ms: Option<u64>,
2119    /// When the daemon last collected this entry, as unix milliseconds.
2120    ///
2121    /// `supervisor.health` answers from the supervisor's STORED record rather
2122    /// than probing, so every field above describes some moment in the past and
2123    /// nothing here said which. That matters most right after a restart, where
2124    /// the surface is used to confirm a deploy: a record collected before the
2125    /// restart reports the OLD process, reads as a failed deploy, and invites a
2126    /// redeploy of something that was already correct.
2127    ///
2128    /// `None` means never probed — distinct from probed-long-ago, and the reader
2129    /// must not collapse them. Absent on modules that advertise no health
2130    /// capability, which is why it is optional rather than defaulted to zero.
2131    #[serde(default, skip_serializing_if = "Option::is_none")]
2132    pub last_probe_ms: Option<u64>,
2133}
2134
2135#[cfg(test)]
2136mod tests {
2137    use super::*;
2138    use subc_protocol::{BindIdentity, RouteTarget};
2139
2140    #[test]
2141    fn legacy_terminal_decoder_ignores_deliberate_severance_kind() {
2142        let entry = TerminalEntry {
2143            daemon_incarnation: Some("daemon-before-restart".into()),
2144            exit_code: Some(1),
2145            exit_signal: None,
2146            at_ms: 1_700_000_000_123,
2147            disposition: TerminalDisposition::Restarting,
2148            exit_kind: Some(TerminalExitKind::DeliberateSeverance),
2149            disposition_detail: None,
2150        };
2151        let wire = serde_json::to_string(&entry).expect("terminal entry serializes");
2152        assert_eq!(
2153            serde_json::from_str::<serde_json::Value>(&wire).expect("terminal entry is JSON")
2154                ["exit_kind"],
2155            "deliberate_severance"
2156        );
2157
2158        #[derive(serde::Deserialize)]
2159        struct LegacyTerminalEntry {
2160            exit_code: Option<i32>,
2161            exit_signal: Option<i32>,
2162            at_ms: u64,
2163            disposition: TerminalDisposition,
2164        }
2165
2166        let decoded: LegacyTerminalEntry =
2167            serde_json::from_str(&wire).expect("legacy decoder keeps the terminal record");
2168        assert_eq!(decoded.exit_code, Some(1));
2169        assert_eq!(decoded.exit_signal, None);
2170        assert_eq!(decoded.at_ms, 1_700_000_000_123);
2171        assert_eq!(decoded.disposition, TerminalDisposition::Restarting);
2172
2173        let future_wire = wire.replace("deliberate_severance", "future_exit_kind");
2174        let future: TerminalEntry =
2175            serde_json::from_str(&future_wire).expect("new decoder keeps a future terminal kind");
2176        assert_eq!(
2177            future.exit_kind,
2178            Some(TerminalExitKind::Unknown("future_exit_kind".to_string()))
2179        );
2180    }
2181
2182    #[test]
2183    fn terminal_incarnation_is_optional_for_older_daemons() {
2184        let entry: TerminalEntry = serde_json::from_value(serde_json::json!({
2185            "at_ms": 123,
2186            "disposition": "stopped"
2187        }))
2188        .unwrap();
2189        let encoded = serde_json::to_value(&entry).unwrap();
2190        assert_eq!(
2191            (entry.daemon_incarnation, encoded.get("daemon_incarnation")),
2192            (None, None)
2193        );
2194    }
2195
2196    #[test]
2197    fn route_poll_uses_kind_field() {
2198        let body = serde_json::to_value(ClientControlRequest::RoutePoll {
2199            route_channel: 7,
2200            route_epoch: 11,
2201            kind: PollKind::Status,
2202        })
2203        .unwrap();
2204
2205        assert_eq!(body["op"], "route.poll");
2206        assert_eq!(body["route_epoch"], 11);
2207        assert_eq!(body["kind"], "status");
2208        assert!(body.get("op").is_some());
2209    }
2210
2211    #[test]
2212    fn route_open_is_internally_tagged() {
2213        let request = ClientControlRequest::RouteOpen {
2214            target: RouteTarget::ToolProvider {
2215                module_id: "aft".to_string(),
2216            },
2217            identity: BindIdentity::new("/tmp/project", "opencode", "session-1"),
2218            consumer_identity: None,
2219            consumer_capabilities: None,
2220            admission_facts: None,
2221        };
2222
2223        let body = serde_json::to_value(request).unwrap();
2224        assert_eq!(body["op"], "route.open");
2225        assert_eq!(body["target"]["kind"], "tool_provider");
2226        assert!(body.get("consumer_identity").is_none());
2227        assert!(body.get("consumer_capabilities").is_none());
2228    }
2229
2230    #[test]
2231    fn route_open_without_optional_fields_still_decodes() {
2232        let body = serde_json::json!({
2233            "op": "route.open",
2234            "target": { "kind": "tool_provider", "module_id": "aft" },
2235            "identity": {
2236                "project_root": "/tmp/project",
2237                "harness": "opencode",
2238                "session": "session-1"
2239            }
2240        });
2241
2242        let decoded: ClientControlRequest = serde_json::from_value(body).unwrap();
2243        let ClientControlRequest::RouteOpen {
2244            consumer_identity,
2245            consumer_capabilities,
2246            admission_facts,
2247            ..
2248        } = decoded
2249        else {
2250            panic!("decoded wrong request variant");
2251        };
2252        assert_eq!(consumer_identity, None);
2253        assert_eq!(consumer_capabilities, None);
2254        assert_eq!(admission_facts, None);
2255    }
2256
2257    #[test]
2258    fn new_route_closed_decoder_defaults_fields_absent_from_old_daemon() {
2259        let old_wire = r#"{"op":"route.closed","module_id":"aft-tools","reason":"crash","drained":false,"abandoned":0}"#;
2260        let decoded: ClientControlPush = serde_json::from_str(old_wire).unwrap();
2261        match decoded {
2262            ClientControlPush::RouteClosed {
2263                excluded_subscriptions,
2264                terminal,
2265                ..
2266            } => {
2267                assert_eq!(excluded_subscriptions, 0);
2268                assert_eq!(terminal, None);
2269            }
2270            other => panic!("unexpected push: {other:?}"),
2271        }
2272        assert!(!serde_json::to_string(&decoded)
2273            .unwrap()
2274            .contains("terminal"));
2275    }
2276
2277    #[test]
2278    fn old_route_closed_decoder_ignores_new_terminal_field() {
2279        #[derive(serde::Deserialize)]
2280        #[serde(tag = "op")]
2281        enum LegacyClientControlPush {
2282            #[serde(rename = "route.closed")]
2283            RouteClosed {
2284                module_id: String,
2285                reason: RouteCloseReason,
2286                drained: bool,
2287                abandoned: u32,
2288            },
2289        }
2290
2291        let wire = r#"{"op":"route.closed","module_id":"aft-tools","reason":"crash","drained":false,"abandoned":0,"excluded_subscriptions":3,"terminal":true}"#;
2292        let decoded: LegacyClientControlPush = serde_json::from_str(wire).unwrap();
2293        match decoded {
2294            LegacyClientControlPush::RouteClosed {
2295                module_id,
2296                reason,
2297                drained,
2298                abandoned,
2299            } => {
2300                assert_eq!(module_id, "aft-tools");
2301                assert_eq!(reason, RouteCloseReason::Crash);
2302                assert!(!drained);
2303                assert_eq!(abandoned, 0);
2304            }
2305        }
2306    }
2307
2308    #[test]
2309    fn supervisor_routes_is_a_control_plane_request() {
2310        let body = serde_json::json!({
2311            "op": "supervisor.routes",
2312            "module_id": "aft"
2313        });
2314
2315        let request: ClientControlRequest = serde_json::from_value(body.clone()).unwrap();
2316        assert_eq!(serde_json::to_value(request).unwrap(), body);
2317    }
2318
2319    #[test]
2320    fn diagnostic_string_enums_retain_unknown_wire_values() {
2321        let reason: RunningImageUnavailableReason =
2322            serde_json::from_str("\"future_reason\"").unwrap();
2323        let disposition: TerminalDisposition =
2324            serde_json::from_str("\"future_disposition\"").unwrap();
2325
2326        assert_eq!(
2327            reason,
2328            RunningImageUnavailableReason::Unknown("future_reason".to_string())
2329        );
2330        assert_eq!(
2331            disposition,
2332            TerminalDisposition::Unknown("future_disposition".to_string())
2333        );
2334    }
2335
2336    #[test]
2337    fn diagnostic_string_enums_preserve_existing_wire_names() {
2338        let names = [
2339            (RunningImageUnavailableReason::NotRunning, "not_running"),
2340            (
2341                RunningImageUnavailableReason::UnsupportedPlatform,
2342                "unsupported_platform",
2343            ),
2344            (
2345                RunningImageUnavailableReason::RunningExecutableUnreadable,
2346                "running_executable_unreadable",
2347            ),
2348            (
2349                RunningImageUnavailableReason::SpawnedPathUnreadable,
2350                "spawned_path_unreadable",
2351            ),
2352            (RunningImageUnavailableReason::HashFailed, "hash_failed"),
2353            (
2354                RunningImageUnavailableReason::ProcessIdentityUnconfirmed,
2355                "process_identity_unconfirmed",
2356            ),
2357        ];
2358        for (value, expected) in names {
2359            let wire = serde_json::to_string(&value).unwrap();
2360            assert_eq!(wire, format!("\"{expected}\""));
2361            let decoded: RunningImageUnavailableReason = serde_json::from_str(&wire).unwrap();
2362            assert_eq!(decoded, value);
2363        }
2364
2365        for (value, expected) in [
2366            (TerminalDisposition::Stopped, "stopped"),
2367            (TerminalDisposition::Disabled, "disabled"),
2368            (TerminalDisposition::Failed, "failed"),
2369            (TerminalDisposition::Restarting, "restarting"),
2370            (TerminalDisposition::DaemonShutdown, "daemon_shutdown"),
2371        ] {
2372            let wire = serde_json::to_string(&value).unwrap();
2373            assert_eq!(wire, format!("\"{expected}\""));
2374            let decoded: TerminalDisposition = serde_json::from_str(&wire).unwrap();
2375            assert_eq!(decoded, value);
2376        }
2377    }
2378
2379    #[test]
2380    fn diagnostic_string_enums_reject_non_string_bodies() {
2381        assert!(serde_json::from_str::<RunningImageUnavailableReason>("42").is_err());
2382        assert!(serde_json::from_str::<TerminalDisposition>("{\"value\":\"failed\"}").is_err());
2383    }
2384
2385    #[test]
2386    fn unknown_provenance_reason_does_not_discard_healthy_siblings() {
2387        let body = serde_json::json!({
2388            "op": "supervisor.provenance",
2389            "daemon": {
2390                "daemon_build": {},
2391                "daemon_observed": {
2392                    "running_image": {
2393                        "status": "unavailable",
2394                        "reason": "not_running"
2395                    }
2396                }
2397            },
2398            "modules": [
2399                {
2400                    "module_id": "future",
2401                    "module_declared": { "status": "unverifiable" },
2402                    "daemon_observed": {
2403                        "running_image": {
2404                            "status": "unavailable",
2405                            "reason": "future_reason"
2406                        }
2407                    }
2408                },
2409                {
2410                    "module_id": "healthy-a",
2411                    "module_declared": { "status": "unverifiable" },
2412                    "daemon_observed": {
2413                        "running_image": {
2414                            "status": "match",
2415                            "evidence": {
2416                                "method": "linux_proc_sha256",
2417                                "digest": "aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa"
2418                            }
2419                        }
2420                    }
2421                },
2422                {
2423                    "module_id": "healthy-b",
2424                    "module_declared": { "status": "unverifiable" },
2425                    "daemon_observed": {
2426                        "running_image": {
2427                            "status": "unavailable",
2428                            "reason": "unsupported_platform"
2429                        }
2430                    }
2431                }
2432            ]
2433        });
2434
2435        let decoded: ClientControlResponse = serde_json::from_value(body).unwrap();
2436        let ClientControlResponse::SupervisorProvenance { modules, .. } = decoded else {
2437            panic!("decoded wrong response variant");
2438        };
2439        assert_eq!(modules.len(), 3);
2440        assert_eq!(modules[0].module_id, "future");
2441        assert_eq!(
2442            modules[0].daemon_observed.running_image,
2443            RunningImageAgreement::Unavailable {
2444                reason: RunningImageUnavailableReason::Unknown("future_reason".to_string())
2445            }
2446        );
2447        assert_eq!(modules[1].module_id, "healthy-a");
2448        assert_eq!(modules[2].module_id, "healthy-b");
2449    }
2450
2451    #[test]
2452    fn tagged_unknown_values_retain_tag_and_body() {
2453        macro_rules! assert_unknown_round_trip {
2454            ($ty:ident, $field:literal, $value:expr) => {
2455                let value = $value;
2456                let wire = serde_json::to_string(&value).unwrap();
2457                let decoded: $ty = serde_json::from_str(&wire).unwrap();
2458                match decoded {
2459                    $ty::Unknown { tag, body } => {
2460                        assert_eq!(tag, value[$field].as_str().unwrap());
2461                        assert_eq!(serde_json::to_value(&body).unwrap(), value);
2462                    }
2463                    _ => panic!("decoded known variant"),
2464                }
2465            };
2466        }
2467
2468        assert_unknown_round_trip!(
2469            ModuleDeclaredProvenance,
2470            "status",
2471            serde_json::json!({"status": "future", "build": {"version": 7}})
2472        );
2473        assert_unknown_round_trip!(
2474            RunningImageAgreement,
2475            "status",
2476            serde_json::json!({"status": "future", "evidence": {"digest": "abc"}})
2477        );
2478        assert_unknown_round_trip!(
2479            RunningImageEvidence,
2480            "method",
2481            serde_json::json!({"method": "future", "digest": "abc"})
2482        );
2483        assert_unknown_round_trip!(
2484            SupervisorRouteConsumer,
2485            "kind",
2486            serde_json::json!({"kind": "future", "module_id": "m"})
2487        );
2488        assert_unknown_round_trip!(
2489            StderrCaptureState,
2490            "state",
2491            serde_json::json!({"state": "future", "reason": "because"})
2492        );
2493        assert_unknown_round_trip!(
2494            StderrTailEntry,
2495            "kind",
2496            serde_json::json!({"kind": "future", "text": "line"})
2497        );
2498        assert_unknown_round_trip!(
2499            ChildResourceUsage,
2500            "status",
2501            serde_json::json!({"status": "future", "memory_bytes": 1})
2502        );
2503    }
2504
2505    #[test]
2506    fn child_resource_usage_round_trips_both_known_states() {
2507        let measured = ChildResourceUsage::Measured(ChildResourceReading {
2508            memory_bytes: 0,
2509            memory_kind: ChildMemoryKind::ResidentSet,
2510            swap_bytes: Some(0),
2511            cpu_user_ms: 0,
2512            cpu_system_ms: 0,
2513        });
2514        let wire = serde_json::to_value(&measured).unwrap();
2515        assert_eq!(
2516            wire,
2517            serde_json::json!({
2518                "status": "measured",
2519                "memory_bytes": 0,
2520                "memory_kind": "resident_set",
2521                "swap_bytes": 0,
2522                "cpu_user_ms": 0,
2523                "cpu_system_ms": 0
2524            })
2525        );
2526        assert_eq!(
2527            serde_json::from_value::<ChildResourceUsage>(wire).unwrap(),
2528            measured
2529        );
2530
2531        let unavailable = ChildResourceUsage::Unavailable {
2532            reason: ChildResourceUnavailableReason::NotRunning,
2533        };
2534        let wire = serde_json::to_value(&unavailable).unwrap();
2535        assert_eq!(
2536            wire,
2537            serde_json::json!({"status": "unavailable", "reason": "not_running"})
2538        );
2539        assert_eq!(
2540            serde_json::from_value::<ChildResourceUsage>(wire).unwrap(),
2541            unavailable
2542        );
2543    }
2544
2545    #[test]
2546    fn a_stderr_line_decodes_with_and_without_its_capture_time() {
2547        // A current daemon stamps each line; an older one sends no `at_ms`.
2548        // Both must decode, and the absent case must stay absent rather than
2549        // turn into a time nobody recorded.
2550        let stamped: StderrTailEntry = serde_json::from_str(
2551            r#"{"kind":"line","text":"boom","truncated":true,"at_ms":1789801440685}"#,
2552        )
2553        .unwrap();
2554        assert_eq!(
2555            stamped,
2556            StderrTailEntry::Line {
2557                text: "boom".to_string(),
2558                truncated: true,
2559                at_ms: Some(1_789_801_440_685),
2560            }
2561        );
2562        let unstamped: StderrTailEntry =
2563            serde_json::from_str(r#"{"kind":"line","text":"boom"}"#).unwrap();
2564        assert_eq!(
2565            unstamped,
2566            StderrTailEntry::Line {
2567                text: "boom".to_string(),
2568                truncated: false,
2569                at_ms: None,
2570            }
2571        );
2572        // Absent stays absent on the way out, so a reply without stamps is
2573        // exactly what an older daemon would have sent.
2574        assert_eq!(
2575            serde_json::to_string(&unstamped).unwrap(),
2576            r#"{"kind":"line","text":"boom"}"#
2577        );
2578        assert_eq!(
2579            serde_json::to_value(&stamped).unwrap()["at_ms"],
2580            serde_json::json!(1_789_801_440_685u64)
2581        );
2582    }
2583
2584    #[test]
2585    fn tagged_unknown_values_round_trip_the_original_json() {
2586        let wire = r#"{"kind":"future_consumer","detail":{"z":1}}"#;
2587        let decoded: SupervisorRouteConsumer = serde_json::from_str(wire).unwrap();
2588        assert_eq!(serde_json::to_string(&decoded).unwrap(), wire);
2589    }
2590
2591    #[test]
2592    fn tagged_unknown_values_round_trip_trailing_tag() {
2593        let route_wire = r#"{"detail":{"z":1},"kind":"future_consumer"}"#;
2594        let route: SupervisorRouteConsumer = serde_json::from_str(route_wire).unwrap();
2595        assert_eq!(serde_json::to_string(&route).unwrap(), route_wire);
2596
2597        let stderr_wire = r#"{"reason":"because","state":"future_state"}"#;
2598        let stderr: StderrCaptureState = serde_json::from_str(stderr_wire).unwrap();
2599        assert_eq!(serde_json::to_string(&stderr).unwrap(), stderr_wire);
2600    }
2601
2602    #[test]
2603    fn tagged_unknown_values_round_trip_middle_tag() {
2604        let route_wire = r#"{"a":1,"kind":"future_x","b":2}"#;
2605        let route: SupervisorRouteConsumer = serde_json::from_str(route_wire).unwrap();
2606        assert_eq!(serde_json::to_string(&route).unwrap(), route_wire);
2607
2608        let stderr_wire = r#"{"a":1,"state":"future_state","b":2}"#;
2609        let stderr: StderrCaptureState = serde_json::from_str(stderr_wire).unwrap();
2610        assert_eq!(serde_json::to_string(&stderr).unwrap(), stderr_wire);
2611    }
2612
2613    #[test]
2614    fn tagged_unknown_values_round_trip_deep_payload() {
2615        let route_wire = r#"{"a":{"n":[1,2]},"kind":"future_x","zz":"s","b":null}"#;
2616        let route: SupervisorRouteConsumer = serde_json::from_str(route_wire).unwrap();
2617        assert_eq!(serde_json::to_string(&route).unwrap(), route_wire);
2618
2619        let stderr_wire = r#"{"a":{"n":[1,2]},"state":"future_state","zz":"s","b":null}"#;
2620        let stderr: StderrCaptureState = serde_json::from_str(stderr_wire).unwrap();
2621        assert_eq!(serde_json::to_string(&stderr).unwrap(), stderr_wire);
2622    }
2623
2624    #[test]
2625    fn tagged_unknown_values_reject_non_object_bodies() {
2626        for wire in ["42", r#""future""#, "[]"] {
2627            assert!(serde_json::from_str::<SupervisorRouteConsumer>(wire).is_err());
2628            assert!(serde_json::from_str::<StderrCaptureState>(wire).is_err());
2629        }
2630    }
2631
2632    #[test]
2633    fn duplicate_discriminators_reject_without_panicking() {
2634        assert_eq!(
2635            serde_json::from_str::<ModuleDeclaredProvenance>(r#"{"status":"unverifiable"}"#)
2636                .unwrap(),
2637            ModuleDeclaredProvenance::Unverifiable
2638        );
2639        match serde_json::from_str::<ModuleDeclaredProvenance>(r#"{"status":"future_thing"}"#)
2640            .unwrap()
2641        {
2642            ModuleDeclaredProvenance::Unknown { tag, .. } => assert_eq!(tag, "future_thing"),
2643            _ => panic!("future discriminator decoded as a known variant"),
2644        }
2645
2646        let wires = [
2647            r#"{"status":"reported","status":"unverifiable"}"#,
2648            r#"{"status":"unverifiable","status":"reported"}"#,
2649            r#"{"status":"reported","build":{},"status":"unverifiable"}"#,
2650            r#"{"status":"unverifiable","build":{},"status":"reported"}"#,
2651        ];
2652
2653        for wire in wires {
2654            let result =
2655                std::panic::catch_unwind(|| serde_json::from_str::<ModuleDeclaredProvenance>(wire));
2656            assert!(result.is_ok(), "duplicate discriminator panicked: {wire}");
2657            assert!(
2658                result.unwrap().is_err(),
2659                "duplicate discriminator decoded: {wire}"
2660            );
2661        }
2662
2663        let wire = r#"{"state":"captured","state":"incomplete","reason":"x"}"#;
2664        let result = std::panic::catch_unwind(|| serde_json::from_str::<StderrCaptureState>(wire));
2665        assert!(result.is_ok(), "duplicate discriminator panicked: {wire}");
2666        assert!(
2667            result.unwrap().is_err(),
2668            "duplicate discriminator decoded: {wire}"
2669        );
2670    }
2671
2672    #[test]
2673    fn nested_unknown_values_round_trip_without_normalizing_member_order() {
2674        let known_wire =
2675            r#"{"status":"match","evidence":{"method":"linux_proc_sha256","digest":"abc"}}"#;
2676        let known: RunningImageAgreement = serde_json::from_str(known_wire).unwrap();
2677        assert_eq!(serde_json::to_string(&known).unwrap(), known_wire);
2678
2679        for wire in [
2680            r#"{"kind":"future_x","detail":{"zeta":1,"alpha":2}}"#,
2681            r#"{"kind":"future_x","d":{"b":{"zz":1,"aa":2}}}"#,
2682        ] {
2683            let decoded: SupervisorRouteConsumer = serde_json::from_str(wire).unwrap();
2684            assert_eq!(serde_json::to_string(&decoded).unwrap(), wire);
2685        }
2686
2687        for wire in [
2688            r#"{"status":"match","evidence":{"method":"future_probe","zz":1,"aa":2}}"#,
2689            r#"{"status":"match","evidence":{"method":"future_probe","d":{"zz":1,"aa":2}}}"#,
2690        ] {
2691            let decoded: RunningImageAgreement = serde_json::from_str(wire).unwrap();
2692            assert_eq!(serde_json::to_string(&decoded).unwrap(), wire);
2693        }
2694
2695        let wire = r#"{"status":"mismatch","running":{"detail":{"z":1},"method":"future_running"},"disk":{"method":"future_disk","detail":{"z":1}}}"#;
2696        let decoded: RunningImageAgreement = serde_json::from_str(wire).unwrap();
2697        assert_eq!(serde_json::to_string(&decoded).unwrap(), wire);
2698
2699        let wire = r#"{"capture":{"state":"captured"},"entries":[{"detail":{"z":1,"a":2},"kind":"future_line"},{"kind":"future_restart","meta":{"b":{"zz":1,"aa":2}}}]}"#;
2700        let decoded: StderrTail = serde_json::from_str(wire).unwrap();
2701        assert_eq!(serde_json::to_string(&decoded).unwrap(), wire);
2702    }
2703
2704    #[test]
2705    fn tagged_unknown_member_does_not_discard_known_siblings() {
2706        let body = serde_json::json!({
2707            "modules": [{
2708                "module_id": "target",
2709                "routes": [
2710                    {"consumer": {"kind": "future_consumer", "module_id": "m", "detail": {"retry": true}}, "age_ms": 0, "draining": false},
2711                    {"consumer": {"kind": "direct", "connection_id": 7}, "age_ms": 0, "draining": false}
2712                ]
2713            }]
2714        });
2715        let decoded: ClientControlResponse = serde_json::from_value(
2716            serde_json::json!({"op": "supervisor.routes", "modules": body["modules"]}),
2717        )
2718        .unwrap();
2719        let ClientControlResponse::SupervisorRoutes { modules } = decoded else {
2720            panic!("decoded wrong response variant");
2721        };
2722        assert_eq!(modules[0].routes.len(), 2);
2723        assert_eq!(
2724            modules[0].routes[1].consumer,
2725            SupervisorRouteConsumer::Direct { connection_id: 7 }
2726        );
2727    }
2728}
2729
2730#[cfg(test)]
2731mod launch_nonce_redaction_tests {
2732    use super::*;
2733
2734    const NONCE: &str = "nonce-f00dfeed1234abcd";
2735
2736    fn identity() -> ConsumerIdentity {
2737        ConsumerIdentity {
2738            module_id: "wernicke".to_string(),
2739            launch_nonce: NONCE.to_string(),
2740        }
2741    }
2742
2743    #[test]
2744    fn consumer_identity_debug_names_the_module_and_never_the_nonce() {
2745        let printed = format!("{:?}", identity());
2746        assert!(printed.contains("wernicke"), "{printed}");
2747        assert!(!printed.contains(NONCE), "launch nonce printed: {printed}");
2748    }
2749
2750    #[test]
2751    fn route_open_request_debug_never_prints_the_nonce() {
2752        let request = ClientControlRequest::RouteOpen {
2753            target: subc_protocol::RouteTarget::ToolProvider {
2754                module_id: "broca".to_string(),
2755            },
2756            identity: subc_protocol::BindIdentity::new(
2757                PathBuf::from("/tmp/project"),
2758                "test".to_string(),
2759                "session".to_string(),
2760            ),
2761            consumer_identity: Some(identity()),
2762            consumer_capabilities: None,
2763            admission_facts: None,
2764        };
2765        let printed = format!("{request:?}");
2766        assert!(!printed.contains(NONCE), "launch nonce printed: {printed}");
2767    }
2768}