Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
83 changes: 83 additions & 0 deletions kernel/relayflowd-core/src/entry.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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() {

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2: The test's name and doc claim it proves ALL enumerates every variant, but it only checks that the listed labels are unique. Both this test and every_journal_label_matches_serialized iterate over CompletionReason::ALL, so if a new variant is added to the enum and journal_label but not to ALL, neither test observes it. The claim '(b) a test failure if you skip ALL' is therefore false — the missing variant is silently undetected and can drift the journal label out of sync with the serialized form. Either add a real coverage guard or correct the misleading name/doc.

Prompt for AI agents
Check if this issue is valid — if so, understand the root cause and fix it. At kernel/relayflowd-core/src/entry.rs, line 316:

<comment>The test's name and doc claim it proves `ALL` enumerates every variant, but it only checks that the listed labels are unique. Both this test and `every_journal_label_matches_serialized` iterate over `CompletionReason::ALL`, so if a new variant is added to the enum and `journal_label` but not to `ALL`, neither test observes it. The claim '(b) a test failure if you skip ALL' is therefore false — the missing variant is silently undetected and can drift the journal label out of sync with the serialized form. Either add a real coverage guard or correct the misleading name/doc.</comment>

<file context>
@@ -241,6 +241,89 @@ pub enum CompletionReason {
+    /// 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();
</file context>

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 {
Expand Down
26 changes: 1 addition & 25 deletions kernel/relayflowd-core/src/machine.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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())
}),
}),
};
Expand Down
24 changes: 6 additions & 18 deletions kernel/relayflowd-core/src/machine/tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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 {

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P3: This test duplicates the owner-side drift test in entry.rs without exercising the machine path. Remove the wrapper or make it assert through completion_actions so it provides behavior coverage rather than a second identical assertion.

Prompt for AI agents
Check if this issue is valid — if so, understand the root cause and fix it. At kernel/relayflowd-core/src/machine/tests.rs, line 646:

<comment>This test duplicates the owner-side drift test in `entry.rs` without exercising the machine path. Remove the wrapper or make it assert through `completion_actions` so it provides behavior coverage rather than a second identical assertion.</comment>

<file context>
@@ -637,29 +637,17 @@ fn worker_reported_failure_without_detail_still_records_a_verification() {
-        CompletionReason::BudgetExceeded,
-        CompletionReason::Canceled,
-    ] {
+    for reason in CompletionReason::ALL {
         let serialized = serde_json::to_value(reason).unwrap();
         assert_eq!(
</file context>

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(),
);
}
}
Expand Down
166 changes: 151 additions & 15 deletions kernel/relayflowd/src/engine/remote.rs
Original file line number Diff line number Diff line change
Expand Up @@ -77,10 +77,19 @@ impl Engine<WallClock> {
// 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;
Expand Down Expand Up @@ -378,20 +387,90 @@ fn next_stream_offset(journal: &SqliteJournal, stream: &str) -> Result<u64> {
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<String> {
// 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 {

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2: When a failure string has more than 8,256 leading whitespace bytes, this prefix cap discards the useful diagnostic and returns None. Trim the string before applying the render cap so bounded rendering preserves the diagnostic.

Prompt for AI agents
Check if this issue is valid — if so, understand the root cause and fix it. At kernel/relayflowd/src/engine/remote.rs, line 416:

<comment>When a failure string has more than 8,256 leading whitespace bytes, this prefix cap discards the useful diagnostic and returns `None`. Trim the string before applying the render cap so bounded rendering preserves the diagnostic.</comment>

<file context>
@@ -378,20 +387,90 @@ fn next_stream_offset(journal: &SqliteJournal, stream: &str) -> Result<u64> {
-        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.
</file context>

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<u8>,
budget: usize,
truncated: bool,
}
impl std::io::Write for BoundedWriter {
fn write(&mut self, chunk: &[u8]) -> std::io::Result<usize> {
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;

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2: When the render cap and character cap both apply, this flag is ignored and the detail reports only the remainder of the capped prefix. Mark the detail as render-bounded, or report the byte count as a lower bound, whenever render_truncated is true.

Prompt for AI agents
Check if this issue is valid — if so, understand the root cause and fix it. At kernel/relayflowd/src/engine/remote.rs, line 464:

<comment>When the render cap and character cap both apply, this flag is ignored and the detail reports only the remainder of the capped prefix. Mark the detail as render-bounded, or report the byte count as a lower bound, whenever `render_truncated` is true.</comment>

<file context>
@@ -378,20 +387,90 @@ fn next_stream_offset(journal: &SqliteJournal, stream: &str) -> Result<u64> {
+                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() {
</file context>

// 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() {
Expand All @@ -400,7 +479,13 @@ fn worker_failure_detail(output: &Value) -> Option<String> {
// 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()
}
}

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Truncation count ignores render cap

Low Severity

The bytes truncated remainder uses the length of the already-capped render. When worker output exceeds MAX_RENDER_BYTES and still has more than MAX_CHARS, the journaled figure is the leftover of the cap, not how much output was actually dropped.

Additional Locations (1)
Fix in Cursor Fix in Web

Reviewed by Cursor Bugbot for commit c32351c. Configure here.

Some((cut, _)) => format!(
"{}… ({} bytes truncated)",
&trimmed[..cut],
Expand Down Expand Up @@ -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);
Expand All @@ -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<Value> = (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()
);
}
}
Loading