diff --git a/lib/apps/fabro-cli/src/commands/run/petri_worker.rs b/lib/apps/fabro-cli/src/commands/run/petri_worker.rs index f215143a4..23576c3ce 100644 --- a/lib/apps/fabro-cli/src/commands/run/petri_worker.rs +++ b/lib/apps/fabro-cli/src/commands/run/petri_worker.rs @@ -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, ) } }; diff --git a/lib/apps/fabro-cli/src/commands/run/runner.rs b/lib/apps/fabro-cli/src/commands/run/runner.rs index fa6c32d69..9facc453f 100644 --- a/lib/apps/fabro-cli/src/commands/run/runner.rs +++ b/lib/apps/fabro-cli/src/commands/run/runner.rs @@ -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, diff --git a/lib/apps/fabro-cli/tests/it/cmd/attach.rs b/lib/apps/fabro-cli/tests/it/cmd/attach.rs index 7f74114d0..81df5cc9a 100644 --- a/lib/apps/fabro-cli/tests/it/cmd/attach.rs +++ b/lib/apps/fabro-cli/tests/it/cmd/attach.rs @@ -66,39 +66,43 @@ fn format_output_snapshot(output: &Output, filters: &[(String, String)]) -> Stri } fn normalize_attach_json_progress_event(mut event: Value) -> Value { - // Definition and spec blob hashes are already rewritten to - // [BLOB_HASH] by the shared json_snapshot_filters regexes. - // Strip v2-shape server/version fields that the bridge emits, - // since the test fixture's socket path is randomised per run. - if let Some(settings) = event - .pointer_mut("/properties/settings") - .and_then(Value::as_object_mut) - { - settings.remove("_version"); - settings.remove("server"); - settings.remove("version"); - } - if let Some(target) = event - .pointer_mut("/properties/settings/cli/target") - .and_then(Value::as_object_mut) - { - if target.contains_key("path") { - target.insert( - "path".to_string(), - Value::String("[CLI_SOCKET]".to_string()), - ); + // The `run.created` record carries the whole run spec, whose graph + // and settings vary with the fixture's socket path and node order; + // the test does not check it. + if event.pointer("/item/record/kind") == Some(&Value::String("run.created".to_string())) { + if let Some(spec) = event.pointer_mut("/item/record/spec") { + *spec = Value::String("[RUN_SPEC]".to_string()); } } - if let Some(model_name) = event.pointer_mut("/properties/settings/run/model/name") { - assert!( - model_name.is_string(), - "default model should serialize as a string" - ); - *model_name = Value::String("[DEFAULT_MODEL]".to_string()); - } + redact_volatile_fields(&mut event); event } +/// Replace, at every depth, the wall-clock epoch milliseconds a stream +/// item carries, the content digests of the graph, which hashes the +/// fixture's temporary path, and the commit shas of the fixture's repository. +fn redact_volatile_fields(value: &mut Value) { + match value { + Value::Object(fields) => { + for (key, field) in fields.iter_mut() { + if key == "recorded_at" { + *field = Value::String("[EPOCH_MS]".to_string()); + } else { + redact_volatile_fields(field); + } + } + } + Value::Array(items) => items.iter_mut().for_each(redact_volatile_fields), + Value::String(text) + if matches!(text.len(), 40 | 64) + && text.bytes().all(|byte| byte.is_ascii_hexdigit()) => + { + *text = "[DIGEST]".to_string(); + } + _ => {} + } +} + fn wait_for_output_signal( child: &mut std::process::Child, stdout: &mut impl Read, @@ -388,7 +392,6 @@ fn attach_replays_completed_detached_run() { ----- stdout ----- ----- stderr ----- Web UI: http://localhost:3000/runs/[ULID] - Sandbox: local (ready in [TIME]) ✓ Start [TIME] ✓ Run Tests [TIME] ✓ Report [TIME] @@ -502,7 +505,8 @@ fn attach_advances_when_pending_question_is_answered_elsewhere() { runtime.block_on(async { 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() @@ -632,7 +636,6 @@ fn attach_before_completion_streams_to_finished_state() { ----- stdout ----- ----- stderr ----- Web UI: http://localhost:3000/runs/[ULID] - Sandbox: local (ready in [TIME]) ✓ start [DURATION] ✓ wait [DURATION] ✓ exit [DURATION] @@ -705,10 +708,9 @@ fn attach_json_errors_without_prompting_for_human_input() { .filter(|line| !line.trim().is_empty()) .map(|line| serde_json::from_str(line).expect("log line should be valid JSON")) .collect(); - if log_events.iter().any(|event| { - event["event"] == "stage.started" - && event["node_id"] == "approve" - && event["properties"]["handler_type"] == "human" + // The gate's question: Petri records it as a parsed step progress. + if log_events.iter().any(|item| { + item.pointer("/item/derived/parsed/kind") == Some(&Value::String("question".into())) }) { break; } @@ -746,10 +748,8 @@ fn attach_json_errors_without_prompting_for_human_input() { .map(|line| serde_json::from_str(line).expect("log line should be valid JSON")) .collect(); assert!( - log_events.iter().any(|event| { - event["event"] == "stage.started" - && event["node_id"] == "approve" - && event["properties"]["handler_type"] == "human" + log_events.iter().any(|item| { + item.pointer("/item/derived/parsed/kind") == Some(&Value::String("question".into())) }), "the run should still be waiting on the human gate" ); @@ -774,660 +774,1824 @@ fn attach_json_errors_without_prompting_for_human_input() { fabro_json_snapshot!(context, &progress, @r#" [ { - "actor": { - "auth_method": "dev_token", - "identity": { - "issuer": "fabro:dev", - "subject": "dev" - }, - "kind": "user", - "login": "dev" - }, - "event": "run.created", - "id": "[EVENT_ID]", - "properties": { - "graph": { - "attrs": { - "goal": { - "String": "Wait for approval" - } - }, - "edges": [ - { - "attrs": {}, - "from": "start", - "to": "approve" - }, - { - "attrs": { - "label": { - "String": "[A] Approve" - } - }, - "from": "approve", - "to": "ship" - }, - { - "attrs": { - "label": { - "String": "[R] Revise" - } - }, - "from": "approve", - "to": "revise" - }, - { - "attrs": {}, - "from": "ship", - "to": "exit" - }, - { - "attrs": {}, - "from": "revise", - "to": "exit" - } - ], - "name": "HumanGate", - "nodes": { - "approve": { - "attrs": { - "label": { - "String": "Approve?" - }, - "shape": { - "String": "hexagon" - } - }, - "id": "approve" - }, - "exit": { - "attrs": { - "label": { - "String": "Exit" - }, - "shape": { - "String": "Msquare" - } - }, - "id": "exit" - }, - "revise": { - "attrs": { - "script": { - "String": "echo revised" - }, - "shape": { - "String": "parallelogram" - } - }, - "id": "revise" - }, - "ship": { - "attrs": { - "script": { - "String": "echo shipped" - }, - "shape": { - "String": "parallelogram" - } - }, - "id": "ship" - }, - "start": { - "attrs": { - "label": { - "String": "Start" - }, - "shape": { - "String": "Mdiamond" - } - }, - "id": "start" - } - } - }, - "provenance": { - "client": { - "name": "fabro-cli", - "user_agent": "fabro-cli/[VERSION]", - "version": "[VERSION]" - }, - "server": { - "version": "[VERSION]" - }, - "subject": { - "auth_method": "dev_token", - "identity": { - "issuer": "fabro:dev", - "subject": "dev" - }, - "kind": "user", - "login": "dev" - } - }, - "settings": { - "project": { - "description": null, - "metadata": {}, - "name": null - }, - "run": { - "agent": { - "fabro_tools": false, - "mcps": {} - }, - "artifacts": { - "include": [] - }, - "checkpoint": { - "commit_timeout_ms": 30000, - "exclude_globs": [], - "skip_git_hooks": false - }, - "clone": { - "depth": 100, - "enabled": true - }, - "environment": { - "env": {}, - "id": "local", - "image": { - "docker": null, - "dockerfile": null - }, - "labels": {}, - "lifecycle": { - "auto_stop": null, - "preserve": false, - "stop_on_terminal": true - }, - "network": { - "allow": [], - "mode": "allow_all" - }, - "provider": "local", - "resources": { - "cpu": null, - "disk": null, - "memory": null - } - }, - "execution": { - "approval": "prompt", - "mode": "normal" - }, - "git": { - "author": null - }, - "goal": { - "type": "inline", - "value": "Wait for approval" - }, - "hooks": [], - "inputs": {}, - "integrations": { - "github": { - "permissions": {} - } - }, - "interviews": { - "provider": null, - "slack": null - }, - "metadata": {}, - "model": { - "controls": { - "reasoning_effort": null, - "speed": null - }, - "fallbacks": {}, - "name": "[DEFAULT_MODEL]", - "provider": "openai" - }, - "notifications": {}, - "prepare": { - "steps": [], - "timeout_ms": 300000 - }, - "pull_request": null, - "run_branch": { - "enabled": true, - "push": true - }, - "scm": { - "github": null, - "owner": null, - "provider": null, - "repository": null - }, - "working_dir": null - }, - "workflow": { - "description": null, - "graph": "workflow.fabro", - "metadata": {}, - "name": null - } - }, - "source_directory": "[TEMP_DIR]", - "spec_blob": "[BLOB_HASH]", - "target": { - "kind": "folder", - "path": "[TEMP_DIR]" - }, - "title": "Wait for approval", - "web_url": "http://localhost:3000/runs/[ULID]", - "workflow_slug": "human-gate", - "workflow_source": "digraph HumanGate {/n graph [goal=\"Wait for approval\"]/n start [shape=Mdiamond, label=\"Start\"]/n exit [shape=Msquare, label=\"Exit\"]/n approve [shape=hexagon, label=\"Approve?\"]/n ship [shape=parallelogram, script=\"echo shipped\"]/n revise [shape=parallelogram, script=\"echo revised\"]/n start -> approve/n approve -> ship [label=\"[A] Approve\"]/n approve -> revise [label=\"[R] Revise\"]/n ship -> exit/n revise -> exit/n}/n", - "workflow_version_id": "fc1611d3be115f2db472e4ac05a5034f449743089259566b18f204ff961a0c18" - }, "run_id": "[ULID]", - "ts": "[TIMESTAMP]" - }, - { - "event": "run.submitted", + "stream_seq": 1, + "kind": "platform", "id": "[EVENT_ID]", - "properties": { - "definition_blob": "[BLOB_HASH]" - }, - "run_id": "[ULID]", - "ts": "[TIMESTAMP]" - }, - { - "actor": { - "auth_method": "dev_token", - "identity": { - "issuer": "fabro:dev", - "subject": "dev" - }, - "kind": "user", - "login": "dev" - }, - "event": "run.start_requested", - "id": "[EVENT_ID]", - "properties": { - "resume": false - }, - "run_id": "[ULID]", - "ts": "[TIMESTAMP]" - }, - { - "actor": { - "auth_method": "dev_token", - "identity": { - "issuer": "fabro:dev", - "subject": "dev" - }, - "kind": "user", - "login": "dev" - }, - "event": "run.runnable", - "id": "[EVENT_ID]", - "properties": { - "source": "start_requested" - }, - "run_id": "[ULID]", - "ts": "[TIMESTAMP]" - }, - { - "actor": { - "kind": "worker", - "run_id": "[ULID]" - }, - "event": "run.starting", - "id": "[EVENT_ID]", - "properties": {}, - "run_id": "[ULID]", - "ts": "[TIMESTAMP]" - }, - { - "actor": { - "kind": "worker", - "run_id": "[ULID]" - }, - "event": "sandbox.initializing", - "id": "[EVENT_ID]", - "properties": { - "provider": "local" - }, - "run_id": "[ULID]", - "ts": "[TIMESTAMP]" - }, - { - "actor": { - "kind": "worker", - "run_id": "[ULID]" - }, - "event": "sandbox.create.started", - "id": "[EVENT_ID]", - "properties": { - "action": "create", - "correlation_id": "[ULID]", - "id": { - "sequence": 1, - "source_id": "[HEX]" - }, - "occurred_at": "[TIMESTAMP]", - "operation_id": "[HEX]", - "provider": "host", - "subject": { - "id": "host-dir-[HEX]", - "type": "sandbox" - }, - "type": "operation_started" - }, - "run_id": "[ULID]", - "ts": "[TIMESTAMP]" - }, - { - "actor": { - "kind": "worker", - "run_id": "[ULID]" - }, - "event": "sandbox.create.progress", - "id": "[EVENT_ID]", - "properties": { - "action": "create", - "correlation_id": "[ULID]", - "id": { - "sequence": 2, - "source_id": "[HEX]" - }, - "occurred_at": "[TIMESTAMP]", - "operation_id": "[HEX]", - "progress": { - "code": "sandbox.provision" - }, - "provider": "host", - "subject": { - "id": "host-dir-[HEX]", - "type": "sandbox" - }, - "type": "operation_progress" - }, - "run_id": "[ULID]", - "ts": "[TIMESTAMP]" - }, - { - "actor": { - "kind": "worker", - "run_id": "[ULID]" - }, - "event": "sandbox.create.completed", - "id": "[EVENT_ID]", - "properties": { - "action": "create", - "correlation_id": "[ULID]", - "duration": { - "nanos": "[NANOS]", - "secs": 0 - }, - "id": { - "sequence": 3, - "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]" - }, - { - "actor": { - "kind": "worker", - "run_id": "[ULID]" - }, - "event": "sandbox.ready", - "id": "[EVENT_ID]", - "properties": { - "duration_ms": "[DURATION_MS]", - "provider": "local" - }, - "run_id": "[ULID]", - "ts": "[TIMESTAMP]" - }, - { - "actor": { - "kind": "worker", - "run_id": "[ULID]" - }, - "event": "sandbox.initialized", - "id": "[EVENT_ID]", - "properties": { - "id": "host-dir-[HEX]", - "provider": "local", - "repo_cloned": false, - "repos_root": "[TEMP_DIR]/.repos", - "working_directory": "[TEMP_DIR]", - "workspace_root": "[TEMP_DIR]" - }, - "run_id": "[ULID]", - "ts": "[TIMESTAMP]" - }, - { - "actor": { - "kind": "worker", - "run_id": "[ULID]" - }, - "event": "git.identity.resolved", - "id": "[EVENT_ID]", - "properties": { - "email": "noreply@fabro.sh", - "name": "Fabro", - "source": "default" - }, - "run_id": "[ULID]", - "ts": "[TIMESTAMP]" - }, - { - "actor": { - "kind": "worker", - "run_id": "[ULID]" - }, - "event": "run.started", - "id": "[EVENT_ID]", - "properties": { - "goal": "Wait for approval", - "name": "HumanGate" - }, - "run_id": "[ULID]", - "ts": "[TIMESTAMP]" - }, - { - "actor": { - "kind": "worker", - "run_id": "[ULID]" - }, - "event": "run.running", - "id": "[EVENT_ID]", - "properties": {}, - "run_id": "[ULID]", - "ts": "[TIMESTAMP]" - }, - { - "actor": { - "kind": "worker", - "run_id": "[ULID]" - }, - "event": "stage.started", - "id": "[EVENT_ID]", - "node_id": "start", - "node_label": "Start", - "properties": { - "attempt": 1, - "graph_visit": 1, - "handler_type": "start", - "index": 0, - "max_attempts": 1 - }, - "run_id": "[ULID]", - "stage_id": "start@1", - "ts": "[TIMESTAMP]" - }, - { - "actor": { - "kind": "worker", - "run_id": "[ULID]" - }, - "event": "stage.completed", - "id": "[EVENT_ID]", - "node_id": "start", - "node_label": "Start", - "properties": { - "attempt": 1, - "context_values": { - "current_node": "start", - "graph.goal": "Wait for approval", - "internal.fidelity": "compact", - "internal.node_visit_count": 1, - "internal.run_id": "[ULID]", - "internal.thread_id": null - }, - "index": 0, - "max_attempts": 1, - "node_visits": { - "start": 1 - }, - "status": "succeeded", - "timing": { - "active_time_ms": "[ACTIVE_TIME_MS]", - "inference_time_ms": "[INFERENCE_TIME_MS]", - "tool_time_ms": "[TOOL_TIME_MS]", - "wall_time_ms": "[WALL_TIME_MS]" + "recorded_at": "[EPOCH_MS]", + "item": { + "seq": 1, + "recorded_at": "[EPOCH_MS]", + "record": { + "kind": "run.created", + "spec": "[RUN_SPEC]", + "title": "Wait for approval", + "web_url": "http://localhost:3000/runs/[ULID]" } - }, - "run_id": "[ULID]", - "stage_id": "start@1", - "ts": "[TIMESTAMP]" + } }, { - "actor": { - "kind": "worker", - "run_id": "[ULID]" - }, - "event": "edge.selected", - "id": "[EVENT_ID]", - "properties": { - "from_node": "start", - "is_jump": false, - "reason": "unconditional", - "stage_status": "succeeded", - "to_node": "approve" - }, "run_id": "[ULID]", - "ts": "[TIMESTAMP]" - }, - { - "actor": { - "kind": "worker", - "run_id": "[ULID]" - }, - "event": "checkpoint.completed", + "stream_seq": 2, + "kind": "platform", "id": "[EVENT_ID]", - "node_id": "start", - "node_label": "start", - "properties": { - "completed_nodes": [ - "start" - ], - "context_values": { - "current_node": "start", - "failure_class": "", - "failure_signature": "", - "graph.goal": "Wait for approval", - "internal.fidelity": "compact", - "internal.node_visit_count": 1, - "internal.retry_count.start": 0, - "internal.run_id": "[ULID]", - "internal.thread_id": null, - "outcome": "succeeded" - }, - "current_node": "start", - "graph_visit": 1, - "next_node_id": "approve", - "node_outcomes": { - "start": { - "status": "succeeded", - "usage": null + "recorded_at": "[EPOCH_MS]", + "item": { + "seq": 2, + "recorded_at": "[EPOCH_MS]", + "record": { + "kind": "run.lifecycle", + "transition": "submitted", + "status": { + "kind": "submitted" } - }, - "node_visits": { - "start": 1 - }, - "status": "succeeded" - }, - "run_id": "[ULID]", - "stage_id": "start@1", - "ts": "[TIMESTAMP]" + } + } }, { - "actor": { - "kind": "worker", - "run_id": "[ULID]" - }, - "event": "stage.started", - "id": "[EVENT_ID]", - "node_id": "approve", - "node_label": "Approve?", - "properties": { - "attempt": 1, - "graph_visit": 1, - "handler_type": "human", - "index": 1, - "max_attempts": 1 - }, "run_id": "[ULID]", - "stage_id": "approve@1", - "ts": "[TIMESTAMP]" + "stream_seq": 3, + "kind": "platform", + "id": "[EVENT_ID]", + "recorded_at": "[EPOCH_MS]", + "item": { + "seq": 3, + "recorded_at": "[EPOCH_MS]", + "record": { + "kind": "run.lifecycle", + "transition": "start_requested", + "source": "start" + } + } }, { - "actor": { - "kind": "worker", - "run_id": "[ULID]" - }, - "event": "interview.started", + "run_id": "[ULID]", + "stream_seq": 4, + "kind": "platform", "id": "[EVENT_ID]", - "node_id": "approve", - "node_label": "approve", - "properties": { - "allow_freeform": false, - "options": [ - { - "key": "A", - "label": "[A] Approve" + "recorded_at": "[EPOCH_MS]", + "item": { + "seq": 4, + "recorded_at": "[EPOCH_MS]", + "record": { + "kind": "run.lifecycle", + "transition": "runnable", + "status": { + "kind": "runnable" }, - { - "key": "R", - "label": "[R] Revise" - } - ], - "question": "Approve?", - "question_id": "[ULID]", - "question_type": "multiple_choice", - "stage": "approve" - }, - "run_id": "[ULID]", - "stage_id": "approve@1", - "ts": "[TIMESTAMP]" + "source": "start_requested" + } + } }, { - "actor": { - "kind": "worker", - "run_id": "[ULID]" - }, - "event": "run.blocked", - "id": "[EVENT_ID]", - "properties": { - "blocked_reason": "human_input_required" - }, "run_id": "[ULID]", - "ts": "[TIMESTAMP]" + "stream_seq": 5, + "kind": "platform", + "id": "[EVENT_ID]", + "recorded_at": "[EPOCH_MS]", + "item": { + "seq": 5, + "recorded_at": "[EPOCH_MS]", + "record": { + "kind": "run.lifecycle", + "transition": "starting", + "status": { + "kind": "starting" + } + } + } + }, + { + "run_id": "[ULID]", + "stream_seq": 6, + "kind": "platform", + "id": "[EVENT_ID]", + "recorded_at": "[EPOCH_MS]", + "item": { + "seq": 6, + "recorded_at": "[EPOCH_MS]", + "record": { + "kind": "run.lifecycle", + "transition": "running", + "status": { + "kind": "running" + } + } + } + }, + { + "run_id": "[ULID]", + "stream_seq": 7, + "kind": "petri", + "id": "coordinator/0/0", + "recorded_at": "[EPOCH_MS]", + "item": { + "id": { + "log": "coordinator", + "seq": 0, + "index": 0 + }, + "origin": "external", + "context": { + "invocation": 0 + }, + "recorded_at": "[EPOCH_MS]", + "record": { + "seq": 0, + "origin": "external", + "recorded_at": "[EPOCH_MS]", + "body": { + "event": "run.started", + "format_version": 5, + "key": "[ULID]", + "root": 0, + "middleware_chain": [ + "circuit-breaker" + ] + } + } + } + }, + { + "run_id": "[ULID]", + "stream_seq": 8, + "kind": "petri", + "id": "coordinator/1/0", + "recorded_at": "[EPOCH_MS]", + "item": { + "id": { + "log": "coordinator", + "seq": 1, + "index": 0 + }, + "origin": "external", + "context": {}, + "recorded_at": "[EPOCH_MS]", + "record": { + "seq": 1, + "origin": "external", + "recorded_at": "[EPOCH_MS]", + "body": { + "event": "graph.registered", + "digest": "[DIGEST]" + } + } + } + }, + { + "run_id": "[ULID]", + "stream_seq": 9, + "kind": "petri", + "id": "coordinator/2/0", + "recorded_at": "[EPOCH_MS]", + "item": { + "id": { + "log": "coordinator", + "seq": 2, + "index": 0 + }, + "origin": "external", + "context": { + "invocation": 0 + }, + "recorded_at": "[EPOCH_MS]", + "record": { + "seq": 2, + "origin": "external", + "recorded_at": "[EPOCH_MS]", + "body": { + "event": "invocation.declared", + "invocation": 0, + "call": null, + "graph": "[DIGEST]", + "context": {}, + "secret_bindings": "none", + "sandbox": "isolated" + } + } + } + }, + { + "run_id": "[ULID]", + "stream_seq": 10, + "kind": "petri", + "id": "coordinator/3/0", + "recorded_at": "[EPOCH_MS]", + "item": { + "id": { + "log": "coordinator", + "seq": 3, + "index": 0 + }, + "origin": "external", + "context": { + "invocation": 0, + "execution": 0 + }, + "recorded_at": "[EPOCH_MS]", + "record": { + "seq": 3, + "origin": "external", + "recorded_at": "[EPOCH_MS]", + "body": { + "event": "execution.declared", + "execution": 0, + "invocation": 0, + "predecessor": null, + "start": { + "entry": "graph_entries", + "context": {}, + "prior_firings": {}, + "execution_index": 0, + "max_executions": 32 + }, + "middleware_state": { + "circuit-breaker": [ + 1, + { + "loop_signatures": {}, + "restart_signatures": {}, + "pending": {} + } + ] + } + } + } + } + }, + { + "run_id": "[ULID]", + "stream_seq": 11, + "kind": "petri", + "id": "execution 0/0/0", + "recorded_at": "[EPOCH_MS]", + "item": { + "id": { + "log": "execution", + "execution": 0, + "seq": 0, + "index": 0 + }, + "origin": "external", + "context": { + "invocation": 0, + "execution": 0 + }, + "recorded_at": "[EPOCH_MS]", + "record": { + "seq": 0, + "origin": "external", + "recorded_at": "[EPOCH_MS]", + "body": { + "event": "execution.started", + "entry": "graph_entries", + "context": {}, + "prior_firings": {}, + "execution_index": 0, + "max_executions": 32 + } + } + } + }, + { + "run_id": "[ULID]", + "stream_seq": 12, + "kind": "petri", + "id": "execution 0/1/0", + "recorded_at": "[EPOCH_MS]", + "item": { + "id": { + "log": "execution", + "execution": 0, + "seq": 1, + "index": 0 + }, + "origin": "external", + "context": { + "invocation": 0, + "execution": 0 + }, + "recorded_at": "[EPOCH_MS]", + "record": { + "seq": 1, + "origin": "external", + "recorded_at": "[EPOCH_MS]", + "body": { + "event": "admission.decided", + "decision_id": "execution_start", + "decision": "admit", + "trace": [] + } + } + } + }, + { + "run_id": "[ULID]", + "stream_seq": 13, + "kind": "petri", + "id": "execution 0/1/1", + "recorded_at": "[EPOCH_MS]", + "item": { + "id": { + "log": "execution", + "execution": 0, + "seq": 1, + "index": 1 + }, + "origin": "derived", + "context": { + "invocation": 0, + "execution": 0 + }, + "subject": { + "node": { + "id": 0, + "name": "start", + "kind": "attractor/stage", + "meta": { + "label": "Start", + "shape": "Mdiamond", + "kind": "start", + "classes": [], + "span": { + "line": 3, + "column": 3 + }, + "admission_hooks": "step", + "edges": { + "0": { + "to": "approve", + "label": null + } + } + } + }, + "firing": 1, + "visit": 1, + "attempt": 1, + "generation": 0, + "branch": { + "role": "none" + } + }, + "recorded_at": "[EPOCH_MS]", + "derived": { + "event": "visit.started", + "inputs": [ + { + "edge": 5, + "generation": 0, + "payload": null, + "from": 0 + } + ] + } + } + }, + { + "run_id": "[ULID]", + "stream_seq": 14, + "kind": "petri", + "id": "execution 0/1/2", + "recorded_at": "[EPOCH_MS]", + "item": { + "id": { + "log": "execution", + "execution": 0, + "seq": 1, + "index": 2 + }, + "origin": "derived", + "context": { + "invocation": 0, + "execution": 0 + }, + "subject": { + "node": { + "id": 0, + "name": "start", + "kind": "attractor/stage", + "meta": { + "label": "Start", + "shape": "Mdiamond", + "kind": "start", + "classes": [], + "span": { + "line": 3, + "column": 3 + }, + "admission_hooks": "step", + "edges": { + "0": { + "to": "approve", + "label": null + } + } + } + }, + "firing": 1, + "visit": 1, + "attempt": 1, + "generation": 0, + "branch": { + "role": "none" + } + }, + "recorded_at": "[EPOCH_MS]", + "derived": { + "event": "wait.state.changed", + "state": "awaiting_admission" + } + } + }, + { + "run_id": "[ULID]", + "stream_seq": 15, + "kind": "petri", + "id": "execution 0/2/0", + "recorded_at": "[EPOCH_MS]", + "item": { + "id": { + "log": "execution", + "execution": 0, + "seq": 2, + "index": 0 + }, + "origin": "core", + "context": { + "invocation": 0, + "execution": 0 + }, + "recorded_at": "[EPOCH_MS]", + "record": { + "seq": 2, + "origin": "core", + "recorded_at": "[EPOCH_MS]", + "body": { + "event": "token.emitted", + "edge": 5, + "generation": 0, + "payload": null, + "from": 0 + } + } + } + }, + { + "run_id": "[ULID]", + "stream_seq": 16, + "kind": "petri", + "id": "execution 0/3/0", + "recorded_at": "[EPOCH_MS]", + "item": { + "id": { + "log": "execution", + "execution": 0, + "seq": 3, + "index": 0 + }, + "origin": "external", + "context": { + "invocation": 0, + "execution": 0 + }, + "subject": { + "node": { + "id": 0, + "name": "start", + "kind": "attractor/stage", + "meta": { + "label": "Start", + "shape": "Mdiamond", + "kind": "start", + "classes": [], + "span": { + "line": 3, + "column": 3 + }, + "admission_hooks": "step", + "edges": { + "0": { + "to": "approve", + "label": null + } + } + } + }, + "firing": 1, + "visit": 1, + "attempt": 1, + "generation": 0, + "branch": { + "role": "none" + } + }, + "recorded_at": "[EPOCH_MS]", + "record": { + "seq": 3, + "origin": "external", + "recorded_at": "[EPOCH_MS]", + "body": { + "event": "admission.decided", + "decision_id": { + "attempt_start": { + "firing": 1, + "attempt": 1 + } + }, + "decision": "admit", + "trace": [] + } + } + } + }, + { + "run_id": "[ULID]", + "stream_seq": 17, + "kind": "petri", + "id": "execution 0/4/0", + "recorded_at": "[EPOCH_MS]", + "item": { + "id": { + "log": "execution", + "execution": 0, + "seq": 4, + "index": 0 + }, + "origin": "external", + "context": { + "invocation": 0, + "execution": 0 + }, + "subject": { + "node": { + "id": 0, + "name": "start", + "kind": "attractor/stage", + "meta": { + "label": "Start", + "shape": "Mdiamond", + "kind": "start", + "classes": [], + "span": { + "line": 3, + "column": 3 + }, + "admission_hooks": "step", + "edges": { + "0": { + "to": "approve", + "label": null + } + } + } + }, + "firing": 1, + "visit": 1, + "attempt": 1, + "generation": 0, + "branch": { + "role": "none" + } + }, + "recorded_at": "[EPOCH_MS]", + "record": { + "seq": 4, + "origin": "external", + "recorded_at": "[EPOCH_MS]", + "body": { + "event": "step.started", + "firing": 1, + "attempt": 1 + } + } + } + }, + { + "run_id": "[ULID]", + "stream_seq": 18, + "kind": "petri", + "id": "execution 0/4/1", + "recorded_at": "[EPOCH_MS]", + "item": { + "id": { + "log": "execution", + "execution": 0, + "seq": 4, + "index": 1 + }, + "origin": "derived", + "context": { + "invocation": 0, + "execution": 0 + }, + "subject": { + "node": { + "id": 0, + "name": "start", + "kind": "attractor/stage", + "meta": { + "label": "Start", + "shape": "Mdiamond", + "kind": "start", + "classes": [], + "span": { + "line": 3, + "column": 3 + }, + "admission_hooks": "step", + "edges": { + "0": { + "to": "approve", + "label": null + } + } + } + }, + "firing": 1, + "visit": 1, + "attempt": 1, + "generation": 0, + "branch": { + "role": "none" + } + }, + "recorded_at": "[EPOCH_MS]", + "derived": { + "event": "wait.state.changed", + "state": "running" + } + } + }, + { + "run_id": "[ULID]", + "stream_seq": 19, + "kind": "petri", + "id": "execution 0/5/0", + "recorded_at": "[EPOCH_MS]", + "item": { + "id": { + "log": "execution", + "execution": 0, + "seq": 5, + "index": 0 + }, + "origin": "external", + "context": { + "invocation": 0, + "execution": 0 + }, + "subject": { + "node": { + "id": 0, + "name": "start", + "kind": "attractor/stage", + "meta": { + "label": "Start", + "shape": "Mdiamond", + "kind": "start", + "classes": [], + "span": { + "line": 3, + "column": 3 + }, + "admission_hooks": "step", + "edges": { + "0": { + "to": "approve", + "label": null + } + } + } + }, + "firing": 1, + "visit": 1, + "attempt": 1, + "generation": 0, + "branch": { + "role": "none" + } + }, + "recorded_at": "[EPOCH_MS]", + "record": { + "seq": 5, + "origin": "external", + "recorded_at": "[EPOCH_MS]", + "body": { + "event": "step.progress.recorded", + "firing": 1, + "ev": { + "log": { + "stream": "stderr", + "line": "checkout: [TEMP_DIR] is not a Git repository; the workspace starts empty" + } + } + } + } + } + }, + { + "run_id": "[ULID]", + "stream_seq": 20, + "kind": "petri", + "id": "execution 0/6/0", + "recorded_at": "[EPOCH_MS]", + "item": { + "id": { + "log": "execution", + "execution": 0, + "seq": 6, + "index": 0 + }, + "origin": "external", + "context": { + "invocation": 0, + "execution": 0 + }, + "subject": { + "node": { + "id": 0, + "name": "start", + "kind": "attractor/stage", + "meta": { + "label": "Start", + "shape": "Mdiamond", + "kind": "start", + "classes": [], + "span": { + "line": 3, + "column": 3 + }, + "admission_hooks": "step", + "edges": { + "0": { + "to": "approve", + "label": null + } + } + } + }, + "firing": 1, + "visit": 1, + "attempt": 1, + "generation": 0, + "branch": { + "role": "none" + } + }, + "recorded_at": "[EPOCH_MS]", + "record": { + "seq": 6, + "origin": "external", + "recorded_at": "[EPOCH_MS]", + "body": { + "event": "step.progress.recorded", + "firing": 1, + "ev": { + "custom": { + "$note": { + "kind": "fabro.checkpoint", + "payload": { + "execution": 0, + "firing": 1, + "attempt": 1, + "workspace": "invocation-0-scope-0", + "git_commit_sha": "[DIGEST]", + "reused": false + } + } + } + } + } + }, + "derived": { + "parsed": { + "kind": "note", + "note": { + "kind": "fabro.checkpoint", + "payload": { + "execution": 0, + "firing": 1, + "attempt": 1, + "workspace": "invocation-0-scope-0", + "git_commit_sha": "[DIGEST]", + "reused": false + } + } + } + } + } + }, + { + "run_id": "[ULID]", + "stream_seq": 21, + "kind": "petri", + "id": "execution 0/7/0", + "recorded_at": "[EPOCH_MS]", + "item": { + "id": { + "log": "execution", + "execution": 0, + "seq": 7, + "index": 0 + }, + "origin": "external", + "context": { + "invocation": 0, + "execution": 0 + }, + "subject": { + "node": { + "id": 0, + "name": "start", + "kind": "attractor/stage", + "meta": { + "label": "Start", + "shape": "Mdiamond", + "kind": "start", + "classes": [], + "span": { + "line": 3, + "column": 3 + }, + "admission_hooks": "step", + "edges": { + "0": { + "to": "approve", + "label": null + } + } + } + }, + "firing": 1, + "visit": 1, + "attempt": 1, + "generation": 0, + "branch": { + "role": "none" + } + }, + "recorded_at": "[EPOCH_MS]", + "record": { + "seq": 7, + "origin": "external", + "recorded_at": "[EPOCH_MS]", + "body": { + "event": "step.finished", + "firing": 1, + "attempt": 1, + "outcome": { + "status": "success", + "output": { + "outcome": "succeeded", + "failure_class": "" + }, + "metrics": { + "duration_ms": "[DURATION_MS]" + }, + "context_updates": { + "failure_class": "", + "internal.run_id": "petri" + } + } + } + }, + "derived": { + "final": true, + "exhausted": false + } + } + }, + { + "run_id": "[ULID]", + "stream_seq": 22, + "kind": "petri", + "id": "execution 0/7/1", + "recorded_at": "[EPOCH_MS]", + "item": { + "id": { + "log": "execution", + "execution": 0, + "seq": 7, + "index": 1 + }, + "origin": "derived", + "context": { + "invocation": 0, + "execution": 0 + }, + "subject": { + "node": { + "id": 0, + "name": "start", + "kind": "attractor/stage", + "meta": { + "label": "Start", + "shape": "Mdiamond", + "kind": "start", + "classes": [], + "span": { + "line": 3, + "column": 3 + }, + "admission_hooks": "step", + "edges": { + "0": { + "to": "approve", + "label": null + } + } + } + }, + "firing": 1, + "visit": 1, + "attempt": 1, + "generation": 0, + "branch": { + "role": "none" + } + }, + "recorded_at": "[EPOCH_MS]", + "derived": { + "event": "visit.completed", + "outcome": { + "status": "success", + "output": { + "outcome": "succeeded", + "failure_class": "" + }, + "metrics": { + "duration_ms": "[DURATION_MS]" + }, + "context_updates": { + "failure_class": "", + "internal.run_id": "petri" + } + }, + "executed": true, + "attempts": 1 + } + } + }, + { + "run_id": "[ULID]", + "stream_seq": 23, + "kind": "platform", + "id": "[EVENT_ID]", + "recorded_at": "[EPOCH_MS]", + "item": { + "seq": 7, + "recorded_at": "[EPOCH_MS]", + "record": { + "kind": "checkpoint", + "execution": 0, + "firing": 1, + "attempt": 1, + "workspace": "invocation-0-scope-0", + "git_commit_sha": "[DIGEST]", + "operation": { + "execution": 0, + "decision": { + "attempt_start": { + "firing": 1, + "attempt": 1 + } + }, + "effect": "checkpoint" + } + }, + "position": { + "execution": 0, + "firing": 1 + } + } + }, + { + "run_id": "[ULID]", + "stream_seq": 24, + "kind": "petri", + "id": "execution 0/8/0", + "recorded_at": "[EPOCH_MS]", + "item": { + "id": { + "log": "execution", + "execution": 0, + "seq": 8, + "index": 0 + }, + "origin": "external", + "context": { + "invocation": 0, + "execution": 0 + }, + "subject": { + "node": { + "id": 0, + "name": "start", + "kind": "attractor/stage", + "meta": { + "label": "Start", + "shape": "Mdiamond", + "kind": "start", + "classes": [], + "span": { + "line": 3, + "column": 3 + }, + "admission_hooks": "step", + "edges": { + "0": { + "to": "approve", + "label": null + } + } + } + }, + "firing": 1, + "visit": 1, + "attempt": 1, + "generation": 0, + "branch": { + "role": "none" + } + }, + "recorded_at": "[EPOCH_MS]", + "record": { + "seq": 8, + "origin": "external", + "recorded_at": "[EPOCH_MS]", + "body": { + "event": "routing.resolved", + "decision_id": { + "route": { + "firing": 1, + "attempt": 1 + } + }, + "groups": [ + { + "group": 0, + "draw": null, + "trace": [], + "decision": { + "emit": 0 + } + } + ] + } + }, + "derived": { + "groups": [ + { + "group": 0, + "target": { + "id": 2, + "name": "approve", + "kind": "attractor/human", + "meta": { + "label": "Approve?", + "shape": "hexagon", + "kind": "human", + "classes": [], + "span": { + "line": 5, + "column": 3 + }, + "edges": { + "1": { + "to": "ship", + "label": "[A] Approve" + }, + "2": { + "to": "revise", + "label": "[R] Revise" + } + } + } + } + } + ] + } + } + }, + { + "run_id": "[ULID]", + "stream_seq": 25, + "kind": "petri", + "id": "execution 0/8/1", + "recorded_at": "[EPOCH_MS]", + "item": { + "id": { + "log": "execution", + "execution": 0, + "seq": 8, + "index": 1 + }, + "origin": "derived", + "context": { + "invocation": 0, + "execution": 0 + }, + "subject": { + "node": { + "id": 2, + "name": "approve", + "kind": "attractor/human", + "meta": { + "label": "Approve?", + "shape": "hexagon", + "kind": "human", + "classes": [], + "span": { + "line": 5, + "column": 3 + }, + "edges": { + "1": { + "to": "ship", + "label": "[A] Approve" + }, + "2": { + "to": "revise", + "label": "[R] Revise" + } + } + } + }, + "firing": 2, + "visit": 1, + "attempt": 1, + "generation": 0, + "branch": { + "role": "none" + } + }, + "recorded_at": "[EPOCH_MS]", + "derived": { + "event": "visit.started", + "inputs": [ + { + "edge": 0, + "generation": 0, + "payload": { + "outcome": "succeeded", + "failure_class": "" + }, + "from": 1 + } + ] + } + } + }, + { + "run_id": "[ULID]", + "stream_seq": 26, + "kind": "petri", + "id": "execution 0/8/2", + "recorded_at": "[EPOCH_MS]", + "item": { + "id": { + "log": "execution", + "execution": 0, + "seq": 8, + "index": 2 + }, + "origin": "derived", + "context": { + "invocation": 0, + "execution": 0 + }, + "subject": { + "node": { + "id": 2, + "name": "approve", + "kind": "attractor/human", + "meta": { + "label": "Approve?", + "shape": "hexagon", + "kind": "human", + "classes": [], + "span": { + "line": 5, + "column": 3 + }, + "edges": { + "1": { + "to": "ship", + "label": "[A] Approve" + }, + "2": { + "to": "revise", + "label": "[R] Revise" + } + } + } + }, + "firing": 2, + "visit": 1, + "attempt": 1, + "generation": 0, + "branch": { + "role": "none" + } + }, + "recorded_at": "[EPOCH_MS]", + "derived": { + "event": "wait.state.changed", + "state": "awaiting_admission" + } + } + }, + { + "run_id": "[ULID]", + "stream_seq": 27, + "kind": "petri", + "id": "execution 0/9/0", + "recorded_at": "[EPOCH_MS]", + "item": { + "id": { + "log": "execution", + "execution": 0, + "seq": 9, + "index": 0 + }, + "origin": "core", + "context": { + "invocation": 0, + "execution": 0 + }, + "subject": { + "node": { + "id": 0, + "name": "start", + "kind": "attractor/stage", + "meta": { + "label": "Start", + "shape": "Mdiamond", + "kind": "start", + "classes": [], + "span": { + "line": 3, + "column": 3 + }, + "admission_hooks": "step", + "edges": { + "0": { + "to": "approve", + "label": null + } + } + } + }, + "firing": 1, + "visit": 1, + "attempt": 1, + "generation": 0, + "branch": { + "role": "none" + } + }, + "recorded_at": "[EPOCH_MS]", + "record": { + "seq": 9, + "origin": "core", + "recorded_at": "[EPOCH_MS]", + "body": { + "event": "route.applied", + "kind": "edge", + "firing": 1, + "group": 0, + "edge": 0 + } + }, + "derived": { + "target": { + "id": 2, + "name": "approve", + "kind": "attractor/human", + "meta": { + "label": "Approve?", + "shape": "hexagon", + "kind": "human", + "classes": [], + "span": { + "line": 5, + "column": 3 + }, + "edges": { + "1": { + "to": "ship", + "label": "[A] Approve" + }, + "2": { + "to": "revise", + "label": "[R] Revise" + } + } + } + }, + "transition": "Continue", + "back": false + } + } + }, + { + "run_id": "[ULID]", + "stream_seq": 28, + "kind": "petri", + "id": "execution 0/10/0", + "recorded_at": "[EPOCH_MS]", + "item": { + "id": { + "log": "execution", + "execution": 0, + "seq": 10, + "index": 0 + }, + "origin": "core", + "context": { + "invocation": 0, + "execution": 0 + }, + "subject": { + "node": { + "id": 0, + "name": "start", + "kind": "attractor/stage", + "meta": { + "label": "Start", + "shape": "Mdiamond", + "kind": "start", + "classes": [], + "span": { + "line": 3, + "column": 3 + }, + "admission_hooks": "step", + "edges": { + "0": { + "to": "approve", + "label": null + } + } + } + }, + "firing": 1, + "visit": 1, + "attempt": 1, + "generation": 0, + "branch": { + "role": "none" + } + }, + "recorded_at": "[EPOCH_MS]", + "record": { + "seq": 10, + "origin": "core", + "recorded_at": "[EPOCH_MS]", + "body": { + "event": "token.emitted", + "edge": 0, + "generation": 0, + "payload": { + "outcome": "succeeded", + "failure_class": "" + }, + "from": 1 + } + } + } + }, + { + "run_id": "[ULID]", + "stream_seq": 29, + "kind": "petri", + "id": "execution 0/11/0", + "recorded_at": "[EPOCH_MS]", + "item": { + "id": { + "log": "execution", + "execution": 0, + "seq": 11, + "index": 0 + }, + "origin": "external", + "context": { + "invocation": 0, + "execution": 0 + }, + "subject": { + "node": { + "id": 2, + "name": "approve", + "kind": "attractor/human", + "meta": { + "label": "Approve?", + "shape": "hexagon", + "kind": "human", + "classes": [], + "span": { + "line": 5, + "column": 3 + }, + "edges": { + "1": { + "to": "ship", + "label": "[A] Approve" + }, + "2": { + "to": "revise", + "label": "[R] Revise" + } + } + } + }, + "firing": 2, + "visit": 1, + "attempt": 1, + "generation": 0, + "branch": { + "role": "none" + } + }, + "recorded_at": "[EPOCH_MS]", + "record": { + "seq": 11, + "origin": "external", + "recorded_at": "[EPOCH_MS]", + "body": { + "event": "admission.decided", + "decision_id": { + "attempt_start": { + "firing": 2, + "attempt": 1 + } + }, + "decision": "admit", + "trace": [] + } + } + } + }, + { + "run_id": "[ULID]", + "stream_seq": 30, + "kind": "petri", + "id": "execution 0/12/0", + "recorded_at": "[EPOCH_MS]", + "item": { + "id": { + "log": "execution", + "execution": 0, + "seq": 12, + "index": 0 + }, + "origin": "external", + "context": { + "invocation": 0, + "execution": 0 + }, + "subject": { + "node": { + "id": 2, + "name": "approve", + "kind": "attractor/human", + "meta": { + "label": "Approve?", + "shape": "hexagon", + "kind": "human", + "classes": [], + "span": { + "line": 5, + "column": 3 + }, + "edges": { + "1": { + "to": "ship", + "label": "[A] Approve" + }, + "2": { + "to": "revise", + "label": "[R] Revise" + } + } + } + }, + "firing": 2, + "visit": 1, + "attempt": 1, + "generation": 0, + "branch": { + "role": "none" + } + }, + "recorded_at": "[EPOCH_MS]", + "record": { + "seq": 12, + "origin": "external", + "recorded_at": "[EPOCH_MS]", + "body": { + "event": "step.started", + "firing": 2, + "attempt": 1 + } + } + } + }, + { + "run_id": "[ULID]", + "stream_seq": 31, + "kind": "petri", + "id": "execution 0/12/1", + "recorded_at": "[EPOCH_MS]", + "item": { + "id": { + "log": "execution", + "execution": 0, + "seq": 12, + "index": 1 + }, + "origin": "derived", + "context": { + "invocation": 0, + "execution": 0 + }, + "subject": { + "node": { + "id": 2, + "name": "approve", + "kind": "attractor/human", + "meta": { + "label": "Approve?", + "shape": "hexagon", + "kind": "human", + "classes": [], + "span": { + "line": 5, + "column": 3 + }, + "edges": { + "1": { + "to": "ship", + "label": "[A] Approve" + }, + "2": { + "to": "revise", + "label": "[R] Revise" + } + } + } + }, + "firing": 2, + "visit": 1, + "attempt": 1, + "generation": 0, + "branch": { + "role": "none" + } + }, + "recorded_at": "[EPOCH_MS]", + "derived": { + "event": "wait.state.changed", + "state": "running" + } + } + }, + { + "run_id": "[ULID]", + "stream_seq": 32, + "kind": "petri", + "id": "execution 0/13/0", + "recorded_at": "[EPOCH_MS]", + "item": { + "id": { + "log": "execution", + "execution": 0, + "seq": 13, + "index": 0 + }, + "origin": "external", + "context": { + "invocation": 0, + "execution": 0 + }, + "subject": { + "node": { + "id": 2, + "name": "approve", + "kind": "attractor/human", + "meta": { + "label": "Approve?", + "shape": "hexagon", + "kind": "human", + "classes": [], + "span": { + "line": 5, + "column": 3 + }, + "edges": { + "1": { + "to": "ship", + "label": "[A] Approve" + }, + "2": { + "to": "revise", + "label": "[R] Revise" + } + } + } + }, + "firing": 2, + "visit": 1, + "attempt": 1, + "generation": 0, + "branch": { + "role": "none" + } + }, + "recorded_at": "[EPOCH_MS]", + "record": { + "seq": 13, + "origin": "external", + "recorded_at": "[EPOCH_MS]", + "body": { + "event": "step.progress.recorded", + "firing": 2, + "ev": { + "custom": { + "$question": { + "id": "approve#2", + "text": "Approve?", + "options": [ + { + "key": "A", + "label": "[A] Approve" + }, + { + "key": "R", + "label": "[R] Revise" + } + ], + "default": "A", + "freeform": false, + "sensitive": false + } + } + } + } + }, + "derived": { + "parsed": { + "kind": "question", + "question": { + "id": "approve#2", + "text": "Approve?", + "options": [ + { + "key": "A", + "label": "[A] Approve" + }, + { + "key": "R", + "label": "[R] Revise" + } + ], + "default": "A", + "freeform": false, + "sensitive": false + } + } + } + } + }, + { + "run_id": "[ULID]", + "stream_seq": 33, + "kind": "petri", + "id": "execution 0/13/1", + "recorded_at": "[EPOCH_MS]", + "item": { + "id": { + "log": "execution", + "execution": 0, + "seq": 13, + "index": 1 + }, + "origin": "derived", + "context": { + "invocation": 0, + "execution": 0 + }, + "subject": { + "node": { + "id": 2, + "name": "approve", + "kind": "attractor/human", + "meta": { + "label": "Approve?", + "shape": "hexagon", + "kind": "human", + "classes": [], + "span": { + "line": 5, + "column": 3 + }, + "edges": { + "1": { + "to": "ship", + "label": "[A] Approve" + }, + "2": { + "to": "revise", + "label": "[R] Revise" + } + } + } + }, + "firing": 2, + "visit": 1, + "attempt": 1, + "generation": 0, + "branch": { + "role": "none" + } + }, + "recorded_at": "[EPOCH_MS]", + "derived": { + "event": "wait.state.changed", + "state": "awaiting_answer" + } + } + }, + { + "run_id": "[ULID]", + "stream_seq": 34, + "kind": "petri", + "id": "execution 0/14/0", + "recorded_at": "[EPOCH_MS]", + "item": { + "id": { + "log": "execution", + "execution": 0, + "seq": 14, + "index": 0 + }, + "origin": "external", + "context": { + "invocation": 0, + "execution": 0 + }, + "subject": { + "node": { + "id": 2, + "name": "approve", + "kind": "attractor/human", + "meta": { + "label": "Approve?", + "shape": "hexagon", + "kind": "human", + "classes": [], + "span": { + "line": 5, + "column": 3 + }, + "edges": { + "1": { + "to": "ship", + "label": "[A] Approve" + }, + "2": { + "to": "revise", + "label": "[R] Revise" + } + } + } + }, + "firing": 2, + "visit": 1, + "attempt": 1, + "generation": 0, + "branch": { + "role": "none" + } + }, + "recorded_at": "[EPOCH_MS]", + "record": { + "seq": 14, + "origin": "external", + "recorded_at": "[EPOCH_MS]", + "body": { + "event": "step.progress.recorded", + "firing": 2, + "ev": { + "log": { + "stream": "stdout", + "line": "waiting for an answer: Approve?" + } + } + } + } + } } ] "#); @@ -1445,7 +2609,8 @@ fn attach_json_errors_without_prompting_for_human_input() { 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() diff --git a/lib/apps/fabro-cli/tests/it/cmd/create.rs b/lib/apps/fabro-cli/tests/it/cmd/create.rs index 31fdc67da..2b4ee82df 100644 --- a/lib/apps/fabro-cli/tests/it/cmd/create.rs +++ b/lib/apps/fabro-cli/tests/it/cmd/create.rs @@ -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); diff --git a/lib/apps/fabro-cli/tests/it/cmd/diff.rs b/lib/apps/fabro-cli/tests/it/cmd/diff.rs index febbb3f4d..0a8370413 100644 --- a/lib/apps/fabro-cli/tests/it/cmd/diff.rs +++ b/lib/apps/fabro-cli/tests/it/cmd/diff.rs @@ -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 ----- - "); -} diff --git a/lib/apps/fabro-cli/tests/it/cmd/dump.rs b/lib/apps/fabro-cli/tests/it/cmd/dump.rs index 4fb079cf9..e23a1712c 100644 --- a/lib/apps/fabro-cli/tests/it/cmd/dump.rs +++ b/lib/apps/fabro-cli/tests/it/cmd/dump.rs @@ -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 "); diff --git a/lib/apps/fabro-cli/tests/it/cmd/events.rs b/lib/apps/fabro-cli/tests/it/cmd/events.rs index 9de00fde9..8cd7a6155 100644 --- a/lib/apps/fabro-cli/tests/it/cmd/events.rs +++ b/lib/apps/fabro-cli/tests/it/cmd/events.rs @@ -14,11 +14,19 @@ fn parse_ndjson(stdout: &[u8]) -> Vec { .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 ----- "); } diff --git a/lib/apps/fabro-cli/tests/it/cmd/inspect.rs b/lib/apps/fabro-cli/tests/it/cmd/inspect.rs index b6ffc0b71..277691ebe 100644 --- a/lib/apps/fabro-cli/tests/it/cmd/inspect.rs +++ b/lib/apps/fabro-cli/tests/it/cmd/inspect.rs @@ -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]" - } - } - ] - "# - ); -} diff --git a/lib/apps/fabro-cli/tests/it/cmd/json_global.rs b/lib/apps/fabro-cli/tests/it/cmd/json_global.rs index 9a9992806..8e5676afc 100644 --- a/lib/apps/fabro-cli/tests/it/cmd/json_global.rs +++ b/lib/apps/fabro-cli/tests/it/cmd/json_global.rs @@ -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] diff --git a/lib/apps/fabro-cli/tests/it/cmd/run.rs b/lib/apps/fabro-cli/tests/it/cmd/run.rs index 7e7188cb2..dc2212f04 100644 --- a/lib/apps/fabro-cli/tests/it/cmd/run.rs +++ b/lib/apps/fabro-cli/tests/it/cmd/run.rs @@ -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() ); } diff --git a/lib/apps/fabro-cli/tests/it/cmd/runner.rs b/lib/apps/fabro-cli/tests/it/cmd/runner.rs index dfc3b11fd..46a32f3a5 100644 --- a/lib/apps/fabro-cli/tests/it/cmd/runner.rs +++ b/lib/apps/fabro-cli/tests/it/cmd/runner.rs @@ -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); diff --git a/lib/apps/fabro-cli/tests/it/cmd/support.rs b/lib/apps/fabro-cli/tests/it/cmd/support.rs index 27c4cde88..8333bf934 100644 --- a/lib/apps/fabro-cli/tests/it/cmd/support.rs +++ b/lib/apps/fabro-cli/tests/it/cmd/support.rs @@ -12,7 +12,6 @@ use std::collections::BTreeMap; use std::path::{Path, PathBuf}; use std::process::Output; -use std::sync::atomic::{AtomicU64, Ordering}; use std::sync::{Arc, Mutex}; use std::time::{Duration, Instant}; @@ -34,7 +33,6 @@ use shlex::try_quote; const LOCAL_COMMAND_TIMEOUT: Duration = Duration::from_secs(30); const CI_COMMAND_TIMEOUT: Duration = Duration::from_secs(90); -static NEXT_SEEDED_EVENT_ID: AtomicU64 = AtomicU64::new(1); pub(crate) use fabro_store::RunProjection; @@ -55,10 +53,6 @@ pub(crate) struct RunSetup { pub(crate) run_dir: PathBuf, } -pub(crate) struct SeededGitRunSetup { - pub(crate) run: RunSetup, -} - pub(crate) struct ProjectFixture { pub(crate) project_dir: PathBuf, pub(crate) fabro_root: PathBuf, @@ -73,12 +67,6 @@ pub(crate) struct WorkflowGate { gate_path: PathBuf, } -#[derive(Clone, Copy)] -enum SeededRunState { - Submitted, - Completed, -} - fn command_timeout() -> Duration { if std::env::var_os("CI").is_some() { CI_COMMAND_TIMEOUT @@ -386,12 +374,14 @@ pub(crate) fn setup_completed_fast_dry_run(context: &TestContext) -> RunSetup { run_completed_dry_run(context, &workflow) } +/// A completed run of the fast simple workflow: a real dry run, since a +/// run's history is what the engine recorded. pub(crate) fn setup_seeded_completed_dry_run(context: &TestContext) -> RunSetup { - block_on(seed_dry_run(context, SeededRunState::Completed)) + setup_completed_fast_dry_run(context) } pub(crate) fn setup_seeded_created_dry_run(context: &TestContext) -> RunSetup { - block_on(seed_dry_run(context, SeededRunState::Submitted)) + block_on(seed_dry_run(context)) } fn run_completed_dry_run(context: &TestContext, workflow: &Path) -> RunSetup { @@ -409,14 +399,27 @@ fn run_completed_dry_run(context: &TestContext, workflow: &Path) -> RunSetup { stderr(&output) ); } - let run_setup = single_run_setup(context); - wait_for_event_names(&run_setup.run_dir, &[ - "run.completed", - "sandbox.stop.completed", - ]); + let run_id = run_id_from_run_output(&output); + let run_setup = RunSetup { + run_dir: context.find_run_dir(&run_id), + run_id, + }; + wait_for_run_finished(&run_setup.run_dir); run_setup } +/// The run id `fabro run` prints (`Run: `) for the run it created. +fn run_id_from_run_output(output: &Output) -> String { + let text = stderr(output); + text.lines() + .find_map(|line| line.trim().strip_prefix("Run: ")) + .map_or_else( + || panic!("fabro run should print the run id:\n{text}"), + str::trim, + ) + .to_string() +} + fn fast_simple_workflow(context: &TestContext) -> PathBuf { let workflow = context.temp_dir.join("simple.fabro"); if !workflow.exists() { @@ -479,16 +482,8 @@ pub(crate) fn setup_detached_dry_run(context: &TestContext) -> RunSetup { run } -pub(crate) fn setup_seeded_git_backed_changed_run(context: &TestContext) -> SeededGitRunSetup { - block_on(seed_git_backed_changed_run(context)) -} - -pub(crate) fn setup_seeded_git_backed_noop_run(context: &TestContext) -> RunSetup { - block_on(seed_git_backed_noop_run(context)) -} - pub(crate) fn setup_seeded_artifact_run(context: &TestContext) -> RunSetup { - block_on(seed_artifact_run(context)) + seed_artifact_run(context) } pub(crate) fn setup_project_fixture(context: &TestContext) -> ProjectFixture { @@ -538,6 +533,9 @@ goal = "Exercise sandbox commands" [run.environment] id = "local" +[environments.local] +provider = "local" + "#, ); @@ -687,32 +685,6 @@ fn run_dirs_for_test_case(context: &TestContext) -> Vec { .collect() } -pub(crate) fn git_filters(context: &TestContext) -> Vec<(String, String)> { - let mut filters = context.filters(); - filters.push((r"\b[0-9a-f]{7,40}\b".to_string(), "[SHA]".to_string())); - filters.push(( - r"(fabro resume )[0-9A-HJKMNP-TV-Z]{8}\b".to_string(), - "$1[RUN_PREFIX]".to_string(), - )); - filters.push(( - r"(Forked run )[0-9A-HJKMNP-TV-Z]{8}\b".to_string(), - "$1[RUN_PREFIX]".to_string(), - )); - filters.push(( - r"(-> )[0-9A-HJKMNP-TV-Z]{8}\b".to_string(), - "$1[RUN_PREFIX]".to_string(), - )); - filters.push(( - r"(Rewound )[0-9A-HJKMNP-TV-Z]{8}\b".to_string(), - "$1[RUN_PREFIX]".to_string(), - )); - filters.push(( - r"(; new run )[0-9A-HJKMNP-TV-Z]{8}\b".to_string(), - "$1[RUN_PREFIX]".to_string(), - )); - filters -} - #[expect( clippy::disallowed_methods, reason = "This sync integration helper polls for the run directory to appear without requiring a Tokio runtime." @@ -901,36 +873,55 @@ pub(crate) fn command_log_text(run_dir: &Path, stage_id: &StageId) -> String { String::from_utf8(bytes).expect("command log should be UTF-8") } +/// Wait until the run's stream holds the terminal lifecycle record. +pub(crate) fn wait_for_run_finished(run_dir: &Path) { + wait_for_stream_item(run_dir, "the terminal lifecycle record", |item| { + crate::support::is_terminal_lifecycle(item) + }); +} + +/// Wait until the run's stream holds the `run.lifecycle` record of +/// `transition` (`running`, `succeeded`, ...). +pub(crate) fn wait_for_lifecycle(run_dir: &Path, transition: &str) { + wait_for_stream_item( + run_dir, + &format!("the {transition} lifecycle record"), + |item| { + let record = item.item.get("record"); + record + .and_then(|record| record.get("kind")) + .and_then(serde_json::Value::as_str) + == Some("run.lifecycle") + && record + .and_then(|record| record.get("transition")) + .and_then(serde_json::Value::as_str) + == Some(transition) + }, + ); +} + #[expect( clippy::disallowed_methods, - reason = "This sync integration helper polls stored events without requiring a Tokio runtime." + reason = "This sync integration helper polls the run stream without requiring a Tokio runtime." )] -pub(crate) fn wait_for_event_names(run_dir: &Path, expected: &[&str]) { +fn wait_for_stream_item(run_dir: &Path, what: &str, matches: impl Fn(&RunStreamItem) -> bool) { let deadline = std::time::Instant::now() + command_timeout(); - loop { - let event_names = run_events(run_dir) - .into_iter() - .filter_map(|item| item.name().map(str::to_string)) - .collect::>(); - - if expected - .iter() - .all(|expected_name| event_names.iter().any(|name| name == expected_name)) - { + if run_events(run_dir).iter().any(&matches) { return; } - assert!( std::time::Instant::now() < deadline, - "timed out waiting for events {expected:?}; saw {event_names:?}" + "timed out waiting for {what} in {}", + run_dir.display() ); std::thread::sleep(std::time::Duration::from_millis(50)); } } -async fn seed_dry_run(context: &TestContext, state: SeededRunState) -> RunSetup { - let run = create_seeded_run( +/// A created, unstarted dry run of the fast simple workflow. +async fn seed_dry_run(context: &TestContext) -> RunSetup { + create_seeded_run( context, "simple.fabro", fast_simple_workflow_source(), @@ -942,109 +933,40 @@ async fn seed_dry_run(context: &TestContext, state: SeededRunState) -> RunSetup }, false, ) - .await; - - if matches!(state, SeededRunState::Completed) { - let (client, base_url) = server_endpoint(&context.storage_dir) - .expect("test server endpoint should be available for seeded run events"); - append_seeded_simple_completion_events(&client, &base_url, &run, context).await; - } - - run + .await } -async fn seed_git_backed_changed_run(context: &TestContext) -> SeededGitRunSetup { - let step_one_sha = "2222222222222222222222222222222222222222"; - let step_two_sha = "3333333333333333333333333333333333333333"; - let run = create_seeded_run( - context, - "flow.fabro", - changed_git_workflow_source(), - RunIntentArgs { - provider: Some("openai".to_string()), - labels: test_label_map(context), - ..Default::default() - }, - true, - ) - .await; - - let base_sha = run_git(&context.temp_dir, &["rev-parse", "HEAD"]); - let base_sha = base_sha.trim(); - let (client, base_url) = server_endpoint(&context.storage_dir) - .expect("test server endpoint should be available for seeded run events"); - append_seeded_git_completion_events( - &client, - &base_url, - &run, - context, - base_sha, - step_one_sha, - step_two_sha, - ) - .await; - - SeededGitRunSetup { run } -} - -async fn seed_git_backed_noop_run(context: &TestContext) -> RunSetup { - let run = create_seeded_run( - context, - "flow.fabro", - noop_git_workflow_source(), - RunIntentArgs { - provider: Some("openai".to_string()), - labels: test_label_map(context), - ..Default::default() - }, - true, - ) - .await; - - let base_sha = run_git(&context.temp_dir, &["rev-parse", "HEAD"]); - let base_sha = base_sha.trim(); - let (client, base_url) = server_endpoint(&context.storage_dir) - .expect("test server endpoint should be available for seeded run events"); - append_seeded_git_noop_events(&client, &base_url, &run, context, base_sha).await; - run -} - -async fn seed_artifact_run(context: &TestContext) -> RunSetup { - let run = create_seeded_run( - context, - "artifact_run.fabro", - artifact_workflow_source(), - RunIntentArgs { - labels: test_label_map(context), - ..Default::default() - }, - false, - ) - .await; +/// A completed dry run of the artifact workflow, with artifacts uploaded +/// for its stages through the API. +fn seed_artifact_run(context: &TestContext) -> RunSetup { + let workflow = context.temp_dir.join("artifact_run.fabro"); + write_text_file(&workflow, artifact_workflow_source()); + let run = run_completed_dry_run(context, &workflow); let (client, base_url) = server_endpoint(&context.storage_dir) .expect("test server endpoint should be available for seeded artifacts"); - append_seeded_artifact_run_events(&client, &base_url, &run, context).await; - for (stage_id, retry, path, contents) in [ - ("create_assets@1", 1, "assets/node_a/summary.txt", "alpha"), - ("create_assets@1", 1, "assets/shared/report.txt", "one"), - ("create_assets@2", 1, "assets/shared/report.txt", "two"), - ("create_colliding@1", 1, "assets/other/summary.txt", "beta"), - ("create_colliding@1", 1, "assets/retry/report.txt", "second"), - ("retry_assets@1", 1, "assets/retry/report.txt", "first"), - ("retry_assets@1", 2, "assets/retry/report.txt", "second"), - ] { - upload_seeded_artifact( - &client, - &base_url, - &run.run_id, - stage_id, - retry, - path, - contents, - ) - .await; - } + block_on(async { + for (stage_id, retry, path, contents) in [ + ("create_assets@1", 1, "assets/node_a/summary.txt", "alpha"), + ("create_assets@1", 1, "assets/shared/report.txt", "one"), + ("create_assets@2", 1, "assets/shared/report.txt", "two"), + ("create_colliding@1", 1, "assets/other/summary.txt", "beta"), + ("create_colliding@1", 1, "assets/retry/report.txt", "second"), + ("retry_assets@1", 1, "assets/retry/report.txt", "first"), + ("retry_assets@1", 2, "assets/retry/report.txt", "second"), + ] { + upload_seeded_artifact( + &client, + &base_url, + &run.run_id, + stage_id, + retry, + path, + contents, + ) + .await; + } + }); run } @@ -1107,459 +1029,6 @@ async fn create_seeded_run( } } -async fn append_seeded_simple_completion_events( - client: &fabro_http::HttpClient, - base_url: &str, - run: &RunSetup, - context: &TestContext, -) { - append_run_event( - client, - base_url, - &run.run_id, - None, - "sandbox.ready", - serde_json::json!({ - "provider": "local", - "duration_ms": 1, - "name": null, - "cpu": null, - "memory": null, - "url": null, - }), - ) - .await; - append_run_event( - client, - base_url, - &run.run_id, - None, - "sandbox.initialized", - serde_json::json!({ - "working_directory": context.temp_dir.display().to_string(), - "provider": "local", - "id": fabro_sandbox::test_support::local_sandbox_id(&context.temp_dir).await, - "repo_cloned": false, - "clone_origin_url": null, - "clone_branch": null, - }), - ) - .await; - append_run_event( - client, - base_url, - &run.run_id, - None, - "run.started", - serde_json::json!({ - "name": "Simple", - "base_branch": null, - "base_sha": null, - "run_branch": null, - "worktree_dir": null, - "goal": "Run tests and report results", - }), - ) - .await; - append_run_event( - client, - base_url, - &run.run_id, - None, - "run.runnable", - serde_json::json!({ "source": "start_requested" }), - ) - .await; - append_run_event( - client, - base_url, - &run.run_id, - None, - "run.starting", - serde_json::json!({}), - ) - .await; - append_run_event( - client, - base_url, - &run.run_id, - None, - "run.running", - serde_json::json!({}), - ) - .await; - - append_seeded_stage(client, base_url, &run.run_id, "start", "Start", 0, None).await; - append_seeded_edge(client, base_url, &run.run_id, "start", "run_tests").await; - append_seeded_stage( - client, - base_url, - &run.run_id, - "run_tests", - "Run Tests", - 1, - Some("Dry run: would execute `true`."), - ) - .await; - append_seeded_edge(client, base_url, &run.run_id, "run_tests", "report").await; - append_seeded_stage( - client, - base_url, - &run.run_id, - "report", - "Report", - 2, - Some("Dry run: would execute `true`."), - ) - .await; - append_seeded_edge(client, base_url, &run.run_id, "report", "exit").await; - append_seeded_stage(client, base_url, &run.run_id, "exit", "Exit", 3, None).await; - append_run_event( - client, - base_url, - &run.run_id, - Some("report"), - "checkpoint.completed", - checkpoint_properties( - "success", - "report", - &["start", "run_tests", "report"], - Some("exit"), - None, - None, - ), - ) - .await; - append_run_event( - client, - base_url, - &run.run_id, - None, - "run.completed", - serde_json::json!({ - "timing": {"wall_time_ms": 123, "inference_time_ms": 0, "tool_time_ms": 0, "active_time_ms": 0}, - "artifact_count": 0, - "status": "succeeded", - "reason": "completed", - "final_git_commit_sha": null, - "final_patch": null, - "usage": null, - }), - ) - .await; - append_run_event( - client, - base_url, - &run.run_id, - None, - "sandbox.stop.started", - serde_json::json!({ - "provider": "local", - }), - ) - .await; - append_run_event( - client, - base_url, - &run.run_id, - None, - "sandbox.stop.completed", - serde_json::json!({ - "provider": "local", - "duration_ms": 1, - }), - ) - .await; -} - -async fn append_seeded_git_completion_events( - client: &fabro_http::HttpClient, - base_url: &str, - run: &RunSetup, - context: &TestContext, - base_sha: &str, - step_one_sha: &str, - step_two_sha: &str, -) { - append_run_event( - client, - base_url, - &run.run_id, - None, - "sandbox.ready", - serde_json::json!({ - "provider": "local", - "duration_ms": 1, - "name": null, - "cpu": null, - "memory": null, - "url": null, - }), - ) - .await; - append_run_event( - client, - base_url, - &run.run_id, - None, - "sandbox.initialized", - serde_json::json!({ - "working_directory": context.temp_dir.display().to_string(), - "provider": "local", - "id": fabro_sandbox::test_support::local_sandbox_id(&context.temp_dir).await, - "repo_cloned": false, - "clone_origin_url": null, - "clone_branch": null, - }), - ) - .await; - append_run_event( - client, - base_url, - &run.run_id, - None, - "run.started", - serde_json::json!({ - "name": "Flow", - "base_branch": "main", - "base_sha": base_sha, - "run_branch": format!("fabro/run/{}", run.run_id), - "worktree_dir": context.temp_dir.display().to_string(), - "goal": "Edit a tracked file", - }), - ) - .await; - append_run_event( - client, - base_url, - &run.run_id, - None, - "run.runnable", - serde_json::json!({ "source": "start_requested" }), - ) - .await; - append_run_event( - client, - base_url, - &run.run_id, - None, - "run.starting", - serde_json::json!({}), - ) - .await; - append_run_event( - client, - base_url, - &run.run_id, - None, - "run.running", - serde_json::json!({}), - ) - .await; - append_run_event( - client, - base_url, - &run.run_id, - Some("start"), - "checkpoint.completed", - checkpoint_properties( - "succeeded", - "start", - &["start"], - Some("step_one"), - None, - None, - ), - ) - .await; - append_run_event( - client, - base_url, - &run.run_id, - Some("step_one"), - "checkpoint.completed", - checkpoint_properties( - "success", - "step_one", - &["start", "step_one"], - Some("step_two"), - Some(step_one_sha), - Some(step_one_patch()), - ), - ) - .await; - append_run_event( - client, - base_url, - &run.run_id, - Some("step_two"), - "checkpoint.completed", - checkpoint_properties( - "success", - "step_two", - &["start", "step_one", "step_two"], - Some("exit"), - Some(step_two_sha), - Some(step_two_patch()), - ), - ) - .await; - append_run_event( - client, - base_url, - &run.run_id, - None, - "run.completed", - serde_json::json!({ - "timing": {"wall_time_ms": 456, "inference_time_ms": 0, "tool_time_ms": 0, "active_time_ms": 0}, - "artifact_count": 0, - "status": "succeeded", - "reason": "completed", - "final_git_commit_sha": step_two_sha, - "final_patch": final_story_patch(), - "usage": null, - }), - ) - .await; -} - -async fn append_seeded_git_noop_events( - client: &fabro_http::HttpClient, - base_url: &str, - run: &RunSetup, - context: &TestContext, - base_sha: &str, -) { - append_run_event( - client, - base_url, - &run.run_id, - None, - "run.started", - serde_json::json!({ - "name": "Flow", - "base_branch": "main", - "base_sha": base_sha, - "run_branch": format!("fabro/run/{}", run.run_id), - "worktree_dir": context.temp_dir.display().to_string(), - "goal": "Leave tracked files unchanged", - }), - ) - .await; - append_run_event( - client, - base_url, - &run.run_id, - None, - "run.runnable", - serde_json::json!({ "source": "start_requested" }), - ) - .await; - append_run_event( - client, - base_url, - &run.run_id, - None, - "run.starting", - serde_json::json!({}), - ) - .await; - append_run_event( - client, - base_url, - &run.run_id, - None, - "run.running", - serde_json::json!({}), - ) - .await; - append_run_event( - client, - base_url, - &run.run_id, - None, - "run.completed", - serde_json::json!({ - "timing": {"wall_time_ms": 123, "inference_time_ms": 0, "tool_time_ms": 0, "active_time_ms": 0}, - "artifact_count": 0, - "status": "succeeded", - "reason": "completed", - "final_git_commit_sha": base_sha, - "final_patch": null, - "usage": null, - }), - ) - .await; -} - -async fn append_seeded_artifact_run_events( - client: &fabro_http::HttpClient, - base_url: &str, - run: &RunSetup, - context: &TestContext, -) { - append_run_event( - client, - base_url, - &run.run_id, - None, - "run.started", - serde_json::json!({ - "name": "ArtifactRun", - "base_branch": null, - "base_sha": null, - "run_branch": null, - "worktree_dir": context.temp_dir.display().to_string(), - "goal": "Exercise artifact commands", - }), - ) - .await; - append_run_event( - client, - base_url, - &run.run_id, - None, - "run.runnable", - serde_json::json!({ "source": "start_requested" }), - ) - .await; - append_run_event( - client, - base_url, - &run.run_id, - None, - "run.starting", - serde_json::json!({}), - ) - .await; - append_run_event( - client, - base_url, - &run.run_id, - None, - "run.running", - serde_json::json!({}), - ) - .await; - append_run_event( - client, - base_url, - &run.run_id, - None, - "run.completed", - serde_json::json!({ - "timing": {"wall_time_ms": 123, "inference_time_ms": 0, "tool_time_ms": 0, "active_time_ms": 0}, - "artifact_count": 7, - "status": "succeeded", - "reason": "completed", - "final_git_commit_sha": null, - "final_patch": null, - "usage": null, - }), - ) - .await; -} - async fn upload_seeded_artifact( client: &fabro_http::HttpClient, base_url: &str, @@ -1586,109 +1055,6 @@ async fn upload_seeded_artifact( .await; } -async fn append_seeded_stage( - client: &fabro_http::HttpClient, - base_url: &str, - run_id: &str, - node_id: &str, - name: &str, - index: usize, - response: Option<&str>, -) { - append_run_event( - client, - base_url, - run_id, - Some(node_id), - "stage.started", - serde_json::json!({ - "index": index, - "handler_type": "noop", - "attempt": 1, - "max_attempts": 1, - }), - ) - .await; - append_run_event( - client, - base_url, - run_id, - Some(node_id), - "stage.completed", - stage_completed_properties(index, response), - ) - .await; - - let _ = name; -} - -async fn append_seeded_edge( - client: &fabro_http::HttpClient, - base_url: &str, - run_id: &str, - from_node: &str, - to_node: &str, -) { - append_run_event( - client, - base_url, - run_id, - Some(from_node), - "edge.selected", - serde_json::json!({ - "from_node": from_node, - "to_node": to_node, - "label": null, - "condition": null, - "reason": "unconditional", - "preferred_label": null, - "suggested_next_ids": [], - "stage_status": "succeeded", - "is_jump": false, - }), - ) - .await; -} - -async fn append_run_event( - client: &fabro_http::HttpClient, - base_url: &str, - run_id: &str, - node_id: Option<&str>, - event_name: &str, - properties: serde_json::Value, -) { - let event_id = NEXT_SEEDED_EVENT_ID.fetch_add(1, Ordering::Relaxed); - let mut event = serde_json::json!({ - "id": format!("00000000-0000-0000-0000-{event_id:012x}"), - "ts": chrono::Utc::now().to_rfc3339(), - "run_id": run_id, - "event": event_name, - "properties": properties, - "actor": { - "kind": "worker", - "run_id": run_id, - }, - }); - if let Some(node_id) = node_id { - event["node_id"] = serde_json::Value::String(node_id.to_string()); - event["node_label"] = serde_json::Value::String(node_label(node_id).to_string()); - } - - let response = client - .post(format!("{base_url}/api/v1/runs/{run_id}/events")) - .json(&event) - .send() - .await - .unwrap_or_else(|err| panic!("append seeded event {event_name} should execute: {err}")); - expect_reqwest_status( - response, - fabro_http::StatusCode::OK, - format!("POST /api/v1/runs/{run_id}/events ({event_name})"), - ) - .await; -} - fn test_label_map(context: &TestContext) -> std::collections::HashMap { test_labels(context) .into_iter() @@ -1705,67 +1071,6 @@ fn test_labels(context: &TestContext) -> Vec { vec![context.test_run_label(), context.test_case_label()] } -fn stage_completed_properties(index: usize, response: Option<&str>) -> serde_json::Value { - serde_json::json!({ - "index": index, - "timing": {"wall_time_ms": 1, "inference_time_ms": 0, "tool_time_ms": 0, "active_time_ms": 0}, - "status": "succeeded", - "preferred_label": null, - "suggested_next_ids": [], - "usage": null, - "failure": null, - "notes": null, - "files_touched": [], - "context_updates": null, - "jump_to_node": null, - "context_values": null, - "node_visits": null, - "loop_failure_signatures": null, - "restart_failure_signatures": null, - "response": response, - "attempt": 1, - "max_attempts": 1, - }) -} - -fn checkpoint_properties( - status: &str, - current_node: &str, - completed_nodes: &[&str], - next_node_id: Option<&str>, - git_commit_sha: Option<&str>, - diff: Option<&str>, -) -> serde_json::Value { - serde_json::json!({ - "status": status, - "current_node": current_node, - "completed_nodes": completed_nodes, - "node_retries": {}, - "context_values": {}, - "node_outcomes": {}, - "next_node_id": next_node_id, - "git_commit_sha": git_commit_sha, - "loop_failure_signatures": {}, - "restart_failure_signatures": {}, - "node_visits": { - (current_node): 1, - }, - "diff": diff, - }) -} - -fn node_label(node_id: &str) -> &str { - match node_id { - "start" => "Start", - "run_tests" => "Run Tests", - "report" => "Report", - "exit" => "Exit", - "step_one" => "step_one", - "step_two" => "step_two", - other => other, - } -} - fn fast_simple_workflow_source() -> &'static str { r#"digraph Simple { graph [goal="Run tests and report results"] @@ -1782,29 +1087,6 @@ fn fast_simple_workflow_source() -> &'static str { "# } -fn changed_git_workflow_source() -> &'static str { - r#"digraph Flow { - graph [goal="Edit a tracked file"]; - start [shape=Mdiamond]; - exit [shape=Msquare]; - step_one [shape=parallelogram, script="printf 'line 1\nline 2\n' > story.txt"]; - step_two [shape=parallelogram, script="printf 'line 1\nline 2\nline 3\n' > story.txt"]; - start -> step_one -> step_two -> exit; -} -"# -} - -fn noop_git_workflow_source() -> &'static str { - r#"digraph Flow { - graph [goal="Leave tracked files unchanged"]; - start [shape=Mdiamond]; - exit [shape=Msquare]; - check [shape=parallelogram, script="test -f story.txt"]; - start -> check -> exit; -} -"# -} - fn artifact_workflow_source() -> &'static str { r#"digraph ArtifactRun { graph [goal="Exercise artifact commands", default_max_retries=0] @@ -1818,18 +1100,6 @@ fn artifact_workflow_source() -> &'static str { "# } -fn step_one_patch() -> &'static str { - "diff --git a/story.txt b/story.txt\nindex 1111111..2222222 100644\n--- a/story.txt\n+++ b/story.txt\n@@ -1 +1,2 @@\n line 1\n+line 2\n" -} - -fn step_two_patch() -> &'static str { - "diff --git a/story.txt b/story.txt\nindex 2222222..3333333 100644\n--- a/story.txt\n+++ b/story.txt\n@@ -1,2 +1,3 @@\n line 1\n line 2\n+line 3\n" -} - -fn final_story_patch() -> &'static str { - "diff --git a/story.txt b/story.txt\nindex 1111111..3333333 100644\n--- a/story.txt\n+++ b/story.txt\n@@ -1 +1,3 @@\n line 1\n+line 2\n+line 3\n" -} - pub(crate) fn text_tree(root: &Path) -> Vec { fn visit(root: &Path, dir: &Path, entries: &mut Vec) { let mut children: Vec<_> = std::fs::read_dir(dir) @@ -1927,70 +1197,6 @@ pub(crate) fn compact_inspect(output: &Output) -> Value { ) } -pub(crate) fn compact_git_inspect(output: &Output) -> Value { - let items: Vec = - serde_json::from_str(&stdout(output)).expect("inspect output should be valid JSON"); - Value::Array( - items.into_iter() - .map(|item| { - let run_spec = item["run_spec"].clone(); - let start_record = item["start_record"].clone(); - let checkpoint = item["checkpoint"].clone(); - let conclusion = item["conclusion"].clone(); - let sandbox = item["sandbox"].clone(); - serde_json::json!({ - "run_id": "[ULID]", - "status": item["status"], - "run_spec": { - "goal": run_spec.pointer("/settings/run/goal"), - "workflow_name": run_spec.pointer("/graph/name"), - "workflow_slug": run_spec.pointer("/workflow_slug"), - "llm_provider": run_spec.pointer("/settings/run/model/provider"), - "sandbox_provider": run_spec.pointer("/settings/run/sandbox/provider"), - "provenance": run_spec.pointer("/provenance").as_ref().map(|_| { - serde_json::json!({ - "server_version": "[VERSION]", - "client_name": run_spec.pointer("/provenance/client/name"), - "client_version": "[VERSION]", - "subject_auth_method": run_spec.pointer("/provenance/subject/auth_method"), - }) - }), - }, - "start_record": start_record.as_object().map(|_| { - serde_json::json!({ - "has_start_time": true, - "run_branch": "fabro/run/[ULID]", - "base_sha": "[SHA]", - }) - }), - "conclusion": conclusion.as_object().map(|_| { - serde_json::json!({ - "status": conclusion["status"], - "timing": "[TIMING]", - "final_git_commit_sha": "[SHA]", - "stage_count": conclusion["stages"].as_array().map(|stages| stages.len()), - }) - }), - "checkpoint": checkpoint.as_object().map(|_| { - serde_json::json!({ - "current_node": checkpoint["current_node"], - "completed_nodes": checkpoint["completed_nodes"], - "next_node_id": checkpoint["next_node_id"], - "git_commit_sha": "[SHA]", - }) - }), - "sandbox": sandbox.as_object().map(|_| { - serde_json::json!({ - "provider": compact_sandbox_provider(&sandbox), - "working_directory": "[WORKTREE]", - }) - }), - }) - }) - .collect(), - ) -} - fn compact_sandbox_provider(sandbox: &Value) -> Value { sandbox .pointer("/instance/provider") diff --git a/lib/apps/fabro-cli/tests/it/support/mod.rs b/lib/apps/fabro-cli/tests/it/support/mod.rs index df937ad5b..7e52a2126 100644 --- a/lib/apps/fabro-cli/tests/it/support/mod.rs +++ b/lib/apps/fabro-cli/tests/it/support/mod.rs @@ -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(), diff --git a/lib/apps/fabro-cli/tests/it/workflow/dry_run_examples.rs b/lib/apps/fabro-cli/tests/it/workflow/dry_run_examples.rs index a6facf8d0..35630a335 100644 --- a/lib/apps/fabro-cli/tests/it/workflow/dry_run_examples.rs +++ b/lib/apps/fabro-cli/tests/it/workflow/dry_run_examples.rs @@ -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]