From 38ae79b5c2c3d219977c9c27401b6817f3ef49a0 Mon Sep 17 00:00:00 2001 From: Relayflow Date: Sun, 20 Sep 2026 17:47:49 +0000 Subject: [PATCH 1/3] feat: advisory deterministic outcomes so a red check can be repaired, not hidden MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit v2 gated every deterministic step on exit code with no opt-out, so every repair-before-failure flow re-invented ` || true`. That workaround discards the exit code: a later gate has nothing to read, and a forgotten gate is indistinguishable from success, so a flow can ship red work silently. A deterministic step can now declare `onNonZero: record`. The kernel keeps the command's real exit code and output tails in the journal, passes the step so dependents run, and labels the gate `exit_code:recorded` so no reader can mistake it for a green command. `f.run(cmd, { onNonZero: 'record' })` resolves to a `RunResult` read back from that journaled envelope, so the branch below it reads the code and output the command actually produced. A later step asserts freshness with `{ type: steps_green, ids: [...] }` rather than re-running. The policy covers exit codes only. Timeouts, signals, the `-1` no-exit-status sentinel, declared content and schema gates, and budgets all stay fatal under either policy, and the default is normalized out of the canonical bytes so every existing spec keeps its exact hash. Also fixes a pre-existing hole this feature would have widened: an authored `.gate()` lowers to its own kernel step that runs after the producer, and `readCompletedStepOutput` read only the producer's entry — so a failed gate resolved the operation as though the command had passed. It now checks the run's verdict too, which is the same invisibility the recording policy exists to remove, one layer up. Co-Authored-By: Claude --- docs/SURFACE.md | 122 +++++++++ kernel/relayflowd-core/src/spec.rs | 30 ++- kernel/relayflowd-core/src/spec/tests.rs | 86 +++++++ kernel/relayflowd-core/src/verify.rs | 114 ++++++++- kernel/relayflowd-core/tests/spec_parity.rs | 18 ++ kernel/relayflowd/src/exec_det.rs | 3 + .../tests/advisory_deterministic.rs | 135 ++++++++++ packages/schema/flows.schema.json | 58 ++++- packages/sdk/src/authored-flow-executor.ts | 61 ++++- packages/sdk/src/authored-step-output.ts | 178 ++++++++++--- packages/sdk/src/authored-worker-step.ts | 29 ++- packages/sdk/src/cli.ts | 4 +- packages/sdk/src/compile.ts | 28 +- packages/sdk/src/gate-contract.ts | 14 +- packages/sdk/src/index.ts | 2 + packages/sdk/src/named-gates.ts | 3 + packages/sdk/src/spec.ts | 39 ++- packages/sdk/src/step-fields.ts | 2 +- packages/sdk/src/steps-green.ts | 165 ++++++++++++ packages/sdk/src/validate.ts | 17 +- packages/sdk/tests/advisory-outcome.test.ts | 242 ++++++++++++++++++ .../tests/authored-advisory-run-live.test.ts | 217 ++++++++++++++++ packages/sdk/tests/preflight.test.ts | 13 + packages/sdk/tests/spec-parity.test.ts | 42 +++ packages/sdk/tests/verb-field-lint.test.ts | 7 +- packages/surface/src/context.ts | 56 +++- packages/surface/src/index.ts | 2 +- .../surface/tests/run-on-non-zero.test-d.ts | 67 +++++ scripts/schema-constraints.mjs | 3 + testdata/advisory-repair.flow.yaml | 36 +++ testdata/advisory-repair.spec.canonical.json | 1 + testdata/advisory-repair.spec.sha256 | 1 + 32 files changed, 1710 insertions(+), 85 deletions(-) create mode 100644 kernel/relayflowd/tests/advisory_deterministic.rs create mode 100644 packages/sdk/src/steps-green.ts create mode 100644 packages/sdk/tests/advisory-outcome.test.ts create mode 100644 packages/sdk/tests/authored-advisory-run-live.test.ts create mode 100644 packages/surface/tests/run-on-non-zero.test-d.ts create mode 100644 testdata/advisory-repair.flow.yaml create mode 100644 testdata/advisory-repair.spec.canonical.json create mode 100644 testdata/advisory-repair.spec.sha256 diff --git a/docs/SURFACE.md b/docs/SURFACE.md index 4620b56e..b38a1590 100644 --- a/docs/SURFACE.md +++ b/docs/SURFACE.md @@ -475,6 +475,128 @@ the kernel kills the command's process group and journals `completionReason: tim `f.run` refuses with code `lease_exceeded`. The override applies only to that invocation; calls without options retain the default. +### Repair before failure: `onNonZero: 'record'` + +A deterministic step is gated on its exit code by default: a nonzero exit fails +the step and ends the flow. That is the right default, and it is the wrong one +for the single most common shape in practice — run a check, let it be red, and +hand its output to whatever is built to answer it. + +`onNonZero: 'record'` is the opt-in for that shape. The command still runs, its +exit code and output tails are still journaled, and the step still completes; +what changes is that a positive exit code satisfies the exit check instead of +failing it, so dependents run. In TypeScript the result is a value rather than +a throw: + +```ts +const tests = await f.run('npm test', { onNonZero: 'record' }); +if (!tests.ok) { + await f.agent('fixer', { task: `Fix these failures:\n${tests.output}` }); +} +``` + +`RunResult` is `{ ok, exitCode, stdout, stderr, output }`. Every field is read +back from the step's journaled outcome, never re-measured — `ok` is exactly +`exitCode === 0`. `stdout` and `stderr` are the journal's **bounded tails**, not +complete transcripts; `output` is `stdout` then `stderr` joined by a newline +when both are nonempty, a reading convenience rather than a reconstruction of +the real interleaving. Without the policy, `f.run` still resolves to the stdout +string, so existing bodies are unchanged. + +This is what ` || true` cannot do. `|| true` throws the exit code +away, so nothing downstream can tell a passing check from a failing one, and a +flow that forgets a later gate ships red work silently. `onNonZero: 'record'` +keeps the code, which is the whole point: the branch below it is a fact, not a +guess. + +**The policy covers exit codes only.** A timeout, a signal or a command that +never started produced no verdict to record, and still fails the step under +either policy. Declared `output_contains` and `json_schema` gates are still +enforced in record mode, as are budgets and journal-append failures. Recording +is not a retry policy and does not trigger semantic retries. + +**Recording makes red allowed, not invisible.** `flows check` annotates a +recording step with `[onNonZero: record]`, and the kernel labels the satisfied +check `exit_code:recorded` with the numeric code in its detail, so a recorded +red outcome reads as red in the journal. A flow that records without ever +asserting green may still finish successfully — that outcome is inspectable, +not forbidden. + +**Asserting green again.** To demand that recorded commands were green, use the +declarative `steps_green` gate, which reads the journaled outcomes and never +re-runs anything: + +```yaml +version: '0.1.0' +steps: + - id: tests + type: deterministic + command: npm test + onNonZero: record + - id: lint + type: deterministic + command: npm run lint + onNonZero: record + - id: gate + type: deterministic + command: 'true' + verification: + type: steps_green + ids: [tests, lint] +``` + +`steps_green` is deterministic-hosts-only: a worker step records no exit code. +The ids it names must be deterministic steps that precede it; unknown, forward, +self and non-deterministic references are refused at compile time, as are empty +and duplicate id lists. The host keeps its own command and its own exit policy — +the assertion is lowered to a separate, always-fatal gate step that every +dependent of the host waits on. + +A recorded outcome is immutable, so repair does not turn an old red step green. +A flow that repairs must produce **new** post-repair evidence and gate on that +distinct step: + +```yaml +version: '0.1.0' +steps: + - id: tests + type: deterministic + command: npm test + onNonZero: record + # Declaring the envelope schema is what lets a later step bind the evidence. + verification: + type: json_schema + schema: + type: object + properties: + exit_code: { type: integer } + stdout_tail: { type: string } + stderr_tail: { type: string } + - id: repair + type: agent + dependsOn: [tests] + input: + failures: { step: tests } + instruction: Read input.failures and fix the failing tests. + - id: retest + type: deterministic + dependsOn: [repair] + command: npm test + - id: gate + type: deterministic + command: 'true' + verification: + type: steps_green + ids: [retest] +``` + +`retest` is an ordinary gated step, so the final `steps_green` over it is what +decides the run. Binding a recorded outcome as input requires the source to +declare an output schema, as every input binding does (see *Declarative output +binding*); an `output_contains` source is refused, because that gate cannot also +carry the envelope schema the binding reads — declare the constraint as a JSON +Schema instead. + ### The authored operation lifecycle An authored TypeScript body reaches `done()` only if every step it created was diff --git a/kernel/relayflowd-core/src/spec.rs b/kernel/relayflowd-core/src/spec.rs index bc5f4dae..ff4f4dab 100644 --- a/kernel/relayflowd-core/src/spec.rs +++ b/kernel/relayflowd-core/src/spec.rs @@ -307,7 +307,7 @@ const STEP_COMMON_FIELDS: &[&str] = &[ "memory", "requirements", ]; -const STEP_DETERMINISTIC_FIELDS: &[&str] = &["command", "timeout_ms", "lease_ms"]; +const STEP_DETERMINISTIC_FIELDS: &[&str] = &["command", "timeout_ms", "lease_ms", "on_non_zero"]; const STEP_LLM_FIELDS: &[&str] = &["prompt", "model", "cli"]; const STEP_AGENT_FIELDS: &[&str] = &[ "instruction", @@ -397,6 +397,29 @@ impl StepSpec { } } +/// What a deterministic step's nonzero exit code means. +/// +/// `Fail` is the kernel's long-standing implicit gate: `exit_code == 0` or the +/// step failed. `Record` is the repair-before-failure shape — the command is +/// allowed to be red, its exit code and output tails stay in the journal +/// exactly as executed, dependents run, and a later step reads the recorded +/// outcome instead of re-running the command. It is a policy for the step's +/// own exit code only; timeouts, worker errors, and declared content or schema +/// gates keep their existing fatal semantics. +#[derive(Debug, Clone, Copy, Default, Serialize, Deserialize, PartialEq, Eq)] +#[serde(rename_all = "snake_case")] +pub enum OnNonZero { + #[default] + Fail, + Record, +} + +impl OnNonZero { + fn is_default(&self) -> bool { + matches!(self, OnNonZero::Fail) + } +} + #[derive(Debug, Clone, Serialize, Deserialize, PartialEq)] #[serde(tag = "type", rename_all = "snake_case")] pub enum StepKind { @@ -406,6 +429,11 @@ pub enum StepKind { timeout_ms: Option, #[serde(default, skip_serializing_if = "Option::is_none")] lease_ms: Option, + /// What a nonzero command exit means. Omitted on serialization when + /// it is the default, so every spec written before this field exists + /// keeps its exact canonical bytes and its exact hash. + #[serde(default, skip_serializing_if = "OnNonZero::is_default")] + on_non_zero: OnNonZero, }, Llm { prompt: String, diff --git a/kernel/relayflowd-core/src/spec/tests.rs b/kernel/relayflowd-core/src/spec/tests.rs index b8443bff..9b1fd50f 100644 --- a/kernel/relayflowd-core/src/spec/tests.rs +++ b/kernel/relayflowd-core/src/spec/tests.rs @@ -329,3 +329,89 @@ fn preflight_data_is_fail_closed() { Err(SpecError::EmptyStepCli("a".to_owned())) ); } + +/// `on_non_zero` is a declared policy, so it parses fail-closed like every +/// other spec field: an unrecognized value is a refusal, never a silent +/// fallback to the fatal default. +#[test] +fn on_non_zero_parses_the_two_declared_policies_and_refuses_anything_else() { + let recorded = RunSpec::parse(&json!({ + "steps": [{ + "id": "tests", "type": "deterministic", "command": "npm test", + "on_non_zero": "record" + }] + })) + .unwrap(); + assert!(matches!( + recorded.steps[0].kind, + StepKind::Deterministic { + on_non_zero: OnNonZero::Record, + .. + } + )); + assert_eq!(recorded.validate(), Ok(())); + + let explicit_fail = RunSpec::parse(&json!({ + "steps": [{ + "id": "tests", "type": "deterministic", "command": "npm test", + "on_non_zero": "fail" + }] + })) + .unwrap(); + assert!(matches!( + explicit_fail.steps[0].kind, + StepKind::Deterministic { + on_non_zero: OnNonZero::Fail, + .. + } + )); + + assert!(matches!( + RunSpec::parse(&json!({ + "steps": [{ + "id": "tests", "type": "deterministic", "command": "npm test", + "on_non_zero": "ignore" + }] + })), + Err(SpecError::Malformed(_)) + )); + + // The policy belongs to the deterministic rung only; a worker step that + // declares it is an unknown field, not an inert decoration. + assert!(matches!( + RunSpec::parse(&json!({ + "steps": [{ + "id": "ask", "type": "llm", "prompt": "p", "on_non_zero": "record" + }] + })), + Err(SpecError::UnknownField { field, .. }) if field == "on_non_zero" + )); +} + +/// The boundary artifact is hashed. A step that does not declare the policy +/// must serialize to the exact bytes it did before the field existed, or every +/// committed canonical fixture and every memoized step hash moves at once. +#[test] +fn the_default_policy_is_absent_from_the_serialized_boundary_spec() { + let spec = RunSpec::parse(&json!({ + "steps": [{"id": "a", "type": "deterministic", "command": "true"}] + })) + .unwrap(); + let serialized = serde_json::to_value(&spec).unwrap(); + assert!( + serialized["steps"][0].get("on_non_zero").is_none(), + "{serialized}" + ); + + let recorded = RunSpec::parse(&json!({ + "steps": [{ + "id": "a", "type": "deterministic", "command": "true", + "on_non_zero": "record" + }] + })) + .unwrap(); + assert_eq!( + serde_json::to_value(&recorded).unwrap()["steps"][0]["on_non_zero"], + json!("record") + ); +} diff --git a/kernel/relayflowd-core/src/verify.rs b/kernel/relayflowd-core/src/verify.rs index 7f6700af..904fb662 100644 --- a/kernel/relayflowd-core/src/verify.rs +++ b/kernel/relayflowd-core/src/verify.rs @@ -2,19 +2,36 @@ use serde_json::Value; use crate::{ entry::{VerificationRecord, VerificationVerdict}, - spec::{StepKind, StepSpec}, + spec::{OnNonZero, StepKind, StepSpec}, }; pub fn verify(step: &StepSpec, output: &Value) -> VerificationRecord { let mut gates = Vec::new(); let mut failures = Vec::new(); + // Explanations that are not failures. A recorded red exit is a satisfied + // policy, not a green command, and a reader is owed that distinction even + // when a separate content or schema gate fails alongside it. + let mut notes = Vec::new(); - if matches!(step.kind, StepKind::Deterministic { .. }) { - gates.push("exit_code"); + if let StepKind::Deterministic { on_non_zero, .. } = &step.kind { match output.get("exit_code").and_then(Value::as_i64) { - Some(0) => {} - Some(code) => failures.push(format!("exit code was {code}")), - None => failures.push("output omitted integer exit_code".to_owned()), + Some(0) => gates.push("exit_code"), + // Only an ordinary POSITIVE code is the command's own verdict. The + // executor writes -1 when the process produced no exit status at + // all (killed by a signal, or never spawned); `record` must not + // quietly absorb that, so it falls through to the fatal branch. + Some(code) if code > 0 && *on_non_zero == OnNonZero::Record => { + gates.push("exit_code:recorded"); + notes.push(format!("recorded exit code {code} (on_non_zero: record)")); + } + Some(code) => { + gates.push("exit_code"); + failures.push(format!("exit code was {code}")); + } + None => { + gates.push("exit_code"); + failures.push("output omitted integer exit_code".to_owned()); + } } } @@ -53,10 +70,16 @@ pub fn verify(step: &StepSpec, output: &Value) -> VerificationRecord { } else { VerificationVerdict::Fail }, - detail: if passed { - "all gates passed".to_owned() - } else { - failures.join("; ") + detail: match (passed, notes.is_empty()) { + (true, true) => "all gates passed".to_owned(), + (true, false) => notes.join("; "), + // Notes first: why the step is complete-but-red is the context the + // failures below are read in, and it must not be dropped. + (false, _) => notes + .into_iter() + .chain(failures) + .collect::>() + .join("; "), }, } } @@ -122,6 +145,77 @@ mod tests { ); } + fn recording_step(verification: Value) -> StepSpec { + serde_json::from_value(json!({ + "id": "tests", + "type": "deterministic", + "command": "npm test", + "on_non_zero": "record", + "verification": verification + })) + .unwrap() + } + + /// Repair-before-failure: the command is allowed to be red. The verdict + /// passes so dependents run, and the gate label and detail keep the code + /// so a reader can never mistake this for a green command. + #[test] + fn a_recorded_nonzero_exit_passes_its_gate_and_names_the_code() { + let step = recording_step(json!({})); + let record = verify(&step, &json!({"exit_code": 7, "stdout_tail": "2 failing"})); + assert_eq!(record.verdict, VerificationVerdict::Pass); + assert_eq!(record.gate, "exit_code:recorded"); + assert!(record.detail.contains('7'), "{}", record.detail); + + // A green command under the same policy is still an ordinary pass. + let green = verify(&step, &json!({"exit_code": 0, "stdout_tail": "ok"})); + assert_eq!(green.verdict, VerificationVerdict::Pass); + assert_eq!(green.gate, "exit_code"); + assert_eq!(green.detail, "all gates passed"); + } + + /// `record` is a policy for the exit code alone. A declared content or + /// schema gate still fails the step — and the recorded code survives into + /// the detail beside the failure rather than being displaced by it. + #[test] + fn recording_an_exit_code_does_not_relax_a_declared_content_gate() { + let step = recording_step(json!({"output_contains": "PASS"})); + let record = verify(&step, &json!({"exit_code": 7, "stdout_tail": "2 failing"})); + assert_eq!(record.verdict, VerificationVerdict::Fail); + assert_eq!(record.gate, "exit_code:recorded+output_contains"); + assert!(record.detail.contains('7'), "{}", record.detail); + assert!(record.detail.contains("PASS"), "{}", record.detail); + } + + /// A step that never produced an exit status has no verdict to record. + /// `-1` is the executor's sentinel for "killed, or never ran"; absorbing + /// it would let a timed-out command read as a recorded red test run. + #[test] + fn recording_does_not_absorb_a_missing_or_sentinel_exit_status() { + let step = recording_step(json!({})); + for output in [ + json!({"exit_code": -1, "stdout_tail": ""}), + json!({"stdout_tail": ""}), + json!({"exit_code": "7"}), + ] { + let record = verify(&step, &output); + assert_eq!(record.verdict, VerificationVerdict::Fail, "{output}"); + assert_eq!(record.gate, "exit_code", "{output}"); + } + } + + /// Without the declaration nothing moves: the default is still fatal. + #[test] + fn the_default_policy_still_fails_a_nonzero_exit() { + let step: StepSpec = serde_json::from_value(json!({ + "id": "tests", "type": "deterministic", "command": "npm test" + })) + .unwrap(); + let record = verify(&step, &json!({"exit_code": 7, "stdout_tail": ""})); + assert_eq!(record.verdict, VerificationVerdict::Fail); + assert_eq!(record.gate, "exit_code"); + } + #[test] fn json_schema_is_a_control_gate() { let step: StepSpec = serde_json::from_value(json!({ diff --git a/kernel/relayflowd-core/tests/spec_parity.rs b/kernel/relayflowd-core/tests/spec_parity.rs index 7e327879..fd339a0a 100644 --- a/kernel/relayflowd-core/tests/spec_parity.rs +++ b/kernel/relayflowd-core/tests/spec_parity.rs @@ -116,6 +116,24 @@ fn the_kernel_parses_the_event_triggered_spec_and_stamps_the_same_hash() { ); } +/// The advisory-outcome dialect: a `on_non_zero: "record"` step, a binding that +/// reads its recorded envelope, and the deterministic gate step the SDK lowers +/// `steps_green` into. The SDK spelling of that gate never reaches the kernel, +/// so this proves the kernel parses what the SDK actually emits for it. +#[test] +fn the_kernel_parses_the_advisory_repair_spec_and_stamps_the_same_hash() { + assert_parity( + include_str!(concat!( + env!("CARGO_MANIFEST_DIR"), + "/../../testdata/advisory-repair.spec.canonical.json" + )), + include_str!(concat!( + env!("CARGO_MANIFEST_DIR"), + "/../../testdata/advisory-repair.spec.sha256" + )), + ); +} + fn assert_parity(canonical_fixture: &str, expected_hash: &str) { let value: Value = serde_json::from_str(canonical_fixture.trim()).unwrap(); let spec = RunSpec::parse(&value).expect("kernel must parse the SDK's compiled spec"); diff --git a/kernel/relayflowd/src/exec_det.rs b/kernel/relayflowd/src/exec_det.rs index 3a6ef16b..e588e929 100644 --- a/kernel/relayflowd/src/exec_det.rs +++ b/kernel/relayflowd/src/exec_det.rs @@ -39,10 +39,13 @@ pub(crate) fn execute_placed_with_input( if step.input.is_some() && input.is_none() { return worker_error("deterministic step input bindings were not resolved"); } + // `on_non_zero` is read by `verify`, not here: the executor reports what + // the process did, and the gate decides what that means. let StepKind::Deterministic { command, timeout_ms, lease_ms, + .. } = &step.kind else { return worker_error("deterministic executor received a non-deterministic step"); diff --git a/kernel/relayflowd/tests/advisory_deterministic.rs b/kernel/relayflowd/tests/advisory_deterministic.rs new file mode 100644 index 00000000..0586e4bd --- /dev/null +++ b/kernel/relayflowd/tests/advisory_deterministic.rs @@ -0,0 +1,135 @@ +//! Repair before failure, end to end (#500 follow-up). +//! +//! `on_non_zero: "record"` is the shape `|| true` was standing in for. The +//! difference this file exists to pin is that the exit code SURVIVES: the +//! journal holds `{exit_code, stdout_tail, stderr_tail}` exactly as the +//! command produced them, a dependent binds that evidence as input instead of +//! re-running the command, and a reader can still tell red from green. + +use std::{fs, process::Command}; + +use relayflowd_core::{EntryType, RunSpec}; +use relayflowd_journal::SqliteJournal; +use serde_json::{Value, json}; + +/// A red probe whose recorded outcome is handed to a repair step. +fn spec(artifact: &std::path::Path, probe_runs: &std::path::Path, record: bool) -> Value { + let mut probe = json!({ + "id": "probe", + "type": "deterministic", + "command": format!( + "printf x >> '{}'; printf '2 failing'; exit 7", + probe_runs.display() + ), + "verification": {"json_schema": { + "type": "object", + "properties": {"exit_code": {"type": "integer"}, "stdout_tail": {"type": "string"}} + }} + }); + if record { + probe["on_non_zero"] = json!("record"); + } + json!({"steps": [ + probe, + { + "id": "repair", + "type": "deterministic", + "depends_on": ["probe"], + "input": {"failure": {"step": "probe"}}, + "command": format!("printf '%s' \"$FLOWS_INPUT\" > '{}'", artifact.display()) + } + ]}) +} + +fn run(data: &std::path::Path, spec: &Value) -> std::process::Output { + let path = data.parent().unwrap().join("spec.json"); + fs::write(&path, serde_json::to_vec(spec).unwrap()).unwrap(); + Command::new(env!("CARGO_BIN_EXE_relayflowd")) + .args([ + "--data-dir", + data.to_str().unwrap(), + "run", + path.to_str().unwrap(), + ]) + .output() + .unwrap() +} + +fn completion(data: &std::path::Path, step: &str) -> Value { + let journal = fs::read_dir(data.join("runs")) + .unwrap() + .flatten() + .map(|entry| entry.path()) + .find(|path| path.extension().is_some_and(|extension| extension == "sqlite3")) + .expect("the run wrote a journal"); + SqliteJournal::open(&journal) + .unwrap() + .scan_all() + .unwrap() + .into_iter() + .find(|entry| { + entry.entry_type == EntryType::StepCompleted && entry.step_id.as_deref() == Some(step) + }) + .map(|entry| entry.payload) + .unwrap_or(Value::Null) +} + +#[test] +fn a_recorded_red_step_completes_and_hands_its_exit_code_to_the_next_step() { + let root = tempfile::tempdir().unwrap(); + let data = root.path().join("data"); + let artifact = root.path().join("input.json"); + let probe_runs = root.path().join("probe-runs"); + let spec = spec(&artifact, &probe_runs, true); + RunSpec::parse(&spec).unwrap().validate().unwrap(); + + let result = run(&data, &spec); + assert!( + result.status.success(), + "a recorded red step must not fail the run: {}", + String::from_utf8_lossy(&result.stderr) + ); + + // The evidence reached the dependent as journal data, not as a re-run. + let input: Value = serde_json::from_slice(&fs::read(&artifact).unwrap()).unwrap(); + assert_eq!(input["failure"]["exit_code"], json!(7)); + assert_eq!(input["failure"]["stdout_tail"], json!("2 failing")); + assert_eq!(fs::read_to_string(&probe_runs).unwrap(), "x"); + + // Complete-but-red, and legible as such: the journaled verdict passes so + // dependents run, while the gate name and detail keep the exit code. + let probe = completion(&data, "probe"); + assert_eq!(probe["completionReason"], json!("success")); + assert_eq!(probe["output"]["exit_code"], json!(7)); + assert_eq!( + probe["verification"]["gate"], + json!("exit_code:recorded+json_schema") + ); + assert_eq!(probe["verification"]["verdict"], json!("pass")); + assert!( + probe["verification"]["detail"] + .as_str() + .is_some_and(|detail| detail.contains('7')), + "{}", + probe["verification"] + ); +} + +/// The identical flow without the declaration. `record` is opt-in, so an +/// author who does not ask for it keeps the gate that stops a red run. +#[test] +fn the_same_flow_without_the_declaration_still_fails_and_never_reaches_the_dependent() { + let root = tempfile::tempdir().unwrap(); + let data = root.path().join("data"); + let artifact = root.path().join("input.json"); + let probe_runs = root.path().join("probe-runs"); + + let result = run(&data, &spec(&artifact, &probe_runs, false)); + assert!(!result.status.success()); + assert!(!artifact.exists()); + assert_eq!(completion(&data, "repair"), Value::Null); + assert_eq!( + completion(&data, "probe")["verification"]["verdict"], + json!("fail") + ); +} diff --git a/packages/schema/flows.schema.json b/packages/schema/flows.schema.json index ff0742df..89c258a6 100644 --- a/packages/schema/flows.schema.json +++ b/packages/schema/flows.schema.json @@ -15,6 +15,15 @@ "agent" ] }, + "OnNonZero": { + "title": "OnNonZero", + "description": "A deterministic step's exit-code policy. `'fail'` is the kernel's implicit\ngate; `'record'` completes the step with its exit code journaled so a later\nstep can read the outcome instead of re-running the command.", + "type": "string", + "enum": [ + "fail", + "record" + ] + }, "McpServerConfig": { "title": "McpServerConfig", "description": "Project-owned MCP connections. env contains names, never secret values.", @@ -148,6 +157,7 @@ "type": "string", "enum": [ "exit_code", + "steps_green", "output_contains", "json_schema", "references_input", @@ -414,6 +424,37 @@ ], "additionalProperties": false }, + "StepsGreenGate": { + "title": "StepsGreenGate", + "description": "Asserts that every named step recorded a zero exit code, reading the\njournaled outcome instead of re-running the command. This is the read side\nof `onNonZero: 'record'`: red work is allowed to exist, and this gate is\nwhere a flow declares that it stops being allowed.\n\nDeterministic hosts only — a worker step has no recorded exit code to read —\nand the ids it names must be deterministic steps that precede it.", + "type": "object", + "properties": { + "type": { + "title": "type", + "description": "type in the Relayflows spec.", + "type": "string", + "const": "steps_green" + }, + "ids": { + "title": "ids", + "description": "Earlier deterministic step ids whose recorded exit codes must all be zero.", + "type": "array", + "items": { + "type": "string", + "pattern": "\\S", + "title": "ids items", + "description": "ids items value." + }, + "minItems": 1, + "uniqueItems": true + } + }, + "required": [ + "type", + "ids" + ], + "additionalProperties": false + }, "NamedDataGate": { "title": "NamedDataGate", "description": "NamedDataGate in the Relayflows spec.", @@ -476,8 +517,13 @@ "description": "See ExitCodeGate." }, { - "$ref": "#/$defs/OutputVerificationSpec", + "$ref": "#/$defs/StepsGreenGate", "title": "VerificationSpec alternative 2", + "description": "See StepsGreenGate." + }, + { + "$ref": "#/$defs/OutputVerificationSpec", + "title": "VerificationSpec alternative 3", "description": "See OutputVerificationSpec." } ] @@ -1099,6 +1145,11 @@ "description": "Per-invocation deterministic lease in milliseconds (maximum 15 minutes).", "type": "number" }, + "onNonZero": { + "title": "onNonZero", + "description": "What a nonzero exit means. `'fail'` (the default) is the implicit\n`exit_code == 0` gate. `'record'` journals the exit code and output tails\nas the step's outcome, completes the step red, and lets dependents run —\nthe repair-before-failure shape, without discarding the exit code the way\n` || true` does. Declared content and schema gates, timeouts and\nworker errors keep their existing fatal semantics.", + "$ref": "#/$defs/OnNonZero" + }, "verification": { "title": "verification", "description": "Omit to get the implicit `exit_code` gate.", @@ -2087,6 +2138,11 @@ "title": "lease_ms", "description": "lease_ms in the Relayflows spec.", "type": "number" + }, + "on_non_zero": { + "title": "on_non_zero", + "description": "Omitted for the default `'fail'`, exactly as the kernel skips serializing\nit: a step that does not declare the policy keeps the canonical bytes —\nand therefore the step hash — it had before this field existed.", + "$ref": "#/$defs/OnNonZero" } }, "required": [ diff --git a/packages/sdk/src/authored-flow-executor.ts b/packages/sdk/src/authored-flow-executor.ts index 965ad117..b004d428 100644 --- a/packages/sdk/src/authored-flow-executor.ts +++ b/packages/sdk/src/authored-flow-executor.ts @@ -21,6 +21,8 @@ import { type Ctx, type FlowCompletionReason, type RunCompletionReason as SurfaceRunCompletionReason, + type RunOptions, + type RunResult, type Step, } from '@relayflows/surface'; import { createHelpers, helperProviders, type HelperCall, type FlowHandle } from '@relayflows/surface/runtime'; @@ -332,6 +334,50 @@ export async function executeAuthoredFlow( return trackStep(authoredSteps, llmOp); } + /** + * `f.run` (docs/SURFACE.md §1). Declared as overloads like `llmOperation`, + * because the exit-code policy decides the RESULT TYPE: the default keeps + * the stdout string every body is written against, and `'record'` widens it + * to the `RunResult` whose `ok` is the journaled exit code. + */ + function runOperation(command: string, options?: RunOptions & { onNonZero?: 'fail' }): Step; + function runOperation(command: string, options: RunOptions & { onNonZero: 'record' }): Step; + function runOperation(command: string, options: RunOptions): Step; + function runOperation(command: string, runOptions?: RunOptions): Step | Step { + assertOperationAllowed('run', definition.name, requestedCompletion); + const leaseMs = runOptions?.timeout === undefined ? undefined : parseStepTimeout(runOptions.timeout); + // Refused before an ordinal is consumed or anything is journaled. A + // misspelled policy must not quietly fall back to the fatal default and + // take away the very branch the author wrote below this line. + const policy = runOptions?.onNonZero; + if (policy !== undefined && policy !== 'fail' && policy !== 'record') { + throw new AuthoredFlowExecutionError( + 'unsupported_verb', + `f.run options.onNonZero must be 'fail' or 'record' (got ${JSON.stringify(policy)}).`, + ); + } + const id = `run-${nextStep++}`; + if (policy === 'record') { + let recordOp!: AuthoredFlowOperation; + recordOp = new AuthoredFlowOperation( + id, 'run', + () => assertOperationAllowed('run', definition.name, requestedCompletion), + () => observeStep(id, 'deterministic', + () => lowerDeterministic.recording(id, command, leaseMs, recordOp.namedGate), options.onProgress), + lifecycle, + ); + return trackStep(authoredSteps, recordOp); + } + let runOp!: AuthoredFlowOperation; + runOp = new AuthoredFlowOperation( + id, 'run', + () => assertOperationAllowed('run', definition.name, requestedCompletion), + () => observeStep(id, 'deterministic', () => lowerDeterministic(id, command, false, leaseMs, runOp.namedGate), options.onProgress), + lifecycle, + ); + return trackStep(authoredSteps, runOp); + } + const slackRun = options.rootRunId ?? randomUUID(); function slackOperation(call: SlackCall): Step { assertOperationAllowed(`slack.${call.verb}`, definition.name, requestedCompletion); @@ -383,20 +429,7 @@ export async function executeAuthoredFlow( () => assertOperationAllowed('memory', definition.name, requestedCompletion), definition.header.memory?.script !== false, ), - run(command, runOptions) { - assertOperationAllowed('run', definition.name, requestedCompletion); - const leaseMs = runOptions?.timeout === undefined ? undefined : parseStepTimeout(runOptions.timeout); - const id = `run-${nextStep++}`; - let runOp!: AuthoredFlowOperation; - runOp = new AuthoredFlowOperation( - id, - 'run', - () => assertOperationAllowed('run', definition.name, requestedCompletion), - () => observeStep(id, 'deterministic', () => lowerDeterministic(id, command, false, leaseMs, runOp.namedGate), options.onProgress), - lifecycle, - ); - return trackStep(authoredSteps, runOp); - }, + run: runOperation, llm: llmOperation, agent(name, options) { assertOperationAllowed('agent', definition.name, requestedCompletion); diff --git a/packages/sdk/src/authored-step-output.ts b/packages/sdk/src/authored-step-output.ts index c0797527..b7c86d98 100644 --- a/packages/sdk/src/authored-step-output.ts +++ b/packages/sdk/src/authored-step-output.ts @@ -20,19 +20,23 @@ export interface AuthoredStepContext { } /** - * Shared by every `f.*` verb that lowers to one kernel step run in isolation: - * find its `step.completed` entry, record it, and refuse anything but a + * Shared by every `f.*` verb that lowers to its own kernel run: find the + * operation's `step.completed` entry, record it, and refuse anything but a * clean success before handing the raw `output` back for verb-specific * extraction (a plain string for `f.run`, an `AgentResult` for `f.agent`). * - * Callers are responsible for having already established that the RUN - * reached a terminal, successful state before calling this — `f.run`'s - * caller relies on `runStart`'s own immediate response (the kernel drives a - * deterministic step to completion inline, no race); `f.agent`'s caller - * relies on `classifyOutcome` (cli/run.ts) having already polled to a true - * terminal state. Given that, a single read from the start of this run's - * (small, single-step) journal is enough — no polling here, and no run - * outcome ever needs re-checking. + * Callers are responsible for having already waited for the run to reach a + * terminal state — `f.run`'s caller relies on `runStart`'s own immediate + * response (the kernel drives a deterministic step to completion inline, no + * race); `f.agent`'s caller relies on `classifyOutcome` (cli/run.ts) having + * already polled to one. Given that, a single read from the start of this + * run's small journal is enough, and no polling happens here. + * + * What is NOT delegated to the caller is the run's verdict. A child spec has + * more than one step whenever the author declared a `.gate()`, and that gate + * step runs after the producer, so the producer's own `success` is not the + * run's. Both are checked here, and the run-level check is what makes a failed + * gate reject the authored operation. * * A non-success completion used to throw `journal step "run-1" completed with * retries_exhausted` and stop there, discarding the evidence it was holding: @@ -60,25 +64,19 @@ export async function readCompletedStepOutput( const reason = completed.payload.completionReason; journalSteps.push(Object.freeze({ id: stepId, runId, completionReason: reason })); if (reason !== 'success') { - let message = `journal step "${stepId}" completed with ${reason}`; - let details: StepFailedDetails | undefined; - try { - details = await stepFailureDetails(journal, runId); - } catch (error) { - // Inspection must not erase the already known step failure; the two - // failures stay separately visible, as in classifyOutcome. - message += ` Could not inspect the failed step: ${errorMessage(error)}`; - } - if (details !== undefined) message += renderStepEvidence(details); - // Appended whatever inspection found — including nothing. A failure shape - // this reader does not recognise must still end with somewhere to go. - const where = inspectionHint(runId, details?.stepId ?? stepId, context.dataDir); - message += renderInspection(where); - message += await alsoRecord(journal, context.rootRunId, { - step: stepId, runId, state: 'completed', completionReason: reason, - ...(details?.stepId === undefined ? {} : { kernelStep: details.stepId }), - }); - throw new AuthoredFlowExecutionError('step_failed', message, reason, runId, { ...details, ...where }); + throw await stepFailure(journal, runId, stepId, reason, context, + `journal step "${stepId}" completed with ${reason}`); + } + // An authored operation does not always lower to ONE kernel step. A declared + // `.gate()` becomes its own barrier step that runs AFTER the producer + // (`lowerNamedGates`), so the producer's own success is not the run's + // verdict. Reading only the producer's entry let a failed gate resolve as + // though the command had passed — the `|| true` invisibility this module's + // recording policy exists to remove, reappearing one layer up. + const runReason = entries.map(runCompletionReason).find((value) => value !== undefined); + if (runReason !== undefined && runReason !== 'success') { + throw await stepFailure(journal, runId, stepId, runReason, context, + `journal run for step "${stepId}" completed with ${runReason}`); } await recordAuthoredChild(journal, context.rootRunId, { step: stepId, runId, state: 'completed', completionReason: reason, @@ -86,6 +84,40 @@ export async function readCompletedStepOutput( return completed.payload.output; } +/** + * The one `step_failed` grammar both the step-level and the run-level refusal + * report through: inspect the journal for the step that actually failed, render + * its evidence, name where to look, and index the child before throwing. + */ +async function stepFailure( + journal: JournalClient, + runId: string, + stepId: string, + reason: ProtocolCompletionReason | ProtocolRunCompletionReason, + context: AuthoredStepContext, + header: string, +): Promise { + let message = header; + let details: StepFailedDetails | undefined; + try { + details = await stepFailureDetails(journal, runId); + } catch (error) { + // Inspection must not erase the already known step failure; the two + // failures stay separately visible, as in classifyOutcome. + message += ` Could not inspect the failed step: ${errorMessage(error)}`; + } + if (details !== undefined) message += renderStepEvidence(details); + // Appended whatever inspection found — including nothing. A failure shape + // this reader does not recognise must still end with somewhere to go. + const where = inspectionHint(runId, details?.stepId ?? stepId, context.dataDir); + message += renderInspection(where); + message += await alsoRecord(journal, context.rootRunId, { + step: stepId, runId, state: 'completed', completionReason: reason, + ...(details?.stepId === undefined ? {} : { kernelStep: details.stepId }), + }); + return new AuthoredFlowExecutionError('step_failed', message, reason, runId, { ...details, ...where }); +} + export async function readSuccessfulOutput( journal: JournalClient, outcome: RunOutcome, @@ -94,24 +126,74 @@ export async function readSuccessfulOutput( context: AuthoredStepContext = {}, ): Promise { const output = await readCompletedStepOutput(journal, outcome.run_id, stepId, journalSteps, context) - .catch(error => { - if (error instanceof AuthoredFlowExecutionError - && (error.completionReason === 'timeout' || error.completionReason === 'lease_expired')) { - // A timeout is still a failed command, and the evidence extracted for - // it is the same evidence. Re-labelling the error must not delete it. - throw new AuthoredFlowExecutionError('lease_exceeded', - `f.run step "${stepId}" exceeded its command timeout.` - + ` ${error.message.slice(`${error.code}: `.length)}`, - error.completionReason, error.runId, error.details); - } - throw error; - }); + .catch(relabelCommandTimeout(stepId)); if (!isRecord(output) || typeof output['stdout_tail'] !== 'string') { throw protocolViolation(outcome.run_id, `step "${stepId}" has no string stdout_tail`); } return output['stdout_tail']; } +/** What `f.run` resolves to under `onNonZero: 'record'` (surface `RunResult`). */ +export interface RecordedRunOutcome { + ok: boolean; + exitCode: number; + output: string; + stdout: string; + stderr: string; +} + +/** + * The read side of `onNonZero: 'record'`. + * + * A recorded red exit is a SUCCESSFUL step, so this takes the same path as + * `readSuccessfulOutput` and every field comes from the journaled envelope — + * nothing is re-run and nothing is re-measured. What still throws is what the + * policy never covered: a timeout, a signal, a step that never ran, or a + * declared content or schema gate that failed. The `-1` the executor writes + * for "no exit status" is refused rather than handed back as a plausible + * code, mirroring the kernel's own rule, so a malformed envelope can never + * reach an author's `if (!result.ok)` as a real verdict. + */ +export async function readRecordedOutcome( + journal: JournalClient, + outcome: RunOutcome, + stepId: string, + journalSteps: AuthoredFlowJournalStep[], + context: AuthoredStepContext = {}, +): Promise { + const envelope = await readCompletedStepOutput(journal, outcome.run_id, stepId, journalSteps, context) + .catch(relabelCommandTimeout(stepId)); + const exitCode = isRecord(envelope) ? envelope['exit_code'] : undefined; + if (!isRecord(envelope) || typeof envelope['stdout_tail'] !== 'string' + || typeof envelope['stderr_tail'] !== 'string' + || typeof exitCode !== 'number' || !Number.isInteger(exitCode) || exitCode < 0) { + throw protocolViolation(outcome.run_id, + `step "${stepId}" recorded no {exit_code, stdout_tail, stderr_tail} to report`); + } + const stdout = envelope['stdout_tail']; + const stderr = envelope['stderr_tail']; + return { + ok: exitCode === 0, exitCode, stdout, stderr, + // Concatenated, not interleaved: the two tails were captured separately, + // so joining them is a reading convenience and is documented as one. + output: stdout !== '' && stderr !== '' ? `${stdout}\n${stderr}` : stdout + stderr, + }; +} + +/** A timeout is still a failed command; re-labelling it must not delete its evidence. */ +function relabelCommandTimeout(stepId: string): (error: unknown) => never { + return (error) => { + if (error instanceof AuthoredFlowExecutionError + && (error.completionReason === 'timeout' || error.completionReason === 'lease_expired')) { + throw new AuthoredFlowExecutionError('lease_exceeded', + `f.run step "${stepId}" exceeded its command timeout.` + + ` ${error.message.slice(`${error.code}: `.length)}`, + error.completionReason, error.runId, error.details); + } + throw error; + }; +} + interface StepCompletedEntry { entry_type: 'step.completed'; step_id: string; @@ -131,6 +213,20 @@ function isStepCompleted(value: unknown, stepId: string): value is StepCompleted && 'output' in payload; } +/** + * The run's own terminal reason, when this entry is the one carrying it. + * + * Anything else — including a `run.completed` whose payload this reader does + * not recognise — yields `undefined` and leaves the step-level verdict + * untouched, so an unfamiliar journal shape cannot invent a failure. + */ +function runCompletionReason(entry: unknown): ProtocolRunCompletionReason | undefined { + if (!isRecord(entry) || entry['entry_type'] !== 'run.completed') return undefined; + const payload = entry['payload']; + const reason = isRecord(payload) ? payload['completionReason'] : undefined; + return isSurfaceRunCompletionReason(reason) ? reason : undefined; +} + export function isSurfaceCompletionReason(value: unknown): value is ProtocolCompletionReason { return typeof value === 'string' && (COMPLETION_REASONS as readonly string[]).includes(value); diff --git a/packages/sdk/src/authored-worker-step.ts b/packages/sdk/src/authored-worker-step.ts index bd85fac0..770d4a88 100644 --- a/packages/sdk/src/authored-worker-step.ts +++ b/packages/sdk/src/authored-worker-step.ts @@ -8,7 +8,7 @@ import type { PreflightDiagnostic } from './preflight.js'; import { AuthoredFlowExecutionError } from './authored-flow-error.js'; import type { JournalClient } from './journal-client.js'; import { SPEC_SCHEMA_VERSION, type FlowSpec, type PermissionsSpec, type StepSpec } from './spec.js'; -import { isSurfaceCompletionReason, readCompletedStepOutput, readSuccessfulOutput, type AuthoredStepContext } from './authored-step-output.js'; +import { isSurfaceCompletionReason, readCompletedStepOutput, readRecordedOutcome, readSuccessfulOutput, type AuthoredStepContext, type RecordedRunOutcome } from './authored-step-output.js'; import type { AuthoredFlowJournalStep } from './authored-flow-executor.js'; import { snapshotJsonValue } from './json-value.js'; import { authoredChildAdmissionKey } from './authored-admission.js'; @@ -217,24 +217,41 @@ export function authoredDeterministicRunner( context: AuthoredStepContext = {}, ) { const rootRunId = context.rootRunId; - return async (id: string, command: string, terminal = false, leaseMs?: number, verification?: NamedGate): Promise => { + // One lowering, two readers. `record` changes what the child run's verdict + // MEANS, so the reader changes with it; everything before the read — spec, + // admission key, budget, child index — is deliberately shared, because a + // recorded red command is still an ordinary child run. + function lower( + id: string, command: string, terminal: boolean, leaseMs: number | undefined, + verification: NamedGate | undefined, onNonZero: 'record' | undefined, + read: (outcome: import('./protocol.js').RunOutcome) => Promise, + ): Promise { const spec = toKernelSpec(compileSpec({ version: SPEC_SCHEMA_VERSION, name: `${name}/${id}`, steps: [{ id, type: 'deterministic', command, ...(leaseMs === undefined ? {} : { lease_ms: leaseMs }), + ...(onNonZero === undefined ? {} : { onNonZero }), ...(verification === undefined ? {} : { verification }), }], })); const admissionKey = authoredChildAdmissionKey(rootRunId, id); - const consume = async (outcome: import('./protocol.js').RunOutcome): Promise => { + const consume = async (outcome: import('./protocol.js').RunOutcome): Promise => { await recordAuthoredChild(journal, rootRunId, { step: id, runId: outcome.run_id, state: 'admitted' }); - return readSuccessfulOutput(journal, outcome, id, journalSteps, context); + return read(outcome); }; - if (terminal) return consume(await journal.runStart(spec, undefined, admissionKey)); + if (terminal) return journal.runStart(spec, undefined, admissionKey).then(consume); return budget.execute(journal, spec, consume, admissionKey); - }; + } + const lowerRun = (id: string, command: string, terminal = false, leaseMs?: number, verification?: NamedGate): Promise => + lower(id, command, terminal, leaseMs, verification, undefined, + outcome => readSuccessfulOutput(journal, outcome, id, journalSteps, context)); + /** `f.run(..., { onNonZero: 'record' })`: the exit code is the result, not a throw. */ + lowerRun.recording = (id: string, command: string, leaseMs?: number, verification?: NamedGate): Promise => + lower(id, command, false, leaseMs, verification, 'record', + outcome => readRecordedOutcome(journal, outcome, id, journalSteps, context)); + return lowerRun; } /** diff --git a/packages/sdk/src/cli.ts b/packages/sdk/src/cli.ts index 75bfc1b1..30b9f213 100644 --- a/packages/sdk/src/cli.ts +++ b/packages/sdk/src/cli.ts @@ -926,7 +926,9 @@ function emitCheckReport(report: CheckReport, json: boolean, io: CliIo): void { // A gate that accepts every output is legal, but it must not read like a // gate that judges something. const vacuous = gate.acceptsAnyOutput === true ? ' [json_schema accepts any output]' : ''; - io.stdout(`GATE step "${gate.stepId}" ${gate.checks.join('+')} from data (kernel, journal-replayable)${vacuous}`); + // Red-but-complete is a declared shape, so it is a visible one. + const recorded = gate.recordsNonZeroExit === true ? ' [onNonZero: record]' : ''; + io.stdout(`GATE step "${gate.stepId}" ${gate.checks.join('+')} from data (kernel, journal-replayable)${vacuous}${recorded}`); } for (const schedule of report.schedules ?? []) { const declared = schedule.cron !== undefined diff --git a/packages/sdk/src/compile.ts b/packages/sdk/src/compile.ts index 27368a10..bda02ed3 100644 --- a/packages/sdk/src/compile.ts +++ b/packages/sdk/src/compile.ts @@ -19,6 +19,7 @@ import { parseBudget, toKernelBudget } from './budget.js'; import { bindingDependencies } from './input-binding.js'; import { namedGateFailure } from './named-gates.js'; import { lowerNamedGates } from './named-gate-lowering.js'; +import { lowerStepsGreenGates } from './steps-green.js'; import type { AgentStepSpec, DeterministicStepSpec, @@ -208,6 +209,10 @@ function compileStep(step: StepSpec): StepSpec { // their dispatch timeout. It must be spread HERE and nowhere in `base`. ...(s.timeoutMs !== undefined ? { timeoutMs: s.timeoutMs } : {}), ...(s.lease_ms !== undefined ? { lease_ms: parseStepTimeout(s.lease_ms) } : {}), + // Normalized away when it names the default, so an author who spells + // out `fail` gets the same canonical spec — and the same hash — as one + // who omits the field entirely. + ...(s.onNonZero !== undefined && s.onNonZero !== 'fail' ? { onNonZero: s.onNonZero } : {}), }; } case 'llm': { @@ -259,6 +264,11 @@ function outputGate(gate: VerificationSpec, stepId: string): OutputVerificationS `step "${stepId}": exit_code is supported only on deterministic steps`, ]); } + if (gate.type === 'steps_green') { + throw new CompileError([ + `step "${stepId}": gate_host_unsupported: steps_green is supported only on deterministic steps`, + ], 'gate_host_unsupported'); + } return gate; } @@ -331,7 +341,13 @@ export function toKernelSpec(flow: FlowSpec): KernelRunSpec { // reverts the other, and `validateSpec` and `flows check` would both still // look correct. See ops/reviews/20260903-pr139-repair-0903.md section 10. ...(compiled.triggers?.length ? { triggers: compiled.triggers.map(toKernelTrigger) } : {}), - steps: lowerNamedGates(compiled.steps).map((step) => toKernelStep(resolveNamedAgent(step, compiled.agents))), + // `steps_green` lowers FIRST: it rewrites the steps it reads to the + // permissive envelope schema its bindings need, and the named-gate pass + // then adds each source's own gate barrier to the generated step. The + // reverse order would hand the green gate a source whose barrier it does + // not wait for. + steps: lowerNamedGates(lowerStepsGreenGates(compiled.steps)) + .map((step) => toKernelStep(resolveNamedAgent(step, compiled.agents))), ...(compiled.budget !== undefined ? { budget: toKernelBudget(compiled.budget) } : {}), }; } @@ -426,14 +442,14 @@ function kernelTriggerToAuthoring(value: unknown, at: string): unknown { function kernelStepToAuthoring(value: unknown, at: string): unknown { const unionKeys = [ 'id', 'type', 'depends_on', 'max_iterations', 'retry', 'verification', 'memory', 'requirements', 'input', - 'command', 'timeout_ms', 'lease_ms', 'prompt', 'model', 'cli', 'instruction', + 'command', 'timeout_ms', 'lease_ms', 'on_non_zero', 'prompt', 'model', 'cli', 'instruction', 'recovery_mode', 'surfaces', 'permissions', ] as const; const step = requireKernelObject(value, unionKeys, at); const type = step['type']; const commonKeys = ['id', 'type', 'depends_on', 'max_iterations', 'retry', 'verification', 'memory', 'requirements', 'input'] as const; const typeKeys = type === 'deterministic' - ? ['command', 'timeout_ms', 'lease_ms'] as const + ? ['command', 'timeout_ms', 'lease_ms', 'on_non_zero'] as const : type === 'llm' ? ['prompt', 'model', 'cli'] as const : type === 'agent' @@ -460,6 +476,7 @@ function kernelStepToAuthoring(value: unknown, at: string): unknown { command: step['command'], ...(step['timeout_ms'] !== undefined ? { timeoutMs: step['timeout_ms'] } : {}), ...(step['lease_ms'] !== undefined ? { lease_ms: step['lease_ms'] } : {}), + ...(step['on_non_zero'] !== undefined ? { onNonZero: step['on_non_zero'] } : {}), }; } if (type === 'llm') { @@ -610,6 +627,11 @@ function toKernelStep(step: StepSpec): KernelStepSpec { command: step.command, ...(step.timeoutMs !== undefined ? { timeout_ms: step.timeoutMs } : {}), ...(step.lease_ms !== undefined ? { lease_ms: parseStepTimeout(step.lease_ms) } : {}), + // Omitted for the default, exactly as the kernel skips serializing it: + // a spec written before this field existed keeps its canonical bytes + // and therefore its spec hash. + ...(step.onNonZero !== undefined && step.onNonZero !== 'fail' + ? { on_non_zero: step.onNonZero } : {}), }; case 'llm': { return { diff --git a/packages/sdk/src/gate-contract.ts b/packages/sdk/src/gate-contract.ts index cc45be5a..6bf85526 100644 --- a/packages/sdk/src/gate-contract.ts +++ b/packages/sdk/src/gate-contract.ts @@ -24,6 +24,13 @@ export interface StepGateInspection extends DataGateClassification { * identically to a gate that judges something. */ acceptsAnyOutput?: true; + /** + * Present when the step declares `onNonZero: 'record'`. The `exit_code` + * check is still applied and still journaled, but a positive code satisfies + * it instead of failing the step. Without this a recording step reads + * exactly like a gated one — the `|| true` ambiguity, moved into the report. + */ + recordsNonZeroExit?: true; } /** Annotation-only keywords: present, they still constrain no instance. */ @@ -55,11 +62,15 @@ export function acceptsAnyOutput(schema: unknown): boolean { /** Describe the exact named checks the existing kernel applies to a step. */ export function inspectStepGate(step: StepSpec): StepGateInspection { const checks: JournalGateCheck[] = []; + // A policy for the step's own exit code, not a gate of its own: it never + // adds or removes a check, it changes what a positive code means to one. + const recording = step.type === 'deterministic' && step.onNonZero === 'record' + ? { recordsNonZeroExit: true as const } : {}; if (isNamedGate(step.verification)) { checks.push('exit_code'); if (step.verification.type === 'references_input') checks.push('output_contains'); if (step.verification.type === 'word_count_bounds') checks.push('json_schema'); - return { stepId: step.id, kind: 'data', checks, evaluator: 'kernel', preflightable: true, replayable: true }; + return { stepId: step.id, kind: 'data', checks, evaluator: 'kernel', preflightable: true, replayable: true, ...recording }; } if (step.type === 'deterministic') checks.push('exit_code'); if (step.verification?.type === 'output_contains') checks.push('output_contains'); @@ -75,5 +86,6 @@ export function inspectStepGate(step: StepSpec): StepGateInspection { preflightable: true, replayable: true, ...(vacuous ? { acceptsAnyOutput: true as const } : {}), + ...recording, }; } diff --git a/packages/sdk/src/index.ts b/packages/sdk/src/index.ts index f65fa12d..00760fcb 100644 --- a/packages/sdk/src/index.ts +++ b/packages/sdk/src/index.ts @@ -34,6 +34,8 @@ export type { LlmStepSpec, NamedAgentSpec, NamedDataGate, + OnNonZero, + StepsGreenGate, ReferencesInputGate, SubprocessGate, WordCountBoundsGate, diff --git a/packages/sdk/src/named-gates.ts b/packages/sdk/src/named-gates.ts index 20a50b1a..87af9600 100644 --- a/packages/sdk/src/named-gates.ts +++ b/packages/sdk/src/named-gates.ts @@ -3,6 +3,9 @@ import type { NamedDataGate, VerificationSpec } from './spec.js'; export const NAMED_GATE_FAILURE_KINDS = [ 'unknown_gate_kind', 'gate_pattern_invalid', 'gate_command_missing', 'gate_bound_invalid', 'gate_path_invalid', + // steps_green refusals: the gate is declarative like the named gates above, + // so it reports through the same kind channel a caller already switches on. + 'gate_host_unsupported', 'gate_ids_invalid', 'gate_source_unsupported', ] as const; export type NamedGateFailureKind = typeof NAMED_GATE_FAILURE_KINDS[number]; diff --git a/packages/sdk/src/spec.ts b/packages/sdk/src/spec.ts index 6016a380..496b6a09 100644 --- a/packages/sdk/src/spec.ts +++ b/packages/sdk/src/spec.ts @@ -15,6 +15,13 @@ export type { JsonOutputSchema } from './output-schema.js'; /** The three rungs of the ladder (RFC §1; AGENTS.md rule 7). */ export type StepType = 'deterministic' | 'llm' | 'agent'; +/** + * A deterministic step's exit-code policy. `'fail'` is the kernel's implicit + * gate; `'record'` completes the step with its exit code journaled so a later + * step can read the outcome instead of re-running the command. + */ +export type OnNonZero = 'fail' | 'record'; + /** Project-owned MCP connections. env contains names, never secret values. */ export type McpServerConfig = | { command: string; args?: string[]; env?: string[] } @@ -93,9 +100,24 @@ export interface ArtifactExistsGate { path: string; } +/** + * Asserts that every named step recorded a zero exit code, reading the + * journaled outcome instead of re-running the command. This is the read side + * of `onNonZero: 'record'`: red work is allowed to exist, and this gate is + * where a flow declares that it stops being allowed. + * + * Deterministic hosts only — a worker step has no recorded exit code to read — + * and the ids it names must be deterministic steps that precede it. + */ +export interface StepsGreenGate { + type: 'steps_green'; + /** Earlier deterministic step ids whose recorded exit codes must all be zero. */ + ids: string[]; +} + export type NamedDataGate = ReferencesInputGate | SubprocessGate | WordCountBoundsGate | RegexMatchGate | ArtifactExistsGate; export type OutputVerificationSpec = OutputContainsGate | JsonSchemaGate | NamedDataGate; -export type VerificationSpec = ExitCodeGate | OutputVerificationSpec; +export type VerificationSpec = ExitCodeGate | StepsGreenGate | OutputVerificationSpec; /** * Agent-step recovery modes (RFC Appendix A rule 4). Default is `reset`. @@ -226,6 +248,15 @@ export interface DeterministicStepSpec extends BaseStepSpec { timeoutMs?: number; /** Per-invocation deterministic lease in milliseconds (maximum 15 minutes). */ lease_ms?: number; + /** + * What a nonzero exit means. `'fail'` (the default) is the implicit + * `exit_code == 0` gate. `'record'` journals the exit code and output tails + * as the step's outcome, completes the step red, and lets dependents run — + * the repair-before-failure shape, without discarding the exit code the way + * ` || true` does. Declared content and schema gates, timeouts and + * worker errors keep their existing fatal semantics. + */ + onNonZero?: OnNonZero; /** Omit to get the implicit `exit_code` gate. */ verification?: VerificationSpec; } @@ -440,6 +471,12 @@ export interface KernelDeterministicStep extends KernelStepCommon { command: string; timeout_ms?: number; lease_ms?: number; + /** + * Omitted for the default `'fail'`, exactly as the kernel skips serializing + * it: a step that does not declare the policy keeps the canonical bytes — + * and therefore the step hash — it had before this field existed. + */ + on_non_zero?: OnNonZero; } export interface KernelLlmStep extends KernelStepCommon { diff --git a/packages/sdk/src/step-fields.ts b/packages/sdk/src/step-fields.ts index d51e1a27..076dc0bd 100644 --- a/packages/sdk/src/step-fields.ts +++ b/packages/sdk/src/step-fields.ts @@ -35,7 +35,7 @@ export const STEP_COMMON_FIELDS = [ * per-verb boundary through a second, drifting allowlist. */ export const STEP_FIELDS_BY_TYPE = { - deterministic: ['command', 'timeoutMs', 'lease_ms'], + deterministic: ['command', 'timeoutMs', 'lease_ms', 'onNonZero'], llm: ['prompt', 'model', 'cli', 'output'], agent: ['instruction', 'agent', 'cli', 'model', 'cwd', 'transport', 'surfaces', 'recoveryMode', 'permissions', 'output'], } as const satisfies Record; diff --git a/packages/sdk/src/steps-green.ts b/packages/sdk/src/steps-green.ts new file mode 100644 index 00000000..ee944ae9 --- /dev/null +++ b/packages/sdk/src/steps-green.ts @@ -0,0 +1,165 @@ +// The read side of `onNonZero: 'record'`. A `steps_green` gate names earlier +// deterministic steps and asserts every one of them recorded exit code zero. +// The verdict comes from the journal — the recorded outcome each step already +// produced — so asserting it costs nothing and cannot disagree with what ran. + +import { bindingDependencies } from './input-binding.js'; +import type { StepSpec, StepsGreenGate, VerificationSpec } from './spec.js'; + +export const STEPS_GREEN_KEYS = ['type', 'ids'] as const; + +export function isStepsGreenGate(gate: VerificationSpec | undefined): gate is StepsGreenGate { + return gate?.type === 'steps_green'; +} + +/** Shape errors for one gate, without reference to the rest of the flow. */ +export function stepsGreenErrors(gate: Record, at: string, stepType: string): string[] { + const errors: string[] = []; + if (stepType !== 'deterministic') { + errors.push(`${at}: gate_host_unsupported: steps_green is supported only on deterministic steps`); + } + const ids = gate['ids']; + if (!Array.isArray(ids) || ids.length === 0 + || !ids.every(id => typeof id === 'string' && id.trim().length > 0)) { + errors.push(`${at}.ids: gate_ids_invalid: expected a non-empty array of step ids`); + return errors; + } + if (new Set(ids).size !== ids.length) { + errors.push(`${at}.ids: gate_ids_invalid: step ids must be unique`); + } + return errors; +} + +/** + * Flow-level errors: the ids must resolve to earlier deterministic steps whose + * recorded outcome this gate can actually read. Refusing here is the whole + * point — a gate that silently reads nothing is the `|| true` failure mode + * this feature replaces. + */ +export function stepsGreenReferenceErrors(steps: readonly unknown[]): string[] { + const errors: string[] = []; + const sources = new Map(steps.flatMap((step, index) => + isObject(step) && typeof step['id'] === 'string' ? [[step['id'], { step, index }] as const] : [])); + for (const [index, step] of steps.entries()) { + if (!isObject(step)) continue; + const gate = step['verification']; + if (!isObject(gate) || gate['type'] !== 'steps_green') continue; + const ids = gate['ids']; + if (!Array.isArray(ids)) continue; + const at = `spec.steps[${index}].verification.ids`; + for (const id of ids) { + if (typeof id !== 'string') continue; + const source = sources.get(id); + if (source === undefined) { + errors.push(`${at}: gate_source_unsupported: unknown step "${id}"`); + continue; + } + if (source.index >= index) { + errors.push(`${at}: gate_source_unsupported: source "${id}" must precede this step; forward and self references are not supported`); + continue; + } + if (source.step['type'] !== 'deterministic') { + errors.push(`${at}: gate_source_unsupported: source "${id}" is a ${String(source.step['type'])} step and records no exit code`); + continue; + } + // The gate binds the source's whole output envelope, which requires a + // declared output schema. Every source shape is reachable except this + // one: `output_contains` is the single gate the compiler cannot widen + // without dropping the author's own check. + const sourceGate = source.step['verification']; + if (isObject(sourceGate) && sourceGate['type'] === 'output_contains') { + errors.push(`${at}: gate_source_unsupported: source "${id}" declares an output_contains gate, which cannot also carry the output schema this gate reads`); + } + } + } + return errors; +} + +/** + * Lower every `steps_green` gate into a generated deterministic gate step. + * + * The host keeps its own command and its own exit policy; the gate is a + * separate step, always fatal, that binds each named source's journaled output + * envelope and refuses any nonzero or malformed exit code. Dependents of the + * host gain a barrier on the gate, so "these steps were green" is established + * before anything downstream of the assertion runs. + * + * Runs BEFORE `lowerNamedGates`, which then adds each source's own named-gate + * barrier to the generated step and rewrites named-gate hosts to the permissive + * envelope schema this gate's bindings need. + */ +export function lowerStepsGreenGates(steps: readonly StepSpec[]): StepSpec[] { + const used = new Set(steps.map(step => step.id)); + const barriers = new Map(); + for (const step of steps) { + // Deterministic hosts only, which `validateSpec` already enforces; the + // narrowing here is what lets the rewrite below stay type-exact. + if (step.type !== 'deterministic' || !isStepsGreenGate(step.verification)) continue; + let id = `${step.id}.green`; + while (used.has(id)) id += '.green'; + used.add(id); + barriers.set(step.id, id); + } + if (barriers.size === 0) return [...steps]; + + // Every source's envelope has to be bindable. A source that declares no + // gate, or the implicit exit_code gate, gets the permissive whole-output + // schema; its exit policy is untouched, so a recording source stays red. + const bound = new Set(steps.flatMap(step => + isStepsGreenGate(step.verification) ? step.verification.ids : [])); + const output: StepSpec[] = []; + for (const step of steps) { + const dependencies = [...new Set([...(step.dependsOn ?? []), ...bindingDependencies(step.input)])]; + const dependsOn = [...new Set([...dependencies, ...dependencies.flatMap(id => barriers.get(id) ?? [])])]; + let producer: StepSpec = dependsOn.length ? { ...step, dependsOn } : step; + if (bound.has(step.id) && producer.type === 'deterministic' + && (producer.verification === undefined || producer.verification.type === 'exit_code')) { + producer = { ...producer, verification: { type: 'json_schema', schema: true } }; + } + const gate = step.verification; + if (producer.type !== 'deterministic' || !isStepsGreenGate(gate)) { + output.push(producer); + continue; + } + // The host's own gate becomes the implicit exit_code check. Its command + // still runs, and its own `onNonZero` still applies to its own exit. A + // host that is itself read by a later gate keeps the bindable envelope. + output.push({ + ...producer, + verification: bound.has(step.id) ? { type: 'json_schema', schema: true } : { type: 'exit_code' }, + }); + const input = Object.fromEntries(gate.ids.map((id, index) => [String(index), { step: id }])); + output.push({ + id: barriers.get(step.id)!, type: 'deterministic', + dependsOn: [...new Set([step.id, ...gate.ids])], input, + command: greenCommand(gate.ids), + verification: { type: 'exit_code' }, maxIterations: step.maxIterations ?? 1, + ...(step.requirements === undefined ? {} : { requirements: step.requirements }), + }); + } + return output; +} + +/** + * Compiler-owned Node command. Author-supplied ids are JSON literals and the + * recorded envelopes travel only through FLOWS_INPUT, never through shell text. + */ +function greenCommand(ids: readonly string[]): string { + const script = `const input=JSON.parse(process.env.FLOWS_INPUT); +const ids=${JSON.stringify(ids)}; +const red=[]; +for(let i=0;i0){process.stderr.write('steps_green: '+red.join('; ')+'\\n');process.exit(1);} +process.stdout.write('steps_green:pass');`; + return `node -e '${script.replaceAll("'", "'\\''")}'`; +} + +function isObject(value: unknown): value is Record { + return typeof value === 'object' && value !== null && !Array.isArray(value); +} diff --git a/packages/sdk/src/validate.ts b/packages/sdk/src/validate.ts index 09b733c1..a5a9953f 100644 --- a/packages/sdk/src/validate.ts +++ b/packages/sdk/src/validate.ts @@ -23,6 +23,7 @@ import { unknownKeyErrors } from './unknown-keys.js'; import { stepDependencyErrors } from './step-dependencies.js'; import { inputBindingErrors } from './input-binding.js'; import { NAMED_GATE_KEYS, namedGateErrors } from './named-gates.js'; +import { STEPS_GREEN_KEYS, stepsGreenErrors, stepsGreenReferenceErrors } from './steps-green.js'; import { AGENT_DECLARATION_FIELDS, FLOW_FIELDS, @@ -74,6 +75,7 @@ const BUDGET_KEYS = ['maxTokensIn', 'maxTokensOut', 'maxDollars', 'maxTokens', ' const VERIFICATION_KEYS: Record = { ...NAMED_GATE_KEYS, exit_code: ['type', 'expect'], + steps_green: STEPS_GREEN_KEYS, output_contains: ['type', 'value'], json_schema: ['type', 'schema'], }; @@ -195,6 +197,9 @@ class Validator { // Dependents must reference real step ids and form a DAG (no cycles). for (const error of stepDependencyErrors(steps, this.ids)) this.fail(error); for (const error of inputBindingErrors(steps)) this.fail(error); + // A steps_green gate that names nothing readable is the `|| true` failure + // it replaces, so its references are refused at the flow level, not lowered. + for (const error of stepsGreenReferenceErrors(steps)) this.fail(error); return this.result(); } @@ -443,10 +448,14 @@ class Validator { } catch (error) { this.fail(error instanceof Error ? error.message : `${at}.schema: expected JSON-compatible data`); } + // Below here are the gates the COMPILER lowers; above are the three the + // kernel evaluates itself. The enumeration follows the same split. + } else if (gate.type === 'steps_green') { + for (const error of stepsGreenErrors(v, at, stepType)) this.fail(error); } else if (Object.hasOwn(NAMED_GATE_KEYS, gate.type)) { for (const error of namedGateErrors(v, input, at)) this.fail(error); } else { - this.fail(`${at}.type: unknown_gate_kind: expected exit_code | output_contains | json_schema | references_input | subprocess_gate | word_count_bounds | regex_match | artifact_exists`); + this.fail(`${at}.type: unknown_gate_kind: expected exit_code | output_contains | json_schema | steps_green | references_input | subprocess_gate | word_count_bounds | regex_match | artifact_exists`); } } @@ -460,6 +469,12 @@ class Validator { if (st.timeoutMs !== undefined && !isPosInt(st.timeoutMs)) { this.fail(`${at}.timeoutMs: expected a positive integer`); } + // Only the two policies the kernel enum names. A third spelling would + // deserialize-fail in the kernel, so it is refused where the author can + // still read the error. + if (st.onNonZero !== undefined && st.onNonZero !== 'fail' && st.onNonZero !== 'record') { + this.fail(`${at}.onNonZero: expected fail | record`); + } } private validateLlm(st: LlmStepSpec, at: string): void { diff --git a/packages/sdk/tests/advisory-outcome.test.ts b/packages/sdk/tests/advisory-outcome.test.ts new file mode 100644 index 00000000..2043be98 --- /dev/null +++ b/packages/sdk/tests/advisory-outcome.test.ts @@ -0,0 +1,242 @@ +// Advisory deterministic outcomes: `onNonZero: 'record'` and the `steps_green` +// gate that reads what it records. +// +// The point of the feature is that a red command stays LEGIBLE. `|| true` +// discards the exit code, so the two things these tests hold onto are: the +// policy never silently defaults when it is misspelled, and a gate that claims +// to read recorded outcomes is refused unless it actually can. + +import { describe, expect, it } from 'vitest'; +import { CompileError, compileSpec, toKernelSpec } from '../src/compile.js'; +import { inspectStepGate } from '../src/gate-contract.js'; +import { lowerStepsGreenGates } from '../src/steps-green.js'; +import { runCli } from '../src/cli.js'; +import { mkdtempSync, rmSync, writeFileSync } from 'node:fs'; +import { tmpdir } from 'node:os'; +import { join } from 'node:path'; +import type { FlowSpec, StepSpec } from '../src/spec.js'; +import { validateSpec } from '../src/validate.js'; + +function flow(steps: unknown[]): unknown { + return { version: '0.1.0', name: 'advisory', steps }; +} + +function errors(steps: unknown[]): string[] { + return validateSpec(flow(steps)).errors; +} + +const recordingCheck = { + id: 'check', type: 'deterministic', command: 'npm test', onNonZero: 'record', +} as const; + +describe('onNonZero policy', () => { + it('accepts both spellings and refuses every other one', () => { + expect(errors([{ ...recordingCheck, onNonZero: 'fail' }])).toEqual([]); + expect(errors([recordingCheck])).toEqual([]); + for (const value of ['ignore', 'true', 'RECORD', '', 0, null, ['record']]) { + expect(errors([{ ...recordingCheck, onNonZero: value }])) + .toContain('spec.steps[0].onNonZero: expected fail | record'); + } + }); + + // The step-field allowlist is per type, so the policy must not leak onto a + // worker step, which has no exit code of its own to record. + it('refuses the policy on llm and agent steps', () => { + expect(errors([{ id: 'a', type: 'llm', prompt: 'hi', onNonZero: 'record' }]).join('\n')) + .toMatch(/onNonZero/); + expect(errors([{ id: 'a', type: 'agent', instruction: 'hi', onNonZero: 'record' }]).join('\n')) + .toMatch(/onNonZero/); + }); + + it('lowers only the non-default policy into the kernel dialect', () => { + const kernel = toKernelSpec(compileSpec(flow([recordingCheck]) as FlowSpec)); + expect(kernel.steps[0]).toMatchObject({ on_non_zero: 'record' }); + + const explicitDefault = toKernelSpec(compileSpec( + flow([{ ...recordingCheck, onNonZero: 'fail' }]) as FlowSpec)); + expect(explicitDefault.steps[0]).not.toHaveProperty('on_non_zero'); + }); + + // Recording changes what a positive exit code MEANS, not which checks run. + // Reporting it as an ordinary exit_code gate would move the `|| true` + // ambiguity into the inspection output. + it('reports the recording policy alongside the checks it does not change', () => { + expect(inspectStepGate(recordingCheck as StepSpec)).toEqual({ + stepId: 'check', kind: 'data', checks: ['exit_code'], + evaluator: 'kernel', preflightable: true, replayable: true, + recordsNonZeroExit: true, + }); + expect(inspectStepGate({ id: 'check', type: 'deterministic', command: 'npm test' })) + .not.toHaveProperty('recordsNonZeroExit'); + }); + + it('annotates the recording step in flows check output', async () => { + const dir = mkdtempSync(join(tmpdir(), 'advisory-check-')); + try { + const path = join(dir, 'advisory.yaml'); + writeFileSync(path, "version: '0.1.0'\nname: advisory\nsteps:\n" + + ' - id: check\n type: deterministic\n command: npm test\n onNonZero: record\n'); + const out: string[] = []; + const code = await runCli(['check', path], { stdout: (line) => out.push(line), stderr: (line) => out.push(line) }); + expect(code).toBe(0); + expect(out.join('\n')).toContain('GATE step "check" exit_code from data (kernel, journal-replayable) [onNonZero: record]'); + } finally { + rmSync(dir, { recursive: true, force: true }); + } + }); +}); + +describe('steps_green gate shape', () => { + const green = (ids: unknown) => [ + recordingCheck, + { id: 'gate', type: 'deterministic', command: 'true', verification: { type: 'steps_green', ids } }, + ]; + + it('accepts a gate over an earlier recording step', () => { + expect(errors(green(['check']))).toEqual([]); + }); + + it('refuses an empty, non-array, blank or duplicated id list', () => { + for (const ids of [[], 'check', [''], [' '], [1]]) { + expect(errors(green(ids)).join('\n')) + .toContain('gate_ids_invalid: expected a non-empty array of step ids'); + } + expect(errors(green(['check', 'check'])).join('\n')) + .toContain('gate_ids_invalid: step ids must be unique'); + }); + + it('refuses unknown keys on the gate', () => { + const errs = errors([ + recordingCheck, + { id: 'gate', type: 'deterministic', command: 'true', verification: { type: 'steps_green', ids: ['check'], expect: 0 } }, + ]); + expect(errs.join('\n')).toMatch(/expect/); + }); + + it('names steps_green in the unknown-gate refusal', () => { + expect(errors([{ id: 'a', type: 'deterministic', command: 'true', verification: { type: 'vibes' } }]).join('\n')) + .toContain('json_schema | steps_green | references_input'); + }); + + it('refuses a worker host, which has no recorded exit code', () => { + expect(errors([ + recordingCheck, + { id: 'gate', type: 'agent', instruction: 'judge', verification: { type: 'steps_green', ids: ['check'] } }, + ]).join('\n')).toContain('gate_host_unsupported: steps_green is supported only on deterministic steps'); + }); +}); + +describe('steps_green references', () => { + const withGate = (ids: string[], ...before: unknown[]) => [ + ...before, + { id: 'gate', type: 'deterministic', command: 'true', verification: { type: 'steps_green', ids } }, + ]; + + it('refuses an unknown id', () => { + expect(errors(withGate(['nope'], recordingCheck)).join('\n')) + .toContain('gate_source_unsupported: unknown step "nope"'); + }); + + // A forward reference reads an outcome that does not exist yet; a self + // reference reads the gate's own, which is not written until it passes. + it('refuses forward and self references', () => { + expect(errors([ + { id: 'gate', type: 'deterministic', command: 'true', verification: { type: 'steps_green', ids: ['check'] } }, + recordingCheck, + ]).join('\n')).toContain('must precede this step'); + expect(errors(withGate(['gate'])).join('\n')).toContain('must precede this step'); + }); + + it('refuses a worker source, which records no exit code', () => { + expect(errors(withGate(['write'], { id: 'write', type: 'agent', instruction: 'work' })).join('\n')) + .toContain('gate_source_unsupported: source "write" is a agent step and records no exit code'); + }); + + // The gate binds the source's whole envelope, and `output_contains` is the + // one gate the compiler cannot widen without deleting the author's check. + it('refuses an output_contains source rather than silently weakening it', () => { + expect(errors(withGate(['check'], { + ...recordingCheck, verification: { type: 'output_contains', value: 'ok' }, + })).join('\n')).toContain('declares an output_contains gate'); + }); + + it('accepts a source that already declares its own json_schema', () => { + expect(errors(withGate(['check'], { + ...recordingCheck, + verification: { type: 'json_schema', schema: { type: 'object' } }, + }))).toEqual([]); + }); +}); + +describe('steps_green lowering', () => { + const lowered = (steps: unknown[]) => toKernelSpec(compileSpec(flow(steps) as FlowSpec)).steps; + + it('keeps the host command and puts the assertion in its own fatal step', () => { + const steps = lowered([ + recordingCheck, + { id: 'gate', type: 'deterministic', command: 'echo done', verification: { type: 'steps_green', ids: ['check'] } }, + ]); + const host = steps.find((step) => step.id === 'gate')!; + expect(host.command).toBe('echo done'); + // The host's own recording policy is its own; the generated gate never + // inherits it, or a red gate could not fail anything. + const generated = steps.find((step) => step.id === 'gate.green')!; + expect(generated).not.toHaveProperty('on_non_zero'); + expect(generated.depends_on).toEqual(expect.arrayContaining(['gate', 'check'])); + expect(generated.input).toEqual({ '0': { step: 'check' } }); + expect(generated.command).toContain('FLOWS_INPUT'); + // Author ids travel as JSON literals, never as shell text. + expect(generated.command).toContain('["check"]'); + }); + + it('gives a bound source the envelope schema its binding needs', () => { + const steps = lowered([ + recordingCheck, + { id: 'gate', type: 'deterministic', command: 'true', verification: { type: 'steps_green', ids: ['check'] } }, + ]); + expect(steps.find((step) => step.id === 'check')!.verification).toEqual({ json_schema: true }); + // The recording policy survives the rewrite; the source stays red-capable. + expect(steps.find((step) => step.id === 'check')).toMatchObject({ on_non_zero: 'record' }); + }); + + it('preserves a source json_schema exactly rather than widening it', () => { + const schema = { type: 'object', required: ['exit_code'], properties: { exit_code: { type: 'integer' } } }; + const steps = lowered([ + { ...recordingCheck, verification: { type: 'json_schema', schema } }, + { id: 'gate', type: 'deterministic', command: 'true', verification: { type: 'steps_green', ids: ['check'] } }, + ]); + expect(steps.find((step) => step.id === 'check')!.verification).toEqual({ json_schema: schema }); + }); + + it('makes dependents of the host wait for the generated gate', () => { + const steps = lowered([ + recordingCheck, + { id: 'gate', type: 'deterministic', command: 'true', verification: { type: 'steps_green', ids: ['check'] } }, + { id: 'ship', type: 'deterministic', dependsOn: ['gate'], command: 'echo ship' }, + ]); + expect(steps.find((step) => step.id === 'ship')!.depends_on) + .toEqual(expect.arrayContaining(['gate', 'gate.green'])); + }); + + it('picks a generated id that cannot collide with an authored one', () => { + const steps = lowerStepsGreenGates([ + { id: 'check', type: 'deterministic', command: 'npm test', onNonZero: 'record' }, + { id: 'gate.green', type: 'deterministic', command: 'echo decoy' }, + { id: 'gate', type: 'deterministic', command: 'true', verification: { type: 'steps_green', ids: ['check'] } }, + ] as StepSpec[]); + expect(steps.map((step) => step.id)).toContain('gate.green.green'); + expect(steps.find((step) => step.id === 'gate.green')!.command).toBe('echo decoy'); + }); + + it('leaves a flow without the gate byte-identical', () => { + const steps: StepSpec[] = [{ id: 'a', type: 'deterministic', command: 'true' }]; + expect(lowerStepsGreenGates(steps)).toEqual(steps); + }); + + it('refuses the gate on a worker host at compile time too', () => { + expect(() => compileSpec(flow([ + recordingCheck, + { id: 'gate', type: 'llm', prompt: 'judge', verification: { type: 'steps_green', ids: ['check'] } }, + ]) as FlowSpec)).toThrow(CompileError); + }); +}); diff --git a/packages/sdk/tests/authored-advisory-run-live.test.ts b/packages/sdk/tests/authored-advisory-run-live.test.ts new file mode 100644 index 00000000..f6392a65 --- /dev/null +++ b/packages/sdk/tests/authored-advisory-run-live.test.ts @@ -0,0 +1,217 @@ +// `f.run(..., { onNonZero: 'record' })` against the live kernel. +// +// The claim under test is the one ` || true` cannot make: a red +// command does not end the flow, AND the branch below it reads the command's +// real exit code and real output, back from the journal. Every assertion here +// is therefore about the values an author actually receives — not about the +// spec that was compiled to produce them. + +import { flow } from '@relayflows/surface'; +import { afterEach, describe, expect, it } from 'vitest'; +import { executeAuthoredFlow } from '../src/authored-flow-executor.js'; +import type { AuthoredFlowExecutionError } from '../src/authored-flow-error.js'; +import { readRecordedOutcome } from '../src/authored-step-output.js'; +import type { JournalClient } from '../src/journal-client.js'; +import type { RunOutcome } from '../src/protocol.js'; +import { chainFixture } from './flow-chain-fixture.js'; + +const cleanup: Array<() => Promise> = []; +afterEach(async () => { for (const close of cleanup.splice(0)) await close(); }); + +describe('a recorded nonzero f.run through the live kernel', () => { + it('resolves to the journaled exit code instead of ending the flow', async () => { + const fixture = chainFixture(); + cleanup.push(() => fixture.close()); + const journal = await fixture.connect(); + + let seen: unknown; + let reachedRepair = false; + const handle = flow('repair-before-failure', async (f) => { + const check = await f.run( + 'printf \'2 failing tests\'; printf \'stack trace\' >&2; exit 7', + { onNonZero: 'record' }); + seen = check; + if (!check.ok) { + // The repair-before-failure shape from the ticket: the red command's + // own output is the input to whatever answers it. + const echoed = await f.run(`printf '%s' ${JSON.stringify(check.stdout)}`); + reachedRepair = echoed === '2 failing tests'; + } + f.done('success'); + }); + + const result = await executeAuthoredFlow(handle, journal, undefined, { dataDir: fixture.data }); + + expect(result.completionReason).toBe('success'); + expect(seen).toEqual({ + ok: false, exitCode: 7, + stdout: '2 failing tests', stderr: 'stack trace', + output: '2 failing tests\nstack trace', + }); + expect(reachedRepair).toBe(true); + // The step is a SUCCESSFUL step whose exit code happens to be 7. Recording + // does not invent a completion reason, and it does not retry. + expect(result.journalSteps.find(step => step.id === 'run-1')) + .toMatchObject({ completionReason: 'success' }); + }, 30_000); + + it('resolves a green command through the same policy', async () => { + const fixture = chainFixture(); + cleanup.push(() => fixture.close()); + const journal = await fixture.connect(); + + let seen: unknown; + const handle = flow('recorded-green', async (f) => { + seen = await f.run('printf all-green', { onNonZero: 'record' }); + f.done('success'); + }); + + await executeAuthoredFlow(handle, journal, undefined, { dataDir: fixture.data }); + + expect(seen).toEqual({ ok: true, exitCode: 0, stdout: 'all-green', stderr: '', output: 'all-green' }); + }, 30_000); + + // The default is what every existing body is written against, and the whole + // feature is opt-in. A policy-free call must be untouched by this change. + it('leaves the default fatal and still resolving to the stdout string', async () => { + const fixture = chainFixture(); + cleanup.push(() => fixture.close()); + const journal = await fixture.connect(); + + const handle = flow('default-still-fatal', async (f) => { + expect(await f.run('printf plain')).toBe('plain'); + await f.run('exit 4'); + f.done('success'); + }); + + await expect(executeAuthoredFlow(handle, journal, undefined, { dataDir: fixture.data })) + .rejects.toThrow(/step "run-2"/u); + }, 30_000); + + // Recording is a policy about EXIT CODES. A command killed by its lease + // produced no exit status to record, so there is nothing to hand an author's + // `if (!result.ok)` and the step still fails. + it('still fails a timeout, which recorded no exit code to report', async () => { + const fixture = chainFixture(); + cleanup.push(() => fixture.close()); + const journal = await fixture.connect(); + + const handle = flow('recorded-timeout', async (f) => { + await f.run('sleep 30', { onNonZero: 'record', timeout: '150ms' }); + f.done('success'); + }); + + const failure = await executeAuthoredFlow(handle, journal, undefined, { dataDir: fixture.data }) + .then(() => undefined, (error) => error as { code?: string; details?: Record }); + + expect(failure?.code).toBe('lease_exceeded'); + expect(failure?.details?.['stepId']).toBe('run-1'); + }, 30_000); + + // A declared gate is the author's own assertion, not the exit-code policy's, + // so recording must not soften it. The producer here completes successfully — + // that is the point. A named gate lowers to its own kernel step that runs + // AFTER the producer, so the producer's `success` entry is not the run's + // verdict, and reading only that entry used to resolve the operation as + // though the gate had passed. + it.each([ + ['a recorded step', { onNonZero: 'record' } as const], + ['a default step', undefined], + ])('still fails a declared named gate on %s whose producer succeeded', async (_label, options) => { + const fixture = chainFixture(); + cleanup.push(() => fixture.close()); + const journal = await fixture.connect(); + + const handle = flow('named-gate-failure', async (f) => { + // `printf` exits 0 under either policy: the only thing that can fail + // this operation is the gate. + await f.run('printf nothing-useful', options) + .gate({ type: 'regex_match', pattern: 'ready', in_output_at: ['stdout_tail'] }); + f.done('success'); + }); + + const failure = await executeAuthoredFlow(handle, journal, undefined, { dataDir: fixture.data }) + .then(() => undefined, (error) => error as AuthoredFlowExecutionError); + + expect(failure?.code).toBe('step_failed'); + // The RUN's reason, and the generated gate step named as the culprit — + // not the producer, which really did succeed. + expect(failure?.completionReason).toBe('step_failed'); + expect(failure?.details?.stepId).toBe('run-1.gate'); + // And the child run itself is terminally failed, so nothing downstream can + // read this operation as a pass. + expect((await journal.runGet(failure!.runId!)).status).toBe('failed'); + }, 30_000); + + // A predicate gate judges the value the author received, so on a recording + // step it judges the `RunResult` — including the exit code the policy kept. + it('hands a predicate gate the recorded result, not the stdout string', async () => { + const fixture = chainFixture(); + cleanup.push(() => fixture.close()); + const journal = await fixture.connect(); + + let judged: unknown; + const handle = flow('recorded-predicate', async (f) => { + await f.run('printf out; exit 5', { onNonZero: 'record' }) + .gate((result) => { judged = result; return result.exitCode === 5; }); + f.done('success'); + }); + + await executeAuthoredFlow(handle, journal, undefined, { dataDir: fixture.data }); + + expect(judged).toMatchObject({ ok: false, exitCode: 5, stdout: 'out' }); + }, 30_000); + + it('refuses a policy value the kernel enum does not name, before journaling anything', async () => { + const fixture = chainFixture(); + cleanup.push(() => fixture.close()); + const journal = await fixture.connect(); + + const handle = flow('bad-policy', async (f) => { + // A JavaScript caller, or a cast, can reach this. Falling back to the + // fatal default would delete the branch the author wrote below it. + await f.run('true', { onNonZero: 'ignore' } as never); + f.done('success'); + }); + + await expect(executeAuthoredFlow(handle, journal, undefined, { dataDir: fixture.data })) + .rejects.toThrow(/onNonZero must be 'fail' or 'record'/u); + }, 30_000); +}); + +describe('reading a recorded outcome from a malformed envelope', () => { + /** A journal whose single completion carries whatever output is given. */ + function completedWith(output: unknown): JournalClient { + return { + async journalRead() { + return { entries: [{ + seq: 1, entry_type: 'step.completed', step_id: 'run-1', + payload: { completionReason: 'success', disposition: 'step_done', output }, + }] }; + }, + } as unknown as JournalClient; + } + + const read = (output: unknown) => readRecordedOutcome( + completedWith(output), { run_id: 'child-1' } as RunOutcome, 'run-1', []); + + it('accepts a well-formed envelope', async () => { + await expect(read({ exit_code: 2, stdout_tail: 'out', stderr_tail: 'err' })) + .resolves.toEqual({ ok: false, exitCode: 2, stdout: 'out', stderr: 'err', output: 'out\nerr' }); + }); + + // `-1` is the executor's "no exit status" sentinel (exec_det.rs), not a + // command's exit code. Handing it back as `ok: false` would present a + // killed process as an ordinary red verdict. + it.each([ + [{ exit_code: -1, stdout_tail: '', stderr_tail: '' }, 'the no-exit-status sentinel'], + [{ exit_code: 1.5, stdout_tail: '', stderr_tail: '' }, 'a non-integer code'], + [{ exit_code: 0, stderr_tail: '' }, 'a missing stdout tail'], + [{ exit_code: 0, stdout_tail: '' }, 'a missing stderr tail'], + [{ stdout_tail: '', stderr_tail: '' }, 'no code at all'], + ['0', 'a non-object envelope'], + [null, 'a null envelope'], + ])('refuses %#: %s', async (output) => { + await expect(read(output)).rejects.toThrow(/recorded no \{exit_code, stdout_tail, stderr_tail\}/u); + }); +}); diff --git a/packages/sdk/tests/preflight.test.ts b/packages/sdk/tests/preflight.test.ts index 6c5596f9..9717b7df 100644 --- a/packages/sdk/tests/preflight.test.ts +++ b/packages/sdk/tests/preflight.test.ts @@ -585,6 +585,19 @@ describe('preflight: CLI resolution and refusal predicates', () => { // `scope_syntax_invalid`: grant string doesn't match "mount/path: mode". preflight({ ...flow({ id: 'a', type: 'deterministic', command: 'x' }), workspace: 'not-a-grant' } as FlowSpec, { probes: probes() }), + // steps_green refusal coverage. All three are compile-time, and all + // three exist because the gate reads a RECORDED exit code: a host that + // never produces one, an id list that names nothing readable, and a + // source whose outcome the gate cannot bind. + preflight(flow({ id: 'a', type: 'llm', prompt: 'p', cli: 'x', + verification: { type: 'steps_green', ids: ['a'] } } as never), { probes: probes() }), + preflight(flow({ id: 'a', type: 'deterministic', command: 'x', + verification: { type: 'steps_green', ids: [] } } as never), { probes: probes() }), + preflight({ version: '0.1.0', steps: [ + { id: 'writer', type: 'agent', instruction: 'i', cli: 'x' }, + { id: 'a', type: 'deterministic', command: 'x', + verification: { type: 'steps_green', ids: ['writer'] } }, + ] } as never, { probes: probes() }), ]; const refusalKinds = scenarios.flatMap((result) => result.diagnostics) .filter((diagnostic) => diagnostic.severity === 'refusal') diff --git a/packages/sdk/tests/spec-parity.test.ts b/packages/sdk/tests/spec-parity.test.ts index 48e25ac2..0bdf999a 100644 --- a/packages/sdk/tests/spec-parity.test.ts +++ b/packages/sdk/tests/spec-parity.test.ts @@ -44,6 +44,48 @@ describe('spec parity: one dialect at the SDK<->kernel boundary', () => { }); } + // The record policy and the gate that reads it, pinned at the same seam. The + // fixture is LOWERED (`steps_green` becomes a generated gate step), so it is + // not round-tripped: `kernelToAuthoring` is the inverse of `toKernelStep`, + // not of gate lowering. The field's own round trip is covered below. + it('compiles advisory-repair to the pinned canonical JSON', () => { + expect(compileYamlToCanonicalJson(fixture('advisory-repair.flow.yaml'))) + .toBe(fixture('advisory-repair.spec.canonical.json').trim()); + }); + + it('hashes advisory-repair to the pinned spec_hash', () => { + expect(compileAndHash(fixture('advisory-repair.flow.yaml')).hash) + .toBe(fixture('advisory-repair.spec.sha256').trim()); + }); + + it('round-trips onNonZero across the kernel dialect', () => { + const flow = compileYaml(` +version: '0.1.0' +name: record-round-trip +steps: + - id: check + type: deterministic + command: npm test + onNonZero: record +`); + const kernel = toKernelSpec(flow); + expect(kernel.steps[0]).toMatchObject({ on_non_zero: 'record' }); + expect(kernelToAuthoring(kernel)).toEqual(flow); + }); + + // The default is normalized away on BOTH sides, so a spec written before this + // field existed keeps its canonical bytes and therefore its spec hash. An + // author who spells out the default must land on those same bytes. + it('leaves the default policy out of the canonical bytes and the hash', () => { + const yaml = fixture('hello-deterministic.flow.yaml'); + const spelled = yaml.replace(' command: "echo hello"', ' command: "echo hello"\n onNonZero: fail'); + expect(spelled).not.toBe(yaml); + expect(compileYamlToCanonicalJson(spelled)) + .toBe(fixture('hello-deterministic.spec.canonical.json').trim()); + expect(compileAndHash(spelled).hash).toBe(fixture('hello-deterministic.spec.sha256').trim()); + expect(JSON.stringify(toKernelSpec(compileYaml(spelled)))).not.toContain('on_non_zero'); + }); + it('normalizes empty triggers exactly as kernel serialization does', () => { const yaml = `${fixture('hello-ladder.flow.yaml')}\ntriggers: []\n`; const flow = compileYaml(yaml); diff --git a/packages/sdk/tests/verb-field-lint.test.ts b/packages/sdk/tests/verb-field-lint.test.ts index b4255713..a853f188 100644 --- a/packages/sdk/tests/verb-field-lint.test.ts +++ b/packages/sdk/tests/verb-field-lint.test.ts @@ -82,6 +82,9 @@ const VERB_FIELD_VALUES: Record = { output: { type: 'object' }, cwd: '/tmp/foreign-cwd', transport: 'direct', + // Not `'fail'`: the default is normalized away on compile, so a default + // value could pass the YAML round trip on a step type that never allows it. + onNonZero: 'record', }; /** @@ -199,13 +202,14 @@ describe('closed per-verb step fields', () => { 'requirements', ]); expect(STEP_FIELDS_BY_TYPE).toEqual({ - deterministic: ['command', 'timeoutMs', 'lease_ms'], + deterministic: ['command', 'timeoutMs', 'lease_ms', 'onNonZero'], llm: ['prompt', 'model', 'cli', 'output'], agent: ['instruction', 'agent', 'cli', 'model', 'cwd', 'transport', 'surfaces', 'recoveryMode', 'permissions', 'output'], }); expect(CROSS_VERB_STEP_FIELDS.map(({ label }) => label).sort()).toEqual([ 'agent foreign command', 'agent foreign lease_ms', + 'agent foreign onNonZero', 'agent foreign prompt', 'agent foreign timeoutMs', 'deterministic foreign agent', @@ -224,6 +228,7 @@ describe('closed per-verb step fields', () => { 'llm foreign cwd', 'llm foreign instruction', 'llm foreign lease_ms', + 'llm foreign onNonZero', 'llm foreign permissions', 'llm foreign recoveryMode', 'llm foreign surfaces', diff --git a/packages/surface/src/context.ts b/packages/surface/src/context.ts index 38b48ce7..df6a30c6 100644 --- a/packages/surface/src/context.ts +++ b/packages/surface/src/context.ts @@ -42,6 +42,46 @@ export interface AgentOptions { transport?: 'direct' | 'relay'; } +/** + * What `f.run` resolves to under `onNonZero: 'record'`. + * + * Every field is read back from the step's journaled output envelope, never + * re-measured: `ok` is the recorded exit code, not a second execution. This is + * the point of the policy — ` || true` discards the code, so the + * branch below it can only ever be taken on faith. + */ +export interface RunResult { + /** Exactly `exitCode === 0`. */ + ok: boolean; + /** The exit code the kernel journaled for this command. */ + exitCode: number; + /** + * `stdout` then `stderr`, joined by a newline when both are nonempty — the + * string to hand an agent asked to repair what the command reported. It is a + * diagnostic convenience, not a reconstruction of the real interleaving: the + * two streams were captured separately and their true ordering is not + * journaled. Read `stdout` when you need to parse the command's own output. + */ + output: string; + /** stdout tail — the same string the default `f.run` resolves to. */ + stdout: string; + /** stderr tail, where a failing check usually explains itself. */ + stderr: string; +} + +export interface RunOptions { + /** Command lease: milliseconds or a duration such as "5m"; default 30s, maximum 15m. */ + timeout?: string | number; + /** + * What a nonzero exit means. `'fail'` (the default) throws, ending the flow + * at this line. `'record'` journals the exit code and output and resolves to + * a `RunResult`, so the body can hand a red check's own output to whatever is + * built to repair it. A timeout, a signal, or a step that never ran still + * throws under either policy: those produced no verdict to record. + */ + onNonZero?: 'fail' | 'record'; +} + export interface LlmOptions { /** JSON Schema checked before the result is accepted by the journal. */ output: Record; @@ -57,8 +97,20 @@ export interface LlmOptions { */ export interface Ctx extends Helpers { readonly mcp: Readonly Step>>>>; - /** Command lease: milliseconds or a duration such as "5m"; default 30s, maximum 15m. */ - run(command: string, options?: { timeout?: string | number }): Step; + /** + * Run a command. Resolves to its stdout tail, and a nonzero exit ends the + * flow — unless `onNonZero: 'record'` is declared, which resolves to a + * `RunResult` carrying the journaled exit code instead. + * + * The literal overloads come first so the default stays the common case: an + * omitted or `'fail'` policy keeps the `Step` every existing body is + * written against, and only the literal `'record'` widens the result. The + * third is for a policy held in a variable, where neither literal applies and + * the author has to narrow the union themselves. + */ + run(command: string, options?: RunOptions & { onNonZero?: 'fail' }): Step; + run(command: string, options: RunOptions & { onNonZero: 'record' }): Step; + run(command: string, options: RunOptions): Step; llm(strings: TemplateStringsArray, ...values: unknown[]): Step; /** JSON Schema validates the value at runtime; narrow unknown in author code. */ llm(prompt: string, options: LlmOptions): Step; diff --git a/packages/surface/src/index.ts b/packages/surface/src/index.ts index 4a6884f6..62b8332d 100644 --- a/packages/surface/src/index.ts +++ b/packages/surface/src/index.ts @@ -7,7 +7,7 @@ export type { WorkerSummary, CloudHelper, } from "./cloud.js"; -export type { AgentOptions, AgentResult, PermissionsSpec, LlmOptions, Ctx } from "./context.js"; +export type { AgentOptions, AgentResult, PermissionsSpec, LlmOptions, RunOptions, RunResult, Ctx } from "./context.js"; export { COMPLETION_REASONS, RUN_COMPLETION_REASONS, diff --git a/packages/surface/tests/run-on-non-zero.test-d.ts b/packages/surface/tests/run-on-non-zero.test-d.ts new file mode 100644 index 00000000..a1160c5f --- /dev/null +++ b/packages/surface/tests/run-on-non-zero.test-d.ts @@ -0,0 +1,67 @@ +// The `f.run` overload set for `onNonZero`, pinned at the type level. +// +// The whole feature is opt-in, so what matters most here is the NEGATIVE +// claim: an omitted or `'fail'` policy must keep resolving to `string`. If an +// overload ever widened the default to `string | RunResult`, every existing +// body that does `(await f.run(...)).trim()` would stop compiling — and the +// compiler is the only thing that can hold that line. + +import type { Ctx, RunOptions, RunResult } from '../src/index.js'; + +export async function runPolicyOverloads(f: Ctx): Promise { + // The default, and the explicit spelling of the default, both stay `string`. + // These annotations are the real assertion: a widened overload would resolve + // to `string | RunResult`, which does not assign to `string`. + const plain: string = await f.run('npm test'); + const timed: string = await f.run('npm test', { timeout: '5m' }); + const explicit: string = await f.run('npm test', { onNonZero: 'fail' }); + void [plain, timed, explicit]; + + // The literal `'record'` — and only it — widens the result. + const recorded: RunResult = await f.run('npm test', { onNonZero: 'record' }); + const ok: boolean = recorded.ok; + const exitCode: number = recorded.exitCode; + const streams: string[] = [recorded.output, recorded.stdout, recorded.stderr]; + void [ok, exitCode, streams]; + + // A policy held in a VARIABLE matches neither literal overload, so the + // author receives the union and has to narrow it. Silently picking one of + // the two branches for them would be a guess about which one they meant. + const chosen: RunOptions['onNonZero'] = Math.random() > 0.5 ? 'record' : 'fail'; + const union: string | RunResult = await f.run('npm test', { onNonZero: chosen }); + if (typeof union !== 'string') { + const narrowed: number = union.exitCode; + void narrowed; + } + + // @ts-expect-error A recorded result is not a string, so the default's callers cannot be silently widened. + const notAString: string = await f.run('npm test', { onNonZero: 'record' }); + void notAString; + // @ts-expect-error The default resolves to a string, which has no exit code to read. + void (await f.run('npm test')).exitCode; + // @ts-expect-error The union must be narrowed before either branch is read. + void (await f.run('npm test', { onNonZero: chosen })).exitCode; + // @ts-expect-error `ignore` is not one of the two policies the kernel enum names. + await f.run('npm test', { onNonZero: 'ignore' }); + // @ts-expect-error The policy is a string literal, not a boolean. + await f.run('npm test', { onNonZero: false }); + // @ts-expect-error Unknown option keys are refused, as everywhere else on the surface. + await f.run('npm test', { onNonzero: 'record' }); + + // A declared gate does not change what the policy resolves to: a named gate + // still yields the step's own result type on both sides. + const gatedDefault: string = await f.run('npm test') + .gate({ type: 'regex_match', pattern: 'ok', in_output_at: ['stdout_tail'] }); + const gatedRecord: RunResult = await f.run('npm test', { onNonZero: 'record' }) + .gate({ type: 'regex_match', pattern: 'ok', in_output_at: ['stdout_tail'] }); + void [gatedDefault, gatedRecord]; + + // A predicate gate judges the value the author received, so its parameter + // type follows the policy rather than being fixed at `string`. + await f.run('npm test').gate((value) => value.trim().length > 0); + await f.run('npm test', { onNonZero: 'record' }).gate((value) => value.exitCode === 0); + // @ts-expect-error On a recording step the predicate sees a RunResult, not a string. + await f.run('npm test', { onNonZero: 'record' }).gate((value) => value.trim() === ''); + // @ts-expect-error On a default step the predicate sees a string, which has no exit code. + await f.run('npm test').gate((value) => value.exitCode === 0); +} diff --git a/scripts/schema-constraints.mjs b/scripts/schema-constraints.mjs index 7c441756..d2d1a2a1 100644 --- a/scripts/schema-constraints.mjs +++ b/scripts/schema-constraints.mjs @@ -47,6 +47,9 @@ export function applyConstraints(defs, version) { Object.assign(defs.AgentSurfaces.properties.external.items, canonicalPath); for (const key of ['fileGlobs', 'networkAllowlist']) Object.assign(defs.PermissionsSpec.properties[key].items, nonempty); property('OutputContainsGate', 'value', nonempty); + // steps-green.ts: a gate naming nothing readable, or the same step twice, is + // the `|| true` failure it replaces, so the editor refuses it too. + property('StepsGreenGate', 'ids', { minItems: 1, uniqueItems: true, items: { type: 'string', pattern: '\\S' } }); // Runtime accepts this legacy spelling although ExitCodeGate omits it. defs.ExitCodeGate.properties.expect = { type: 'integer', const: 0, description: 'Legacy explicit success code. Only zero is supported.' }; defs.JsonSchemaGate.properties.schema = { $ref: '#/$defs/OutputSchema', description: 'JSON Schema object or boolean; references and termination are checked by flows check.' }; diff --git a/testdata/advisory-repair.flow.yaml b/testdata/advisory-repair.flow.yaml new file mode 100644 index 00000000..38850777 --- /dev/null +++ b/testdata/advisory-repair.flow.yaml @@ -0,0 +1,36 @@ +# Advisory deterministic outcomes (docs/SURFACE.md §2, "Repair before failure"). +# `check` is allowed to be red; its recorded envelope is what `report` reads, +# and the run's verdict comes from `steps_green` over the POST-repair `recheck` +# step, never from the immutable red outcome that started the repair. +version: '0.1.0' +name: advisory-repair +description: A recorded nonzero exit, read by a later step, gated on fresh evidence. +steps: + - id: check + type: deterministic + command: "node -e 'process.stderr.write(\"two tests failed\"); process.exit(1)'" + onNonZero: record + # Declaring the envelope schema is what makes the recorded outcome bindable. + verification: + type: json_schema + schema: + type: object + properties: + exit_code: { type: integer } + stdout_tail: { type: string } + stderr_tail: { type: string } + - id: report + type: deterministic + input: + failure: { step: check } + command: printf '%s' "$FLOWS_INPUT" + - id: recheck + type: deterministic + dependsOn: [report] + command: "node -e 'process.exit(0)'" + - id: green + type: deterministic + command: "true" + verification: + type: steps_green + ids: [recheck] diff --git a/testdata/advisory-repair.spec.canonical.json b/testdata/advisory-repair.spec.canonical.json new file mode 100644 index 00000000..d88bacc0 --- /dev/null +++ b/testdata/advisory-repair.spec.canonical.json @@ -0,0 +1 @@ +{"description":"A recorded nonzero exit, read by a later step, gated on fresh evidence.","name":"advisory-repair","steps":[{"command":"node -e 'process.stderr.write(\"two tests failed\"); process.exit(1)'","depends_on":[],"id":"check","max_iterations":1,"on_non_zero":"record","retry":{"initial_backoff_ms":100,"jitter_percent":20,"max_backoff_ms":60000,"multiplier":2},"type":"deterministic","verification":{"json_schema":{"properties":{"exit_code":{"type":"integer"},"stderr_tail":{"type":"string"},"stdout_tail":{"type":"string"}},"type":"object"}}},{"command":"printf '%s' \"$FLOWS_INPUT\"","depends_on":["check"],"id":"report","input":{"failure":{"step":"check"}},"max_iterations":1,"retry":{"initial_backoff_ms":100,"jitter_percent":20,"max_backoff_ms":60000,"multiplier":2},"type":"deterministic","verification":{}},{"command":"node -e 'process.exit(0)'","depends_on":["report"],"id":"recheck","max_iterations":1,"retry":{"initial_backoff_ms":100,"jitter_percent":20,"max_backoff_ms":60000,"multiplier":2},"type":"deterministic","verification":{"json_schema":true}},{"command":"true","depends_on":[],"id":"green","max_iterations":1,"retry":{"initial_backoff_ms":100,"jitter_percent":20,"max_backoff_ms":60000,"multiplier":2},"type":"deterministic","verification":{}},{"command":"node -e 'const input=JSON.parse(process.env.FLOWS_INPUT);\nconst ids=[\"recheck\"];\nconst red=[];\nfor(let i=0;i0){process.stderr.write('\\''steps_green: '\\''+red.join('\\''; '\\'')+'\\''\\n'\\'');process.exit(1);}\nprocess.stdout.write('\\''steps_green:pass'\\'');'","depends_on":["green","recheck"],"id":"green.green","input":{"0":{"step":"recheck"}},"max_iterations":1,"retry":{"initial_backoff_ms":100,"jitter_percent":20,"max_backoff_ms":60000,"multiplier":2},"type":"deterministic","verification":{}}],"version":"0.1.0"} diff --git a/testdata/advisory-repair.spec.sha256 b/testdata/advisory-repair.spec.sha256 new file mode 100644 index 00000000..d2f7887c --- /dev/null +++ b/testdata/advisory-repair.spec.sha256 @@ -0,0 +1 @@ +c08673d759047eb26f02fe0469c8543ca1539b298310b7bcf6faa08492374848 From 8c43cc085eaca027f7688958f4a91ea97e233918 Mon Sep 17 00:00:00 2001 From: Relayflow Date: Sun, 20 Sep 2026 17:48:22 +0000 Subject: [PATCH 2/3] docs: PR summary for advisory deterministic outcomes Co-Authored-By: Claude --- summary.md | 135 +++++++++++++++++++++++++++++++++++++++++++++++++++++ 1 file changed, 135 insertions(+) create mode 100644 summary.md diff --git a/summary.md b/summary.md new file mode 100644 index 00000000..89ab2872 --- /dev/null +++ b/summary.md @@ -0,0 +1,135 @@ +# Advisory deterministic outcomes: `onNonZero: record` and the `steps_green` gate + +Closes the "no advisory deterministic step" gap. v2 gated every `deterministic` +step on exit code with no opt-out, so every repair-before-failure flow +re-invented ` || true` — including this repo's own v1→v2 bridge +(`flows/spec-builder.ts`) and `flows/diagnose/orchestration.spec.ts`'s +`advisoryGate`. + +`|| true` discards the exit code. A later gate has nothing to read, so the flow +either re-runs the command or builds an evidence journal outside the kernel. +Worse, it is indistinguishable from success: a flow that forgets one gate ships +red work silently. + +## What an author writes now + +```ts +const tests = await f.run('npm test', { onNonZero: 'record' }); +if (!tests.ok) { + await f.agent('fixer', { task: `Fix these failures:\n${tests.output}` }); +} +``` + +```yaml +- id: check + type: deterministic + command: npm test + onNonZero: record # 'fail' (default) | 'record' + +- id: green + type: deterministic + command: ./ship.sh + verification: { type: steps_green, ids: [recheck] } +``` + +`tests.ok`, `tests.exitCode`, `tests.stdout` and `tests.stderr` are all read back +from the step's journaled `{exit_code, stdout_tail, stderr_tail}` envelope. +Nothing is re-run and nothing is re-measured — that is the claim `|| true` +cannot make. + +## Acceptance + +| Ticket clause | Where it is held | +| --- | --- | +| A deterministic step can be red without failing the run | `verify.rs` `a_recorded_nonzero_exit_passes_its_gate_and_names_the_code`; `relayflowd/tests/advisory_deterministic.rs` | +| Its exit code and output are readable by a later step from the journal | `authored-advisory-run-live.test.ts` "resolves to the journaled exit code instead of ending the flow" (asserts the repair branch received `2 failing tests`) | +| A gate can assert "these recorded steps were all green" without re-running | `steps-green.ts` + `advisory-outcome.test.ts` "steps_green lowering" | +| `|| true` is no longer the documented answer | `docs/SURFACE.md` §2 "Repair before failure: `onNonZero: record`" | + +## Design decisions worth reviewing + +**Recording is a policy about exit codes, and only about exit codes.** A +timeout, a signal, the executor's `-1` no-exit-status sentinel, a declared +content or schema gate, and a budget all stay fatal under either policy. `-1` +in particular is refused rather than handed back as a plausible code, mirroring +the kernel's own rule — otherwise a killed process would reach an author's +`if (!result.ok)` as a real red verdict. + +**Red is allowed, not invisible.** The kernel labels the gate +`exit_code:recorded` and names the code in the verification detail; +`flows check` annotates the step with `[onNonZero: record]`. Reporting it as an +ordinary `exit_code` gate would have moved the `|| true` ambiguity into the +inspection output. + +**The default is normalized away at both boundaries.** `onNonZero: fail` is +dropped in `compileStep` and again in `toKernelStep`, and `OnNonZero::is_default` +skips it on the Rust side — so every spec written before this field keeps its +exact canonical bytes and its exact `spec_hash`. A parity test asserts that +directly. + +**`steps_green` is compiler-lowered, not a kernel gate.** It becomes a separate +fatal deterministic step bound to the sources' envelopes, so the kernel never +sees the SDK spelling. `testdata/advisory-repair.*` pins the *lowered* output, +which is what proves the kernel parses what the SDK actually emits. The +generated gate never inherits its host's recording policy, or a red gate could +not fail anything. + +**Overload order on `f.run`.** The two literal overloads come first so an +omitted or `'fail'` policy keeps the `Step` every existing body is +written against; only the literal `'record'` widens the result. A third +union-returning overload covers a policy held in a variable, where the author +narrows it themselves rather than having one branch guessed for them. +`packages/surface/tests/run-on-non-zero.test-d.ts` pins all of this, including +that a recorded result does *not* assign to `string`. + +**A misspelled policy is refused, never defaulted.** Falling back to the fatal +default would delete the branch the author wrote below the call. Refused in +`validateSpec` for declarative specs and before any ordinal is consumed in +`f.run`. + +## Bug found and fixed along the way + +An authored `.gate()` lowers to its own kernel step (`.gate`) that runs +*after* the producer. `readCompletedStepOutput` read only the producer's +`step.completed` entry, so a failed gate on a successful command resolved the +operation as though it had passed — the child run was terminally `failed` and +the author never heard about it. + +This reproduces on the **default** policy too, so it predates this change, but +it is the same invisibility the recording policy exists to remove, one layer up, +and `record` would have widened its reach. `readCompletedStepOutput` now also +checks the run's own terminal reason from the entries it already reads. +`authored-advisory-run-live.test.ts` covers it under both policies and asserts +the failed step is named as `run-1.gate`, not the producer. + +## Evidence + +- `cd kernel && sh ../ops/cargo.sh test` — all suites pass, 0 failures + (includes 6 new `verify.rs` tests, 2 new `spec/tests.rs` tests, 2 new + `relayflowd/tests/advisory_deterministic.rs` integration tests, and the new + `spec_parity.rs` fixture test). +- `npm --prefix packages/schema test` — 78 pass (the new `advisory-repair` + fixture is picked up by the existing parity walk). +- `npm --prefix packages/surface test` — 46 pass, 9 files. +- `npm --prefix packages/surface run typecheck:regressions` — clean, including + the new `run-on-non-zero.test-d.ts`. +- `npm --prefix packages/sdk test` — 2401 pass. 7 files fail, **all confirmed + failing identically at pristine HEAD** (verified by stashing and rebuilding): + `live-kernel`, `webhook-live`, `mcp`, `provider-trigger-executor`, + `communication-mixed-resume` need `kernel/target/{debug,release}/relayflowd`, + which `ops/cargo.sh` deliberately builds outside the repo; + `stuck-run-triage` and `authored-node-runtime` hit a vitest module-duplication + issue for flow files imported from outside the package root. +- Schema regeneration is reproducible: regenerating over the tracked file is + byte-identical. + +Four suites needed updating for this change rather than being broken by it: +`verb-field-lint` (its fail-closed pin demands a sample value for every new +step field), `validate` and `gate-contract` (the `unknown_gate_kind` +enumeration), and `preflight` (its converse test requires every declared +refusal kind to be reachable through the public boundary — the three +`steps_green` kinds now have scenarios). + +## Not in scope + +Retry policy (`maxIterations`) and agent-step failure semantics, per the ticket. From 88f6d179812fce6345c5afa40ac79fff41b2011a Mon Sep 17 00:00:00 2001 From: khaliqgant Date: Tue, 22 Sep 2026 23:21:06 -0700 Subject: [PATCH 3/3] fix(sdk): register run operations in the authored step DAG The merge hoisted runOperation for onNonZero but dropped the node argument, so f.run steps never entered the step graph and status --json reported after: undefined. Restore the {} node (a command is not a display label) on both the record and fail paths, matching main. --- packages/sdk/src/authored-flow-executor.ts | 6 ++++++ 1 file changed, 6 insertions(+) diff --git a/packages/sdk/src/authored-flow-executor.ts b/packages/sdk/src/authored-flow-executor.ts index 83a75f2b..e24383bd 100644 --- a/packages/sdk/src/authored-flow-executor.ts +++ b/packages/sdk/src/authored-flow-executor.ts @@ -406,6 +406,9 @@ export async function executeAuthoredFlow( () => observeStep(id, 'deterministic', () => lowerDeterministic.recording(id, command, leaseMs, recordOp.namedGate), options.onProgress), lifecycle, + // No label: a command is not a display name. It carries literal + // tokens and URLs, and no prefix of it is safe to show (displayLabel). + {}, ); return trackStep(authoredSteps, recordOp); } @@ -415,6 +418,9 @@ export async function executeAuthoredFlow( () => assertOperationAllowed('run', definition.name, requestedCompletion), () => observeStep(id, 'deterministic', () => lowerDeterministic(id, command, false, leaseMs, runOp.namedGate), options.onProgress), lifecycle, + // No label: a command is not a display name. It carries literal + // tokens and URLs, and no prefix of it is safe to show (displayLabel). + {}, ); return trackStep(authoredSteps, runOp); }