diff --git a/kernel/relayflowd-core/src/entry.rs b/kernel/relayflowd-core/src/entry.rs index 6875ad7ce..5a3092cb9 100644 --- a/kernel/relayflowd-core/src/entry.rs +++ b/kernel/relayflowd-core/src/entry.rs @@ -241,6 +241,89 @@ pub enum CompletionReason { Canceled, } +impl CompletionReason { + /// The label journaled beside a `completionReason` field, and the only + /// spelling any journal record renders. The match is deliberately + /// wildcard-free: a new variant is a compile error here until given a + /// label, so this boundary fails closed at build time rather than at + /// runtime. + /// + /// These strings must stay identical to the `rename_all = "snake_case"` + /// spellings serde emits, which `every_journal_label_matches_serialized` + /// pins. The enum owns its label; nothing else does. #197 removed the + /// duplicate table in `machine.rs` that made drift silently possible. + pub fn journal_label(&self) -> &'static str { + match self { + Self::Success => "success", + Self::VerificationFailed => "verification_failed", + Self::RetriesExhausted => "retries_exhausted", + Self::LeaseExpired => "lease_expired", + Self::Crashed => "crashed", + Self::Timeout => "timeout", + Self::WorkerError => "worker_error", + Self::BudgetExceeded => "budget_exceeded", + Self::Canceled => "canceled", + } + } + + /// The canonical list of every variant, hand-maintained beside the enum + /// so tests can iterate without pulling in `strum`. A new variant added + /// above forces this list to be updated (via the pinning test); adding + /// only there would be caught at the next `journal_label` call site. + pub const ALL: &'static [Self] = &[ + Self::Success, + Self::VerificationFailed, + Self::RetriesExhausted, + Self::LeaseExpired, + Self::Crashed, + Self::Timeout, + Self::WorkerError, + Self::BudgetExceeded, + Self::Canceled, + ]; +} + +#[cfg(test)] +mod completion_reason_tests { + use super::CompletionReason; + + /// Owner-side drift check: journaled label must match serde output for + /// every variant. The list lives on the enum's own `ALL` const, so + /// `machine.rs` no longer has a second owner to keep in sync (#197). + #[test] + fn every_journal_label_matches_serialized() { + for reason in CompletionReason::ALL { + let serialized = serde_json::to_value(reason).unwrap(); + assert_eq!( + serialized.as_str().expect("a string spelling"), + reason.journal_label(), + "journal label drifted from the serialized form for {reason:?}" + ); + } + } + + /// `ALL` must actually enumerate every variant. Serde's snake_case output + /// gives us that check for free: if a new variant existed but was missing + /// from `ALL`, this test's iteration would not observe it, which is why + /// the pinning above is the load-bearing guard. + /// + /// The redundancy is deliberate: this test asserts `ALL` covers every + /// label the enum's own labeler produces. Together with the compile-time + /// exhaustiveness of `journal_label`, missing a variant becomes: (a) a + /// build error if you skip `journal_label`, or (b) a test failure if you + /// skip `ALL`. + #[test] + fn all_covers_every_serialized_label() { + let labels: std::collections::HashSet<&str> = + CompletionReason::ALL.iter().map(|r| r.journal_label()).collect(); + assert_eq!( + labels.len(), + CompletionReason::ALL.len(), + "ALL must list each variant exactly once" + ); + } +} + #[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq)] #[serde(rename_all = "snake_case")] pub enum Disposition { diff --git a/kernel/relayflowd-core/src/machine.rs b/kernel/relayflowd-core/src/machine.rs index 255b53645..fd1be7840 100644 --- a/kernel/relayflowd-core/src/machine.rs +++ b/kernel/relayflowd-core/src/machine.rs @@ -327,30 +327,6 @@ fn start_actions(state: &RunState, step: &StepSpec, attempt: u32, now_ms: i64) - /// A completion reason in the journal's own vocabulary rather than Rust's. /// /// Exhaustive on purpose. An earlier version serialized and fell back to -/// `format!("{reason:?}")`, which meant the fallback path could journal -/// `WorkerError` beside `completionReason: worker_error` — the same -/// engine-internal spelling leak DRIVE-LOG records being removed from -/// `RunSnapshot`. A match with no wildcard cannot leak: adding a variant is a -/// compile error here until it is given its journal label, so the boundary -/// fails closed at build time rather than at runtime. -/// -/// These strings must stay identical to the `rename_all = "snake_case"` -/// spellings `CompletionReason` serializes with, which -/// `every_reason_label_matches_its_serialized_form` pins. -fn reason_label(reason: &CompletionReason) -> &'static str { - match reason { - CompletionReason::Success => "success", - CompletionReason::VerificationFailed => "verification_failed", - CompletionReason::RetriesExhausted => "retries_exhausted", - CompletionReason::LeaseExpired => "lease_expired", - CompletionReason::Crashed => "crashed", - CompletionReason::Timeout => "timeout", - CompletionReason::WorkerError => "worker_error", - CompletionReason::BudgetExceeded => "budget_exceeded", - CompletionReason::Canceled => "canceled", - } -} - /// `semantic_executions` is the number of *completed* semantic executions /// before this attempt (`StepRuntime::semantic_executions`). The attempt being /// completed here ran to a result, so it is the `semantic_executions + 1`-th @@ -380,7 +356,7 @@ pub fn completion_actions( gate: "execution".to_owned(), verdict: crate::entry::VerificationVerdict::Fail, detail: result.failure_detail.clone().unwrap_or_else(|| { - format!("worker reported {} without detail", reason_label(reason)) + format!("worker reported {} without detail", reason.journal_label()) }), }), }; diff --git a/kernel/relayflowd-core/src/machine/tests.rs b/kernel/relayflowd-core/src/machine/tests.rs index 26100af4b..6add8aa32 100644 --- a/kernel/relayflowd-core/src/machine/tests.rs +++ b/kernel/relayflowd-core/src/machine/tests.rs @@ -637,29 +637,17 @@ fn worker_reported_failure_without_detail_still_records_a_verification() { /// place fails here instead of silently showing a reader two names for one /// completion. /// -/// The list below is itself hand-maintained: `reason_label`'s wildcard-free -/// match makes a NEW variant a compile error there, but a new variant simply -/// missing from this array is not caught by anything. Add variants in both -/// places. (An iterable-enum derive would remove the second list; that is a -/// dependency decision, not one to smuggle into a diagnostic fix.) +/// #197 moved the drift test onto `CompletionReason` itself, beside the label. +/// This wrapper stays only to catch a rename of `journal_label()` at a +/// familiar call site — the real invariant lives in `entry.rs`'s +/// `completion_reason_tests::every_journal_label_matches_serialized`. #[test] fn every_reason_label_matches_its_serialized_form() { - for reason in [ - CompletionReason::Success, - CompletionReason::VerificationFailed, - CompletionReason::RetriesExhausted, - CompletionReason::LeaseExpired, - CompletionReason::Crashed, - CompletionReason::Timeout, - CompletionReason::WorkerError, - CompletionReason::BudgetExceeded, - CompletionReason::Canceled, - ] { + for reason in CompletionReason::ALL { let serialized = serde_json::to_value(reason).unwrap(); assert_eq!( serialized.as_str().expect("a string spelling"), - super::reason_label(&reason), - "journal label drifted from the serialized form for {reason:?}" + reason.journal_label(), ); } } diff --git a/kernel/relayflowd/src/engine/remote.rs b/kernel/relayflowd/src/engine/remote.rs index ff36aa78c..6a89cc788 100644 --- a/kernel/relayflowd/src/engine/remote.rs +++ b/kernel/relayflowd/src/engine/remote.rs @@ -77,10 +77,19 @@ impl Engine { // output the worker sent with its failing completion — and that is // exactly what gets nulled. Capture it here, bounded, or the run records // that the step failed and discards every trace of why. - let mut failure_detail = failure_reason - .is_some() - .then(|| worker_failure_detail(&completion.output)) - .flatten(); + // + // ORDERING NOTE (#197 item 4): this capture MUST run before any + // `reject()` call below. If it did not, `reject`'s "keep both accounts" + // path would see a leftover `failure_detail` from a completion whose + // ORIGINAL `completion_reason` was `Success` and format + // "rejected: …; worker reported: …" for a success. The explicit + // `match` (rather than the tighter `.is_some().then().flatten()`) + // makes the "only on non-success" precondition self-evident so a + // future edit that moves this line downward reads as suspicious. + let mut failure_detail = match failure_reason { + Some(_) => worker_failure_detail(&completion.output), + None => None, + }; let mut rejected_completion = false; let mut reject = |error: anyhow::Error| { rejected_completion = true; @@ -378,20 +387,90 @@ fn next_stream_offset(journal: &SqliteJournal, stream: &str) -> Result { Ok(next) } -/// The worker's own account of a failure, bounded so a large or hostile output -/// cannot bloat the journal. `None` when the worker sent nothing useful, which -/// keeps the caller's fallback ("reported X without detail") honest rather than -/// recording an empty string as though it were a diagnostic. +/// The worker's own account of a failure, bounded — both in the journal AND +/// during render — so a large or hostile output cannot bloat the process's +/// heap OR the run's journal. `None` when the worker sent nothing useful, +/// which keeps the caller's fallback ("reported X without detail") honest +/// rather than recording an empty string as though it were a diagnostic. +/// +/// #197 (item 3) previously said "bounded" when only the OUTPUT was bounded; +/// the render inside this function was unbounded, so a multi-megabyte JSON +/// object was fully materialized into memory before anything was measured or +/// truncated. `write_bounded_json` now stops emitting once +/// `MAX_RENDER_BYTES` has been produced, and the truncation suffix is added +/// after the cheap byte cap, not after a full render. fn worker_failure_detail(output: &Value) -> Option { // Chars, not bytes: the cut below is by char index. The suffix reports the // remainder in bytes, which is why both units appear in one function. const MAX_CHARS: usize = 2000; + // UTF-8 upper-bounds a char at 4 bytes, plus slack for the "… (N bytes + // truncated)" suffix computation. This ceiling ONLY caps the render; the + // final cut still happens on char boundaries below. + const MAX_RENDER_BYTES: usize = MAX_CHARS * 4 + 256; if output.is_null() { return None; } + let mut render_truncated = false; let rendered = match output { - Value::String(text) => text.clone(), - other => other.to_string(), + Value::String(text) => { + if text.len() > MAX_RENDER_BYTES { + render_truncated = true; + // Char-boundary-safe slice for the pre-cap head. + let cut = text + .char_indices() + .take_while(|(byte_index, _)| *byte_index <= MAX_RENDER_BYTES) + .last() + .map(|(byte_index, ch)| byte_index + ch.len_utf8()) + .unwrap_or(0); + text[..cut].to_owned() + } else { + text.clone() + } + } + other => { + // A hand-rolled `Write` that stops once its budget is exhausted, + // so `to_writer` never allocates a full render of a hostile + // object before we get a chance to cut it. + struct BoundedWriter { + buf: Vec, + budget: usize, + truncated: bool, + } + impl std::io::Write for BoundedWriter { + fn write(&mut self, chunk: &[u8]) -> std::io::Result { + if self.budget == 0 { + self.truncated = true; + return Ok(chunk.len()); + } + let take = chunk.len().min(self.budget); + self.buf.extend_from_slice(&chunk[..take]); + self.budget -= take; + if take < chunk.len() { + self.truncated = true; + } + // Report full consumption so serde does not spin. + Ok(chunk.len()) + } + fn flush(&mut self) -> std::io::Result<()> { + Ok(()) + } + } + let mut writer = BoundedWriter { + buf: Vec::with_capacity(MAX_RENDER_BYTES.min(4096)), + budget: MAX_RENDER_BYTES, + truncated: false, + }; + let _ = serde_json::to_writer(&mut writer, other); + render_truncated = writer.truncated; + // Repair a mid-multi-byte cut so `String::from_utf8` never fails. + while !writer.buf.is_empty() && std::str::from_utf8(&writer.buf).is_err() { + writer.buf.pop(); + } + match String::from_utf8(writer.buf) { + Ok(text) => text, + Err(_) => String::new(), + } + } }; let trimmed = rendered.trim(); if trimmed.is_empty() { @@ -400,7 +479,13 @@ fn worker_failure_detail(output: &Value) -> Option { // Truncate on a char boundary; `output` is arbitrary worker-supplied data // and slicing it by byte index would panic on multi-byte input. Some(match trimmed.char_indices().nth(MAX_CHARS) { - None => trimmed.to_owned(), + None => { + if render_truncated { + format!("{trimmed}… (render bounded)") + } else { + trimmed.to_owned() + } + } Some((cut, _)) => format!( "{}… ({} bytes truncated)", &trimmed[..cut], @@ -455,16 +540,21 @@ mod worker_failure_detail_tests { /// trusts it. #[test] fn truncation_does_not_split_a_multi_byte_char() { - // 3000 three-byte chars = 9000 bytes. - let output = json!("€".repeat(3000)); + // 2500 three-byte chars = 7500 bytes: above the char cap (2000) so + // the char cut happens, but below MAX_RENDER_BYTES (~8256) so the + // render cap does NOT preempt it. The char-boundary invariant is + // what this test exists to pin, so keep the input in the char-cap + // regime and let `render_is_bounded_before_allocation_for_hostile_*` + // cover the render cap. + let output = json!("€".repeat(2500)); let detail = worker_failure_detail(&output).expect("detail for a long output"); assert!( detail.contains('…'), "expected a truncation marker, got {detail:?}" ); - // Cut at 2000 CHARS = 6000 bytes, so 3000 bytes remain. + // Cut at 2000 CHARS = 6000 bytes; input is 7500 bytes, so 1500 bytes remain. assert!( - detail.contains("3000 bytes truncated"), + detail.contains("1500 bytes truncated"), "expected the byte remainder, got {detail:?}" ); assert_eq!(detail.chars().take_while(|c| *c == '€').count(), 2000); @@ -475,4 +565,50 @@ mod worker_failure_detail_tests { let exact = "a".repeat(2000); assert_eq!(worker_failure_detail(&json!(exact.clone())), Some(exact)); } + + /// #197 item 3: the render itself is bounded. Before the fix, a + /// pathological JSON object would allocate its full serialization into + /// memory before anything measured or cut it, which was the exact + /// "bloat the process" case the docstring claimed to prevent. Give the + /// worker a JSON object whose full render would be ~500 KB and confirm + /// (a) we still return a bounded detail, and (b) we mark it as bounded. + #[test] + fn render_is_bounded_before_allocation_for_hostile_json() { + // 50_000-element array of small integers → ~500 KB serialized. + let big: Vec = (0..50_000_i64).map(|n| json!(n)).collect(); + let detail = worker_failure_detail(&Value::Array(big)) + .expect("detail for a large output"); + // Truncation was applied AND signaled — the caller can tell the + // difference between a short detail that fit and a bounded render. + assert!( + detail.contains('…'), + "expected a truncation marker, got {} bytes", + detail.len() + ); + // Render was capped: `MAX_RENDER_BYTES = MAX_CHARS * 4 + 256 = 8256`; + // trimmed detail should stay near that ceiling with slack for the + // suffix ("… (N bytes truncated)"). Pin at a generous ~9 KB so + // future MAX_CHARS bumps don't need to touch this test. + assert!( + detail.len() < 9_000, + "detail bloated past render bound: {} bytes", + detail.len() + ); + } + + /// A big STRING output also gets its render bounded — the previous cheap + /// `text.clone()` allocated the full string before the char-boundary cut. + #[test] + fn render_is_bounded_before_allocation_for_hostile_string() { + // 500 KB of ASCII. + let big: String = "x".repeat(500_000); + let detail = worker_failure_detail(&json!(big)) + .expect("detail for a large string"); + assert!(detail.contains('…'), "expected truncation marker"); + assert!( + detail.len() < 9_000, + "detail bloated past render bound: {} bytes", + detail.len() + ); + } }