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
50 changes: 42 additions & 8 deletions kernel/relayflowd-core/src/machine.rs
Original file line number Diff line number Diff line change
Expand Up @@ -302,6 +302,33 @@ fn start_actions(state: &RunState, step: &StepSpec, attempt: u32, now_ms: i64) -
vec![Action::Append(started), execute]
}

/// 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 All @@ -320,14 +347,21 @@ pub fn completion_actions(
// taxonomy label and the diagnostic is gone.
let verification = match &result.failure_reason {
None => Some(verify(step, &result.output)),
Some(_) => result
.failure_detail
.as_ref()
.map(|detail| crate::entry::VerificationRecord {
gate: "execution".to_owned(),
verdict: crate::entry::VerificationVerdict::Fail,
detail: detail.clone(),
}),
// `failure_detail` is populated only for kernel-side rejections, so
// mapping over it dropped the record entirely whenever a WORKER
// reported the failure — leaving the reason in the taxonomy label
// alone, which is the outcome this branch exists to prevent. The
// record is now unconditional: a failure always names itself, and the
// fallback marks that no detail accompanied the report rather than
// implying one was given.
Some(reason) => Some(crate::entry::VerificationRecord {
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))),
}),
};
let verified = verification
.as_ref()
Expand Down
104 changes: 104 additions & 0 deletions kernel/relayflowd-core/src/machine/tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -489,3 +489,107 @@ impl AppendAction for Action {
}
}
}

/// Regression, 2026-09-06 (#195): a WORKER-reported failure carries no
/// `failure_detail` — that field is set only for kernel-side rejections — and
/// the completion used to map over it, journaling `verification: null`. The
/// reason then survived only as the taxonomy label, which is precisely what the
/// branch was written to prevent, and `output` is nulled for every non-success
/// so nothing else carried it either.
///
/// Every other row in this file sets `failure_detail: Some(..)`, so the whole
/// suite exercised the arm that worked and none of it touched the arm that did
/// not. This asserts the arm that did not.
#[test]
fn worker_reported_failure_without_detail_still_records_a_verification() {
let spec: crate::RunSpec = serde_json::from_value(json!({
"version": "0.1.0",
"steps": [
{
"id": "only",
"type": "deterministic",
"command": "false",
"max_iterations": 1
}
]
}))
.unwrap();
let result = AttemptResult {
output: Value::Null,
budget: Budget::default(),
completed_by: "worker".to_owned(),
end_pins: None,
effects: Vec::new(),
trajectory_tail: None,
failure_reason: Some(CompletionReason::WorkerError),
// The point of the case: the worker reported a failure and sent no
// detail with it.
failure_detail: None,
};
let entries: Vec<_> = completion_actions("run", &spec.steps[0], 1, 0, result, 1_000)
.into_iter()
.filter_map(|action| match action {
Action::Append(entry) => Some(entry),
_ => None,
})
.collect();

let completed = entries
.iter()
.find(|entry| entry.entry_type == EntryType::StepCompleted)
.expect("a step.completed entry");
let payload: StepCompletedPayload =
serde_json::from_value(serde_json::to_value(&completed.payload).unwrap()).unwrap();

let record = payload
.verification
.expect("a worker-reported failure must journal WHY, not just its taxonomy label");
assert_eq!(record.verdict, crate::VerificationVerdict::Fail);
assert_eq!(record.gate, "execution");
// The reason itself has to appear, or the record is present but empty of
// information and the diagnostic is still gone.
// The journal's own vocabulary, not Rust's Debug spelling: the same
// `worker_error` a reader sees in `completionReason`. Pinning the string
// keeps the two from drifting apart.
assert!(
record.detail.contains("worker_error"),
"detail should name the reported reason in journal vocabulary, got {:?}",
record.detail
);
// And the taxonomy label must still be the reason the worker gave, not
// overwritten by the verification bookkeeping.
assert_eq!(payload.completion_reason, CompletionReason::WorkerError);
}

/// `reason_label` hand-writes the journal spellings, so nothing but a test stops
/// it drifting from what `CompletionReason` actually serializes. Pins every
/// variant against serde rather than spot-checking one, so a rename in either
/// 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.)
#[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,
] {
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:?}"
);
}
}
115 changes: 113 additions & 2 deletions kernel/relayflowd/src/engine/remote.rs
Original file line number Diff line number Diff line change
Expand Up @@ -65,12 +65,31 @@ impl Engine<WallClock> {
// A rejected completion names the mistake in the journal. `output` is
// nulled for every non-success, so the detail rides the completion's
// verification record — the same channel a failed gate uses.
let mut failure_detail = None;
//
// A WORKER-reported failure gets that treatment too. `OutOfBandCompletion`
// carries no error field, so the only account of what went wrong is the
// 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();
let mut rejected_completion = false;
let mut reject = |error: anyhow::Error| {
rejected_completion = true;
failure_reason = Some(CompletionReason::WorkerError);
failure_detail = Some(format!("{error:#}"));
// Keep BOTH accounts when a worker reports its own failure and then
// trips validation. The rejection says why the kernel refused the
// completion; the worker's output says what went wrong upstream of
// that, and the two are rarely the same story. Overwriting here
// would discard the worker's account for exactly the completions
// that have the most gone wrong — the loss this whole change exists
// to stop, reintroduced one layer up.
failure_detail = Some(match failure_detail.take() {
Some(reported) => format!("rejected: {error:#}; worker reported: {reported}"),
None => format!("{error:#}"),
});
};
let effects = if matches!(step.kind, StepKind::Agent { .. }) {
let recorded = recorded_effects(&journal, &step, completion.attempt)?;
Expand Down Expand Up @@ -337,3 +356,95 @@ 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.
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;
if output.is_null() {
return None;
}
let rendered = match output {
Value::String(text) => text.clone(),
other => other.to_string(),
};
let trimmed = rendered.trim();
if trimmed.is_empty() {
return None;
}
// 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(),
Some((cut, _)) => format!("{}… ({} bytes truncated)", &trimmed[..cut], trimmed.len() - cut),
})
}

#[cfg(test)]
mod worker_failure_detail_tests {
use super::worker_failure_detail;
use serde_json::{Value, json};

#[test]
fn a_null_or_blank_output_yields_no_detail() {
// The caller's fallback ("reported X without detail") is only honest if
// this returns None rather than an empty string dressed as a diagnostic.
assert_eq!(worker_failure_detail(&Value::Null), None);
assert_eq!(worker_failure_detail(&json!("")), None);
assert_eq!(worker_failure_detail(&json!(" \n\t ")), None);
}

#[test]
fn a_string_output_is_carried_verbatim_and_trimmed() {
assert_eq!(
worker_failure_detail(&json!(" analyzer exited 1: no such model ")),
Some("analyzer exited 1: no such model".to_owned())
);
}

#[test]
fn a_non_string_output_is_rendered_rather_than_dropped() {
// A worker that reports structured failure data must not have it
// discarded just because it is not a bare string.
assert_eq!(
worker_failure_detail(&json!({"code": 2})),
Some(r#"{"code":2}"#.to_owned())
);
}

/// The truncation comment names a panic mode — byte slicing on multi-byte
/// input — and this pins it. The character matters: `€` is THREE bytes, so
/// byte index `MAX_CHARS` (2000) falls at 666 chars + 2 bytes, mid-character,
/// and `&trimmed[..MAX_CHARS]` panics on it.
///
/// An earlier version of this test used `é` and claimed the same thing. That
/// was wrong: `é` is two bytes, so byte 2000 is a valid boundary and the
/// byte-slice mutation does NOT panic there — it silently returns half the
/// intended characters. The test still failed, but on a length assertion,
/// which is a much weaker signal than the panic it advertised. A test whose
/// stated rationale is false is worse than no test, because the next reader
/// trusts it.
#[test]
fn truncation_does_not_split_a_multi_byte_char() {
// 3000 three-byte chars = 9000 bytes.
let output = json!("€".repeat(3000));
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.
assert!(
detail.contains("3000 bytes truncated"),
"expected the byte remainder, got {detail:?}"
);
assert_eq!(detail.chars().take_while(|c| *c == '€').count(), 2000);
}

#[test]
fn an_output_at_the_boundary_is_not_truncated() {
let exact = "a".repeat(2000);
assert_eq!(worker_failure_detail(&json!(exact.clone())), Some(exact));
}
}
Loading