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///
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}