diff --git a/crates/subc-control/src/lib.rs b/crates/subc-control/src/lib.rs index 5ce38b2e..ff80f6bf 100644 --- a/crates/subc-control/src/lib.rs +++ b/crates/subc-control/src/lib.rs @@ -38,6 +38,7 @@ pub mod ops { pub const SUPERVISOR_SET_ENABLED: &str = "supervisor.set_enabled"; pub const SUPERVISOR_HEALTH_PROBE: &str = "supervisor.health_probe"; pub const SUPERVISOR_HEALTH: &str = "supervisor.health"; + pub const SUPERVISOR_STDERR_TAIL: &str = "supervisor.stderr_tail"; } /// Client-originated channel-0 control RPC body. @@ -127,6 +128,21 @@ pub enum ClientControlRequest { SupervisorHealthProbe { module_id: String }, #[serde(rename = "supervisor.health")] SupervisorHealth {}, + /// Retained stderr for one module. + /// + /// A separate op rather than a field on `supervisor.list`: the tail is + /// kilobytes per module and `list` renders every module, so carrying it in + /// the snapshot would charge every status read for a payload almost no + /// caller wants. Caps ride on the REQUEST so a caller wanting twenty lines + /// and one wanting the whole ring need no separate fields anywhere. + #[serde(rename = "supervisor.stderr_tail")] + SupervisorStderrTail { + module_id: String, + #[serde(default, skip_serializing_if = "Option::is_none")] + max_lines: Option, + #[serde(default, skip_serializing_if = "Option::is_none")] + max_bytes: Option, + }, } /// subc's channel-0 response body for client control RPCs. @@ -186,6 +202,70 @@ pub enum ClientControlResponse { generation: u64, modules: Vec, }, + #[serde(rename = "supervisor.stderr_tail")] + SupervisorStderrTail { + module_id: String, + #[serde(flatten)] + tail: StderrTail, + }, +} + +/// A module's retained stderr, oldest entry first. +#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)] +pub struct StderrTail { + pub capture: StderrCaptureState, + pub entries: Vec, + /// Lines not present above: evicted by the ring, or held back by this + /// request's own caps. + /// + /// Non-zero means the first entry is not the first line the module wrote. A + /// reader hunting a cause needs that, or an absent explanation reads as a + /// module that never gave one. + /// + /// Zero is skipped so the common complete-tail case stays compact. + #[serde(default, skip_serializing_if = "is_zero_u64")] + pub dropped_lines: u64, +} + +/// Whether stderr is being captured for a module, and if not, why not. +/// +/// A typed state rather than an empty-tail convention. "The module printed +/// nothing before dying" and "nobody was capturing" send an operator in opposite +/// directions, and rendering them alike is the defect this op exists to fix -- +/// the same shape as a `detail -` that means both no-detail and never-probed. +#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)] +#[serde(tag = "state", rename_all = "snake_case")] +pub enum StderrCaptureState { + /// A reader is attached, or was attached and saw clean EOF. An empty + /// `entries` under this state means the module genuinely wrote nothing. + Captured, + /// Retained entries are valid, but the stderr reader ended before clean EOF. + Incomplete { reason: String }, + /// No reader was attached. `entries` says nothing about what the module wrote. + NotCaptured { reason: String }, +} + +#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)] +#[serde(tag = "kind", rename_all = "snake_case")] +pub enum StderrTailEntry { + Line { + text: String, + /// The line was cut at the per-line cap and `text` is a prefix. + /// + /// Carried as a field rather than left to a marker in `text` so a + /// consumer can branch on it without string matching. + #[serde(default, skip_serializing_if = "std::ops::Not::not")] + truncated: bool, + }, + /// The supervisor spawned a new process. Entries after this came from it. + /// + /// In-band because position is the information: which side of the restart a + /// line falls on is unanswerable from a count. + ProcessStart, +} + +fn is_zero_u64(value: &u64) -> bool { + *value == 0 } #[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq)] diff --git a/crates/subc-control/tests/golden/client_control_response_catalog_list.json b/crates/subc-control/tests/golden/client_control_response_catalog_list.json index 79210fa0..0de8f00c 100644 --- a/crates/subc-control/tests/golden/client_control_response_catalog_list.json +++ b/crates/subc-control/tests/golden/client_control_response_catalog_list.json @@ -89,6 +89,7 @@ "supervisor.rescan", "supervisor.set_enabled", "supervisor.health_probe", - "supervisor.health" + "supervisor.health", + "supervisor.stderr_tail" ] } diff --git a/crates/subc-control/tests/golden/client_control_response_server_describe.json b/crates/subc-control/tests/golden/client_control_response_server_describe.json index 6cb3a300..ae4aaab2 100644 --- a/crates/subc-control/tests/golden/client_control_response_server_describe.json +++ b/crates/subc-control/tests/golden/client_control_response_server_describe.json @@ -16,6 +16,7 @@ "supervisor.rescan", "supervisor.set_enabled", "supervisor.health_probe", - "supervisor.health" + "supervisor.health", + "supervisor.stderr_tail" ] } diff --git a/crates/subc-control/tests/golden/client_control_response_server_describe_with_counters.json b/crates/subc-control/tests/golden/client_control_response_server_describe_with_counters.json index b5712716..750b87e5 100644 --- a/crates/subc-control/tests/golden/client_control_response_server_describe_with_counters.json +++ b/crates/subc-control/tests/golden/client_control_response_server_describe_with_counters.json @@ -25,6 +25,7 @@ "supervisor.rescan", "supervisor.set_enabled", "supervisor.health_probe", - "supervisor.health" + "supervisor.health", + "supervisor.stderr_tail" ] } diff --git a/crates/subc-control/tests/golden/client_control_response_supervisor_stderr_tail.json b/crates/subc-control/tests/golden/client_control_response_supervisor_stderr_tail.json new file mode 100644 index 00000000..72b0ab1c --- /dev/null +++ b/crates/subc-control/tests/golden/client_control_response_supervisor_stderr_tail.json @@ -0,0 +1,22 @@ +{ + "capture": { + "state": "captured" + }, + "dropped_lines": 12, + "entries": [ + { + "kind": "line", + "text": "config error: missing top-level `storage`" + }, + { + "kind": "process_start" + }, + { + "kind": "line", + "text": "config error: missing top-level `stor", + "truncated": true + } + ], + "module_id": "aft-tools", + "op": "supervisor.stderr_tail" +} diff --git a/crates/subc-control/tests/golden/client_control_response_supervisor_stderr_tail_incomplete.json b/crates/subc-control/tests/golden/client_control_response_supervisor_stderr_tail_incomplete.json new file mode 100644 index 00000000..aefe5792 --- /dev/null +++ b/crates/subc-control/tests/golden/client_control_response_supervisor_stderr_tail_incomplete.json @@ -0,0 +1,14 @@ +{ + "capture": { + "reason": "stderr read failed: reader failed", + "state": "incomplete" + }, + "entries": [ + { + "kind": "line", + "text": "config error: missing top-level `storage`" + } + ], + "module_id": "aft-tools", + "op": "supervisor.stderr_tail" +} diff --git a/crates/subc-control/tests/golden/client_control_response_supervisor_stderr_tail_not_captured.json b/crates/subc-control/tests/golden/client_control_response_supervisor_stderr_tail_not_captured.json new file mode 100644 index 00000000..ea152dd8 --- /dev/null +++ b/crates/subc-control/tests/golden/client_control_response_supervisor_stderr_tail_not_captured.json @@ -0,0 +1,9 @@ +{ + "capture": { + "reason": "stderr pipe was not available on spawn", + "state": "not_captured" + }, + "entries": [], + "module_id": "aft-tools", + "op": "supervisor.stderr_tail" +} diff --git a/crates/subc-control/tests/golden_json.rs b/crates/subc-control/tests/golden_json.rs index 362ed2a5..19b1b667 100644 --- a/crates/subc-control/tests/golden_json.rs +++ b/crates/subc-control/tests/golden_json.rs @@ -4,7 +4,8 @@ use serde::{de::DeserializeOwned, Serialize}; use serde_json::Value; use subc_control::{ CatalogEntry, ClientControlRequest, ClientControlResponse, ConsumerIdentity, PollKind, - SupervisorEntry, SupervisorHealthEntry, SupervisorHealthStatus, SupervisorRescanResult, + StderrCaptureState, StderrTail, StderrTailEntry, SupervisorEntry, SupervisorHealthEntry, + SupervisorHealthStatus, SupervisorRescanResult, }; use subc_protocol::{ manifest::{ @@ -264,6 +265,59 @@ fn client_control_responses() -> Vec<(&'static str, ClientControlResponse)> { modules: vec![supervisor_health_entry()], }, ), + ( + "client_control_response_supervisor_stderr_tail", + ClientControlResponse::SupervisorStderrTail { + module_id: "aft-tools".to_string(), + tail: StderrTail { + capture: StderrCaptureState::Captured, + entries: vec![ + StderrTailEntry::Line { + text: "config error: missing top-level `storage`".to_string(), + truncated: false, + }, + StderrTailEntry::ProcessStart, + StderrTailEntry::Line { + text: "config error: missing top-level `stor".to_string(), + truncated: true, + }, + ], + dropped_lines: 12, + }, + }, + ), + ( + // Pinned separately because it is the state the empty-tail convention + // could not express, and a fixture is the only thing that keeps the + // distinction from being collapsed back into an empty list later. + "client_control_response_supervisor_stderr_tail_not_captured", + ClientControlResponse::SupervisorStderrTail { + module_id: "aft-tools".to_string(), + tail: StderrTail { + capture: StderrCaptureState::NotCaptured { + reason: "stderr pipe was not available on spawn".to_string(), + }, + entries: Vec::new(), + dropped_lines: 0, + }, + }, + ), + ( + "client_control_response_supervisor_stderr_tail_incomplete", + ClientControlResponse::SupervisorStderrTail { + module_id: "aft-tools".to_string(), + tail: StderrTail { + capture: StderrCaptureState::Incomplete { + reason: "stderr read failed: reader failed".to_string(), + }, + entries: vec![StderrTailEntry::Line { + text: "config error: missing top-level `storage`".to_string(), + truncated: false, + }], + dropped_lines: 0, + }, + }, + ), ] } @@ -280,6 +334,7 @@ fn thin_core_ops() -> Vec { "supervisor.set_enabled".to_string(), "supervisor.health_probe".to_string(), "supervisor.health".to_string(), + "supervisor.stderr_tail".to_string(), ] } diff --git a/crates/subc-core/src/bin/ck.rs b/crates/subc-core/src/bin/ck.rs index e4b27e3d..febd607c 100644 --- a/crates/subc-core/src/bin/ck.rs +++ b/crates/subc-core/src/bin/ck.rs @@ -41,7 +41,7 @@ const PROD_CONNECTION_RELATIVE_PATH: &[&str] = const QUOTA_MODULE_ID: &str = "insula"; const CK_HARNESS: &str = "ck"; -const TOP_HELP_BASE: &str = "ck — CortexKit operator CLI\n\nusage:\n ck [--subc ] [--json] [] []\n\ndomains:\n module supervised modules: list, status, restart, stop, start, rescan\n health one-line health for every supervised module\n quota AI-provider quota and usage windows\n daemon daemon version, uptime, and connection info"; +const TOP_HELP_BASE: &str = "ck — CortexKit operator CLI\n\nusage:\n ck [--subc ] [--json] [] []\n\ndomains:\n module supervised modules: list, status, stderr, restart, stop, start, rescan\n health one-line health for every supervised module\n quota AI-provider quota and usage windows\n daemon daemon version, uptime, and connection info"; const TOP_HELP_TAIL: &str = "flags:\n --subc use a specific connection file (default: auto-discover)\n --json raw JSON output instead of tables\n\nrun 'ck ' with no verb to see that domain's commands"; @@ -107,7 +107,8 @@ fn discover_external_domains() -> Vec { domains } -const MODULE_HELP: &str = "ck module — inspect and control supervised modules\n\nusage: ck [--json] module []\n\nverbs:\n ck module list all modules with state and health\n ck module status one module in detail\n ck module restart drain-restart a module\n ck module stop disable and stop a module (persists until start)\n ck module start enable and spawn a module\n ck module rescan re-read subc.jsonc and reconcile the module set\n ck module rescan --dry-run show what a rescan would change, without changing it"; +const MODULE_HELP: &str = "ck module — inspect and control supervised modules\n\nusage: ck [--json] module []\n\nverbs:\n ck module list all modules with state and health\n ck module status one module in detail + ck module stderr retained stderr for a module (-n to limit)\n ck module restart drain-restart a module\n ck module stop disable and stop a module (persists until start)\n ck module start enable and spawn a module\n ck module rescan re-read subc.jsonc and reconcile the module set\n ck module rescan --dry-run show what a rescan would change, without changing it"; const QUOTA_HELP: &str = "ck quota - AI-provider quota and usage windows\n\nusage: ck [--json] quota [--verbose] []\n\n ck quota connected providers and their usage windows\n ck quota --verbose all tracked providers, including unavailable ones\n ck quota claude one provider's windows and status in detail"; @@ -148,6 +149,10 @@ async fn run(argv: impl IntoIterator) -> Result<(), CkError> { Command::Module(ModuleCommand::Status { module_id }) => { module_status(&mut client, &module_id, args.json).await } + Command::Module(ModuleCommand::StderrTail { + module_id, + max_lines, + }) => module_stderr_tail(&mut client, &module_id, max_lines, args.json).await, Command::Module(ModuleCommand::Restart { module_id }) => { module_restart(&mut client, &module_id, args.json).await } @@ -219,11 +224,25 @@ enum Command { enum ModuleCommand { List, - Status { module_id: String }, - Restart { module_id: String }, - Rescan { preview: bool }, - Stop { module_id: String }, - Start { module_id: String }, + Status { + module_id: String, + }, + Restart { + module_id: String, + }, + Rescan { + preview: bool, + }, + Stop { + module_id: String, + }, + Start { + module_id: String, + }, + StderrTail { + module_id: String, + max_lines: Option, + }, } struct ResolvedConnection { @@ -469,6 +488,118 @@ async fn module_status( Ok(()) } +/// Read `-n ` from a verb's own tail. +/// +/// Scoped to the verb rather than the global argument set, like `--dry-run` on +/// rescan, so it cannot silently apply somewhere else. +fn parse_tail_count(tail: &[std::ffi::OsString]) -> Result, CkError> { + let Some(position) = tail.iter().position(|arg| arg == "-n") else { + return Ok(None); + }; + let raw = tail + .get(position + 1) + .map(|arg| arg.to_string_lossy().into_owned()) + .ok_or_else(|| { + CkError::Usage(format!( + "ck module stderr -n needs a count\n\n{MODULE_HELP}" + )) + })?; + raw.parse::() + .map(Some) + .map_err(|_| CkError::Usage(format!("ck module stderr -n needs a number, got '{raw}'"))) +} + +async fn module_stderr_tail( + client: &mut CkClient, + module_id: &str, + max_lines: Option, + json_output: bool, +) -> Result<(), CkError> { + let response = client + .rpc_value(ClientControlRequest::SupervisorStderrTail { + module_id: module_id.to_string(), + max_lines, + max_bytes: None, + }) + .await?; + + if json_output { + print_json(&response)?; + return Ok(()); + } + + // An uncaptured tail is reported instead of the lines, never alongside them: + // the entries under that state carry no information about what the module + // wrote, and printing them under a warning invites reading them as complete. + let capture = response.get("capture"); + if capture + .and_then(|capture| capture.get("state")) + .and_then(Value::as_str) + .is_some_and(|state| state == "not_captured") + { + let reason = capture + .and_then(|capture| capture.get("reason")) + .and_then(Value::as_str) + .unwrap_or("unknown"); + println!("stderr not captured for {module_id}: {reason}"); + return Ok(()); + } + let incomplete_reason = capture + .filter(|capture| { + capture + .get("state") + .and_then(Value::as_str) + .is_some_and(|state| state == "incomplete") + }) + .and_then(|capture| capture.get("reason")) + .and_then(Value::as_str); + + let dropped = response + .get("dropped_lines") + .and_then(Value::as_u64) + .unwrap_or(0); + if dropped > 0 { + // Printed BEFORE the lines: a reader scanning for a cause needs to know + // the first line shown is not the first line written, and a footer after + // a long tail is read too late to change how the tail is read. + println!("... {dropped} earlier line(s) dropped"); + } + + let entries = response + .get("entries") + .and_then(Value::as_array) + .cloned() + .unwrap_or_default(); + if entries.is_empty() { + println!("(no stderr output captured)"); + } else { + for entry in entries { + match entry.get("kind").and_then(Value::as_str) { + Some("process_start") => println!("--- process start ---"), + _ => { + let text = entry + .get("text") + .and_then(Value::as_str) + .unwrap_or_default(); + let truncated = entry + .get("truncated") + .and_then(Value::as_bool) + .unwrap_or(false); + if truncated { + println!("{text} [truncated]"); + } else { + println!("{text}"); + } + } + } + } + } + if let Some(reason) = incomplete_reason { + println!("stderr capture incomplete for {module_id}: {reason}"); + } + Ok(()) +} + async fn module_restart( client: &mut CkClient, module_id: &str, @@ -2416,6 +2547,13 @@ fn parse_command(domain: &str, tail: &[OsString]) -> Result { preview: tail.iter().any(|t| t == "--dry-run"), }, "status" => ModuleCommand::Status { module_id: id(1)? }, + // `-n ` narrows the tail daemon-side rather than here, so + // a caller asking for 20 lines is not shipped the whole ring to + // discard most of it. + "stderr" => ModuleCommand::StderrTail { + module_id: id(1)?, + max_lines: parse_tail_count(tail)?, + }, "restart" => ModuleCommand::Restart { module_id: id(1)? }, "stop" => ModuleCommand::Stop { module_id: id(1)? }, "start" => ModuleCommand::Start { module_id: id(1)? }, diff --git a/crates/subc-core/src/bin/fake-aft-stub.rs b/crates/subc-core/src/bin/fake-aft-stub.rs index 3cac1bbe..fcdb0030 100644 --- a/crates/subc-core/src/bin/fake-aft-stub.rs +++ b/crates/subc-core/src/bin/fake-aft-stub.rs @@ -9,6 +9,7 @@ use std::{ io::{self, Write as _}, net::{IpAddr, SocketAddr}, path::{Path, PathBuf}, + process::{Command, Stdio}, sync::{Arc, Mutex}, time::Duration, }; @@ -71,6 +72,29 @@ const FAKE_AFT_HEALTH_NEVER_REPLY_FIRST_PATH_ENV: &str = "FAKE_AFT_HEALTH_NEVER_ const FAKE_AFT_HEALTH_STATUS_ENV: &str = "FAKE_AFT_HEALTH_STATUS"; const FAKE_AFT_HEALTH_DETAIL_ENV: &str = "FAKE_AFT_HEALTH_DETAIL"; const FAKE_AFT_HEALTH_METRICS_ENV: &str = "FAKE_AFT_HEALTH_METRICS"; +/// Presence (not value) is the trigger: when set, the stub writes +/// `FAKE_AFT_STDERR_LINE` (if any) to stderr and exits with this code before +/// ever touching the connection file. Lets a portable spawn stand in for a +/// freestanding `/bin/sh -c '...; exit N'` script, which never dialled subc +/// either -- a normal stub run would connect, HELLO, and register, and any +/// noise from that path would land in the very stderr ring these tests +/// assert on. +const FAKE_AFT_EXIT_CODE_ENV: &str = "FAKE_AFT_EXIT_CODE"; +/// Text written to stderr before the `FAKE_AFT_EXIT_CODE` exit. Supports a +/// `{pid}` token, substituted with this process's pid, for tests that must +/// distinguish which generation across a restart produced a line. +const FAKE_AFT_STDERR_LINE_ENV: &str = "FAKE_AFT_STDERR_LINE"; +/// Milliseconds the detached orphan writer (below) sleeps before writing +/// `FAKE_AFT_ORPHAN_WRITER_LINE` to stderr. Set alongside `FAKE_AFT_EXIT_CODE` +/// to reproduce a wedged pump: a child that inherits this process's stderr +/// pipe and keeps its write end open after the parent has already exited. +const FAKE_AFT_ORPHAN_WRITER_DELAY_MS_ENV: &str = "FAKE_AFT_ORPHAN_WRITER_DELAY_MS"; +/// Text the detached orphan writer emits after its delay. +const FAKE_AFT_ORPHAN_WRITER_LINE_ENV: &str = "FAKE_AFT_ORPHAN_WRITER_LINE"; +/// Internal marker set only on the re-exec'd orphan child, never by a test. +/// Distinguishes "I am the orphan, sleep then write" from "spawn an orphan" +/// on the same binary. +const FAKE_AFT_ORPHAN_WRITER_MODE_ENV: &str = "FAKE_AFT_ORPHAN_WRITER_MODE"; /// Id used when `FAKE_AFT_MODULE_ID` is absent. /// /// TESTS THAT ASSERT A MODULE APPEARS IN THE CATALOG MUST CONFIGURE AN ID THAT @@ -94,10 +118,61 @@ type InFlightRegistry = Arc>>>; #[tokio::main] async fn main() -> Result<(), StubError> { + // Checked before StubConfig::from_env(), which requires a `--subc + // ` argument neither of these paths receives: the + // orphan re-exec is spawned by `spawn_orphan_writer` with no args at all, + // and a caller using FAKE_AFT_EXIT_CODE as a portable stand-in for + // `/bin/sh -c '...; exit N'` may spawn this binary the same way a raw + // shell script would -- with no `--subc` argument, because a shell script + // never dialled subc either. Checking first also means a connect/HELLO + // failure can never land its own noise in the very stderr ring this knob + // is configured to control. + if env_flag(FAKE_AFT_ORPHAN_WRITER_MODE_ENV) { + return run_detached_orphan_writer().await; + } + if let Some(exit_code) = exit_code_from_env()? { + run_exit_only(exit_code).await?; + unreachable!("run_exit_only always exits the process"); + } + let config = StubConfig::from_env()?; run(config).await } +/// Writes the configured stderr line (if any), optionally spawns the orphan +/// writer, then exits with `exit_code`. Never touches `--subc`, the +/// connection file, or the network. +async fn run_exit_only(exit_code: i32) -> Result<(), StubError> { + if let Ok(line) = env::var(FAKE_AFT_STDERR_LINE_ENV) { + if !line.is_empty() { + eprintln!("{}", substitute_pid_token(&line)); + } + } + if let Some(delay) = env::var(FAKE_AFT_ORPHAN_WRITER_DELAY_MS_ENV) + .ok() + .map(|raw| { + raw.parse::() + .map(Duration::from_millis) + .map_err(|source| StubError::InvalidOrphanWriterDelay { raw, source }) + }) + .transpose()? + { + let line = env::var(FAKE_AFT_ORPHAN_WRITER_LINE_ENV).ok(); + spawn_orphan_writer(delay, line)?; + } + std::process::exit(exit_code); +} + +fn exit_code_from_env() -> Result, StubError> { + env::var(FAKE_AFT_EXIT_CODE_ENV) + .ok() + .map(|raw| { + raw.parse::() + .map_err(|source| StubError::InvalidExitCode { raw, source }) + }) + .transpose() +} + async fn run(config: StubConfig) -> Result<(), StubError> { if config.fail_registration { std::process::exit(2); @@ -120,6 +195,57 @@ async fn run(config: StubConfig) -> Result<(), StubError> { } } +/// Entry point for the re-exec'd orphan: sleep, then write one line to +/// (inherited) stderr and exit. Never dials subc, never reads `--subc`. +async fn run_detached_orphan_writer() -> Result<(), StubError> { + let delay = env::var(FAKE_AFT_ORPHAN_WRITER_DELAY_MS_ENV) + .ok() + .map(|raw| { + raw.parse::() + .map(Duration::from_millis) + .map_err(|source| StubError::InvalidOrphanWriterDelay { raw, source }) + }) + .transpose()? + .unwrap_or(Duration::ZERO); + sleep(delay).await; + if let Ok(line) = env::var(FAKE_AFT_ORPHAN_WRITER_LINE_ENV) { + eprintln!("{line}"); + } + Ok(()) +} + +/// Spawns a copy of this binary in orphan-writer mode with `Stdio::inherit()` +/// on stderr, so the child holds a duplicate of THIS process's stderr write +/// end -- the same handle the supervisor's stderr pump is reading from the +/// other side of. The child is spawned and dropped without a wait: dropping a +/// `std::process::Child` does not kill it, so it keeps that handle open past +/// this process's own exit, reproducing a wedged pump on both Unix and +/// Windows without a shell. +fn spawn_orphan_writer(delay: Duration, line: Option) -> Result<(), StubError> { + let exe = env::current_exe().map_err(StubError::Io)?; + let mut command = Command::new(exe); + command + .env(FAKE_AFT_ORPHAN_WRITER_MODE_ENV, "1") + .env( + FAKE_AFT_ORPHAN_WRITER_DELAY_MS_ENV, + delay.as_millis().to_string(), + ) + .stdin(Stdio::null()) + .stdout(Stdio::null()) + .stderr(Stdio::inherit()); + if let Some(line) = line { + command.env(FAKE_AFT_ORPHAN_WRITER_LINE_ENV, line); + } + command.spawn().map_err(StubError::Io)?; + Ok(()) +} + +/// Substitutes a `{pid}` token in a configured stderr line with this +/// process's real pid, so successive restart generations are distinguishable. +fn substitute_pid_token(line: &str) -> String { + line.replace("{pid}", &std::process::id().to_string()) +} + async fn connect_to_subc(connection_file_path: &Path) -> Result { // Any future reconnect loop must call this helper for every reconnect, so key // rotation is observed by re-reading the connection file each time. @@ -1517,6 +1643,14 @@ enum StubError { raw: String, source: std::num::ParseIntError, }, + InvalidExitCode { + raw: String, + source: std::num::ParseIntError, + }, + InvalidOrphanWriterDelay { + raw: String, + source: std::num::ParseIntError, + }, InvalidHealthStatus(String), InvalidConcurrency { raw: String, @@ -1578,6 +1712,14 @@ impl fmt::Display for StubError { f, "invalid {FAKE_AFT_TOOLCALL_DELAY_MS_ENV} value '{raw}': {source}" ), + Self::InvalidExitCode { raw, source } => write!( + f, + "invalid {FAKE_AFT_EXIT_CODE_ENV} value '{raw}': {source}" + ), + Self::InvalidOrphanWriterDelay { raw, source } => write!( + f, + "invalid {FAKE_AFT_ORPHAN_WRITER_DELAY_MS_ENV} value '{raw}': {source}" + ), Self::InvalidHealthStatus(raw) => write!( f, "invalid {FAKE_AFT_HEALTH_STATUS_ENV} value '{raw}': expected ok, degraded, or failing" @@ -1644,6 +1786,8 @@ impl Error for StubError { match self { Self::InvalidCrashAfter { source, .. } => Some(source), Self::InvalidToolcallDelay { source, .. } => Some(source), + Self::InvalidExitCode { source, .. } => Some(source), + Self::InvalidOrphanWriterDelay { source, .. } => Some(source), Self::ConnectionFile { source, .. } => Some(source), Self::Connect { source, .. } => Some(source), Self::Auth { source, .. } => Some(source), diff --git a/crates/subc-core/src/control.rs b/crates/subc-core/src/control.rs index 2bf95c5f..2f1cc39b 100644 --- a/crates/subc-core/src/control.rs +++ b/crates/subc-core/src/control.rs @@ -9,7 +9,8 @@ use std::{ use serde::{Deserialize, Serialize}; use subc_control::{ ops, CatalogEntry, ClientControlRequest, ClientControlResponse, ConsumerIdentity, PollKind, - SupervisorEntry, SupervisorHealthEntry, SupervisorRescanResult, + StderrCaptureState, StderrTail, StderrTailEntry, SupervisorEntry, SupervisorHealthEntry, + SupervisorRescanResult, }; use subc_protocol::{ manifest::{Concurrency, ModuleManifest, ProviderRole}, @@ -32,6 +33,7 @@ use crate::{ }, registry::{ChannelState, ConnectionId, Registry, RegistryError}, router::{RouteCtx, RouterError}, + stderr_tail::{CaptureState, TailEntry}, supervise::{validate_spec, ModuleProcessLiveness, ReservedHelloRejection, SupervisorHandle}, ConnectedClients, DaemonCounters, Frame, ProjectRootId, Supervisor, }; @@ -61,6 +63,7 @@ const SUBC_CONTROL_OPS: &[&str] = &[ ops::SUPERVISOR_SET_ENABLED, ops::SUPERVISOR_HEALTH_PROBE, ops::SUPERVISOR_HEALTH, + ops::SUPERVISOR_STDERR_TAIL, ]; const MODULE_TO_SUBC_CONTROL_OPS: &[&str] = &[MODULE_TO_SUBC_OP_CATALOG_UPDATE]; @@ -838,6 +841,11 @@ impl ControlHandler { self.handle_supervisor_health_probe(frame, module_id).await } ClientControlRequest::SupervisorHealth {} => self.handle_supervisor_health(frame), + ClientControlRequest::SupervisorStderrTail { + module_id, + max_lines, + max_bytes, + } => self.handle_supervisor_stderr_tail(frame, module_id, max_lines, max_bytes), } } @@ -1362,6 +1370,58 @@ impl ControlHandler { )?]) } + fn handle_supervisor_stderr_tail( + &self, + frame: Frame, + module_id: String, + max_lines: Option, + max_bytes: Option, + ) -> Result, RouterError> { + let Some(module) = self.supervisor.get(&module_id) else { + return Ok(vec![control_error_frame( + &frame, + "unknown_module", + format!("module_id '{module_id}' is not supervised"), + )?]); + }; + + let snapshot = module.stderr_tail( + max_lines.map(|value| value as usize), + max_bytes.map(|value| value as usize), + ); + + let response = ClientControlResponse::SupervisorStderrTail { + module_id, + tail: StderrTail { + capture: match snapshot.capture { + CaptureState::Captured => StderrCaptureState::Captured, + CaptureState::Incomplete { reason } => { + StderrCaptureState::Incomplete { reason } + } + CaptureState::NotCaptured { reason } => { + StderrCaptureState::NotCaptured { reason } + } + }, + entries: snapshot + .entries + .into_iter() + .map(|entry| match entry { + TailEntry::Line { text, truncated } => { + StderrTailEntry::Line { text, truncated } + } + TailEntry::ProcessStart => StderrTailEntry::ProcessStart, + }) + .collect(), + dropped_lines: snapshot.dropped_lines, + }, + }; + Ok(vec![control_response_body_frame( + &frame, + &response, + "ClientControlResponse::SupervisorStderrTail", + )?]) + } + fn handle_supervisor_health(&self, frame: Frame) -> Result, RouterError> { let generation = self .registry @@ -2751,7 +2811,7 @@ fn send_goodbye_target_best_effort(target: &GoodbyeTarget, context: &str) { #[cfg(test)] mod tests { - use std::{path::PathBuf, sync::Arc}; + use std::{path::PathBuf, sync::Arc, time::Duration}; use serde_json::{json, Value}; use subc_protocol::{ @@ -2765,8 +2825,38 @@ mod tests { }; use super::*; - use crate::{registry::ChannelState, router::FrameSink, RouteCtx, Router}; - use tokio::sync::mpsc; + use crate::{ + registry::ChannelState, + router::FrameSink, + stderr_tail::DEFAULT_MAX_LINE_BYTES, + supervise::{ModuleSpec, RestartPolicy, Supervisor}, + RouteCtx, Router, + }; + use tokio::{ + sync::mpsc, + time::{sleep, Instant}, + }; + + /// Locates the `fake-aft-stub` binary from a `src/lib.rs` unit test. + /// + /// `CARGO_BIN_EXE_*` (compile-time `env!` and runtime `std::env::var` alike) + /// is only populated for `tests/*.rs` integration test binaries -- this file + /// compiles as part of the library target, which gets neither. Cargo still + /// builds every `[[bin]]` target before running the library's tests, so the + /// binary is on disk; this test's own executable path is + /// `//deps/subc_core-`, and the sibling binary + /// lives two directories up at `//fake-aft-stub`. + fn fake_aft_stub_path() -> PathBuf { + let mut path = std::env::current_exe().expect("current_exe available in tests"); + path.pop(); // .../deps/ + path.pop(); // ...// + path.push(if cfg!(windows) { + "fake-aft-stub.exe" + } else { + "fake-aft-stub" + }); + path + } /// The retryable set both SDKs branch on, kept byte-identical to /// `is_retryable_route_open_code` in subc-client-rs and subc-client. @@ -3252,6 +3342,101 @@ mod tests { } } + #[tokio::test(flavor = "multi_thread", worker_threads = 2)] + async fn supervisor_stderr_tail_converts_a_real_truncated_ring_entry_to_prefix_only_wire_data() + { + let registry = Arc::new(Registry::default()); + let supervisor_handle = SupervisorHandle::new(); + let supervisor = Supervisor::new( + Arc::clone(®istry), + RestartPolicy::new(1, Duration::from_millis(10)), + ) + .with_handle(supervisor_handle.clone()); + let source_line = format!("config error: {}", "x".repeat(DEFAULT_MAX_LINE_BYTES)); + let module = supervisor + .spawn(ModuleSpec { + module_id: "stderr-tail-wire".to_string(), + program: fake_aft_stub_path(), + args: Vec::new(), + env: vec![ + ("FAKE_AFT_STDERR_LINE".to_string(), source_line.clone()), + ("FAKE_AFT_EXIT_CODE".to_string(), "1".to_string()), + ], + reserved: false, + reserved_prefixes: Vec::new(), + }) + .unwrap(); + + let deadline = Instant::now() + Duration::from_secs(5); + loop { + let tail = module.stderr_tail(None, None); + if tail + .entries + .iter() + .any(|entry| matches!(entry, TailEntry::ProcessStart)) + && tail.entries.iter().any(|entry| { + matches!( + entry, + TailEntry::Line { + truncated: true, + .. + } + ) + }) + { + break; + } + assert!( + Instant::now() < deadline, + "module did not produce a truncated line and restart boundary: {tail:?}" + ); + sleep(Duration::from_millis(10)).await; + } + + let handler = ControlHandler::new(Arc::clone(®istry)).with_supervisor(supervisor_handle); + let request = ClientControlRequest::SupervisorStderrTail { + module_id: "stderr-tail-wire".to_string(), + max_lines: None, + max_bytes: None, + }; + let frame = Frame::build( + FrameType::Request, + control_flags(), + 0, + 0, + 1, + serde_json::to_vec(&request).unwrap(), + ) + .unwrap(); + let (ctx, _egress) = route_ctx(ConnectionId::new(1)); + let responses = handler.handle_control_frame(&ctx, frame).await.unwrap(); + let ClientControlResponse::SupervisorStderrTail { tail, .. } = + serde_json::from_slice(&responses[0].body).unwrap() + else { + panic!("expected supervisor.stderr_tail response"); + }; + + assert!( + tail.entries + .iter() + .any(|entry| matches!(entry, StderrTailEntry::ProcessStart)), + "the control response lost the restart boundary" + ); + let Some(StderrTailEntry::Line { text, truncated }) = tail.entries.iter().find(|entry| { + matches!( + entry, + StderrTailEntry::Line { + truncated: true, + .. + } + ) + }) else { + panic!("the control response lost the truncated line"); + }; + assert_eq!(text, &source_line[..DEFAULT_MAX_LINE_BYTES]); + assert!(*truncated); + } + #[test] fn hello_registers_manifest_and_returns_ack() { let registry = Arc::new(Registry::default()); diff --git a/crates/subc-core/src/lib.rs b/crates/subc-core/src/lib.rs index 54e58795..3dad197c 100644 --- a/crates/subc-core/src/lib.rs +++ b/crates/subc-core/src/lib.rs @@ -17,6 +17,7 @@ pub mod observability; pub mod registry; pub mod router; pub mod server; +pub mod stderr_tail; pub mod supervise; pub mod watchdog; diff --git a/crates/subc-core/src/stderr_tail.rs b/crates/subc-core/src/stderr_tail.rs new file mode 100644 index 00000000..cc20a5ca --- /dev/null +++ b/crates/subc-core/src/stderr_tail.rs @@ -0,0 +1,923 @@ +//! Bounded per-module stderr capture. +//! +//! # Why this exists +//! +//! `last_exit_code` survives a respawn because the supervisor holds it in memory. +//! Stderr had no such path: it went to the daemon's inherited fd, and from there +//! to whatever rotates or evicts it. On the box this was written for, that window +//! was about three hours; on another it was bounded by a log file reaching 908 MB +//! with one module accounting for 98% of it. Two hosts, two mechanisms, the same +//! outcome -- the text explaining a crash is gone by the time anyone asks. +//! +//! So this keeps the last few lines where `last_exit` already lives: in supervisor +//! memory, immune to whatever happens to the log. +//! +//! # What it is not +//! +//! Not a log. The ring is deliberately small and lossy, and callers are expected +//! to know they are reading a tail rather than a history. The daemon log keeps +//! doing its job; this exists because that job has a time limit. + +use std::collections::VecDeque; +use std::io::Write; +use std::sync::{Arc, Mutex}; + +use tokio::io::AsyncReadExt; + +/// Longest single line admitted to the ring before truncation. +/// +/// A module emitting a 40 MB backtrace on one line satisfies any line-count cap +/// while evicting everything else -- the pathological emitter wins twice, once by +/// filling the ring and once by being unreadable itself. Truncating on the way in +/// costs that emitter one line instead of the whole tail. +pub const DEFAULT_MAX_LINE_BYTES: usize = 2048; + +/// Lines retained per module. +pub const DEFAULT_MAX_LINES: usize = 200; + +/// Total bytes retained per module, across all lines. +/// +/// Both this and [`DEFAULT_MAX_LINES`] apply; whichever binds first wins. A line +/// cap alone is satisfied by 200 lines of 2 KB, which is not a budget worth +/// holding for fifteen modules. +pub const DEFAULT_MAX_BYTES: usize = 64 * 1024; + +/// Whether a module's stderr is being captured, and if not, why not. +/// +/// Typed rather than nullable so `NotCaptured` has to be handled rather than +/// defaulted past. An empty tail and an uncaptured one send an operator in +/// opposite directions -- one says the module printed nothing before dying, the +/// other says nobody was listening -- and rendering them alike is the defect this +/// module exists to fix, reproduced one layer up. +#[derive(Debug, Clone, PartialEq, Eq)] +pub enum CaptureState { + /// A reader is attached, or was attached and reached clean EOF. + Captured, + /// Retained entries are valid, but capture ended before clean EOF. + Incomplete { reason: String }, + /// No reader was attached. The tail says nothing about what the module wrote. + NotCaptured { reason: String }, +} + +/// One retained entry. +/// +/// Boundaries are in-band rather than a separate field because their position +/// relative to the lines is the whole point: "these three lines came from the +/// process that died, those came from its replacement" is unanswerable from a +/// count. +#[derive(Debug, Clone, PartialEq, Eq)] +pub enum TailEntry { + Line { + text: String, + /// This line was cut at the per-line cap. + truncated: bool, + }, + /// The supervisor spawned a new process for this module. Lines after this + /// entry come from the new one. + ProcessStart, +} + +impl TailEntry { + fn cost(&self) -> usize { + match self { + Self::Line { text, .. } => text.len(), + Self::ProcessStart => 0, + } + } +} + +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub struct StderrTailConfig { + max_lines: usize, + max_bytes: usize, + max_line_bytes: usize, +} + +impl StderrTailConfig { + /// Keeps every retained line within the ring's total byte budget. + /// Clamp rather than reject so diagnostics degrade without blocking supervisor startup. + pub const fn new(max_lines: usize, max_bytes: usize, max_line_bytes: usize) -> Self { + Self { + max_lines, + max_bytes, + max_line_bytes: if max_line_bytes > max_bytes { + max_bytes + } else { + max_line_bytes + }, + } + } +} + +impl Default for StderrTailConfig { + fn default() -> Self { + Self::new(DEFAULT_MAX_LINES, DEFAULT_MAX_BYTES, DEFAULT_MAX_LINE_BYTES) + } +} + +/// A module's retained stderr, oldest first. +#[derive(Debug, Clone, PartialEq, Eq)] +pub struct StderrTailSnapshot { + pub capture: CaptureState, + pub entries: Vec, + /// Lines evicted since the module was first supervised. + /// + /// Non-zero means the tail starts mid-stream. That is the ring working as + /// intended, but a reader diagnosing a crash needs to know the first retained + /// line is not the first line the module wrote -- otherwise an absent cause + /// reads as a module that never explained itself. + pub dropped_lines: u64, +} + +impl StderrTailSnapshot { + /// The uncaptured case, for a module whose stderr was never piped. + pub fn not_captured(reason: impl Into) -> Self { + Self { + capture: CaptureState::NotCaptured { + reason: reason.into(), + }, + entries: Vec::new(), + dropped_lines: 0, + } + } +} + +/// Bounded ring of a single module's stderr lines. +/// +/// Survives respawn deliberately. The stderr explaining an exit is written +/// *before* that exit, so clearing on restart would discard the lines exactly +/// when they become the thing being asked for. [`TailEntry::ProcessStart`] keeps +/// the generations distinguishable instead. +#[derive(Debug)] +pub struct StderrRing { + config: StderrTailConfig, + entries: VecDeque, + bytes: usize, + dropped_lines: u64, + capture: CaptureState, +} + +impl StderrRing { + pub fn new(config: StderrTailConfig) -> Self { + Self { + config, + entries: VecDeque::new(), + bytes: 0, + dropped_lines: 0, + // Until a reader attaches, the honest answer is that nothing is + // listening -- not that the module has been quiet. + capture: CaptureState::NotCaptured { + reason: "stderr reader has not started".to_string(), + }, + } + } + + pub fn mark_captured(&mut self) { + if matches!(self.capture, CaptureState::NotCaptured { .. }) { + self.capture = CaptureState::Captured; + } + } + + pub fn mark_incomplete(&mut self, reason: impl Into) { + self.capture = CaptureState::Incomplete { + reason: reason.into(), + }; + } + + pub fn mark_not_captured(&mut self, reason: impl Into) { + self.capture = CaptureState::NotCaptured { + reason: reason.into(), + }; + } + + /// Record that a new process was spawned for this module. + /// + /// A boundary separates output on either side of it, so one with nothing + /// before it separates nothing: on the FIRST spawn it would make a module + /// that printed nothing render as a marker rather than as empty, and the + /// caller then has to decide whether a one-marker tail counts as silence. + /// Recording it only once there is something to divide keeps "captured and + /// empty" literally empty. + pub fn push_process_start(&mut self) { + if self.entries.is_empty() && self.dropped_lines == 0 { + return; + } + if matches!(self.entries.back(), Some(TailEntry::ProcessStart)) { + return; + } + self.push_entry(TailEntry::ProcessStart); + } + + /// Admit one complete line, truncating it if it exceeds the per-line cap. + /// + /// `line` must not contain a trailing newline; the reader strips it so the + /// stored text and the byte accounting agree. + pub fn push_line(&mut self, line: &str) { + let (text, truncated) = truncate_line(line, self.config.max_line_bytes); + self.push_entry(TailEntry::Line { text, truncated }); + } + + fn push_entry(&mut self, entry: TailEntry) { + self.bytes += entry.cost(); + self.entries.push_back(entry); + self.evict_to_fit(); + } + + fn evict_to_fit(&mut self) { + while self + .entries + .iter() + .filter(|entry| matches!(entry, TailEntry::Line { .. })) + .count() + > self.config.max_lines + || (self.bytes > self.config.max_bytes && self.entries.len() > 1) + { + let Some(evicted) = self.entries.pop_front() else { + break; + }; + self.bytes -= evicted.cost(); + if matches!(evicted, TailEntry::Line { .. }) { + self.dropped_lines += 1; + } + } + } + + /// The most recent entries, oldest first, bounded by the caller's limits. + /// + /// `max_lines`/`max_bytes` narrow the ring's own caps; they cannot widen them. + pub fn snapshot( + &self, + max_lines: Option, + max_bytes: Option, + ) -> StderrTailSnapshot { + let line_limit = max_lines.unwrap_or(self.config.max_lines); + let byte_limit = max_bytes.unwrap_or(self.config.max_bytes); + + let mut taken: Vec = Vec::new(); + let mut bytes = 0usize; + let mut lines = 0usize; + // Walk backwards: a tail is anchored at the newest end, so a caller + // asking for 20 lines wants the last 20, not the first 20. + for entry in self.entries.iter().rev() { + match entry { + TailEntry::Line { .. } => { + if lines >= line_limit { + break; + } + let cost = entry.cost(); + if lines > 0 && bytes + cost > byte_limit { + break; + } + bytes += cost; + lines += 1; + taken.push(entry.clone()); + } + TailEntry::ProcessStart if lines > 0 => taken.push(entry.clone()), + TailEntry::ProcessStart => {} + } + } + taken.reverse(); + + let withheld = self + .entries + .iter() + .filter(|entry| matches!(entry, TailEntry::Line { .. })) + .count() + .saturating_sub( + taken + .iter() + .filter(|entry| matches!(entry, TailEntry::Line { .. })) + .count(), + ); + + StderrTailSnapshot { + capture: self.capture.clone(), + entries: taken, + // Lines the ring evicted plus lines this request's own limits held + // back. Both mean the same thing to the reader -- the text above is + // not the beginning -- and separating them would invite treating a + // narrow request as evidence of a quiet module. + dropped_lines: self.dropped_lines + withheld as u64, + } + } +} + +/// Reassembly buffer ceiling for a line with no newline in sight. +/// +/// The ring truncates what it stores, but the READER has to hold the bytes until +/// it finds a delimiter. A module emitting a gigabyte with no newline would grow +/// this buffer without bound and take the daemon down with it -- a module fault +/// escalating into a fleet fault, which is exactly what supervision exists to +/// prevent. At this ceiling the pending bytes are flushed as a line and +/// reassembly restarts. +const MAX_PENDING_LINE_BYTES: usize = 1024 * 1024; + +/// Read a child's stderr to EOF: retain a bounded tail, and forward every line on. +/// +/// # Forwarding is mandatory, not a courtesy +/// +/// Measured on the live daemon log: 4727 of the last 5000 lines carried a module +/// tag. That file is a module log with some daemon lines in it, not the reverse. +/// A tap that captured without forwarding would leave it nearly empty, and every +/// existing reader -- including fleet scripts -- would report clean on nothing. +/// An absence that reads as calm is worse than the interleaving this replaces. +/// +/// # Bytes are forwarded verbatim +/// +/// The ring stores lossy UTF-8 because it renders into JSON; the forward writes +/// the ORIGINAL bytes. Anything else silently rewrites a log other tools parse. +/// +/// # One write per complete line +/// +/// Inheriting the daemon's fd gave line atomicity for free: a module's own write +/// reached the fd in one syscall. Reading a pipe and re-emitting can split a line +/// that used to be atomic, so this reassembles first and writes each complete +/// line in a single call -- otherwise the fix introduces a defect the previous +/// design did not have. +pub async fn pump_stderr(source: R, ring: Arc>) +where + R: AsyncReadExt + Unpin, +{ + pump_stderr_into(source, ring, &mut StderrSink).await +} + +/// Where forwarded lines go. Exists so tests can assert that forwarding HAPPENS +/// and that each line arrives in one write -- the property that makes this a +/// replacement for inherited stdio rather than a regression from it. +pub trait LineSink { + fn write_line(&mut self, line: &[u8]); +} + +struct StderrSink; + +impl LineSink for StderrSink { + fn write_line(&mut self, line: &[u8]) { + let stderr = std::io::stderr(); + let mut handle = stderr.lock(); + // Best-effort: a failed forward must not stop capture. Losing a log line + // is recoverable; losing the tail that explains a crash is the thing + // being fixed. + let _ = handle.write_all(line); + } +} + +async fn pump_stderr_into(mut source: R, ring: Arc>, sink: &mut S) +where + R: AsyncReadExt + Unpin, + S: LineSink, +{ + { + let mut guard = lock_ring(&ring); + guard.mark_captured(); + } + + let mut pending: Vec = Vec::new(); + let mut chunk = [0u8; 8192]; + + loop { + let read = match source.read(&mut chunk).await { + Ok(0) => break, + Ok(n) => n, + Err(err) => { + let mut guard = lock_ring(&ring); + // The tail up to this point stays valid and readable; what changes + // is that it is no longer complete, and saying so beats letting a + // truncated capture read as a module that stopped talking. + guard.mark_incomplete(format!("stderr read failed: {err}")); + return; + } + }; + + pending.extend_from_slice(&chunk[..read]); + + while let Some(newline) = pending.iter().position(|byte| *byte == b'\n') { + let line: Vec = pending.drain(..=newline).collect(); + emit_line(&ring, sink, &line[..line.len() - 1], true); + } + + if pending.len() >= MAX_PENDING_LINE_BYTES { + let line = std::mem::take(&mut pending); + emit_line(&ring, sink, &line, false); + } + } + + // A process that dies mid-line still wrote the bytes, and on a crash that + // fragment is disproportionately likely to be the message worth reading. + if !pending.is_empty() { + emit_line(&ring, sink, &pending, false); + } +} + +fn emit_line( + ring: &Arc>, + sink: &mut S, + raw: &[u8], + terminated: bool, +) { + { + let mut guard = lock_ring(ring); + guard.push_line(&String::from_utf8_lossy(raw)); + } + + // Framed and written in ONE call. Two writes -- body then newline -- would + // reintroduce exactly the interleaving that inheriting the fd avoided. + if terminated { + let mut framed = Vec::with_capacity(raw.len() + 1); + framed.extend_from_slice(raw); + framed.push(b'\n'); + sink.write_line(&framed); + } else { + sink.write_line(raw); + } +} + +fn lock_ring(ring: &Arc>) -> std::sync::MutexGuard<'_, StderrRing> { + ring.lock().unwrap_or_else(|poisoned| poisoned.into_inner()) +} + +/// Cut `line` to at most `max_bytes`, reporting whether it was shortened. +/// +/// Cuts on a char boundary: slicing a multi-byte sequence would produce invalid +/// UTF-8, and a panic while capturing a crash message is the worst possible time +/// to discover that. +fn truncate_line(line: &str, max_bytes: usize) -> (String, bool) { + if line.len() <= max_bytes { + return (line.to_string(), false); + } + let mut end = max_bytes; + while end > 0 && !line.is_char_boundary(end) { + end -= 1; + } + (line[..end].to_string(), true) +} + +#[cfg(test)] +mod tests { + use std::{ + io, + pin::Pin, + task::{Context, Poll}, + }; + + use super::*; + use tokio::io::{AsyncRead, ReadBuf}; + + fn ring(max_lines: usize, max_bytes: usize, max_line_bytes: usize) -> StderrRing { + StderrRing::new(StderrTailConfig::new(max_lines, max_bytes, max_line_bytes)) + } + + fn lines(snapshot: &StderrTailSnapshot) -> Vec { + snapshot + .entries + .iter() + .filter_map(|entry| match entry { + TailEntry::Line { text, .. } => Some(text.clone()), + TailEntry::ProcessStart => None, + }) + .collect() + } + + #[test] + fn a_fresh_ring_reports_not_captured_rather_than_empty() { + // The distinction this whole module exists for: "nobody was listening" + // must not render as "the module said nothing". + let ring = ring(10, 1024, 128); + let snapshot = ring.snapshot(None, None); + assert!(matches!(snapshot.capture, CaptureState::NotCaptured { .. })); + assert!(snapshot.entries.is_empty()); + } + + #[test] + fn a_captured_module_that_printed_nothing_is_distinguishable_from_an_uncaptured_one() { + let mut captured = ring(10, 1024, 128); + captured.mark_captured(); + let uncaptured = ring(10, 1024, 128); + + let captured = captured.snapshot(None, None); + let uncaptured = uncaptured.snapshot(None, None); + + // Both are empty. Only the capture state separates them, which is the + // point -- an assertion on emptiness alone would pass either way. + assert!(captured.entries.is_empty()); + assert!(uncaptured.entries.is_empty()); + assert_eq!(captured.capture, CaptureState::Captured); + assert!(matches!( + uncaptured.capture, + CaptureState::NotCaptured { .. } + )); + } + + #[test] + fn the_line_cap_evicts_oldest_first_and_counts_what_it_dropped() { + let mut ring = ring(3, 10_000, 128); + ring.mark_captured(); + for i in 0..6 { + ring.push_line(&format!("line{i}")); + } + let snapshot = ring.snapshot(None, None); + assert_eq!(lines(&snapshot), vec!["line3", "line4", "line5"]); + // Without this the tail silently becomes "the last lines that happened + // to survive" and reads as complete. + assert_eq!(snapshot.dropped_lines, 3); + } + + #[test] + fn the_byte_cap_binds_before_the_line_cap_when_lines_are_large() { + // 100 lines allowed, but only ~30 bytes of them. + let mut ring = ring(100, 30, 128); + ring.mark_captured(); + for i in 0..10 { + ring.push_line(&format!("{i}--------")); // 9 bytes each + } + let snapshot = ring.snapshot(None, None); + assert!( + snapshot.entries.len() < 10, + "byte cap did not bind: {} entries retained", + snapshot.entries.len() + ); + let retained: usize = lines(&snapshot).iter().map(String::len).sum(); + assert!( + retained <= 30, + "retained {retained} bytes over a 30 byte cap" + ); + assert!(snapshot.dropped_lines > 0); + } + + #[test] + fn one_enormous_line_is_truncated_rather_than_evicting_the_tail() { + // The pathological-emitter case: without per-line truncation this single + // line would evict every other line AND be unreadable itself. + let mut ring = ring(10, 10_000, 64); + ring.mark_captured(); + ring.push_line("context line that must survive"); + ring.push_line(&"x".repeat(40_000)); + + let snapshot = ring.snapshot(None, None); + let kept = &snapshot.entries; + assert!(matches!( + &kept[0], + TailEntry::Line { text, truncated: false } + if text == "context line that must survive" + )); + let TailEntry::Line { text, truncated } = &kept[1] else { + panic!("expected a truncated line"); + }; + assert_eq!(text, &"x".repeat(64)); + assert!(*truncated); + } + + #[test] + fn truncation_is_visible_so_a_cut_line_is_not_mistaken_for_a_short_one() { + let mut ring = ring(10, 10_000, 16); + ring.mark_captured(); + ring.push_line("0123456789abcdefghij"); + ring.push_line("short"); + + let snapshot = ring.snapshot(None, None); + let TailEntry::Line { truncated, .. } = &snapshot.entries[0] else { + panic!("expected a line"); + }; + assert!(truncated); + let TailEntry::Line { truncated, .. } = &snapshot.entries[1] else { + panic!("expected a line"); + }; + assert!(!truncated, "a short line must not be reported as truncated"); + } + + #[test] + fn truncation_cuts_on_a_char_boundary_rather_than_splitting_utf8() { + // A panic message with non-ASCII in it is not exotic, and slicing mid + // sequence would panic while capturing a crash. + let mut ring = ring(10, 10_000, 5); + ring.mark_captured(); + ring.push_line("aa€€€€"); + let snapshot = ring.snapshot(None, None); + let TailEntry::Line { text, truncated } = &snapshot.entries[0] else { + panic!("expected a line"); + }; + assert!(truncated); + assert!(text.starts_with("aa")); + } + + #[test] + fn a_restart_boundary_keeps_generations_distinguishable() { + let mut ring = ring(10, 10_000, 128); + ring.mark_captured(); + ring.push_line("before the crash"); + ring.push_process_start(); + ring.push_line("after the respawn"); + + let snapshot = ring.snapshot(None, None); + assert_eq!( + snapshot.entries, + vec![ + TailEntry::Line { + text: "before the crash".to_string(), + truncated: false + }, + TailEntry::ProcessStart, + TailEntry::Line { + text: "after the respawn".to_string(), + truncated: false + }, + ] + ); + } + + #[test] + fn the_ring_survives_respawn_because_the_cause_is_written_before_the_exit() { + // Clearing on restart would discard the lines at the exact moment they + // become the thing being asked for. + let mut ring = ring(10, 10_000, 128); + ring.mark_captured(); + ring.push_line("Error: storage section missing"); + ring.push_process_start(); + + let snapshot = ring.snapshot(None, None); + assert!(lines(&snapshot).contains(&"Error: storage section missing".to_string())); + } + + #[test] + fn a_caller_limit_returns_the_newest_lines_not_the_oldest() { + let mut ring = ring(100, 100_000, 128); + ring.mark_captured(); + for i in 0..10 { + ring.push_line(&format!("line{i}")); + } + let snapshot = ring.snapshot(Some(3), None); + assert_eq!(lines(&snapshot), vec!["line7", "line8", "line9"]); + } + + #[test] + fn a_caller_line_limit_keeps_the_boundary_before_the_selected_line() { + let mut ring = ring(100, 100_000, 128); + ring.mark_captured(); + ring.push_line("before restart"); + ring.push_process_start(); + ring.push_line("after restart"); + + let snapshot = ring.snapshot(Some(1), None); + assert_eq!( + snapshot.entries, + vec![ + TailEntry::ProcessStart, + TailEntry::Line { + text: "after restart".to_string(), + truncated: false, + }, + ] + ); + } + + #[test] + fn a_caller_line_limit_omits_a_trailing_boundary_after_the_selected_line() { + let mut ring = ring(100, 100_000, 128); + ring.mark_captured(); + ring.push_line("before restart"); + ring.push_process_start(); + + let snapshot = ring.snapshot(Some(1), None); + assert_eq!( + snapshot.entries, + vec![TailEntry::Line { + text: "before restart".to_string(), + truncated: false, + }] + ); + } + + #[test] + fn a_caller_limit_reports_what_it_withheld_rather_than_looking_complete() { + let mut ring = ring(100, 100_000, 128); + ring.mark_captured(); + for i in 0..10 { + ring.push_line(&format!("line{i}")); + } + // Nothing was evicted; the narrowing is the caller's own. It still has to + // be reported, or a 3-line request reads as a module that wrote 3 lines. + assert_eq!(ring.snapshot(Some(3), None).dropped_lines, 7); + assert_eq!(ring.snapshot(None, None).dropped_lines, 0); + } + + #[test] + fn a_caller_limit_cannot_widen_the_rings_own_caps() { + let mut ring = ring(2, 10_000, 128); + ring.mark_captured(); + for i in 0..5 { + ring.push_line(&format!("line{i}")); + } + let snapshot = ring.snapshot(Some(1000), Some(1_000_000)); + assert_eq!(lines(&snapshot).len(), 2); + } + + fn shared(max_lines: usize, max_bytes: usize, max_line_bytes: usize) -> Arc> { + Arc::new(Mutex::new(ring(max_lines, max_bytes, max_line_bytes))) + } + + /// Records each forwarded write separately, so a test can tell one write of + /// `b"abc\n"` from two writes of `b"abc"` and `b"\n"`. + #[derive(Default)] + struct RecordingSink { + writes: Vec>, + } + + impl LineSink for RecordingSink { + fn write_line(&mut self, line: &[u8]) { + self.writes.push(line.to_vec()); + } + } + + struct FailingReader { + bytes: Vec, + emitted: bool, + } + + impl AsyncRead for FailingReader { + fn poll_read( + mut self: Pin<&mut Self>, + _cx: &mut Context<'_>, + buf: &mut ReadBuf<'_>, + ) -> Poll> { + if self.emitted { + return Poll::Ready(Err(io::Error::other("reader failed"))); + } + self.emitted = true; + buf.put_slice(&self.bytes); + Poll::Ready(Ok(())) + } + } + + #[tokio::test] + async fn the_pump_splits_on_newlines_and_keeps_a_trailing_fragment() { + let ring = shared(10, 10_000, 128); + // No trailing newline on the last line: a crashing process routinely dies + // mid-line, and that fragment is often the message worth reading. + let source = std::io::Cursor::new(b"one\ntwo\nthree".to_vec()); + let mut sink = RecordingSink::default(); + pump_stderr_into(source, Arc::clone(&ring), &mut sink).await; + + let snapshot = lock_ring(&ring).snapshot(None, None); + assert_eq!(lines(&snapshot), vec!["one", "two", "three"]); + assert_eq!(snapshot.capture, CaptureState::Captured); + assert_eq!( + sink.writes, + vec![b"one\n".to_vec(), b"two\n".to_vec(), b"three".to_vec()] + ); + } + + #[tokio::test] + async fn a_read_failure_keeps_prior_lines_and_marks_the_capture_incomplete() { + let ring = shared(10, 10_000, 128); + let source = FailingReader { + bytes: b"crash cause\n".to_vec(), + emitted: false, + }; + let mut sink = RecordingSink::default(); + pump_stderr_into(source, Arc::clone(&ring), &mut sink).await; + + let snapshot = lock_ring(&ring).snapshot(None, None); + assert_eq!(lines(&snapshot), vec!["crash cause"]); + assert!(matches!( + snapshot.capture, + CaptureState::Incomplete { ref reason } if reason.contains("reader failed") + )); + assert_eq!(sink.writes, vec![b"crash cause\n".to_vec()]); + } + + #[tokio::test] + async fn every_captured_line_is_also_forwarded() { + // Forwarding is not optional. The daemon log is overwhelmingly module + // output; a tap that captured without forwarding would leave it nearly + // empty and every existing reader would report clean on nothing. + let ring = shared(10, 10_000, 128); + let source = std::io::Cursor::new(b"alpha\nbeta\n".to_vec()); + let mut sink = RecordingSink::default(); + pump_stderr_into(source, Arc::clone(&ring), &mut sink).await; + + assert_eq!(sink.writes, vec![b"alpha\n".to_vec(), b"beta\n".to_vec()]); + } + + #[tokio::test] + async fn each_forwarded_line_is_exactly_one_write() { + // Inheriting the fd gave line atomicity for free. Reading a pipe and + // re-emitting can split a line that used to be atomic, so the framing + // must be one syscall per complete line -- asserted as one write per + // line, not merely as correct bytes. + let ring = shared(10, 10_000, 128); + let source = std::io::Cursor::new(b"first\nsecond\nthird\n".to_vec()); + let mut sink = RecordingSink::default(); + pump_stderr_into(source, Arc::clone(&ring), &mut sink).await; + + assert_eq!(sink.writes.len(), 3); + for write in &sink.writes { + assert_eq!( + write.iter().filter(|byte| **byte == b'\n').count(), + 1, + "a write carried something other than exactly one complete line" + ); + assert_eq!(*write.last().unwrap(), b'\n'); + } + } + + #[test] + fn the_first_process_start_is_not_recorded_because_it_divides_nothing() { + // Otherwise a module that printed nothing renders as a lone boundary + // marker, and every caller has to decide whether that counts as silence. + let mut ring = ring(10, 10_000, 128); + ring.push_process_start(); + assert!(ring.snapshot(None, None).entries.is_empty()); + + ring.push_line("first process said this"); + ring.push_process_start(); + assert!( + matches!(ring.entries.back(), Some(TailEntry::ProcessStart)), + "a boundary with output before it must be recorded" + ); + } + + #[test] + fn a_process_start_is_recorded_when_only_dropped_lines_precede_it() { + // The ring can be non-empty in the sense that matters -- lines were + // written and evicted -- while `entries` is empty. Suppressing the + // boundary there would attribute surviving output to the wrong process. + let mut ring = ring(1, 10_000, 128); + ring.push_line("evicted"); + ring.push_line("also evicted"); + ring.entries.clear(); + ring.push_process_start(); + assert!(matches!( + ring.entries.front(), + Some(TailEntry::ProcessStart) + )); + } + + #[tokio::test] + async fn the_pump_marks_captured_even_when_the_module_writes_nothing() { + // Clean EOF with no output is a module that was quiet, not one nobody + // listened to -- and the two must not render alike. + let ring = shared(10, 10_000, 128); + let source = std::io::Cursor::new(Vec::new()); + let mut sink = RecordingSink::default(); + pump_stderr_into(source, Arc::clone(&ring), &mut sink).await; + + let snapshot = lock_ring(&ring).snapshot(None, None); + assert!(snapshot.entries.is_empty()); + assert_eq!(snapshot.capture, CaptureState::Captured); + assert!(sink.writes.is_empty()); + } + + #[tokio::test] + async fn a_line_with_no_newline_cannot_grow_the_reader_without_bound() { + // A module fault must not become a daemon fault: without the pending + // ceiling this buffer grows to whatever the module writes. + let ring = shared(10, 10_000_000, 4 * 1024 * 1024); + let source = std::io::Cursor::new(vec![b'x'; MAX_PENDING_LINE_BYTES + 4096]); + let mut sink = RecordingSink::default(); + pump_stderr_into(source, Arc::clone(&ring), &mut sink).await; + + let snapshot = lock_ring(&ring).snapshot(None, None); + assert_eq!( + lines(&snapshot).len(), + 2, + "expected a forced flush at the ceiling plus the remainder" + ); + assert_eq!( + sink.writes, + vec![vec![b'x'; MAX_PENDING_LINE_BYTES], vec![b'x'; 4096],], + "forced flushes and EOF fragments must not invent delimiters" + ); + } + + #[test] + fn a_byte_limit_smaller_than_one_line_still_returns_that_line() { + // Returning nothing would be indistinguishable from a quiet module, which + // is the failure this module exists to prevent. + let mut ring = ring(10, 10_000, 128); + ring.mark_captured(); + ring.push_line("a line considerably longer than the request limit"); + let snapshot = ring.snapshot(None, Some(4)); + assert_eq!(snapshot.entries.len(), 1); + } + + #[test] + fn an_incoherent_config_clamps_the_line_cap_and_keeps_its_restart_boundary() { + let config = StderrTailConfig::new(2, 10, 100); + assert_eq!(config.max_line_bytes, config.max_bytes); + let mut ring = StderrRing::new(config); + ring.mark_captured(); + ring.push_line("old"); + ring.push_process_start(); + ring.push_line("new process line longer than the ring byte cap"); + + assert_eq!( + ring.snapshot(None, None).entries, + vec![ + TailEntry::ProcessStart, + TailEntry::Line { + text: "new proces".to_string(), + truncated: true, + }, + ] + ); + } +} diff --git a/crates/subc-core/src/supervise.rs b/crates/subc-core/src/supervise.rs index ea3a965c..185aac51 100644 --- a/crates/subc-core/src/supervise.rs +++ b/crates/subc-core/src/supervise.rs @@ -3,7 +3,7 @@ use std::{ error::Error, fmt, io, path::PathBuf, - process::ExitStatus, + process::{ExitStatus, Stdio}, sync::{Arc, Mutex, OnceLock}, time::{Duration, SystemTime, UNIX_EPOCH}, }; @@ -28,6 +28,7 @@ use crate::{ ModuleDrainTarget, PendingModuleControlRpc, }, registry::RegistryError, + stderr_tail::{pump_stderr, StderrRing, StderrTailConfig, StderrTailSnapshot}, Frame, Registry, }; @@ -43,6 +44,58 @@ const DEFAULT_BACKOFF: Duration = Duration::from_millis(100); const DEFAULT_DRAIN_TIMEOUT: Duration = Duration::from_secs(2); const REGISTRY_RELEASE_TIMEOUT: Duration = Duration::from_secs(1); const REGISTRY_RELEASE_POLL: Duration = Duration::from_millis(10); +const STDERR_PUMP_DRAIN_TIMEOUT: Duration = Duration::from_millis(250); + +struct SupervisedChild { + child: Child, + stderr_pump: Option>, + stderr_ring: Arc>, +} + +impl SupervisedChild { + fn id(&self) -> Option { + self.child.id() + } + + async fn wait(&mut self) -> io::Result { + self.child.wait().await + } + + fn start_kill(&mut self) -> io::Result<()> { + self.child.start_kill() + } + + async fn drain_stderr(&mut self, module_id: &str) { + let Some(mut pump) = self.stderr_pump.take() else { + return; + }; + match timeout(STDERR_PUMP_DRAIN_TIMEOUT, &mut pump).await { + Ok(Ok(())) => {} + Ok(Err(err)) => { + self.stderr_ring + .lock() + .unwrap_or_else(|poisoned| poisoned.into_inner()) + .mark_incomplete(format!("stderr pump ended unexpectedly: {err}")); + warn!(module_id, error = %err, "stderr pump ended before clean EOF"); + } + Err(_) => { + pump.abort(); + self.stderr_ring + .lock() + .unwrap_or_else(|poisoned| poisoned.into_inner()) + .mark_incomplete(format!( + "stderr pump did not reach EOF within {:?} before restart", + STDERR_PUMP_DRAIN_TIMEOUT + )); + warn!( + module_id, + waited = ?STDERR_PUMP_DRAIN_TIMEOUT, + "stderr pump did not drain before restart; stopped it before marking the new process" + ); + } + } + } +} fn registration_release_events() -> &'static watch::Sender { static EVENTS: OnceLock> = OnceLock::new(); @@ -366,6 +419,13 @@ struct SupervisorRuntimeConfig { /// The shared handle, so every spawn path (initial, restart, reload) records the /// reserved-module launch nonce the HELLO verifier checks against. supervisor_handle: Option, + /// This module's stderr tail, shared with the [`SupervisedModule`] that answers + /// status queries. + /// + /// One ring per module, held across every respawn. The lines explaining an exit + /// are written BEFORE that exit, so a ring recreated per process would be empty + /// exactly when it is asked for. + stderr_ring: Arc>, } #[derive(Debug, Clone, PartialEq, Eq)] @@ -700,6 +760,7 @@ impl Supervisor { &spec, runtime.connection_file_path.as_deref(), self.supervisor_handle.as_ref(), + &runtime.stderr_ring, )?; set_running(&snapshot, child.id())?; self.process_liveness @@ -731,6 +792,7 @@ impl Supervisor { &spec, runtime.connection_file_path.as_deref(), self.supervisor_handle.as_ref(), + &runtime.stderr_ring, ) { Ok(child) => { set_running(&snapshot, child.id())?; @@ -771,6 +833,7 @@ impl Supervisor { &spec, runtime.connection_file_path.as_deref(), self.supervisor_handle.as_ref(), + &runtime.stderr_ring, ) { Ok(child) => { set_running(&snapshot, child.id())?; @@ -808,6 +871,7 @@ impl Supervisor { connection_file_path: self.connection_file_path.clone(), forwarding: self.forwarding.clone(), supervisor_handle: self.supervisor_handle.clone(), + stderr_ring: Arc::new(Mutex::new(StderrRing::new(StderrTailConfig::default()))), } } @@ -816,12 +880,13 @@ impl Supervisor { spec: ModuleSpec, runtime: SupervisorRuntimeConfig, snapshot: SharedSnapshot, - child: Option, + child: Option, ) -> SupervisedModule { let configuration = Arc::new(Mutex::new(SupervisedConfiguration { spec: spec.clone(), health: runtime.health, })); + let stderr_ring = Arc::clone(&runtime.stderr_ring); let (tx, rx) = mpsc::channel(4); let monitor = tokio::spawn(supervise_loop( spec.clone(), @@ -840,6 +905,7 @@ impl Supervisor { registry: Arc::clone(&self.registry), snapshot, configuration, + stderr_ring, commands: tx, monitor: Mutex::new(Some(monitor)), }), @@ -869,6 +935,7 @@ struct SupervisedModuleInner { registry: Arc, snapshot: SharedSnapshot, configuration: Arc>, + stderr_ring: Arc>, commands: mpsc::Sender, monitor: Mutex>>, } @@ -891,6 +958,24 @@ impl SupervisedModule { Ok(lock_snapshot(&self.inner.snapshot)?.state) } + /// The module's retained stderr, newest lines last. + /// + /// Deliberately NOT on [`Self::status`]: a bounded tail is kilobytes per + /// module, `supervisor.list` renders every module, and putting it in the + /// shared snapshot would make each status read carry a payload almost nobody + /// asked for. Callers that want the text ask for it. + pub fn stderr_tail( + &self, + max_lines: Option, + max_bytes: Option, + ) -> StderrTailSnapshot { + self.inner + .stderr_ring + .lock() + .unwrap_or_else(|poisoned| poisoned.into_inner()) + .snapshot(max_lines, max_bytes) + } + pub fn status(&self) -> Result { let snapshot = lock_snapshot(&self.inner.snapshot)?.clone(); let registration_active = self @@ -1454,7 +1539,7 @@ async fn run_health_probe_cycle( registry: &Registry, process_liveness: &SupervisorProcessLiveness, snapshot: &SharedSnapshot, - child: &mut Option, + child: &mut Option, ) { let now_ms = unix_ms_now(); match probe_module_health(spec, runtime).await { @@ -1608,7 +1693,7 @@ async fn handle_health_report( registry: &Registry, process_liveness: &SupervisorProcessLiveness, snapshot: &SharedSnapshot, - child: &mut Option, + child: &mut Option, report: HealthReport, now_ms: u64, ) { @@ -1650,7 +1735,7 @@ async fn handle_health_probe_failure( registry: &Registry, process_liveness: &SupervisorProcessLiveness, snapshot: &SharedSnapshot, - child: &mut Option, + child: &mut Option, err: HealthProbeError, now_ms: u64, ) { @@ -1729,7 +1814,7 @@ async fn apply_l3_health_action( registry: &Registry, process_liveness: &SupervisorProcessLiveness, snapshot: &SharedSnapshot, - child: &mut Option, + child: &mut Option, status: SupervisorHealthStatus, detail: Option<&str>, action: HealthAction, @@ -1780,7 +1865,7 @@ async fn health_restart_child( registry: &Registry, process_liveness: &SupervisorProcessLiveness, snapshot: &SharedSnapshot, - child: &mut Option, + child: &mut Option, status: SupervisorHealthStatus, detail: Option<&str>, now_ms: u64, @@ -1968,7 +2053,7 @@ async fn supervise_loop( registry: Arc, process_liveness: Arc, snapshot: SharedSnapshot, - mut child: Option, + mut child: Option, mut commands: mpsc::Receiver, ) { let mut health_probe = HealthProbeRuntime::default(); @@ -1991,6 +2076,7 @@ async fn supervise_loop( let exit_report = match wait_result { Ok(status) => classify_exit(&status), Err(err) => { + active_child.drain_stderr(&spec.module_id).await; fail_snapshot(&snapshot, Some(&spec.module_id), None); untrack_if_registration_released( &process_liveness, @@ -2003,6 +2089,7 @@ async fn supervise_loop( continue; } }; + active_child.drain_stderr(&spec.module_id).await; match on_child_exit( &spec, @@ -2110,7 +2197,7 @@ async fn handle_supervisor_command( registry: &Registry, process_liveness: &SupervisorProcessLiveness, snapshot: &SharedSnapshot, - child: &mut Option, + child: &mut Option, ) -> bool { match command { SupervisorCommand::Drain { reply } => { @@ -2208,7 +2295,7 @@ async fn restart_child( registry: &Registry, process_liveness: &SupervisorProcessLiveness, snapshot: &SharedSnapshot, - child: &mut Option, + child: &mut Option, ) -> Result<(), SuperviseError> { // Restart cycles a running module; it must not silently start a disabled one. if !lock_snapshot(snapshot)?.enabled { @@ -2254,7 +2341,7 @@ async fn reload_child( registry: &Registry, process_liveness: &SupervisorProcessLiveness, snapshot: &SharedSnapshot, - child: &mut Option, + child: &mut Option, ) -> Result<(), SuperviseError> { // Reload cycles a running module; it must not silently start a disabled one. if !lock_snapshot(snapshot)?.enabled { @@ -2321,6 +2408,9 @@ async fn reload_child( Ok(()) } RegistrationWaitOutcome::Exited(exit_report) => { + if let Some(active_child) = child.as_mut() { + active_child.drain_stderr(&spec.module_id).await; + } *child = None; handle_reload_child_registration_failure( spec, @@ -2353,6 +2443,7 @@ async fn reload_child( module_id: spec.module_id.clone(), source, })?; + timed_out_child.drain_stderr(&spec.module_id).await; handle_reload_child_registration_failure( spec, runtime, @@ -2379,7 +2470,7 @@ async fn set_child_enabled( registry: &Registry, process_liveness: &SupervisorProcessLiveness, snapshot: &SharedSnapshot, - child: &mut Option, + child: &mut Option, enabled: bool, ) -> Result { let (current_enabled, current_state) = { @@ -2557,7 +2648,8 @@ fn spawn_child( spec: &ModuleSpec, connection_file_path: Option<&std::path::Path>, handle: Option<&SupervisorHandle>, -) -> Result { + ring: &Arc>, +) -> Result { let mut command = Command::new(&spec.program); command.args(&spec.args); if let Some(connection_file_path) = connection_file_path { @@ -2580,35 +2672,74 @@ fn spawn_child( } command.env(SUBC_LAUNCH_NONCE_ENV, nonce); - // STDIO IS DELIBERATELY LEFT INHERITED, which is a decision no line of this - // function states and therefore worth stating here: every supervised child - // gets the daemon's own stdout/stderr, so all of them share ONE file - // descriptor and one file offset. That is what puts every module's output in - // a single daemon log without a per-child reader task. + // STDOUT STAYS INHERITED; STDERR IS PIPED. The asymmetry is the whole design, + // so it is worth saying why rather than leaving it to be inferred. + // + // Inheriting both was the original choice: every child wrote to the daemon's + // own descriptors, which put all module output in one log with no per-child + // reader task. That was correct about interleaving and silent about DURABILITY, + // and durability is the axis that decides whether a crash can be diagnosed. + // A module's stderr is the only diagnostic input with no in-memory path -- + // `last_exit` survives a respawn because the supervisor holds it, while the + // text explaining that exit went to a sink that rotates or fills. Measured on + // two hosts: a systemd journal at its size cap retaining ~3.2 hours, and a + // plain log file reaching 908 MB with one module accounting for 98% of it. In + // both, the noisiest module sets everyone else's retention and the victim has + // no way to know its window shrank. // - // THE COST IS INTERLEAVING, and the mechanism is finer than a shared fd. A - // module owner measured it: an emitter that formats INCREMENTALLY issues one - // write syscall per format fragment, and while a process-local lock serialises - // those within one process, nothing serialises them ACROSS processes. Two - // processes on one inherited fd, 1500 lines each: 212 of 3000 lines came out - // spliced with an incremental emitter, 0 of 3000 when each line was formatted - // first and written in a single call. So the inheritance is the exposure and - // the multi-syscall write is what converts it into damage. + // So stderr is piped into a bounded in-memory ring the supervisor owns, and + // every line is forwarded on so the daemon log keeps its current content. That + // forwarding is MANDATORY rather than courteous: the log is overwhelmingly + // module output (4727 of 5000 sampled lines carried a module tag), so a tap + // that captured without forwarding would leave it nearly empty and every + // existing reader would report clean on nothing -- an absence that reads as + // calm, which is worse than the interleaving it replaces. // - // CONSEQUENCE FOR ANYONE SCRAPING THE DAEMON LOG: do not anchor patterns at - // line start. A single-write emitter guarantees its line is WHOLE, not that it - // begins a line -- another process mid-write can still land a fragment ahead of - // it, measured at ~8% of lines. An unanchored match found 1500 of 1500 where - // `^`-anchored found 1382. Strip escape sequences too. + // THE TAP MAKES LINE ATOMICITY THIS DAEMON'S PROBLEM. Previously a module's + // own write reached the fd in one syscall and the splicing came from emitters + // that formatted incrementally: two processes on one inherited fd, 1500 lines + // each, produced 212 spliced lines of 3000 with an incremental emitter and 0 + // of 3000 when each line was formatted first and written once. Reading a pipe + // and re-emitting can split a line that WAS atomic, so the reader reassembles + // to a complete line and writes it in a single call -- otherwise this change + // introduces a defect the previous design did not have. // - // The structural fix is a per-child pipe with a line-atomic writer here, which - // would make the property hold for modules this daemon does not own. Not done: - // it adds a reader task per child and moves where logs land, and the mitigation - // above is free for any module that adopts it. + // Stdout is left inherited: modules use it for ordinary output rather than + // diagnostics, and piping it would double the reader tasks for no diagnostic + // gain. + command.stderr(Stdio::piped()); command.kill_on_drop(true); - command.spawn().map_err(|source| SuperviseError::Spawn { + let mut child = command.spawn().map_err(|source| SuperviseError::Spawn { program: spec.program.clone(), source, + })?; + + let stderr_pump = match child.stderr.take() { + Some(stderr) => { + ring.lock() + .unwrap_or_else(|poisoned| poisoned.into_inner()) + .push_process_start(); + Some(tokio::spawn(pump_stderr(stderr, Arc::clone(ring)))) + } + None => { + // Spawning succeeded but the pipe did not materialise. Recording it as + // uncaptured keeps the tail honest: the alternative is an empty tail + // that reads as a module which printed nothing. + ring.lock() + .unwrap_or_else(|poisoned| poisoned.into_inner()) + .mark_not_captured("stderr pipe was not available on spawn"); + warn!( + module_id = %spec.module_id, + "spawned child exposed no stderr pipe; tail will be unavailable" + ); + None + } + }; + + Ok(SupervisedChild { + child, + stderr_pump, + stderr_ring: Arc::clone(ring), }) } @@ -2644,11 +2775,12 @@ fn spawn_and_mark_running( spec: &ModuleSpec, runtime: &SupervisorRuntimeConfig, snapshot: &SharedSnapshot, -) -> Result { +) -> Result { let child = spawn_child( spec, runtime.connection_file_path.as_deref(), runtime.supervisor_handle.as_ref(), + &runtime.stderr_ring, )?; set_running(snapshot, child.id())?; Ok(child) @@ -2851,7 +2983,7 @@ async fn begin_forwarding_drain_with( async fn wait_for_registration_after_reload( registry: &Registry, module_id: &str, - child: &mut Child, + child: &mut SupervisedChild, wait: Duration, ) -> Result { let deadline = Instant::now() + wait; @@ -2897,7 +3029,7 @@ async fn handle_reload_child_registration_failure( registry: &Registry, process_liveness: &SupervisorProcessLiveness, snapshot: &SharedSnapshot, - child: &mut Option, + child: &mut Option, failure: ReloadRegistrationFailure, ) -> Result<(), SuperviseError> { let ReloadRegistrationFailure { @@ -2963,7 +3095,7 @@ async fn handle_reload_spawn_failure( runtime: &SupervisorRuntimeConfig, process_liveness: &SupervisorProcessLiveness, snapshot: &SharedSnapshot, - child: &mut Option, + child: &mut Option, reason: String, ) -> Result<(), SuperviseError> { let mut should_retry = false; @@ -3015,7 +3147,7 @@ async fn drain_optional_child( module_id: &str, registry: &Registry, snapshot: &SharedSnapshot, - child: &mut Option, + child: &mut Option, drain_timeout: Duration, final_state: ModuleState, enabled: Option, @@ -3048,7 +3180,7 @@ async fn drain_child_to_state( module_id: &str, registry: &Registry, snapshot: &SharedSnapshot, - mut child: Child, + mut child: SupervisedChild, drain_timeout: Duration, final_state: ModuleState, enabled: Option, @@ -3091,6 +3223,7 @@ async fn drain_child_to_state( state.pid = None; state.last_exit = Some(exit_report); })?; + child.drain_stderr(module_id).await; wait_for_registration_release(registry, module_id, REGISTRY_RELEASE_TIMEOUT).await } diff --git a/crates/subc-core/tests/closure.rs b/crates/subc-core/tests/closure.rs index 55f4ccb2..e6367c94 100644 --- a/crates/subc-core/tests/closure.rs +++ b/crates/subc-core/tests/closure.rs @@ -700,6 +700,7 @@ fn thin_core_ops() -> BTreeSet<&'static str> { ops::SUPERVISOR_SET_ENABLED, ops::SUPERVISOR_HEALTH_PROBE, ops::SUPERVISOR_HEALTH, + ops::SUPERVISOR_STDERR_TAIL, ]) } diff --git a/crates/subc-core/tests/supervision.rs b/crates/subc-core/tests/supervision.rs index bce3d993..824ec2d7 100644 --- a/crates/subc-core/tests/supervision.rs +++ b/crates/subc-core/tests/supervision.rs @@ -1,6 +1,7 @@ use std::{ops::Deref, path::PathBuf, sync::Arc, time::Duration}; use subc_core::{ + stderr_tail::{CaptureState, StderrTailSnapshot, TailEntry}, ModuleSpec, ModuleState, ModuleStatus, Registry, RestartPolicy, SuperviseError, SupervisedModule, Supervisor, }; @@ -525,6 +526,246 @@ async fn wait_for_registration( } } +/// The case #7 was filed about: a module that dies with its cause on stderr. +/// +/// The claustrum incident had `exit_code: 1` and the reason -- a missing config +/// section -- only in the text, which was gone from the journal by the time +/// anyone looked. This asserts the text is recoverable from the supervisor after +/// the process is dead, with no log file in the path. +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn a_dead_module_leaves_its_stderr_readable_from_the_supervisor() { + let server = TestServer::start().await; + let supervisor = supervisor(&server, 0, Duration::from_millis(10)); + let module = supervisor + .spawn(ModuleSpec { + module_id: "stderr-tail-crasher".to_string(), + program: PathBuf::from(env!("CARGO_BIN_EXE_fake-aft-stub")), + args: Vec::new(), + env: vec![ + ( + "FAKE_AFT_STDERR_LINE".to_string(), + "config error: missing top-level `storage`".to_string(), + ), + ("FAKE_AFT_EXIT_CODE".to_string(), "1".to_string()), + ], + reserved: false, + reserved_prefixes: Vec::new(), + }) + .unwrap(); + + let status = wait_for_status(&module, Duration::from_secs(5), |status| { + status.state == ModuleState::Failed + }) + .await; + assert_eq!( + status.last_exit.as_ref().and_then(|exit| exit.code), + Some(1), + "precondition: the module should have exited non-zero" + ); + + let tail = wait_for_tail(&module, Duration::from_secs(5), |tail| { + matches!(tail.capture, CaptureState::Captured) + && tail.entries.iter().any( + |entry| matches!(entry, TailEntry::Line { text, .. } if text.contains("missing top-level `storage`")), + ) + }) + .await; + assert!( + matches!(tail.capture, CaptureState::Captured), + "a spawned module must report captured, not an empty tail that reads as silence" + ); + + let lines: Vec<&str> = tail + .entries + .iter() + .filter_map(|entry| match entry { + TailEntry::Line { text, .. } => Some(text.as_str()), + TailEntry::ProcessStart => None, + }) + .collect(); + assert!( + lines + .iter() + .any(|line| line.contains("missing top-level `storage`")), + "the cause of the exit was not recoverable from the tail; got {lines:?}" + ); +} + +/// A module that exits cleanly having printed nothing must not look like one +/// nobody was listening to. +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn a_silent_module_reports_captured_and_empty_rather_than_uncaptured() { + let server = TestServer::start().await; + let supervisor = supervisor(&server, 0, Duration::from_millis(10)); + let module = supervisor + .spawn(ModuleSpec { + module_id: "stderr-tail-silent".to_string(), + program: PathBuf::from(env!("CARGO_BIN_EXE_fake-aft-stub")), + args: Vec::new(), + env: vec![("FAKE_AFT_EXIT_CODE".to_string(), "3".to_string())], + reserved: false, + reserved_prefixes: Vec::new(), + }) + .unwrap(); + + wait_for_status(&module, Duration::from_secs(5), |status| { + status.state == ModuleState::Failed + }) + .await; + + let tail = wait_for_tail(&module, Duration::from_secs(5), |tail| { + matches!(tail.capture, CaptureState::Captured) && tail.entries.is_empty() + }) + .await; + assert!( + matches!(tail.capture, CaptureState::Captured), + "silence and absence must be distinguishable; got {:?}", + tail.capture + ); + assert!( + tail.entries.is_empty(), + "expected no lines from a module that printed nothing; got {:?}", + tail.entries + ); +} + +/// The tail has to outlive the process whose death it explains. +/// +/// A ring recreated per spawn would be empty exactly when asked, and the restart +/// boundary has to be visible in-band -- which side of a restart a line falls on +/// is unanswerable from a count. +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn stderr_from_before_a_restart_survives_with_a_marked_boundary() { + let server = TestServer::start().await; + let supervisor = supervisor(&server, 3, Duration::from_millis(10)); + let module = supervisor + .spawn(ModuleSpec { + module_id: "stderr-tail-looper".to_string(), + program: PathBuf::from(env!("CARGO_BIN_EXE_fake-aft-stub")), + args: Vec::new(), + env: vec![ + ("FAKE_AFT_STDERR_LINE".to_string(), "boot {pid}".to_string()), + ("FAKE_AFT_EXIT_CODE".to_string(), "1".to_string()), + ], + reserved: false, + reserved_prefixes: Vec::new(), + }) + .unwrap(); + + let tail = wait_for_tail(&module, Duration::from_secs(5), |tail| { + let boots = tail + .entries + .iter() + .filter( + |entry| matches!(entry, TailEntry::Line { text, .. } if text.starts_with("boot ")), + ) + .count(); + let process_starts = tail + .entries + .iter() + .filter(|entry| matches!(entry, TailEntry::ProcessStart)) + .count() + >= 2; + boots >= 2 && process_starts + }) + .await; + + let boots = tail + .entries + .iter() + .filter(|entry| matches!(entry, TailEntry::Line { text, .. } if text.starts_with("boot "))) + .count(); + assert!( + boots >= 2, + "output from before the restart was lost; got {:?}", + tail.entries + ); + assert!( + tail.entries + .iter() + .filter(|entry| matches!(entry, TailEntry::ProcessStart)) + .count() + >= 2, + "restart boundaries were lost; got {:?}", + tail.entries + ); +} + +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn a_wedged_old_stderr_pump_is_stopped_before_the_next_restart_boundary() { + let server = TestServer::start().await; + let supervisor = supervisor(&server, 1, Duration::from_millis(10)); + let module = supervisor + .spawn(ModuleSpec { + module_id: "stderr-tail-wedged-pump".to_string(), + program: PathBuf::from(env!("CARGO_BIN_EXE_fake-aft-stub")), + args: Vec::new(), + env: vec![ + ("FAKE_AFT_STDERR_LINE".to_string(), "old-start".to_string()), + ("FAKE_AFT_EXIT_CODE".to_string(), "1".to_string()), + ( + "FAKE_AFT_ORPHAN_WRITER_DELAY_MS".to_string(), + "1000".to_string(), + ), + ( + "FAKE_AFT_ORPHAN_WRITER_LINE".to_string(), + "old-trailing".to_string(), + ), + ], + reserved: false, + reserved_prefixes: Vec::new(), + }) + .unwrap(); + + let tail = wait_for_tail(&module, Duration::from_secs(5), |tail| { + matches!(tail.capture, CaptureState::Incomplete { .. }) + && tail + .entries + .iter() + .any(|entry| matches!(entry, TailEntry::ProcessStart)) + }) + .await; + assert!( + tail.entries + .iter() + .any(|entry| matches!(entry, TailEntry::Line { text, .. } if text == "old-start"),), + "the initial process output was not retained: {:?}", + tail.entries + ); + + sleep(Duration::from_millis(1200)).await; + let tail = module.stderr_tail(None, None); + assert!( + !tail + .entries + .iter() + .any(|entry| matches!(entry, TailEntry::Line { text, .. } if text == "old-trailing"),), + "old output crossed the restart boundary: {:?}", + tail.entries + ); +} + +async fn wait_for_tail( + module: &SupervisedModule, + wait: Duration, + matches: impl Fn(&StderrTailSnapshot) -> bool, +) -> StderrTailSnapshot { + let deadline = Instant::now() + wait; + loop { + let tail = module.stderr_tail(None, None); + if matches(&tail) { + return tail; + } + if Instant::now() >= deadline { + panic!( + "module {} did not reach the expected stderr tail within {wait:?}; last: {tail:?}", + module.module_id() + ); + } + sleep(Duration::from_millis(10)).await; + } +} + async fn wait_for_status( module: &SupervisedModule, wait: Duration,