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}