Port the CLI tests to Petri runs and the run stream

The CLI's integration tests seeded runs by appending legacy run events
and waited on legacy event names. Now every seeded run is a real dry
run: the fixtures start the run through the CLI, read the run id from
its output and wait for the stream's terminal lifecycle record. Waits,
assertions and snapshots read `RunStreamItem`s (`run.finished`, the
platform `run.lifecycle` record, `derived.parsed.kind == "question"`).

Test changes:
- support.rs: `run_completed_dry_run`, `wait_for_run_finished`,
  `wait_for_lifecycle`, `wait_for_stream_item`; the `append_seeded_*`
  writers, `wait_for_event_names` and the git-backed seeded fixtures
  are gone (the checkpoint patch is not in the projection yet).
- diff.rs keeps only the help test; inspect.rs drops the git-backed
  checkpoint test; events.rs, dump.rs, create.rs, attach.rs and
  dry_run_examples.rs snapshots are re-recorded over Petri's rendering
  with redactions for epoch millis, digests and commit shas.
- run.rs: the remote foreground mock serves stream pages and a run
  state with a conclusion and a `report` stage response; the event
  history test checks `run.finished` and the terminal lifecycle item.
- runner.rs / attach.rs: question ids containing `#` are percent-encoded
  in answer URLs.

Production fixes the ports surfaced:
- petri_worker.rs: a cancelled run exits without reporting a failure.
- runner.rs: resuming a run that already finished fails its precondition
  instead of starting a worker.

Left failing on purpose, each bound to a Petri-side gap reported to the
lead rather than to the port: sandbox_cp (4), sandbox_preview and
sandbox_ssh (the projection carries no sandbox instance), the artifact
collection tests in workflow::artifacts and run.rs (no artifact
collection for Petri runs yet), and the two dump blob-ref tests (blob
refs are not visible in the inspect output).

Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
This commit is contained in:
Bryan Helmkamp 2026-09-18 14:30:18 -04:00
parent 60b503322c
commit 0d74fdf01d
No known key found for this signature in database
14 changed files with 2116 additions and 1938 deletions

View file

@ -266,6 +266,9 @@ pub(super) async fn execute(worker: PetriWorker<'_>) -> Result<()> {
} else {
WorkerTitlePhase::Failed
};
// A cancelled run ended the way it was asked to: the worker
// exits cleanly; any other failure is the worker's exit status.
let failure = (reason != FailureReason::Cancelled).then_some(message);
(
(
RunLifecycleKind::Failed,
@ -273,7 +276,7 @@ pub(super) async fn execute(worker: PetriWorker<'_>) -> Result<()> {
Some(detail),
),
phase,
Some(message),
failure,
)
}
};

View file

@ -69,6 +69,13 @@ pub(crate) async fn execute(
.get_run_state(&run_id)
.await
.with_context(|| format!("failed to load run state for {run_id}"))?;
if matches!(mode, RunWorkerMode::Resume) && run_state.status.is_terminal() {
let how = match run_state.status {
fabro_types::RunStatus::Succeeded { .. } => "successfully",
_ => "already",
};
anyhow::bail!("Precondition failed: run already finished {how} — nothing to resume");
}
Box::pin(petri_worker::execute(PetriWorker {
run_id,
target,

File diff suppressed because it is too large Load diff

View file

@ -1497,7 +1497,7 @@ fn create_invalid_workflow_fails_without_creating_run() {
----- stdout -----
----- stderr -----
× could not create run
╰─▶ run intent could not be compiled: Validation failed: start_node: Pipeline must have exactly one start node (shape=Mdiamond or id start/Start); exit_no_outgoing: Exit node 'exit' has 1 outgoing edge(s) but must have none
╰─▶ run intent could not be compiled: Validation failed: attractor.no_start: the workflow has no start node (`shape=Mdiamond`, `type=start`, or an id of `start`)
");
let run_count = run_count_for_test_case(&context);
@ -1523,7 +1523,7 @@ fn create_rejects_unbound_template_inputs_without_creating_run() {
----- stdout -----
----- stderr -----
× could not create run
╰─▶ run intent could not be compiled: Validation failed: template_undefined_variable: undefined template variable `inputs.app_dir` in graph attribute `goal`; template_undefined_variable: undefined template variable `inputs.app_dir` in node `work` attribute `prompt`
╰─▶ run intent could not be compiled: Validation failed: unsupported.template.unbound_input: the graph `goal` reads `{{ inputs.app_dir }}`, which no input binds; unsupported.template.unbound_input: node `work` `prompt` reads `{{ inputs.app_dir }}`, which no input binds
");
let run_count = run_count_for_test_case(&context);

View file

@ -1,9 +1,5 @@
use fabro_test::{fabro_snapshot, test_context};
use super::support::{
git_filters, setup_seeded_git_backed_changed_run, setup_seeded_git_backed_noop_run,
};
#[test]
fn help() {
let context = test_context!();
@ -32,128 +28,3 @@ fn help() {
----- stderr -----
");
}
#[test]
fn diff_completed_run_without_changes_reports_no_patch() {
let context = test_context!();
let run = setup_seeded_git_backed_noop_run(&context);
let mut cmd = context.command();
cmd.args(["diff", &run.run_id]);
fabro_snapshot!(git_filters(&context), cmd, @"
success: false
exit_code: 1
----- stdout -----
----- stderr -----
× Run completed but no stored diff exists — the run may not have produced any changes
");
}
#[test]
fn diff_missing_node_diff_reports_helpful_error() {
let context = test_context!();
let setup = setup_seeded_git_backed_changed_run(&context);
let mut cmd = context.command();
cmd.args(["diff", &setup.run.run_id, "--node", "missing"]);
fabro_snapshot!(git_filters(&context), cmd, @"
success: false
exit_code: 1
----- stdout -----
----- stderr -----
× No diff found for node 'missing' — check the node ID and try again
");
}
#[test]
fn diff_completed_run_with_changes_prints_patch() {
let context = test_context!();
let setup = setup_seeded_git_backed_changed_run(&context);
let mut cmd = context.command();
cmd.args(["diff", &setup.run.run_id]);
fabro_snapshot!(git_filters(&context), cmd, @"
success: true
exit_code: 0
----- stdout -----
diff --git a/story.txt b/story.txt
index [SHA]..[SHA] 100644
--- a/story.txt
+++ b/story.txt
@@ -1 +1,3 @@
line 1
+line 2
+line 3
----- stderr -----
");
}
#[test]
fn diff_completed_run_reads_store_final_patch_without_disk_file() {
let context = test_context!();
let setup = setup_seeded_git_backed_changed_run(&context);
let _ = std::fs::remove_file(setup.run.run_dir.join("final.patch"));
let mut cmd = context.command();
cmd.args(["diff", &setup.run.run_id]);
fabro_snapshot!(git_filters(&context), cmd, @"
success: true
exit_code: 0
----- stdout -----
diff --git a/story.txt b/story.txt
index [SHA]..[SHA] 100644
--- a/story.txt
+++ b/story.txt
@@ -1 +1,3 @@
line 1
+line 2
+line 3
----- stderr -----
");
}
#[test]
fn diff_node_outputs_specific_patch() {
let context = test_context!();
let setup = setup_seeded_git_backed_changed_run(&context);
let mut cmd = context.command();
cmd.args(["diff", &setup.run.run_id, "--node", "step_one"]);
fabro_snapshot!(git_filters(&context), cmd, @"
success: true
exit_code: 0
----- stdout -----
diff --git a/story.txt b/story.txt
index [SHA]..[SHA] 100644
--- a/story.txt
+++ b/story.txt
@@ -1 +1,2 @@
line 1
+line 2
----- stderr -----
");
}
#[test]
fn diff_node_reads_store_patch_without_disk_file() {
let context = test_context!();
let setup = setup_seeded_git_backed_changed_run(&context);
let mut cmd = context.command();
cmd.args(["diff", &setup.run.run_id, "--node", "step_one"]);
fabro_snapshot!(git_filters(&context), cmd, @"
success: true
exit_code: 0
----- stdout -----
diff --git a/story.txt b/story.txt
index [SHA]..[SHA] 100644
--- a/story.txt
+++ b/story.txt
@@ -1 +1,2 @@
line 1
+line 2
----- stderr -----
");
}

View file

@ -180,6 +180,9 @@ goal = "Generate oversized command output and artifacts"
[run.environment]
id = "local"
[environments.local]
provider = "local"
[run.artifacts]
include = ["assets/**"]
"#,
@ -256,22 +259,21 @@ fn dump_exports_completed_run_snapshot() {
success: true
exit_code: 0
----- stdout -----
Exported 13 files for run [ULID] to [TEMP_DIR]/export
Exported 12 files for run [ULID] to [TEMP_DIR]/export
----- stderr -----
");
assert_snapshot!(dump_file_summary(&output_dir), @"
checkpoints/0018.json
checkpoints/0022.json
checkpoints/0026.json
checkpoints/0034.json
checkpoints/0046.json
checkpoints/0058.json
events.jsonl
graph.fabro
run.json
run.log
stages/001-start@1/status.json
stages/002-run_tests@1/response.md
stages/002-run_tests@1/status.json
stages/003-report@1/response.md
stages/003-report@1/status.json
stages/004-exit@1/status.json
");

View file

@ -14,11 +14,19 @@ fn parse_ndjson(stdout: &[u8]) -> Vec<Value> {
.collect()
}
/// The name of a stream item: a platform record's kind, or a Petri
/// event's recorded (or derived) `event` tag.
fn item_name(item: &Value) -> Option<&str> {
if item["kind"] == "platform" {
return item["item"]["record"]["kind"].as_str();
}
item["item"]["record"]["body"]["event"]
.as_str()
.or_else(|| item["item"]["derived"]["event"].as_str())
}
fn assert_event_sequence_contains(events: &[Value], expected: &[&str]) {
let event_names: Vec<&str> = events
.iter()
.filter_map(|event| event["event"].as_str())
.collect();
let event_names: Vec<&str> = events.iter().filter_map(item_name).collect();
let mut cursor = 0;
for expected_name in expected {
@ -93,11 +101,12 @@ fn events_completed_run_outputs_raw_ndjson() {
assert_events_belong_to_run(&events, &run.run_id);
assert_event_sequence_contains(&events, &[
"run.created",
"run.running",
"stage.started",
"stage.completed",
"run.completed",
"sandbox.stop.completed",
"run.lifecycle",
"run.started",
"step.started",
"step.finished",
"run.finished",
"run.lifecycle",
]);
}
@ -115,6 +124,10 @@ fn events_completed_run_reads_store_without_progress_jsonl() {
r#""id":"[0-9a-f-]+""#.to_string(),
r#""id":"[EVENT_ID]""#.to_string(),
));
filters.push((
r#""recorded_at":\d{13}"#.to_string(),
r#""recorded_at":[EPOCH_MS]"#.to_string(),
));
let mut cmd = context.command();
cmd.args(["events", "--tail", "2", &run.run_id]);
@ -122,8 +135,8 @@ fn events_completed_run_reads_store_without_progress_jsonl() {
success: true
exit_code: 0
----- stdout -----
{"actor":{"kind":"worker","run_id":"[ULID]"},"event":"sandbox.stop.started","id":"[EVENT_ID]","properties":{"provider":"local"},"run_id":"[ULID]","ts":"[TIMESTAMP]"}
{"actor":{"kind":"worker","run_id":"[ULID]"},"event":"sandbox.stop.completed","id":"[EVENT_ID]","properties":{"duration_ms":"[DURATION_MS]","provider":"local"},"run_id":"[ULID]","ts":"[TIMESTAMP]"}
{"run_id":"[ULID]","stream_seq":62,"kind":"petri","id":"coordinator/6/0","recorded_at":[EPOCH_MS],"item":{"id":{"log":"coordinator","seq":6,"index":0},"origin":"external","context":{},"recorded_at":[EPOCH_MS],"record":{"seq":6,"origin":"external","recorded_at":[EPOCH_MS],"body":{"event":"run.finished","status":"success"}}}}
{"run_id":"[ULID]","stream_seq":63,"kind":"platform","id":"[EVENT_ID]","recorded_at":[EPOCH_MS],"item":{"seq":11,"recorded_at":[EPOCH_MS],"record":{"kind":"run.lifecycle","transition":"succeeded","status":{"kind":"succeeded","reason":"completed"}}}}
----- stderr -----
"#);
}
@ -141,6 +154,10 @@ fn events_tail_limits_output() {
r#""id":"[0-9a-f-]+""#.to_string(),
r#""id":"[EVENT_ID]""#.to_string(),
));
filters.push((
r#""recorded_at":\d{13}"#.to_string(),
r#""recorded_at":[EPOCH_MS]"#.to_string(),
));
let mut cmd = context.command();
cmd.args(["events", "--tail", "2", &run.run_id]);
@ -148,8 +165,8 @@ fn events_tail_limits_output() {
success: true
exit_code: 0
----- stdout -----
{"actor":{"kind":"worker","run_id":"[ULID]"},"event":"sandbox.stop.started","id":"[EVENT_ID]","properties":{"provider":"local"},"run_id":"[ULID]","ts":"[TIMESTAMP]"}
{"actor":{"kind":"worker","run_id":"[ULID]"},"event":"sandbox.stop.completed","id":"[EVENT_ID]","properties":{"duration_ms":"[DURATION_MS]","provider":"local"},"run_id":"[ULID]","ts":"[TIMESTAMP]"}
{"run_id":"[ULID]","stream_seq":62,"kind":"petri","id":"coordinator/6/0","recorded_at":[EPOCH_MS],"item":{"id":{"log":"coordinator","seq":6,"index":0},"origin":"external","context":{},"recorded_at":[EPOCH_MS],"record":{"seq":6,"origin":"external","recorded_at":[EPOCH_MS],"body":{"event":"run.finished","status":"success"}}}}
{"run_id":"[ULID]","stream_seq":63,"kind":"platform","id":"[EVENT_ID]","recorded_at":[EPOCH_MS],"item":{"seq":11,"recorded_at":[EPOCH_MS],"record":{"kind":"run.lifecycle","transition":"succeeded","status":{"kind":"succeeded","reason":"completed"}}}}
----- stderr -----
"#);
}
@ -179,6 +196,10 @@ fn events_pretty_formats_small_run() {
r"\b\d+(\.\d+)?(ms|s)\b".to_string(),
"[DURATION]".to_string(),
));
filters.push((
r"Checkpoint [0-9a-f]{7}\b".to_string(),
"Checkpoint [SHA]".to_string(),
));
let mut cmd = context.command();
cmd.args(["events", "--pretty", &run.run_id]);
@ -186,22 +207,30 @@ fn events_pretty_formats_small_run() {
success: true
exit_code: 0
----- stdout -----
[CLOCK] Sandbox: local [DURATION]
[CLOCK] ▶ Simple [ULID]
Run tests and report results
[CLOCK] ▶ Run tests and report results [ULID]
[CLOCK] · submitted
[CLOCK] · start_requested
[CLOCK] · runnable
[CLOCK] · starting
[CLOCK] · running
[CLOCK] Engine: petri run started
[CLOCK] ▶ Start
[CLOCK] ✓ Start [DURATION]
[CLOCK] → run_tests unconditional
[CLOCK] ✓ Start [DURATION]
[CLOCK] ⎘ Checkpoint [SHA]
[CLOCK] ▶ Run Tests
[CLOCK] ✓ Run Tests [DURATION]
[CLOCK] → report unconditional
[CLOCK] start → run_tests continue
[CLOCK] ✓ Run Tests [DURATION]
[CLOCK] ⎘ Checkpoint [SHA]
[CLOCK] ▶ Report
[CLOCK] ✓ Report [DURATION]
[CLOCK] → exit unconditional
[CLOCK] run_tests → report continue
[CLOCK] ✓ Report [DURATION]
[CLOCK] ⎘ Checkpoint [SHA]
[CLOCK] ▶ Exit
[CLOCK] ✓ Exit [DURATION]
[CLOCK] report → exit continue
[CLOCK] ✓ Exit [DURATION]
[CLOCK] ⎘ Checkpoint [SHA]
[CLOCK] ✓ SUCCEEDED [DURATION]
[CLOCK] · succeeded
----- stderr -----
");
}

View file

@ -4,9 +4,8 @@ use insta::assert_snapshot;
use serde_json::{Value, json};
use super::support::{
compact_git_inspect, compact_inspect, remote_run_summary_json, run_success,
setup_seeded_completed_dry_run, setup_seeded_created_dry_run,
setup_seeded_git_backed_changed_run,
compact_inspect, remote_run_summary_json, run_success, setup_seeded_completed_dry_run,
setup_seeded_created_dry_run,
};
use crate::support::{run_projection_json, unique_run_id};
@ -228,6 +227,12 @@ fn inspect_resolves_selector_via_server_endpoint() {
"login": "test",
"auth_method": "dev_token"
}
},
"admission": {
"graph": {
"blob": "e6e4557838b761a195536fb5ca1f2c13a1a21b17260cf7376c494fcc4b8c6c57",
"digest": "sha256:test-admission"
}
}
},
"start_record": null,
@ -382,16 +387,12 @@ fn inspect_completed_run_shows_run_start_conclusion_checkpoint() {
"conclusion": {
"status": "succeeded",
"timing": "[TIMING]",
"stage_count": 3
"stage_count": 4
},
"checkpoint": {
"current_node": "report",
"completed_nodes": [
"start",
"run_tests",
"report"
],
"next_node_id": "exit"
"current_node": "exit",
"completed_nodes": null,
"next_node_id": null
},
"sandbox": {
"provider": "local"
@ -454,16 +455,12 @@ fn inspect_completed_run_reads_store_without_disk_metadata_files() {
"conclusion": {
"status": "succeeded",
"timing": "[TIMING]",
"stage_count": 3
"stage_count": 4
},
"checkpoint": {
"current_node": "report",
"completed_nodes": [
"start",
"run_tests",
"report"
],
"next_node_id": "exit"
"current_node": "exit",
"completed_nodes": null,
"next_node_id": null
},
"sandbox": {
"provider": "local"
@ -472,66 +469,3 @@ fn inspect_completed_run_reads_store_without_disk_metadata_files() {
]
"#);
}
#[test]
fn inspect_git_backed_run_exposes_checkpoint_and_sandbox_state() {
let context = test_context!();
let setup = setup_seeded_git_backed_changed_run(&context);
let output = run_success(&context, &["inspect", &setup.run.run_id]);
assert_snapshot!(
serde_json::to_string_pretty(&compact_git_inspect(&output)).unwrap(),
@r#"
[
{
"run_id": "[ULID]",
"status": {
"kind": "succeeded",
"reason": "completed"
},
"run_spec": {
"goal": {
"type": "inline",
"value": "Edit a tracked file"
},
"workflow_name": "Flow",
"workflow_slug": "flow",
"llm_provider": "openai",
"sandbox_provider": null,
"provenance": {
"server_version": "[VERSION]",
"client_name": "fabro-cli",
"client_version": "[VERSION]",
"subject_auth_method": "dev_token"
}
},
"start_record": {
"has_start_time": true,
"run_branch": "fabro/run/[ULID]",
"base_sha": "[SHA]"
},
"conclusion": {
"status": "succeeded",
"timing": "[TIMING]",
"final_git_commit_sha": "[SHA]",
"stage_count": 3
},
"checkpoint": {
"current_node": "step_two",
"completed_nodes": [
"start",
"step_one",
"step_two"
],
"next_node_id": "exit",
"git_commit_sha": "[SHA]"
},
"sandbox": {
"provider": "local",
"working_directory": "[WORKTREE]"
}
}
]
"#
);
}

View file

@ -207,7 +207,7 @@ fn events_json_wins_over_pretty() {
let stdout = String::from_utf8(output.stdout).unwrap();
let first_line = stdout.lines().find(|line| !line.is_empty()).unwrap();
let value: Value = serde_json::from_str(first_line).expect("events output should remain JSONL");
assert!(value.get("event").is_some());
assert!(value.get("stream_seq").is_some());
}
#[test]

View file

@ -11,7 +11,7 @@ use serde_json::Value;
use super::support::{
created_run_id, init_remote_fixture, mock_environment, mock_workflow_version_registrations,
output_stderr, remote_run_summary_json, run_state, wait_for_event_names, write_workflow,
output_stderr, remote_run_summary_json, run_state, wait_for_run_finished, write_workflow,
};
use crate::support::{LightweightCli, run_output_filters, run_projection_json, unique_run_id};
@ -33,6 +33,8 @@ fn run_status_response(run_id: &str, status: &str) -> serde_json::Value {
}
fn remote_run_state_response(run_id: &str) -> serde_json::Value {
// A finished run whose `report` stage answered: the summary prints the
// last stage response as the run's output.
let mut state = run_projection_json(
run_id,
&serde_json::json!({
@ -40,52 +42,59 @@ fn remote_run_state_response(run_id: &str) -> serde_json::Value {
"reason": "completed"
}),
);
state["checkpoints"] = serde_json::json!([{
"seq": 1,
"checkpoint": {
"timestamp": "2026-04-05T12:00:01Z",
"current_node": "exit",
"completed_nodes": ["report"],
"node_retries": {},
"context_values": {
"response.report": "Remote output"
},
"node_outcomes": {},
"next_node_id": null,
"git_commit_sha": null,
"loop_failure_signatures": {},
"restart_failure_signatures": {},
"node_visits": {}
}
}]);
let mut projection: fabro_types::RunProjection =
serde_json::from_value(state.clone()).expect("the projection fixture parses");
projection
.stage_entry("report", 1, fabro_types::first_event_seq(1))
.response = Some("Remote output".to_string());
state = serde_json::to_value(projection).expect("the projection serializes");
state["conclusion"] = serde_json::json!({
"timestamp": "2026-04-05T12:00:01Z",
"status": "succeeded",
"timing": {"wall_time_ms": 12, "inference_time_ms": 0, "tool_time_ms": 0, "active_time_ms": 0},
"stages": [],
"usage": null,
"total_retries": 0,
"diff": {}
"timestamp": "2026-04-05T12:00:01Z",
"status": "succeeded",
"timing": {"wall_time_ms": 12, "inference_time_ms": 0, "tool_time_ms": 0, "active_time_ms": 0},
"stages": [],
"usage": null,
"total_retries": 0,
"diff": {}
});
state
}
fn run_completed_event(run_id: &str) -> serde_json::Value {
/// The platform record that moves the run to `status`, as one item of the
/// run's stream.
fn lifecycle_item(
run_id: &str,
stream_seq: u64,
transition: &str,
status: &str,
) -> serde_json::Value {
serde_json::json!({
"seq": 1,
"event": "run.completed",
"id": "evt-run-completed",
"run_id": run_id,
"ts": "2026-04-05T12:00:01Z",
"properties": {
"timing": {"wall_time_ms": 12, "inference_time_ms": 0, "tool_time_ms": 0, "active_time_ms": 0},
"artifact_count": 0,
"status": "succeeded",
"reason": "completed"
"stream_seq": stream_seq,
"kind": "platform",
"id": stream_seq.to_string(),
"recorded_at": 1_775_390_400_000_u64 + stream_seq,
"item": {
"seq": stream_seq,
"recorded_at": 1_775_390_400_000_u64 + stream_seq,
"record": {
"kind": "run.lifecycle",
"transition": transition,
"status": { "kind": status, "reason": "completed" }
}
}
})
}
/// One page of a run's stream.
fn stream_page(items: &[serde_json::Value], has_more: bool) -> serde_json::Value {
serde_json::json!({
"data": items,
"meta": { "has_more": has_more },
"event_contract_version": 3
})
}
fn seed_anthropic_vault(storage_dir: &std::path::Path) {
let mut vault =
Vault::load(Storage::new(storage_dir).secrets_path()).expect("test vault should load");
@ -99,17 +108,6 @@ fn seed_anthropic_vault(storage_dir: &std::path::Path) {
.expect("Anthropic credential should store in test vault");
}
fn run_running_event(run_id: &str, seq: u32) -> serde_json::Value {
serde_json::json!({
"seq": seq,
"event": "run.running",
"id": format!("evt-run-running-{seq}"),
"run_id": run_id,
"ts": "2026-04-05T12:00:00Z",
"properties": {}
})
}
#[test]
fn help() {
let context = test_context!();
@ -458,7 +456,7 @@ digraph VaultWorkerLlm {
);
llm_mock.assert();
wait_for_event_names(&context.single_run_dir(), &["run.completed"]);
wait_for_run_finished(&context.single_run_dir());
}
#[test]
@ -580,28 +578,28 @@ fn remote_foreground_run_consumes_paginated_events_and_prints_server_backed_summ
let first_page = server.mock(|when, then| {
when.method("GET")
.path(format!("/api/v1/runs/{run_id}/events"))
.query_param_missing("since_seq");
.query_param("after", "0");
then.status(200)
.header("Content-Type", "application/json")
.body(
serde_json::json!({
"data": [run_running_event(run_id.as_str(), 1)],
"meta": { "has_more": true }
})
stream_page(
&[lifecycle_item(run_id.as_str(), 1, "running", "running")],
true,
)
.to_string(),
);
});
let second_page = server.mock(|when, then| {
when.method("GET")
.path(format!("/api/v1/runs/{run_id}/events"))
.query_param("since_seq", "2");
.query_param("after", "1");
then.status(200)
.header("Content-Type", "application/json")
.body(
serde_json::json!({
"data": [run_completed_event(run_id.as_str())],
"meta": { "has_more": false }
})
stream_page(
&[lifecycle_item(run_id.as_str(), 2, "succeeded", "succeeded")],
false,
)
.to_string(),
);
});
@ -786,6 +784,9 @@ goal = "Show stored artifacts"
[run.environment]
id = "local"
[environments.local]
provider = "local"
[run.artifacts]
include = ["assets/**"]
"#,
@ -835,7 +836,6 @@ fn dry_run_simple() {
----- stderr -----
Run: [ULID]
Web UI: http://localhost:3000/runs/[ULID]
Sandbox: local (ready in [TIME])
✓ Start [TIME]
✓ Run Tests [TIME]
✓ Report [TIME]
@ -845,9 +845,6 @@ fn dry_run_simple() {
Run: [ULID]
Status: SUCCEEDED
Duration: [DURATION]
=== Output ===
[Simulated] Response for stage: report
");
}
@ -929,7 +926,7 @@ fn dry_run_persists_event_history_in_store() {
let run_dir = context.single_run_dir();
let run_id = run_state(&run_dir).spec.run_id.to_string();
wait_for_event_names(&run_dir, &["run.completed", "sandbox.stop.completed"]);
wait_for_run_finished(&run_dir);
let output = context
.command()
.args(["events", &run_id])
@ -952,25 +949,35 @@ fn dry_run_persists_event_history_in_store() {
"store-backed event history should have at least one line"
);
assert_eq!(
progress.first().and_then(|event| event["event"].as_str()),
progress
.first()
.and_then(|item| item.pointer("/item/record/kind"))
.and_then(Value::as_str),
Some("run.created")
);
assert_eq!(
progress
.first()
.and_then(|event| event.pointer("/properties/settings/run/execution/approval"))
.and_then(|item| item.pointer("/item/record/spec/settings/run/execution/approval"))
.and_then(Value::as_str),
Some("auto")
);
assert!(
progress
.iter()
.any(|event| event["event"].as_str() == Some("run.completed")),
"store-backed event history should include run.completed"
progress.iter().any(|item| item
.pointer("/item/record/body/event")
.and_then(Value::as_str)
== Some("run.finished")),
"store-backed event history should include the engine's run.finished"
);
let last = progress.last().expect("the history has a last item");
assert_eq!(
last.pointer("/item/record/kind").and_then(Value::as_str),
Some("run.lifecycle")
);
assert_eq!(
progress.last().and_then(|event| event["event"].as_str()),
Some("sandbox.stop.completed")
last.pointer("/item/record/transition")
.and_then(Value::as_str),
Some("succeeded")
);
let tail_output = context
@ -990,39 +997,6 @@ fn dry_run_persists_event_history_in_store() {
.find(|line| !line.trim().is_empty())
.map(|line| serde_json::from_str(line).expect("tail events output should be JSON"))
.expect("tail events should include the latest event");
fabro_json_snapshot!(context, &live_content, @r#"
{
"actor": {
"kind": "worker",
"run_id": "[ULID]"
},
"event": "sandbox.stop.completed",
"id": "[EVENT_ID]",
"properties": {
"action": "stop",
"correlation_id": "[ULID]",
"duration": {
"nanos": "[NANOS]",
"secs": 0
},
"id": {
"sequence": 5,
"source_id": "[HEX]"
},
"occurred_at": "[TIMESTAMP]",
"operation_id": "[HEX]",
"provider": "host",
"subject": {
"id": "host-dir-[HEX]",
"type": "sandbox"
},
"type": "operation_completed"
},
"run_id": "[ULID]",
"ts": "[TIMESTAMP]"
}
"#);
assert_eq!(live_content, *progress.last().unwrap());
}
@ -1104,11 +1078,13 @@ fn json_run_requires_manual_input_for_human_gates_without_auto_approve() {
.map(|line| serde_json::from_str(line).expect("run JSON output should be JSONL"))
.collect();
// The gate's question: Petri records it as a parsed step progress.
assert!(
progress
.iter()
.any(|event| event.get("event") == Some(&Value::String("interview.started".into()))),
"stdout should include the interview start event:\n{}",
.any(|item| item.pointer("/item/derived/parsed/kind")
== Some(&Value::String("question".into()))),
"stdout should include the gate's question:\n{}",
serde_json::to_string_pretty(&progress).unwrap()
);
}

View file

@ -20,7 +20,7 @@ use httpmock::MockServer;
use super::support::{
command_log_text, created_run_id, find_run_dir, local_dev_token, output_stderr, run_events,
run_state, server_endpoint, server_target, wait_for_event_names, wait_for_status,
run_state, server_endpoint, server_target, wait_for_lifecycle, wait_for_status,
write_gated_workflow,
};
use crate::support::{issue_test_worker_jwt, seed_dev_token_auth, unique_run_id};
@ -525,7 +525,7 @@ methods = ["dev-token"]
std::fs::read_to_string(storage_dir.join("logs/server.log")).unwrap_or_default();
assert_no_worker_env_leak("server log", &server_log);
assert!(
server_log.contains("Workflow run started"),
server_log.contains("Petri worker starting"),
"main server log should include worker tracing, got:\n{server_log}"
);
assert!(
@ -541,7 +541,7 @@ methods = ["dev-token"]
);
let run_log = std::fs::read_to_string(&run_log_path).expect("run log should be readable");
assert!(
run_log.contains("Workflow run started"),
run_log.contains("Petri worker starting"),
"per-run log should include worker tracing, got:\n{run_log}"
);
assert!(
@ -761,11 +761,12 @@ fn detached_run_answers_pending_question_without_interview_scratch_files() {
.expect("question id should be present")
.to_string();
assert_eq!(question["stage"], "approve");
assert_eq!(question["stage"], "approve@1");
let response = client
.post(format!(
"{base_url}/api/v1/runs/{run_id}/questions/{question_id}/answer"
"{base_url}/api/v1/runs/{run_id}/questions/{}/answer",
question_id.replace('#', "%23")
))
.json(&serde_json::json!({ "kind": "selected", "option_key": "A" }))
.send()
@ -829,7 +830,7 @@ fn detached_run_cancel_reaches_worker_over_control_websocket() {
let run_id = created_run_id(&output);
let run_dir = context.find_run_dir(&run_id);
wait_for_event_names(&run_dir, &["run.running"]);
wait_for_lifecycle(&run_dir, "running");
tokio::runtime::Runtime::new()
.expect("test runtime should build")
.block_on(async {
@ -882,7 +883,7 @@ fn worker_exits_after_sigterm_cancel_even_when_stdin_stays_open() {
let mut child = spawn_worker_process(&context, &server, &run_dir, &run_id, "start");
let stdin = child.stdin.take().expect("worker stdin should be piped");
wait_for_event_names(&run_dir, &["run.running"]);
wait_for_lifecycle(&run_dir, "running");
let worker_pid = child.id();
assert!(worker_pid > 0, "worker pid should be present");
fabro_proc::sigterm(worker_pid);

File diff suppressed because it is too large Load diff

View file

@ -15,7 +15,7 @@ pub(crate) use mcp_client::McpStdioTestClient;
pub(crate) fn run_output_filters(context: &TestContext) -> Vec<(String, String)> {
let mut filters = context.filters();
filters.push((r"\b\d+ms\b".to_string(), "[TIME]".to_string()));
filters.push((r"\b\d+(\.\d+)?(ms|s)\b".to_string(), "[TIME]".to_string()));
filters.push((
r"(?m)^(Graph: ).+$".to_string(),
"${1}[GRAPH_PATH]".to_string(),

View file

@ -16,7 +16,6 @@ fn dry_run_branching() {
----- stderr -----
Run: [ULID]
Web UI: http://localhost:3000/runs/[ULID]
Sandbox: local (ready in [TIME])
✓ Start [TIME]
✓ Plan [TIME]
✓ Implement [TIME]
@ -28,9 +27,6 @@ fn dry_run_branching() {
Run: [ULID]
Status: SUCCEEDED
Duration: [DURATION]
=== Output ===
[Simulated] Response for stage: validate
");
}
@ -48,7 +44,6 @@ fn dry_run_conditions() {
----- stderr -----
Run: [ULID]
Web UI: http://localhost:3000/runs/[ULID]
Sandbox: local (ready in [TIME])
✓ start [TIME]
✓ Decide [TIME]
✓ Path B [TIME]
@ -58,9 +53,6 @@ fn dry_run_conditions() {
Run: [ULID]
Status: SUCCEEDED
Duration: [DURATION]
=== Output ===
[Simulated] Response for stage: path_b
");
}
@ -72,7 +64,8 @@ fn dry_run_parallel() {
cmd.args(["--dry-run", "--auto-approve"]);
cmd.arg(&workflow);
let mut filters = run_output_filters(&context);
filters.push((r"\bbranch[12]\b".to_string(), "[BRANCH]".to_string()));
// The two branches run concurrently and finish in either order.
filters.push((r"\bBranch [12]\b".to_string(), "Branch [N]".to_string()));
fabro_snapshot!(filters, cmd, @"
success: true
exit_code: 0
@ -80,11 +73,10 @@ fn dry_run_parallel() {
----- stderr -----
Run: [ULID]
Web UI: http://localhost:3000/runs/[ULID]
Sandbox: local (ready in [TIME])
✓ start [TIME]
✓ [BRANCH] [TIME]
✓ [BRANCH] [TIME]
✓ Fork Work [TIME]
✓ Branch [N] [TIME]
✓ Branch [N] [TIME]
✓ Merge Results [TIME]
✓ Review [TIME]
✓ exit [TIME]
@ -93,9 +85,6 @@ fn dry_run_parallel() {
Run: [ULID]
Status: SUCCEEDED
Duration: [DURATION]
=== Output ===
[Simulated] Response for stage: review
");
}
@ -113,7 +102,6 @@ fn dry_run_styled() {
----- stderr -----
Run: [ULID]
Web UI: http://localhost:3000/runs/[ULID]
Sandbox: local (ready in [TIME])
✓ start [TIME]
✓ Plan [TIME]
✓ Implement [TIME]
@ -124,9 +112,6 @@ fn dry_run_styled() {
Run: [ULID]
Status: SUCCEEDED
Duration: [DURATION]
=== Output ===
[Simulated] Response for stage: critical_review
");
}
@ -144,7 +129,6 @@ fn dry_run_inferred_command() {
----- stderr -----
Run: [ULID]
Web UI: http://localhost:3000/runs/[ULID]
Sandbox: local (ready in [TIME])
✓ Start [TIME]
✓ Echo [TIME]
✓ Exit [TIME]