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///
155/// @yah:relay(R866, "Deployed-credential drift: apply-time value fingerprints reported back over the ProcStatus rail")
156/// @yah:at(2026-09-05T08:47:33Z)
157/// @yah:status(open)
158/// @yah:assignee(agent:bundle-anthropic-ashguard)
159/// @arch:see(.yah/docs/working/W337-credential-health-and-rotation.md)
160/// @yah:depends_on(R556-F6)
161/// @yah:next("THE PROBLEM. A credential used by a cloud workload is injected at `yah cloud apply` time from the local vault via resolve_serve_env (grep `fn resolve_serve_env` in app/yah/cli/src/cloud.rs — :7979 on 2026-09-05, but it has moved three times). That splits rot in two: (a) the vault copy rotted, which a local `yah keys doctor` probe already catches, and (b) THE DEPLOYED COPY DRIFTED FROM THE VAULT COPY — someone rotated the vault and never re-applied, or applied and the workload never restarted. In case (b) the vault probe is GREEN and the service is DOWN. This relay is (b) and only (b).")
162/// @yah:next("SCOPE THE DELTA. The apply-time PRESENCE half is already built and must NOT be rebuilt: resolve_serve_env is FATAL on an empty resolution, with an error naming the slot and the `yah keys set` fix, so a serve process is never forked with a blank credential. What it cannot catch is a value that resolves fine and is DEAD, or one that DRIFTED after apply.")
163/// @yah:next("SHAPE. At apply time store a SALTED hash of each injected value; have the workload report the hash of what it is actually running with; compare. That detects (b) without moving a secret anywhere. HARD CONSTRAINT: no secret value, and no secret LENGTH, may appear in the fingerprint path, in the reported status document, or in any log.")
164/// @yah:next("USE THE EXISTING ProcStatus RAIL — do not build a bespoke per-workload health endpoint. Right here in this file: workloads publish a `ProcStatus` self-description document (:155), fetched by `fetch_status` (:234) over a `ControlEndpoint` that is `Socket(PathBuf)` OR `Http(String)` (:211, doc: \\\"for a process that already serves HTTP, including every cloud-tier workload\\\"), at conventional path DEFAULT_HTTP_PATH = \\\"/_yah/status\\\" (:111, overridable per-workload via `[process.control] http_path`). Producer side is a published two-line helper crate, oss/kamaji/crates/procctl (serve_env / serve_at / ControlServer, lib.rs:78). Riding this rail makes the feature work for EVERY procctl-conforming workload rather than mesofact alone.")
165/// @yah:next("THE ONE GAP, and it is the first edit: ProcStatus has NO field a hex fingerprint fits. `metrics` is `BTreeMap<String, f64>` (numeric only), `endpoints` is addresses, `detail` is documented as ONE HUMAN LINE. Add a string-valued field — `env_fingerprint`, or a general free-form `labels: BTreeMap<String, String>`; THAT CHOICE IS A NAMING CALL, make it deliberately — carrying #[serde(default)] so workloads built before this change keep deserializing.")
166/// @yah:next("SURFACE IT IN THE EXISTING TABLE, not a new command. `yah cloud mirror-status` already does declared-vs-observed comparison for replicas and already has a --drift filter: `handle_mirror_status` at app/yah/cli/src/cloud.rs:11862, row type `MirrorStatusRow` at :6774, --drift applied at :11929. Add a row type; do not add a command. (Older prose cites :9684 for this — that was never a mirror-status line.)")
167/// @yah:gotcha("THERE IS NO LIVE CONSUMER YET, AND THAT GATES THIS RELAY — it is why depends_on(R556-F6) is set. Measured 2026-09-05: ZERO uncommented `vault:` declarations in any tracked TOML. All four hits are commented out — .yah/services/yah-analytics/mirrors/cloud.toml:616-618 (the `#!` cut-over block) and .yah/qed/gha-actions.toml:25 — so today there is nothing for an apply-time hash to hash and the reporting half would be dead code on both ends. Step (4) of R556-F6's cut-over uncomments that block and creates the first live declaration. Confirm the field shape against what R556-T12 actually SHIPPED, not against the comment: the comment predates it and names `cloudflare-r2-endpoint`, a slot the vault does not have.")
168/// @yah:gotcha("RIPGREP TRAP that has already cost two sessions a false reading: rg skips hidden directories by default, and every `vault:` declaration in this tree lives under .yah/. So `rg '=\\s*\"vault:\"' --glob '*.toml'` WITHOUT --hidden returns clean over a tree that is not clean. Always pass --hidden when re-measuring the trigger.")
169/// @yah:gotcha("BLAST RADIUS IS WIDER THAN IT LOOKS — weigh it before starting. This touches oss/yubaba AND oss/kamaji, which are INDEPENDENT Cargo workspaces excluded from the yah root workspace (so no `workspace = true` inheritance from the root inside them), and procctl is a crates.io PUBLISH surface, meaning a ProcStatus field change is a wire-format change for external consumers (noisetable, in the entambi repo, is named as one in this file's own R-notes). #[serde(default)] on the new field is not optional politeness; it is what keeps already-deployed workloads deserializing.")
170/// @yah:gotcha("HISTORY: this was R856-F8, deferred unbuilt across three sessions (2026-09-03/04/05) because the trigger never fired. Operator decision 2026-09-05 re-filed it here as its own relay rather than holding R856 open — R856's remaining work is vault-local and finished, while this spans two oss workspaces and a publish surface. R856-F8 is archived; its design record is W337 §5.")
171/// @yah:verify("Rotating a vault slot without re-applying shows as drift in `yah cloud mirror-status --drift`, while the local `yah keys doctor` probe still reports Valid on the same slot")
172/// @yah:verify("No secret value and no secret LENGTH appears in the fingerprint path, in the reported ProcStatus document, or in any log")
173/// @yah:verify("A workload built before the new ProcStatus field still deserializes (pin it with a test that feeds the pre-change JSON through serde), so a partial fleet roll cannot break status reporting")
174/// @yah:verify("cargo test --manifest-path oss/yubaba/Cargo.toml -p yah-cloud proc_control # the rail's existing serde round-trip tests still pass (see proc_control.rs:493)")
175#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
176pub struct ProcStatus {
177 /// Lifecycle state, in kamaji's vocabulary.
178 pub state: ProcState,
179 /// Redundant convenience mirror of `state == running`, accepted from
180 /// producers that emit it. Never trusted over `state` — a document
181 /// claiming `{"state":"starting","ready":true}` is a producer bug, and
182 /// believing the optimistic half of it is how a supervisor reports a
183 /// half-booted process as up.
184 #[serde(default, skip_serializing_if = "Option::is_none")]
185 pub ready: Option<bool>,
186 /// Process id, when the process knows and cares to say.
187 #[serde(default, skip_serializing_if = "Option::is_none")]
188 pub pid: Option<u32>,
189 /// Seconds since the process considered itself started.
190 #[serde(default, skip_serializing_if = "Option::is_none")]
191 pub uptime_secs: Option<u64>,
192 /// Build/version string, for an operator staring at two of these.
193 #[serde(default, skip_serializing_if = "Option::is_none")]
194 pub version: Option<String>,
195 /// One human line elaborating on `state` — "replaying WAL 3/7",
196 /// "waiting for GPU". This is the field that replaces log-grepping.
197 #[serde(default, skip_serializing_if = "Option::is_none")]
198 pub detail: Option<String>,
199 /// Named addresses the process serves — `{"http":"http://127.0.0.1:4325"}`.
200 /// A portless process may legitimately name a non-URL surface here.
201 #[serde(default, skip_serializing_if = "std::collections::BTreeMap::is_empty")]
202 pub endpoints: std::collections::BTreeMap<String, String>,
203 /// Numeric gauges the process wants surfaced. Free-form on purpose: this
204 /// is a status channel, not a metrics pipeline.
205 #[serde(default, skip_serializing_if = "std::collections::BTreeMap::is_empty")]
206 pub metrics: std::collections::BTreeMap<String, f64>,
207}
208
209impl ProcStatus {
210 /// Ready iff the *state* says so. See [`Self::ready`] for why the
211 /// producer-supplied boolean does not get a vote.
212 pub fn is_ready(&self) -> bool {
213 self.state.is_ready()
214 }
215
216 /// One-line rendering for an operator-facing note or log line.
217 pub fn summary(&self) -> String {
218 let mut s = format!("{:?}", self.state).to_lowercase();
219 if let Some(detail) = &self.detail {
220 s.push_str(" — ");
221 s.push_str(detail);
222 }
223 if let Some(v) = &self.version {
224 s.push_str(&format!(" (v{v})"));
225 }
226 s
227 }
228}
229
230/// Where to ask for the status document.
231#[derive(Debug, Clone, PartialEq, Eq)]
232pub enum ControlEndpoint {
233 /// Newline-JSON over a unix domain socket — the dev-tier default.
234 Socket(PathBuf),
235 /// `GET <url>` returning the status document — for a process that already
236 /// serves HTTP, including every cloud-tier workload.
237 Http(String),
238}
239
240impl std::fmt::Display for ControlEndpoint {
241 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
242 match self {
243 ControlEndpoint::Socket(p) => write!(f, "unix:{}", p.display()),
244 ControlEndpoint::Http(u) => write!(f, "{u}"),
245 }
246 }
247}
248
249/// Ask a process for its status document once.
250///
251/// Errors mean *unreachable or unparseable*, which is not the same as
252/// unhealthy: a process that has not yet bound its socket is indistinguishable
253/// here from one that never will. Callers deciding readiness should poll with
254/// [`wait_ready`] rather than treating one error as a verdict.
255pub async fn fetch_status(endpoint: &ControlEndpoint) -> anyhow::Result<ProcStatus> {
256 match endpoint {
257 ControlEndpoint::Socket(path) => fetch_status_uds(path).await,
258 ControlEndpoint::Http(url) => {
259 let body = reqwest::Client::new()
260 .get(url)
261 .timeout(Duration::from_secs(2))
262 .send()
263 .await?
264 .error_for_status()?
265 .text()
266 .await?;
267 Ok(serde_json::from_str(&body)?)
268 }
269 }
270}
271
272/// Newline-JSON round trip: write one request line, read one response line.
273///
274/// The connection is not reused. A status poll happens every few seconds at
275/// most, and a per-call connection means a wedged reader on the producer side
276/// cannot poison later polls — worth far more here than the syscalls saved.
277async fn fetch_status_uds(path: &Path) -> anyhow::Result<ProcStatus> {
278 let mut stream = tokio::net::UnixStream::connect(path).await?;
279 stream.write_all(b"{\"cmd\":\"status\"}\n").await?;
280 stream.flush().await?;
281
282 let mut line = String::new();
283 let read = tokio::time::timeout(
284 Duration::from_secs(2),
285 BufReader::new(stream).read_line(&mut line),
286 )
287 .await
288 .map_err(|_| anyhow::anyhow!("control socket {} did not answer within 2s", path.display()))??;
289 if read == 0 {
290 anyhow::bail!(
291 "control socket {} closed without answering",
292 path.display()
293 );
294 }
295 Ok(serde_json::from_str(line.trim())?)
296}
297
298/// Outcome of waiting for a process to report itself ready.
299#[derive(Debug)]
300pub enum ReadyOutcome {
301 /// The process reported [`ProcState::Running`].
302 Ready(ProcStatus),
303 /// The process reported a terminal state — it is not coming up. Failing
304 /// here rather than burning the whole timeout is the practical difference
305 /// between a five-second and a twenty-second edit loop.
306 Terminal(ProcStatus),
307 /// The deadline passed. `last` is the most recent document read, or `None`
308 /// when the endpoint never answered at all — a distinction worth keeping
309 /// in the error message, since "never bound its socket" and "stuck in
310 /// starting" are different bugs with different fixes.
311 TimedOut { last: Option<ProcStatus> },
312}
313
314/// Poll `endpoint` until the process reports ready, reports terminal, or the
315/// timeout expires.
316pub async fn wait_ready(endpoint: &ControlEndpoint, timeout: Duration) -> ReadyOutcome {
317 let deadline = tokio::time::Instant::now() + timeout;
318 let mut last: Option<ProcStatus> = None;
319 loop {
320 if let Ok(status) = fetch_status(endpoint).await {
321 if status.is_ready() {
322 return ReadyOutcome::Ready(status);
323 }
324 if status.state.is_terminal() {
325 return ReadyOutcome::Terminal(status);
326 }
327 last = Some(status);
328 }
329 if tokio::time::Instant::now() >= deadline {
330 return ReadyOutcome::TimedOut { last };
331 }
332 tokio::time::sleep(Duration::from_millis(100)).await;
333 }
334}
335
336#[cfg(test)]
337mod tests {
338 use super::*;
339 use std::collections::BTreeMap;
340
341 /// Minimal conforming producer: accept, read a line, answer one document.
342 /// This is also the reference for how little a workload has to implement.
343 fn serve_once(path: PathBuf, docs: Vec<String>) -> tokio::task::JoinHandle<()> {
344 tokio::spawn(async move {
345 let listener = tokio::net::UnixListener::bind(&path).unwrap();
346 for doc in docs {
347 let Ok((stream, _)) = listener.accept().await else {
348 return;
349 };
350 let (read_half, mut write_half) = stream.into_split();
351 let mut line = String::new();
352 BufReader::new(read_half).read_line(&mut line).await.ok();
353 write_half.write_all(doc.as_bytes()).await.ok();
354 write_half.write_all(b"\n").await.ok();
355 write_half.flush().await.ok();
356 }
357 })
358 }
359
360 #[test]
361 fn the_state_vocabulary_matches_kamajis_workload_state_on_the_wire() {
362 // If this ever drifts, a workload's own report can no longer be handed
363 // to the supervisor verbatim — which is the entire compatibility claim
364 // this module makes.
365 for (state, wire) in [
366 (ProcState::Pending, "\"pending\""),
367 (ProcState::Starting, "\"starting\""),
368 (ProcState::Running, "\"running\""),
369 (ProcState::Draining, "\"draining\""),
370 (ProcState::Exited, "\"exited\""),
371 (ProcState::Failed, "\"failed\""),
372 ] {
373 assert_eq!(serde_json::to_string(&state).unwrap(), wire);
374 assert_eq!(
375 serde_json::from_str::<ProcState>(wire).unwrap(),
376 state,
377 "{wire} must round-trip"
378 );
379 }
380 }
381
382 #[test]
383 fn state_is_the_only_required_field() {
384 let s: ProcStatus = serde_json::from_str(r#"{"state":"running"}"#).unwrap();
385 assert!(s.is_ready());
386 assert_eq!(s.pid, None);
387 assert!(s.endpoints.is_empty());
388 assert!(s.metrics.is_empty());
389 }
390
391 #[test]
392 fn a_full_document_parses() {
393 let s: ProcStatus = serde_json::from_str(
394 r#"{"state":"starting","ready":false,"pid":71455,"uptime_secs":41,
395 "version":"0.8.20","detail":"replaying WAL 3/7",
396 "endpoints":{"gui":"winit://main"},"metrics":{"fps":59.9}}"#,
397 )
398 .unwrap();
399 assert_eq!(s.state, ProcState::Starting);
400 assert!(!s.is_ready(), "starting is not ready");
401 assert_eq!(s.pid, Some(71455));
402 assert_eq!(s.detail.as_deref(), Some("replaying WAL 3/7"));
403 assert_eq!(s.endpoints.get("gui").map(String::as_str), Some("winit://main"));
404 assert_eq!(s.metrics.get("fps"), Some(&59.9));
405 assert_eq!(s.summary(), "starting — replaying WAL 3/7 (v0.8.20)");
406 }
407
408 /// A producer that contradicts itself must not be believed on the
409 /// optimistic half — that is precisely how a half-booted process gets
410 /// reported as up, which is the failure this channel exists to end.
411 #[test]
412 fn a_ready_flag_never_overrides_a_not_running_state() {
413 let s: ProcStatus =
414 serde_json::from_str(r#"{"state":"starting","ready":true}"#).unwrap();
415 assert_eq!(s.ready, Some(true), "the claim is preserved verbatim");
416 assert!(!s.is_ready(), "but state decides");
417 }
418
419 #[tokio::test]
420 async fn fetch_status_reads_a_document_over_a_unix_socket() {
421 let tmp = tempfile::tempdir().unwrap();
422 let sock = tmp.path().join("control.sock");
423 let server = serve_once(
424 sock.clone(),
425 vec![r#"{"state":"running","detail":"3 windows"}"#.to_string()],
426 );
427 // Give the listener a moment to bind.
428 tokio::time::sleep(Duration::from_millis(50)).await;
429
430 let got = fetch_status(&ControlEndpoint::Socket(sock)).await.unwrap();
431 assert_eq!(got.state, ProcState::Running);
432 assert_eq!(got.detail.as_deref(), Some("3 windows"));
433 server.abort();
434 }
435
436 #[tokio::test]
437 async fn wait_ready_polls_through_starting_to_running() {
438 let tmp = tempfile::tempdir().unwrap();
439 let sock = tmp.path().join("control.sock");
440 let server = serve_once(
441 sock.clone(),
442 vec![
443 r#"{"state":"starting"}"#.to_string(),
444 r#"{"state":"starting"}"#.to_string(),
445 r#"{"state":"running"}"#.to_string(),
446 ],
447 );
448 tokio::time::sleep(Duration::from_millis(50)).await;
449
450 let outcome = wait_ready(&ControlEndpoint::Socket(sock), Duration::from_secs(5)).await;
451 assert!(
452 matches!(&outcome, ReadyOutcome::Ready(s) if s.state == ProcState::Running),
453 "{outcome:?}"
454 );
455 server.abort();
456 }
457
458 /// `failed` must short-circuit. Burning the full readiness timeout on a
459 /// process that has already said it is not coming up is the slow-edit-loop
460 /// failure this outcome exists to prevent.
461 #[tokio::test]
462 async fn wait_ready_fails_fast_on_a_terminal_state() {
463 let tmp = tempfile::tempdir().unwrap();
464 let sock = tmp.path().join("control.sock");
465 let server = serve_once(
466 sock.clone(),
467 vec![r#"{"state":"failed","detail":"no GPU"}"#.to_string()],
468 );
469 tokio::time::sleep(Duration::from_millis(50)).await;
470
471 let started = std::time::Instant::now();
472 let outcome = wait_ready(&ControlEndpoint::Socket(sock), Duration::from_secs(30)).await;
473 assert!(
474 matches!(&outcome, ReadyOutcome::Terminal(s) if s.detail.as_deref() == Some("no GPU")),
475 "{outcome:?}"
476 );
477 assert!(
478 started.elapsed() < Duration::from_secs(5),
479 "must not burn the 30s timeout on a terminal state"
480 );
481 server.abort();
482 }
483
484 #[tokio::test]
485 async fn wait_ready_times_out_when_nothing_is_listening() {
486 let tmp = tempfile::tempdir().unwrap();
487 let outcome = wait_ready(
488 &ControlEndpoint::Socket(tmp.path().join("never-bound.sock")),
489 Duration::from_millis(300),
490 )
491 .await;
492 assert!(
493 matches!(outcome, ReadyOutcome::TimedOut { last: None }),
494 "an endpoint that never answered must report no last document"
495 );
496 }
497
498 #[test]
499 fn endpoints_render_distinguishably() {
500 assert_eq!(
501 ControlEndpoint::Socket(PathBuf::from("/tmp/c.sock")).to_string(),
502 "unix:/tmp/c.sock"
503 );
504 assert_eq!(
505 ControlEndpoint::Http("http://127.0.0.1:4325/_yah/status".into()).to_string(),
506 "http://127.0.0.1:4325/_yah/status"
507 );
508 }
509
510 #[test]
511 fn a_status_document_round_trips_through_serialization() {
512 let mut endpoints = BTreeMap::new();
513 endpoints.insert("http".to_string(), "http://127.0.0.1:4325".to_string());
514 let original = ProcStatus {
515 state: ProcState::Running,
516 ready: None,
517 pid: Some(9),
518 uptime_secs: Some(3),
519 version: None,
520 detail: None,
521 endpoints,
522 metrics: BTreeMap::new(),
523 };
524 let json = serde_json::to_string(&original).unwrap();
525 assert_eq!(serde_json::from_str::<ProcStatus>(&json).unwrap(), original);
526 assert!(
527 !json.contains("\"ready\""),
528 "absent optionals must not be emitted: {json}"
529 );
530 }
531}