Skip to main content

cloud/
proc_control.rs

1//! The **process-control channel** — how a yah-supervised process describes
2//! itself to its supervisor and to an agent, instead of being guessed at from
3//! the outside.
4//!
5//! ## Why this is opinionated
6//!
7//! Everything yah runs is supervised by something in the kamaji family, and
8//! until now the supervisor's only questions were "is the pid alive?" and "is
9//! the port open?". Both are proxies. A process that has finished booting, a
10//! process still replaying a WAL, and a process wedged on a lock all answer
11//! them identically — so the only richer signal available to an operator or
12//! an agent was grepping the log tail for a line somebody hopefully logged.
13//!
14//! So: **any process built to run under a yah camp SHOULD expose a control
15//! channel**, dev tier or cloud tier, port or no port. It is the difference
16//! between an agent reading `state = "starting", detail = "migrating 3/7"`
17//! and an agent tailing stdout hoping for a sentence.
18//!
19//! ## The contract
20//!
21//! One required verb. A conforming process answers a `status` request with a
22//! **status document**:
23//!
24//! ```json
25//! {"state":"running","ready":true,"pid":71455,"uptime_secs":41,
26//!  "detail":"3 windows open","endpoints":{"gui":"winit://main"},
27//!  "metrics":{"frames_per_sec":59.9}}
28//! ```
29//!
30//! `state` is the only required field, and its vocabulary is *exactly*
31//! kamaji's [`WorkloadState`](https://docs.rs/kamaji-proto) —
32//! `pending | starting | running | draining | exited | failed`. That is the
33//! compatibility rule that matters: the supervisor already has this enum in
34//! its wire protocol and already answers a `Probe` verb with it, so a
35//! workload that reports in the same words can be believed verbatim rather
36//! than translated. Everything else in the document is optional.
37//!
38//! ## Two transports, one document
39//!
40//! | Transport | Where | How |
41//! |---|---|---|
42//! | Unix socket | dev tier, portless or not | newline-delimited JSON: write `{"cmd":"status"}\n`, read one JSON line back |
43//! | HTTP | any tier that already serves HTTP | `GET <base><path>` (conventionally `/_yah/status`) returning the same document |
44//!
45//! The cloud tier gets this for free: a `Healthcheck { probe: Http { path } }`
46//! in `workload-spec` pointed at the status path is the *same* endpoint the
47//! dev tier reads over a socket. One document, two transports, no per-tier
48//! fork — which is the same rule W265 applies to everything else here.
49//!
50//! ## Why newline-JSON and not the kamaji postcard wire
51//!
52//! Because the producer side has to be implementable in twenty lines with no
53//! dependency, in any language, by someone whose actual job that day is their
54//! own app. kamaji's `kamaji-proto` is postcard over a framed UDS: excellent
55//! between two Rust processes that both link it, a non-starter as a thing you
56//! ask every workload in the fleet to adopt. A strict-subset JSON document
57//! that a Bun script or a Python daemon can emit is the version that actually
58//! gets adopted, and it is trivially bridged into `kamaji-proto`'s
59//! `WorkloadState` because the vocabulary was chosen to match.
60//!
61//! ## Where the socket path comes from
62//!
63//! The supervisor picks it and hands it over in the environment as
64//! **`YAH_CONTROL_SOCK`**. A conforming process binds `$YAH_CONTROL_SOCK` if
65//! it is set and does nothing if it isn't — so the same binary runs unchanged
66//! outside a camp.
67//!
68//! @arch:see(.yah/docs/working/W315-process-control-channel.md)
69//! @arch:see(.yah/docs/working/W265-service-capabilities-and-drivers.md)
70//!
71//! @yah:ticket(R715-F3, "Process-control channel phase 2: producer helper crate, run.spawn injection, kamaji Probe bridge")
72//! @yah:status(review)
73//! @yah:at(2026-08-14T22:30:16Z)
74//! @yah:assignee(agent:bundle-anthropic-ashguard)
75//! @yah:parent(R715)
76//! @arch:see(.yah/docs/working/W315-process-control-channel.md)
77//! @yah:next("Producer-side helper crate so conforming is two lines for a Rust workload (bind the socket, answer status from a Fn() -> ProcStatus). CRATE HOME IS AN OPERATOR CALL: oss/kamaji/crates/* (nearest owner, but an independent workspace and a publish surface) versus a standalone oss/procctl. It cannot live in yah-cloud where the client is, because external consumers (noisetable, in the entambi repo) need it from crates.io.")
78//! @yah:next("run.spawn does not inject YAH_CONTROL_SOCK, so agent-spawned processes have no channel even if they speak it. Injecting the var is trivial; the value only appears once the camp daemon exposes a run.status RPC to read it back. Do both together.")
79//! @yah:next("kamaji bridge: kamaji already answers a Probe verb with WorkloadState, and ProcState was chosen to match it word for word, but nothing wires the two. A workload that reports `starting` is currently believed by the reconciler and invisible to kamaji.")
80//! @yah:next("Run tab polls the status once at bring-up (surfaced via RunningWorkloadSummary.notes). Live polling - a status line that moves starting -> running while you watch - needs an RPC, not new protocol.")
81//! @yah:next("Nothing enforces the SHOULD. A lint over workload.toml (long-running component, no [process.control], no healthcheck) would turn W315 into a gate. One implementation was judged too little evidence to start failing builds over.")
82//! @yah:gotcha("The protocol is deliberately dependency-free (one newline-delimited JSON verb), so nothing REQUIRES a crate to conform. The helper is ergonomics, not a gate - do not let its crate-home question block anyone from implementing the channel by hand.")
83//! @yah:gotcha("Do not switch the wire to kamaji-proto's postcard framing to unify them. That was considered and rejected in W315: the producer side has to be implementable in twenty lines in any language, and postcard-over-framed-UDS is a Rust-links-the-crate contract.")
84//! @yah:next("DECIDED 2026-08-14 by the operator, do not re-litigate: the producer helper crate lives at oss/kamaji/crates/procctl. Kamaji already owns the WorkloadState vocabulary this protocol reuses, so helper and enum move together; accepted cost is one more publish surface in kamaji's workspace. Rejected alternative: a standalone oss/procctl. Note kamaji is an INDEPENDENT cargo workspace with an export mirror - no workspace = true inheritance from yah's root, and the crate ships outward via scripts/export-oss.sh.")
85//! @yah:handoff("All three titled items shipped. (1) PRODUCER CRATE: oss/kamaji/crates/procctl, package kamaji-procctl, lib name procctl (dir per the operator's call; package name follows the kamaji-* convention and is now listed in scripts/reserve-crate-names.sh + scripts/set-trusted-publishers.sh). Default build is std + serde only - a winit GUI with no runtime can adopt it, which was the motivating case. serve_env() returns Ok(None) when YAH_CONTROL_SOCK is unset; the server stamps pid and uptime when the producer omits them; a dead predecessor socket is reclaimed but a LIVE one is refused rather than stolen. Optional features: client (async consumer, tokio) and kamaji (ProcState -> WorkloadState).")
86//! @yah:handoff("(2) KAMAJI BRIDGE: procctl's kamaji feature holds the From impl as an exhaustive match, so a state added to either vocabulary stops compiling - that compile error is the entire mechanism keeping W315's believed-verbatim claim true. The reverse direction is deliberately absent (WorkloadState is non_exhaustive, so matching it needs a wildcard, which is the silent drift the design refuses). ProbeTarget gained control: Option<PathBuf> and healthcheck became Option<Healthcheck> (a portless GUI has no port to probe and inventing a healthcheck for it is a lie); constructors ProbeTarget::healthcheck / ::control replace the struct literals. A control socket is the WHOLE probe when present - unreachable reads Starting, never a fallback to the port, per W315. Registered on the native-exec deploy path only, read from the spec's own YAH_CONTROL_SOCK literal, because a container's socket path names a location inside its mount namespace kamaji has no route to.")
87//! @yah:handoff("(3) RUN.SPAWN INJECTION + RUN.STATUS: every spawn is now offered the channel unconditionally (a conforming process binds when set, does nothing when unset, so this costs a declining process nothing). New rpc::method::RUN_STATUS + RunStatusParams/RunStatusResult, run_status_handler in camp.rs, and a read-only run.status agent tool. The result relays the document verbatim (serde_json::Value) so a producer's own metrics/endpoints reach the caller intact. run.stop now unlinks the socket.")
88//! @yah:verify("cargo test --workspace --all-features in oss/kamaji: all green. kamaji-procctl 26 passed + 1 doc-test; kamaji-bin lib 248 passed (was 216, +5 control-channel probe tests, +3 server tests).")
89//! @yah:verify("cargo test -p yah --lib camp:: - 298 passed, including 5 new r715_f3_control_channel_tests that drive real child processes through run_spawn_handler / run_status_handler / run_stop_handler.")
90//! @yah:verify("cargo test -p yah-agent-tools --lib run_tools - 16 passed. cargo test -p yah-rpc --lib - 45 passed. cargo test -p yah-cloud --lib proc_control (oss/yubaba) - 10 passed, untouched.")
91//! @yah:verify("cargo check --workspace (root) and cargo check -p yubaba (oss/yubaba) clean; scripts/check-workspace-members.sh resolves all 58 members. kamaji-procctl also checked with default features (no client, no kamaji) so the producer half stays std-only.")
92//! @yah:verify("The sandbox rule is pinned by a test that binds from INSIDE the sandbox - a_sandboxed_child_can_bind_the_socket_it_is_handed (camp.rs), a ~15-line stdlib-Python producer spawned through run_spawn_handler. Every other test in that module binds from the test process, which is outside the sandbox and proves nothing about it. It skips (does not fail) where no python3 exists.")
93//! @yah:gotcha("MACOS SANDBOX BUG FOUND AND FIXED, and it would have made the whole injection useless on macOS. run.spawn puts the control socket at <workload_dir>/.yah-control.sock because that is the only writable path both sandbox facilities agree on (Seatbelt allows file-write* under WORKLOAD; bwrap binds workload_dir and nothing else - .yah/jit/ is not even present inside the bwrap namespace). But file-write* is NOT sufficient: Seatbelt gates AF_UNIX bind under network-bind, so a child got EPERM binding a path it could otherwise create any file at. MACOS_SANDBOX_PROFILE now carries (allow network-bind (local unix-socket (subpath (param WORKLOAD)))). Reproduced by hand with sandbox-exec before and after. plugin_host.rs shares the same profile constant and passes the same -D WORKLOAD, so it inherits the rule.")
94//! @yah:gotcha("DISCOVERED, NOT FIXED (deliberately): kamaji never registers a ProbeTarget from a spec's own healthcheck on any path except the mesofact-bundle deploy. insert_probe has exactly three call sites (server.rs registry decl, the bundle deploy, and now the native control-channel registration), so a container workload declaring a healthcheck is probed as Ready unconditionally. Fixing it means live fleet workloads that report Ready today would start reporting real probe status - a behaviour change on running infra, not a drive-by. The control registration added here cannot regress that: it only fires when the spec declares YAH_CONTROL_SOCK, so a spec without it registers nothing and behaves exactly as before.")
95//! @yah:cleanup("kamaji-procctl is not on crates.io. In-tree consumers resolve it by path so nothing is blocked, but kamaji-bin now depends on it and cargo publish --workspace publishes in topological order - the name must be reserved (scripts/reserve-crate-names.sh, entry added) before the next kamaji release, or that release fails on an unpublished dep. The external consumers this crate exists for (noisetable, entambi repo) need it published anyway.")
96
97use std::path::{Path, PathBuf};
98use std::time::Duration;
99
100use serde::{Deserialize, Serialize};
101use tokio::io::{AsyncBufReadExt, AsyncWriteExt, BufReader};
102
103/// Environment variable naming the control socket a supervised process should
104/// bind. Absent → the process is not running under a supervisor that wants a
105/// control channel, and MUST NOT fail for its absence.
106pub const CONTROL_SOCK_ENV: &str = "YAH_CONTROL_SOCK";
107
108/// Conventional HTTP path for the status document on a process that already
109/// serves HTTP. Not enforced — `[process.control] http_path` overrides it —
110/// but a service with no reason to differ should use this one.
111pub const DEFAULT_HTTP_PATH: &str = "/_yah/status";
112
113/// Lifecycle vocabulary of a supervised process.
114///
115/// Deliberately identical to `kamaji_proto::WorkloadState` so a workload's own
116/// report can be handed to the supervisor without a translation table that
117/// would rot the first time either side gained a state. Kept as a separate
118/// type rather than a re-export only because this crate must not force a
119/// `kamaji-proto` dependency on every consumer of the status document.
120#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
121#[serde(rename_all = "snake_case")]
122pub enum ProcState {
123    /// Accepted, nothing started yet.
124    Pending,
125    /// Started, not yet serving — booting, migrating, warming a cache.
126    Starting,
127    /// Serving. This is the only state that counts as ready.
128    Running,
129    /// Shutting down gracefully.
130    Draining,
131    /// Exited cleanly.
132    Exited,
133    /// Exited with a failure, or reported itself unrecoverable.
134    Failed,
135}
136
137impl ProcState {
138    /// Whether a process in this state is ready to be used.
139    ///
140    /// `Starting` is deliberately *not* ready: the entire point of the channel
141    /// is to distinguish "the port is open" from "I am serving".
142    pub fn is_ready(self) -> bool {
143        matches!(self, ProcState::Running)
144    }
145
146    /// Whether this state is terminal — no amount of further polling changes
147    /// it, so a readiness wait should fail fast rather than burn its timeout.
148    pub fn is_terminal(self) -> bool {
149        matches!(self, ProcState::Exited | ProcState::Failed)
150    }
151}
152
153/// A workload's self-description. Only [`Self::state`] is required.
154#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
155pub struct ProcStatus {
156    /// Lifecycle state, in kamaji's vocabulary.
157    pub state: ProcState,
158    /// Redundant convenience mirror of `state == running`, accepted from
159    /// producers that emit it. Never trusted over `state` — a document
160    /// claiming `{"state":"starting","ready":true}` is a producer bug, and
161    /// believing the optimistic half of it is how a supervisor reports a
162    /// half-booted process as up.
163    #[serde(default, skip_serializing_if = "Option::is_none")]
164    pub ready: Option<bool>,
165    /// Process id, when the process knows and cares to say.
166    #[serde(default, skip_serializing_if = "Option::is_none")]
167    pub pid: Option<u32>,
168    /// Seconds since the process considered itself started.
169    #[serde(default, skip_serializing_if = "Option::is_none")]
170    pub uptime_secs: Option<u64>,
171    /// Build/version string, for an operator staring at two of these.
172    #[serde(default, skip_serializing_if = "Option::is_none")]
173    pub version: Option<String>,
174    /// One human line elaborating on `state` — "replaying WAL 3/7",
175    /// "waiting for GPU". This is the field that replaces log-grepping.
176    #[serde(default, skip_serializing_if = "Option::is_none")]
177    pub detail: Option<String>,
178    /// Named addresses the process serves — `{"http":"http://127.0.0.1:4325"}`.
179    /// A portless process may legitimately name a non-URL surface here.
180    #[serde(default, skip_serializing_if = "std::collections::BTreeMap::is_empty")]
181    pub endpoints: std::collections::BTreeMap<String, String>,
182    /// Numeric gauges the process wants surfaced. Free-form on purpose: this
183    /// is a status channel, not a metrics pipeline.
184    #[serde(default, skip_serializing_if = "std::collections::BTreeMap::is_empty")]
185    pub metrics: std::collections::BTreeMap<String, f64>,
186}
187
188impl ProcStatus {
189    /// Ready iff the *state* says so. See [`Self::ready`] for why the
190    /// producer-supplied boolean does not get a vote.
191    pub fn is_ready(&self) -> bool {
192        self.state.is_ready()
193    }
194
195    /// One-line rendering for an operator-facing note or log line.
196    pub fn summary(&self) -> String {
197        let mut s = format!("{:?}", self.state).to_lowercase();
198        if let Some(detail) = &self.detail {
199            s.push_str(" — ");
200            s.push_str(detail);
201        }
202        if let Some(v) = &self.version {
203            s.push_str(&format!(" (v{v})"));
204        }
205        s
206    }
207}
208
209/// Where to ask for the status document.
210#[derive(Debug, Clone, PartialEq, Eq)]
211pub enum ControlEndpoint {
212    /// Newline-JSON over a unix domain socket — the dev-tier default.
213    Socket(PathBuf),
214    /// `GET <url>` returning the status document — for a process that already
215    /// serves HTTP, including every cloud-tier workload.
216    Http(String),
217}
218
219impl std::fmt::Display for ControlEndpoint {
220    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
221        match self {
222            ControlEndpoint::Socket(p) => write!(f, "unix:{}", p.display()),
223            ControlEndpoint::Http(u) => write!(f, "{u}"),
224        }
225    }
226}
227
228/// Ask a process for its status document once.
229///
230/// Errors mean *unreachable or unparseable*, which is not the same as
231/// unhealthy: a process that has not yet bound its socket is indistinguishable
232/// here from one that never will. Callers deciding readiness should poll with
233/// [`wait_ready`] rather than treating one error as a verdict.
234pub async fn fetch_status(endpoint: &ControlEndpoint) -> anyhow::Result<ProcStatus> {
235    match endpoint {
236        ControlEndpoint::Socket(path) => fetch_status_uds(path).await,
237        ControlEndpoint::Http(url) => {
238            let body = reqwest::Client::new()
239                .get(url)
240                .timeout(Duration::from_secs(2))
241                .send()
242                .await?
243                .error_for_status()?
244                .text()
245                .await?;
246            Ok(serde_json::from_str(&body)?)
247        }
248    }
249}
250
251/// Newline-JSON round trip: write one request line, read one response line.
252///
253/// The connection is not reused. A status poll happens every few seconds at
254/// most, and a per-call connection means a wedged reader on the producer side
255/// cannot poison later polls — worth far more here than the syscalls saved.
256async fn fetch_status_uds(path: &Path) -> anyhow::Result<ProcStatus> {
257    let mut stream = tokio::net::UnixStream::connect(path).await?;
258    stream.write_all(b"{\"cmd\":\"status\"}\n").await?;
259    stream.flush().await?;
260
261    let mut line = String::new();
262    let read = tokio::time::timeout(
263        Duration::from_secs(2),
264        BufReader::new(stream).read_line(&mut line),
265    )
266    .await
267    .map_err(|_| anyhow::anyhow!("control socket {} did not answer within 2s", path.display()))??;
268    if read == 0 {
269        anyhow::bail!(
270            "control socket {} closed without answering",
271            path.display()
272        );
273    }
274    Ok(serde_json::from_str(line.trim())?)
275}
276
277/// Outcome of waiting for a process to report itself ready.
278#[derive(Debug)]
279pub enum ReadyOutcome {
280    /// The process reported [`ProcState::Running`].
281    Ready(ProcStatus),
282    /// The process reported a terminal state — it is not coming up. Failing
283    /// here rather than burning the whole timeout is the practical difference
284    /// between a five-second and a twenty-second edit loop.
285    Terminal(ProcStatus),
286    /// The deadline passed. `last` is the most recent document read, or `None`
287    /// when the endpoint never answered at all — a distinction worth keeping
288    /// in the error message, since "never bound its socket" and "stuck in
289    /// starting" are different bugs with different fixes.
290    TimedOut { last: Option<ProcStatus> },
291}
292
293/// Poll `endpoint` until the process reports ready, reports terminal, or the
294/// timeout expires.
295pub async fn wait_ready(endpoint: &ControlEndpoint, timeout: Duration) -> ReadyOutcome {
296    let deadline = tokio::time::Instant::now() + timeout;
297    let mut last: Option<ProcStatus> = None;
298    loop {
299        if let Ok(status) = fetch_status(endpoint).await {
300            if status.is_ready() {
301                return ReadyOutcome::Ready(status);
302            }
303            if status.state.is_terminal() {
304                return ReadyOutcome::Terminal(status);
305            }
306            last = Some(status);
307        }
308        if tokio::time::Instant::now() >= deadline {
309            return ReadyOutcome::TimedOut { last };
310        }
311        tokio::time::sleep(Duration::from_millis(100)).await;
312    }
313}
314
315#[cfg(test)]
316mod tests {
317    use super::*;
318    use std::collections::BTreeMap;
319
320    /// Minimal conforming producer: accept, read a line, answer one document.
321    /// This is also the reference for how little a workload has to implement.
322    fn serve_once(path: PathBuf, docs: Vec<String>) -> tokio::task::JoinHandle<()> {
323        tokio::spawn(async move {
324            let listener = tokio::net::UnixListener::bind(&path).unwrap();
325            for doc in docs {
326                let Ok((stream, _)) = listener.accept().await else {
327                    return;
328                };
329                let (read_half, mut write_half) = stream.into_split();
330                let mut line = String::new();
331                BufReader::new(read_half).read_line(&mut line).await.ok();
332                write_half.write_all(doc.as_bytes()).await.ok();
333                write_half.write_all(b"\n").await.ok();
334                write_half.flush().await.ok();
335            }
336        })
337    }
338
339    #[test]
340    fn the_state_vocabulary_matches_kamajis_workload_state_on_the_wire() {
341        // If this ever drifts, a workload's own report can no longer be handed
342        // to the supervisor verbatim — which is the entire compatibility claim
343        // this module makes.
344        for (state, wire) in [
345            (ProcState::Pending, "\"pending\""),
346            (ProcState::Starting, "\"starting\""),
347            (ProcState::Running, "\"running\""),
348            (ProcState::Draining, "\"draining\""),
349            (ProcState::Exited, "\"exited\""),
350            (ProcState::Failed, "\"failed\""),
351        ] {
352            assert_eq!(serde_json::to_string(&state).unwrap(), wire);
353            assert_eq!(
354                serde_json::from_str::<ProcState>(wire).unwrap(),
355                state,
356                "{wire} must round-trip"
357            );
358        }
359    }
360
361    #[test]
362    fn state_is_the_only_required_field() {
363        let s: ProcStatus = serde_json::from_str(r#"{"state":"running"}"#).unwrap();
364        assert!(s.is_ready());
365        assert_eq!(s.pid, None);
366        assert!(s.endpoints.is_empty());
367        assert!(s.metrics.is_empty());
368    }
369
370    #[test]
371    fn a_full_document_parses() {
372        let s: ProcStatus = serde_json::from_str(
373            r#"{"state":"starting","ready":false,"pid":71455,"uptime_secs":41,
374                "version":"0.8.20","detail":"replaying WAL 3/7",
375                "endpoints":{"gui":"winit://main"},"metrics":{"fps":59.9}}"#,
376        )
377        .unwrap();
378        assert_eq!(s.state, ProcState::Starting);
379        assert!(!s.is_ready(), "starting is not ready");
380        assert_eq!(s.pid, Some(71455));
381        assert_eq!(s.detail.as_deref(), Some("replaying WAL 3/7"));
382        assert_eq!(s.endpoints.get("gui").map(String::as_str), Some("winit://main"));
383        assert_eq!(s.metrics.get("fps"), Some(&59.9));
384        assert_eq!(s.summary(), "starting — replaying WAL 3/7 (v0.8.20)");
385    }
386
387    /// A producer that contradicts itself must not be believed on the
388    /// optimistic half — that is precisely how a half-booted process gets
389    /// reported as up, which is the failure this channel exists to end.
390    #[test]
391    fn a_ready_flag_never_overrides_a_not_running_state() {
392        let s: ProcStatus =
393            serde_json::from_str(r#"{"state":"starting","ready":true}"#).unwrap();
394        assert_eq!(s.ready, Some(true), "the claim is preserved verbatim");
395        assert!(!s.is_ready(), "but state decides");
396    }
397
398    #[tokio::test]
399    async fn fetch_status_reads_a_document_over_a_unix_socket() {
400        let tmp = tempfile::tempdir().unwrap();
401        let sock = tmp.path().join("control.sock");
402        let server = serve_once(
403            sock.clone(),
404            vec![r#"{"state":"running","detail":"3 windows"}"#.to_string()],
405        );
406        // Give the listener a moment to bind.
407        tokio::time::sleep(Duration::from_millis(50)).await;
408
409        let got = fetch_status(&ControlEndpoint::Socket(sock)).await.unwrap();
410        assert_eq!(got.state, ProcState::Running);
411        assert_eq!(got.detail.as_deref(), Some("3 windows"));
412        server.abort();
413    }
414
415    #[tokio::test]
416    async fn wait_ready_polls_through_starting_to_running() {
417        let tmp = tempfile::tempdir().unwrap();
418        let sock = tmp.path().join("control.sock");
419        let server = serve_once(
420            sock.clone(),
421            vec![
422                r#"{"state":"starting"}"#.to_string(),
423                r#"{"state":"starting"}"#.to_string(),
424                r#"{"state":"running"}"#.to_string(),
425            ],
426        );
427        tokio::time::sleep(Duration::from_millis(50)).await;
428
429        let outcome = wait_ready(&ControlEndpoint::Socket(sock), Duration::from_secs(5)).await;
430        assert!(
431            matches!(&outcome, ReadyOutcome::Ready(s) if s.state == ProcState::Running),
432            "{outcome:?}"
433        );
434        server.abort();
435    }
436
437    /// `failed` must short-circuit. Burning the full readiness timeout on a
438    /// process that has already said it is not coming up is the slow-edit-loop
439    /// failure this outcome exists to prevent.
440    #[tokio::test]
441    async fn wait_ready_fails_fast_on_a_terminal_state() {
442        let tmp = tempfile::tempdir().unwrap();
443        let sock = tmp.path().join("control.sock");
444        let server = serve_once(
445            sock.clone(),
446            vec![r#"{"state":"failed","detail":"no GPU"}"#.to_string()],
447        );
448        tokio::time::sleep(Duration::from_millis(50)).await;
449
450        let started = std::time::Instant::now();
451        let outcome = wait_ready(&ControlEndpoint::Socket(sock), Duration::from_secs(30)).await;
452        assert!(
453            matches!(&outcome, ReadyOutcome::Terminal(s) if s.detail.as_deref() == Some("no GPU")),
454            "{outcome:?}"
455        );
456        assert!(
457            started.elapsed() < Duration::from_secs(5),
458            "must not burn the 30s timeout on a terminal state"
459        );
460        server.abort();
461    }
462
463    #[tokio::test]
464    async fn wait_ready_times_out_when_nothing_is_listening() {
465        let tmp = tempfile::tempdir().unwrap();
466        let outcome = wait_ready(
467            &ControlEndpoint::Socket(tmp.path().join("never-bound.sock")),
468            Duration::from_millis(300),
469        )
470        .await;
471        assert!(
472            matches!(outcome, ReadyOutcome::TimedOut { last: None }),
473            "an endpoint that never answered must report no last document"
474        );
475    }
476
477    #[test]
478    fn endpoints_render_distinguishably() {
479        assert_eq!(
480            ControlEndpoint::Socket(PathBuf::from("/tmp/c.sock")).to_string(),
481            "unix:/tmp/c.sock"
482        );
483        assert_eq!(
484            ControlEndpoint::Http("http://127.0.0.1:4325/_yah/status".into()).to_string(),
485            "http://127.0.0.1:4325/_yah/status"
486        );
487    }
488
489    #[test]
490    fn a_status_document_round_trips_through_serialization() {
491        let mut endpoints = BTreeMap::new();
492        endpoints.insert("http".to_string(), "http://127.0.0.1:4325".to_string());
493        let original = ProcStatus {
494            state: ProcState::Running,
495            ready: None,
496            pid: Some(9),
497            uptime_secs: Some(3),
498            version: None,
499            detail: None,
500            endpoints,
501            metrics: BTreeMap::new(),
502        };
503        let json = serde_json::to_string(&original).unwrap();
504        assert_eq!(serde_json::from_str::<ProcStatus>(&json).unwrap(), original);
505        assert!(
506            !json.contains("\"ready\""),
507            "absent optionals must not be emitted: {json}"
508        );
509    }
510}