diff --git a/Cargo.lock b/Cargo.lock index 300041ee1..dd9a7ac07 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -1642,7 +1642,7 @@ checksum = "d7a1e2f27636f116493b8b860f5546edb47c8d8f8ea73e1d2a20be88e28d1fea" [[package]] name = "daytona-api-client" version = "0.1.0" -source = "git+https://github.com/brynary/daytona-sdk-rust?branch=main#0e69058c888a6562c70a5b16d707253914cff563" +source = "git+https://github.com/brynary/daytona-sdk-rust?branch=main#df596bad2093f83fc793d1890bd25ba0d2f93425" dependencies = [ "reqwest 0.13.4", "reqwest-middleware", @@ -1656,7 +1656,7 @@ dependencies = [ [[package]] name = "daytona-sdk" version = "0.1.0" -source = "git+https://github.com/brynary/daytona-sdk-rust?branch=main#0e69058c888a6562c70a5b16d707253914cff563" +source = "git+https://github.com/brynary/daytona-sdk-rust?branch=main#df596bad2093f83fc793d1890bd25ba0d2f93425" dependencies = [ "daytona-api-client", "daytona-toolbox-client", @@ -1676,7 +1676,7 @@ dependencies = [ [[package]] name = "daytona-toolbox-client" version = "0.1.0" -source = "git+https://github.com/brynary/daytona-sdk-rust?branch=main#0e69058c888a6562c70a5b16d707253914cff563" +source = "git+https://github.com/brynary/daytona-sdk-rust?branch=main#df596bad2093f83fc793d1890bd25ba0d2f93425" dependencies = [ "reqwest 0.13.4", "reqwest-middleware", @@ -5166,7 +5166,7 @@ checksum = "9b4f627cb1b25917193a259e49bdad08f671f8d9708acfd5fe0a8c1455d87220" [[package]] name = "petri-attractor-steps" version = "0.1.0" -source = "git+https://github.com/lithoscomputer/petri.git?branch=main#e46845bd0139dd04e9e795d3201b5b9b4b8b1026" +source = "git+https://github.com/lithoscomputer/petri.git?branch=main#505a17839987f052ef0a47b799d30c565fd29475" dependencies = [ "async-trait", "globset", @@ -5197,7 +5197,7 @@ dependencies = [ [[package]] name = "petri-driver" version = "0.1.0" -source = "git+https://github.com/lithoscomputer/petri.git?branch=main#e46845bd0139dd04e9e795d3201b5b9b4b8b1026" +source = "git+https://github.com/lithoscomputer/petri.git?branch=main#505a17839987f052ef0a47b799d30c565fd29475" dependencies = [ "async-trait", "getrandom 0.3.4", @@ -5217,7 +5217,7 @@ dependencies = [ [[package]] name = "petri-engine" version = "0.1.0" -source = "git+https://github.com/lithoscomputer/petri.git?branch=main#e46845bd0139dd04e9e795d3201b5b9b4b8b1026" +source = "git+https://github.com/lithoscomputer/petri.git?branch=main#505a17839987f052ef0a47b799d30c565fd29475" dependencies = [ "petri-ir", "serde", @@ -5229,7 +5229,7 @@ dependencies = [ [[package]] name = "petri-execution" version = "0.1.0" -source = "git+https://github.com/lithoscomputer/petri.git?branch=main#e46845bd0139dd04e9e795d3201b5b9b4b8b1026" +source = "git+https://github.com/lithoscomputer/petri.git?branch=main#505a17839987f052ef0a47b799d30c565fd29475" dependencies = [ "async-trait", "petri-driver", @@ -5253,7 +5253,7 @@ dependencies = [ [[package]] name = "petri-executor" version = "0.1.0" -source = "git+https://github.com/lithoscomputer/petri.git?branch=main#e46845bd0139dd04e9e795d3201b5b9b4b8b1026" +source = "git+https://github.com/lithoscomputer/petri.git?branch=main#505a17839987f052ef0a47b799d30c565fd29475" dependencies = [ "async-trait", "libc", @@ -5268,7 +5268,7 @@ dependencies = [ [[package]] name = "petri-executor-sandbox" version = "0.1.0" -source = "git+https://github.com/lithoscomputer/petri.git?branch=main#e46845bd0139dd04e9e795d3201b5b9b4b8b1026" +source = "git+https://github.com/lithoscomputer/petri.git?branch=main#505a17839987f052ef0a47b799d30c565fd29475" dependencies = [ "async-trait", "petri-executor", @@ -5290,7 +5290,7 @@ dependencies = [ [[package]] name = "petri-frontend" version = "0.1.0" -source = "git+https://github.com/lithoscomputer/petri.git?branch=main#e46845bd0139dd04e9e795d3201b5b9b4b8b1026" +source = "git+https://github.com/lithoscomputer/petri.git?branch=main#505a17839987f052ef0a47b799d30c565fd29475" dependencies = [ "marked-yaml", "petri-ir", @@ -5304,7 +5304,7 @@ dependencies = [ [[package]] name = "petri-frontend-attractor" version = "0.1.0" -source = "git+https://github.com/lithoscomputer/petri.git?branch=main#e46845bd0139dd04e9e795d3201b5b9b4b8b1026" +source = "git+https://github.com/lithoscomputer/petri.git?branch=main#505a17839987f052ef0a47b799d30c565fd29475" dependencies = [ "minijinja", "petri-frontend", @@ -5321,7 +5321,7 @@ dependencies = [ [[package]] name = "petri-frontend-fabro" version = "0.1.0" -source = "git+https://github.com/lithoscomputer/petri.git?branch=main#e46845bd0139dd04e9e795d3201b5b9b4b8b1026" +source = "git+https://github.com/lithoscomputer/petri.git?branch=main#505a17839987f052ef0a47b799d30c565fd29475" dependencies = [ "petri-frontend", "petri-frontend-attractor", @@ -5337,7 +5337,7 @@ dependencies = [ [[package]] name = "petri-frontend-native" version = "0.1.0" -source = "git+https://github.com/lithoscomputer/petri.git?branch=main#e46845bd0139dd04e9e795d3201b5b9b4b8b1026" +source = "git+https://github.com/lithoscomputer/petri.git?branch=main#505a17839987f052ef0a47b799d30c565fd29475" dependencies = [ "petri-frontend", "petri-ir", @@ -5348,7 +5348,7 @@ dependencies = [ [[package]] name = "petri-ir" version = "0.1.0" -source = "git+https://github.com/lithoscomputer/petri.git?branch=main#e46845bd0139dd04e9e795d3201b5b9b4b8b1026" +source = "git+https://github.com/lithoscomputer/petri.git?branch=main#505a17839987f052ef0a47b799d30c565fd29475" dependencies = [ "regex", "serde", @@ -5361,7 +5361,7 @@ dependencies = [ [[package]] name = "petri-runtime" version = "0.1.0" -source = "git+https://github.com/lithoscomputer/petri.git?branch=main#e46845bd0139dd04e9e795d3201b5b9b4b8b1026" +source = "git+https://github.com/lithoscomputer/petri.git?branch=main#505a17839987f052ef0a47b799d30c565fd29475" dependencies = [ "async-trait", "petri-driver", @@ -5382,7 +5382,7 @@ dependencies = [ [[package]] name = "petri-steps" version = "0.1.0" -source = "git+https://github.com/lithoscomputer/petri.git?branch=main#e46845bd0139dd04e9e795d3201b5b9b4b8b1026" +source = "git+https://github.com/lithoscomputer/petri.git?branch=main#505a17839987f052ef0a47b799d30c565fd29475" dependencies = [ "async-trait", "petri-executor", @@ -5398,7 +5398,7 @@ dependencies = [ [[package]] name = "petri-store" version = "0.1.0" -source = "git+https://github.com/lithoscomputer/petri.git?branch=main#e46845bd0139dd04e9e795d3201b5b9b4b8b1026" +source = "git+https://github.com/lithoscomputer/petri.git?branch=main#505a17839987f052ef0a47b799d30c565fd29475" dependencies = [ "async-trait", "getrandom 0.3.4", @@ -5413,7 +5413,7 @@ dependencies = [ [[package]] name = "petri-testkit" version = "0.1.0" -source = "git+https://github.com/lithoscomputer/petri.git?branch=main#e46845bd0139dd04e9e795d3201b5b9b4b8b1026" +source = "git+https://github.com/lithoscomputer/petri.git?branch=main#505a17839987f052ef0a47b799d30c565fd29475" dependencies = [ "async-trait", "petri-driver", diff --git a/docs/internal/events-strategy.md b/docs/internal/events-strategy.md index f6c404911..f1e5440c7 100644 --- a/docs/internal/events-strategy.md +++ b/docs/internal/events-strategy.md @@ -75,9 +75,12 @@ When matching items: `run.finished`); its parsed meaning is under `item.derived` (a pending question is `derived.parsed.kind == "question"`). - A platform record's kind is `item.record.kind`. -- The run has ended when a platform `run.lifecycle` record's `transition` - is `succeeded`, `failed` or `dead`. Petri's `run.finished` precedes it and - carries the engine's own status. +- Petri's coordinator `run.finished` commits the overall result after required + finalization and cleanup. Its optional `finalization_failure` carries the + host-defined code and message. Invocation and stage results remain execution + evidence. A subsequent platform `run.lifecycle` terminal transition reports + that same outcome; it cannot replace it. Early worker failures without a + Petri finish end at the platform terminal lifecycle record. Never rebuild an item downstream: pass the `RunStreamItem` through as read. diff --git a/docs/internal/run-finalization.md b/docs/internal/run-finalization.md new file mode 100644 index 000000000..50a35da51 --- /dev/null +++ b/docs/internal/run-finalization.md @@ -0,0 +1,59 @@ +# Required run finalization + +Petri owns the durable run result. Fabro declares required finalization for +every run and implements it in `FabroHooks::finalize_run`, so a resume or a +fork always matches the stored declaration. A failed checkpoint rejects with +`checkpoint_failed` and skips publication. The checkpoint failure cancelled the +run, so Petri may commit `cancelled`; Fabro reports that finish as a workflow +failure with the checkpoint's message. +For successful workflow execution this prepares the final diff, then uses the +existing publisher to retry the push and reconcile or create the pull request. +A preparation or publication error rejects with `publish_failed` and a rendered +message. Missing branch, checkpoint or workspace evidence is a rejection. +Failed or cancelled execution skips publication; its diff remains best effort. +Local workflow run-end hooks remain observational. + +The coordinator awaits finalization and cleanup before committing `run.finished`. +Until that record, the run stays active without a conclusion, even when every +stage has succeeded. The run stream, API projection, managed server status, +worker return and CLI verdict use the committed overall status and structured +finalization failure. Stage outcomes, metrics and invocation results keep their +execution meaning. Successful stages stay successful when publication fails. + +Both HTTP and in-process worker transports store the finish before settling +managed status. A rejected append cannot settle the run. The later platform +terminal lifecycle record acknowledges the same outcome; it cannot replace a +terminal result or conclusion. Early worker failures without a Petri finish +still end through their committed platform lifecycle record. Run scope, +authorization and lease ownership checks remain required for every worker write. + +## Recovery and stored history + +This integration requires a Petri with required run finalization (coordinator +format 9) and event contract 6. Petri retains its strict stored-format policy: version-8 coordinator +logs cannot resume, inspect or replay with this engine. Fabro does not rewrite +source history or relax that policy. Existing materialized views and stream +rows remain stored; a projector replay failure holds Petri positions and reports +incomplete record health. A historical run without a materialized view cannot +recover its old Petri stage evidence through the new engine. Operators requiring +inspection or recovery of old execution logs must keep the prior compatible +binary and its storage backup. + +Existing incorrect historical publication-success projections are not repaired +by this change. Consumed cursors and immutable conclusions remain unchanged; +historical-view repair requires a separate, bounded process over preserved +source records. No publication is invoked by projection or replay. + +A committed resume returns the same overall result without publishing again. +A fork declares required finalization while seeding its records, as every +Fabro run does; its worker restores the actual publication hooks before resume. +Unfinished recovery must restore the same finalization requirement. Petri may +call the restored finalizer again after interruption, including a crash after +the callback returns but before the terminal record commits. Existing GitHub +reconciliation is retained; this contract does not provide independent retry +orchestration or guarantee external-effect deduplication after every crash. +Lost terminal-execution workspaces are not recreated for finalization. Fabro +rejects publication when it cannot establish the required evidence. Push +attempts retain their existing five-minute timeout; this integration adds no +overall finalizer deadline. Cancellation arriving after +execution ends does not interrupt Petri's awaited finalizer. diff --git a/lib/apps/fabro-cli/src/commands/run/attach.rs b/lib/apps/fabro-cli/src/commands/run/attach.rs index 67998a6d8..15473dcf0 100644 --- a/lib/apps/fabro-cli/src/commands/run/attach.rs +++ b/lib/apps/fabro-cli/src/commands/run/attach.rs @@ -228,7 +228,7 @@ async fn attach_petri_run_with_client( for item in &items { emit_stream_item(&mut progress_ui, item, opts.json_output)?; if let Some(code) = petri_stream::exit_code_of(item) { - replayed_exit_code = Some(ExitCode::from(code)); + replayed_exit_code.get_or_insert(ExitCode::from(code)); } } if let Some(exit_code) = replayed_exit_code.or_else(|| { @@ -1114,4 +1114,109 @@ mod tests { assert_eq!(result.unwrap_err().to_string(), "questions unavailable"); assert_eq!(calls, 1); } + + #[tokio::test] + async fn required_publication_keeps_attach_open_until_the_committed_result() { + let run_id = fabro_types::fixtures::RUN_1; + let server = MockServer::start_async().await; + let mut state = terminal_run_state_response(run_id); + state["status"] = serde_json::json!({"kind": "running"}); + state["conclusion"] = serde_json::Value::Null; + let state: server_client::RunProjection = serde_json::from_value(state).unwrap(); + let executed = serde_json::json!({ + "run_id": run_id, "stream_seq": 1, "kind": "petri", "id": "executed", + "recorded_at": 1000, + "item": {"origin": "external", "context": {}, "record": {"seq": 5, + "body": {"event": "invocation.finished", "invocation": 0, + "result": {"status": "success"}}}} + }); + server + .mock_async(|when, then| { + when.method("GET") + .path(format!("/api/v1/runs/{run_id}/events")); + then.status(200).json_body(serde_json::json!({ + "data": [executed], "meta": {"has_more": false}, "event_contract_version": 6 + })); + }) + .await; + server + .mock_async(|when, then| { + when.method("GET") + .path(format!("/api/v1/runs/{run_id}/state")); + then.status(200) + .json_body(serde_json::to_value(&state).unwrap()); + }) + .await; + server + .mock_async(|when, then| { + when.method("GET") + .path(format!("/api/v1/runs/{run_id}/questions")); + then.status(200) + .json_body(serde_json::json!({ "data": [], "meta": { "has_more": false } })); + }) + .await; + let waiting = server + .mock_async(|when, then| { + when.method("GET") + .path(format!("/api/v1/runs/{run_id}/attach")); + then.status(200) + .header("Content-Type", "text/event-stream") + .body(""); + }) + .await; + let client = server_client::Client::new_no_proxy(&server.base_url()).unwrap(); + let task = tokio::spawn(async move { + Box::pin(attach_petri_run_with_client( + &client, + &run_id, + &state, + no_color_styles(), + AttachOptions { + auto_approve: false, + verbose: false, + kill_on_detach: false, + json_output: true, + }, + Printer::Default, + )) + .await + .unwrap() + }); + tokio::time::timeout(Duration::from_secs(5), async { + while waiting.calls_async().await == 0 { + sleep(Duration::from_millis(10)).await; + } + }) + .await + .unwrap(); + assert!( + !task.is_finished(), + "successful execution does not end attach while publication is pending" + ); + let finished = serde_json::json!({ + "run_id": run_id, "stream_seq": 2, "kind": "petri", "id": "finished", + "recorded_at": 2000, + "item": {"origin": "external", "context": {}, "record": {"seq": 6, + "body": {"event": "run.finished", "status": "failed", + "finalization_failure": {"code": "publish_failed", "message": "the push was rejected"}}}} + }); + let final_stream = server + .mock_async(|when, then| { + when.method("GET") + .path(format!("/api/v1/runs/{run_id}/attach")); + then.status(200) + .header("Content-Type", "text/event-stream") + .body(format!("data: {finished}\n\n")); + }) + .await; + waiting.delete_async().await; + assert_eq!( + tokio::time::timeout(Duration::from_secs(5), task) + .await + .unwrap() + .unwrap(), + ExitCode::from(1) + ); + final_stream.assert_async().await; + } } diff --git a/lib/apps/fabro-cli/src/commands/run/petri_stream.rs b/lib/apps/fabro-cli/src/commands/run/petri_stream.rs index 5b4be90d2..6d0241017 100644 --- a/lib/apps/fabro-cli/src/commands/run/petri_stream.rs +++ b/lib/apps/fabro-cli/src/commands/run/petri_stream.rs @@ -523,7 +523,7 @@ pub(crate) fn format_pretty( "run.finished" => { let status = view.body()?.get("status").and_then(Value::as_str)?; let duration = format_duration_ms(elapsed.unwrap_or(0)); - Some(match status { + let verdict = match status { "success" => format!( "{ts} {} {}", styles.bold_green.apply_to("\u{2713} SUCCEEDED"), @@ -539,6 +539,17 @@ pub(crate) fn format_pretty( styles.bold_red.apply_to("\u{2717} FAILED"), styles.bold.apply_to(&duration), ), + }; + let failure = view + .body()? + .get("finalization_failure") + .and_then(|failure| failure.get("message")) + .and_then(Value::as_str); + Some(match failure { + Some(message) => { + format!("{verdict}\n{ts} {}", styles.bold_red.apply_to(message)) + } + None => verdict, }) } _ => None, @@ -1276,6 +1287,32 @@ mod tests { assert!(line.contains("\u{2717} FAILED"), "{line}"); } + #[test] + fn execution_success_does_not_finish_required_publication() { + let executed = petri( + 8, + json!({ "origin": "external", "context": {}, + "record": { "seq": 5, "body": { "event": "invocation.finished", + "invocation": 0, "result": { "status": "success" } } } + }), + ); + assert_eq!(exit_code_of(&executed), None); + let finished = petri( + 9, + json!({ "origin": "external", "context": {}, + "record": { "seq": 6, "body": { "event": "run.finished", "status": "failed", + "finalization_failure": { "code": "publish_failed", "message": "the push was rejected" } } } + }), + ); + assert_eq!(exit_code_of(&finished), Some(1)); + let mut state = PrettyState::default(); + let line = format_pretty(&finished, &Styles::new(false), &mut state).unwrap(); + insta::assert_snapshot!(line, @" + 04:43:08 ✗ FAILED 0ms + 04:43:08 the push was rejected + "); + } + #[test] fn the_raw_line_is_the_envelope_as_json() { let notice = platform( 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 040f33ef2..267c18070 100644 --- a/lib/apps/fabro-cli/src/commands/run/petri_worker.rs +++ b/lib/apps/fabro-cli/src/commands/run/petri_worker.rs @@ -295,6 +295,9 @@ pub(super) async fn execute(worker: PetriWorker<'_>) -> Result<()> { elapsed_ms = u64::try_from(started.elapsed().as_millis()).unwrap_or(u64::MAX), "Petri run ended" ); + // The engine reads Petri's committed overall outcome, including a + // structured publish_failed rejection. The lifecycle acknowledges it; + // invocation success alone cannot decide the worker's exit. let (record, phase, failure) = match engine::conclusion(&result) { Conclusion::Interrupted { message } => { // The run is not over: it continues from its records in the diff --git a/lib/apps/fabro-cli/src/commands/run/publish.rs b/lib/apps/fabro-cli/src/commands/run/publish.rs index deba41843..9832bb55b 100644 --- a/lib/apps/fabro-cli/src/commands/run/publish.rs +++ b/lib/apps/fabro-cli/src/commands/run/publish.rs @@ -13,9 +13,9 @@ //! long run's pushes into repeated failures. Re-minting near expiry keeps a //! run of any length on a live token. //! -//! Publication runs in Fabro's `run_finished` hook, after the last stage and -//! before the run's terminal record, as the legacy publish step did: the -//! final checkpoint is pushed from inside the sandbox to the run +//! Publication runs in Fabro's required `finalize_run` hook, after the last +//! stage and before the run's terminal record: +//! the final checkpoint is pushed from inside the sandbox to the run //! branch on GitHub, and, when the run changed files and its settings ask //! for one, a pull request is opened and recorded. A failure fails the run //! with `publish_failed`. diff --git a/lib/apps/fabro-cli/tests/it/cmd/attach.rs b/lib/apps/fabro-cli/tests/it/cmd/attach.rs index ff4b4bb42..a30df85d8 100644 --- a/lib/apps/fabro-cli/tests/it/cmd/attach.rs +++ b/lib/apps/fabro-cli/tests/it/cmd/attach.rs @@ -1103,12 +1103,13 @@ fn attach_json_errors_without_prompting_for_human_input() { "recorded_at": "[EPOCH_MS]", "body": { "event": "run.started", - "format_version": 8, + "format_version": 9, "key": "[ULID]", "root": 0, "middleware_chain": [ "circuit-breaker" - ] + ], + "required_finalization": true } } } diff --git a/lib/apps/fabro-cli/tests/it/cmd/validate.rs b/lib/apps/fabro-cli/tests/it/cmd/validate.rs index 3a3896401..800cadf8e 100644 --- a/lib/apps/fabro-cli/tests/it/cmd/validate.rs +++ b/lib/apps/fabro-cli/tests/it/cmd/validate.rs @@ -263,6 +263,12 @@ fn bare_fabro_with_unbound_inputs_in_template_partial_validates_structurally_wit // The include error names the partial relative to the run's working // directory, so the `../` run is as long as that directory is deep. let mut filters = context.filters(); + // A checkout outside the user home can be rendered through a relative + // path (including macOS's /tmp alias), rather than its canonical root. + filters.push(( + r#"(?:\.\./)+[^"\s]*?/test/(templated_unbound_partial/)"#.to_string(), + "[UP][FIXTURES]/$1".to_string(), + )); filters.push(( r"(\.\./)*\.\.\[FIXTURES\]".to_string(), "[UP][FIXTURES]".to_string(), @@ -333,6 +339,10 @@ fn validate_reports_missing_template_dependency() { let mut cmd = context.validate(); cmd.arg(fixture("templates/missing_dependency/workflow.fabro")); let mut filters = context.filters(); + filters.push(( + r"(?:\.\./)+[^`\s]*?/test/(templates/missing_dependency/)".to_string(), + "[FIXTURES]/$1".to_string(), + )); filters.push(( r"(?:\.\./)*\.\.\[FIXTURES\]/".to_string(), "[FIXTURES]/".to_string(), diff --git a/lib/apps/fabro-server/src/petri_runs.rs b/lib/apps/fabro-server/src/petri_runs.rs index 28380e8f2..9fd2840a3 100644 --- a/lib/apps/fabro-server/src/petri_runs.rs +++ b/lib/apps/fabro-server/src/petri_runs.rs @@ -574,13 +574,25 @@ mod tests { transition: RunLifecycleKind, status: RunStatus, ) { + let response = append_lifecycle_as_worker(app, run_id, token, transition, status).await; + assert_eq!(response.status(), StatusCode::OK); + } + + /// Append one lifecycle record as the run's worker, and return the + /// server's response whatever it is. + async fn append_lifecycle_as_worker( + app: &axum::Router, + run_id: RunId, + token: &str, + transition: RunLifecycleKind, + status: RunStatus, + ) -> axum::response::Response { let record = PlatformRecord::RunLifecycle(RunLifecycleRecord::new(transition).with_status(status)); let body = json!({ "record": serde_json::to_value(&record).expect("the record encodes"), }); - let response = app - .clone() + app.clone() .oneshot( Request::builder() .method("POST") @@ -591,8 +603,7 @@ mod tests { .expect("the append request builds"), ) .await - .expect("the append request completes"); - assert_eq!(response.status(), StatusCode::OK); + .expect("the append request completes") } /// Petri's own finish of the run, stored the way its worker stores it: @@ -940,4 +951,308 @@ mod tests { "the failure names the worker's exit" ); } + + async fn append_engine_records( + app: &axum::Router, + run_id: RunId, + token: &str, + owner: &str, + log: &fabro_petri::petri::LogId, + records: &[fabro_petri::petri::Record], + ) -> axum::response::Response { + let body = json!({ "owner": owner, "records": records.iter().map(|record| json!({ + "seq": record.seq, "recorded_at": record.recorded_at, "record": record.record, + })).collect::>() }); + app.clone() + .oneshot( + Request::builder() + .method("POST") + .uri( + format!("/api/v1/runs/{run_id}/petri/logs/{log}/records") + .replace(' ', "%20"), + ) + .header(header::AUTHORIZATION, format!("Bearer {token}")) + .header(header::CONTENT_TYPE, "application/json") + .body(Body::from(body.to_string())) + .unwrap(), + ) + .await + .unwrap() + } + + async fn assert_public_result( + state: &AppState, + app: &axum::Router, + run_id: RunId, + expected: RunStatus, + message: Option<&str>, + ) { + state.petri_projector.settle(run_id).await; + let mut values = Vec::new(); + for suffix in ["", "/state"] { + let response = app + .clone() + .oneshot( + Request::builder() + .uri(format!("/api/v1/runs/{run_id}{suffix}")) + .body(Body::empty()) + .unwrap(), + ) + .await + .unwrap(); + values.push( + fabro_test::expect_axum_json( + response, + StatusCode::OK, + "GET run and state at required finalization boundary", + ) + .await, + ); + } + assert_eq!( + values[0]["lifecycle"]["status"], + serde_json::to_value(expected).unwrap(), + "public state: {:#?}", + values[1] + ); + assert_eq!(values[1]["status"], serde_json::to_value(expected).unwrap()); + assert_eq!(state.test_managed_run_status(&run_id), Some(expected)); + if expected.is_terminal() { + assert_eq!( + values[1]["conclusion"]["failure"]["detail"]["message"].as_str(), + message + ); + assert!( + values[1]["conclusion"]["stages"] + .as_array() + .is_some_and(|stages| !stages.is_empty()) + ); + } else { + assert!(values[1]["conclusion"].is_null()); + } + } + + #[tokio::test(flavor = "multi_thread", worker_threads = 2)] + async fn a_rejected_platform_failure_does_not_settle_the_run() { + let runtime = Arc::new(HeldWorkerRuntime::default()); + let (state, app, run_id, token) = held_worker_run(&runtime).await; + run_to_running_as_worker(&app, run_id, &token).await; + let pool = state.stores.run_summaries.pool(); + sqlx::query( + "CREATE TRIGGER reject_platform_record BEFORE INSERT ON platform_records \ + BEGIN SELECT RAISE(FAIL, 'scripted append failure'); END", + ) + .execute(&pool) + .await + .unwrap(); + let response = append_lifecycle_as_worker( + &app, + run_id, + &token, + RunLifecycleKind::Failed, + RunStatus::Failed { + reason: FailureReason::LaunchFailed, + }, + ) + .await; + fabro_test::assert_axum_status( + response, + StatusCode::INTERNAL_SERVER_ERROR, + "rejected platform terminal append", + ) + .await; + assert_public_result(&state, &app, run_id, RunStatus::Running, None).await; + crate::server::persist_run_failure( + &state, + run_id, + FailureReason::LaunchFailed, + "worker launch failed".into(), + ) + .await; + assert_public_result(&state, &app, run_id, RunStatus::Running, None).await; + sqlx::query("DROP TRIGGER reject_platform_record") + .execute(&pool) + .await + .unwrap(); + runtime.end_worker(); + } + + #[tokio::test(flavor = "multi_thread", worker_threads = 2)] + async fn a_host_failure_is_retried_until_storage_accepts_it() { + let runtime = Arc::new(HeldWorkerRuntime::default()); + let (state, app, run_id, token) = held_worker_run(&runtime).await; + run_to_running_as_worker(&app, run_id, &token).await; + let pool = state.stores.run_summaries.pool(); + sqlx::query( + "CREATE TRIGGER reject_platform_record BEFORE INSERT ON platform_records \ + BEGIN SELECT RAISE(FAIL, 'scripted append failure'); END", + ) + .execute(&pool) + .await + .unwrap(); + let failure = tokio::spawn({ + let state = Arc::clone(&state); + async move { + crate::server::persist_run_failure( + &state, + run_id, + FailureReason::LaunchFailed, + "worker launch failed".into(), + ) + .await; + } + }); + tokio::time::sleep(std::time::Duration::from_millis(100)).await; + sqlx::query("DROP TRIGGER reject_platform_record") + .execute(&pool) + .await + .unwrap(); + failure.await.unwrap(); + let failed = RunStatus::Failed { + reason: FailureReason::LaunchFailed, + }; + let stored = crate::server::run_records::projection(&state, run_id) + .await + .unwrap() + .expect("the run is stored"); + assert_eq!(stored.status, failed, "the retried failure is durable"); + assert_eq!(state.test_managed_run_status(&run_id), Some(failed)); + runtime.end_worker(); + } + + #[tokio::test(flavor = "multi_thread", worker_threads = 2)] + async fn required_publication_result_agrees_across_worker_api_projection_and_cleanup() { + use fabro_petri::petri::LogId; + use fabro_petri::test_support::finalization; + for rejection in [None, Some("the push was rejected")] { + let runtime = Arc::new(HeldWorkerRuntime::default()); + let (state, app, run_id, token) = held_worker_run(&runtime).await; + run_to_running_as_worker(&app, run_id, &token).await; + let fixture = finalization::test_run_records(run_id, rejection).await; + for bytes in fixture.blobs { + let response = app + .clone() + .oneshot( + Request::builder() + .method("POST") + .uri(format!("/api/v1/runs/{run_id}/petri/blobs?owner=worker-1")) + .header(header::AUTHORIZATION, format!("Bearer {token}")) + .body(Body::from(bytes)) + .unwrap(), + ) + .await + .unwrap(); + fabro_test::assert_axum_status(response, StatusCode::OK, "worker graph blob").await; + } + let mut terminal = None; + for (log, mut records) in fixture.logs { + if log == LogId::Coordinator { + let finish = records.pop().unwrap(); + assert_eq!(finish.record["body"]["event"], "run.finished"); + terminal = Some(finish); + } + let response = + append_engine_records(&app, run_id, &token, "worker-1", &log, &records).await; + fabro_test::assert_axum_status( + response, + StatusCode::NO_CONTENT, + "worker execution records", + ) + .await; + } + assert_public_result(&state, &app, run_id, RunStatus::Running, None).await; + let terminal = terminal.unwrap(); + // A rejected writer cannot settle the managed run or store its + // result. This exercises the same endpoint as the valid finish. + let response = append_engine_records( + &app, + run_id, + &token, + "superseded-worker", + &LogId::Coordinator, + std::slice::from_ref(&terminal), + ) + .await; + fabro_test::assert_axum_status( + response, + StatusCode::CONFLICT, + "finish from a stale lease owner", + ) + .await; + assert_public_result(&state, &app, run_id, RunStatus::Running, None).await; + let response = + append_engine_records(&app, run_id, &token, "worker-1", &LogId::Coordinator, &[ + terminal, + ]) + .await; + fabro_test::assert_axum_status( + response, + StatusCode::NO_CONTENT, + "authoritative worker finish", + ) + .await; + let expected = match rejection { + Some(_) => RunStatus::Failed { + reason: FailureReason::PublishFailed, + }, + None => RunStatus::Succeeded { + reason: SuccessReason::Completed, + }, + }; + assert_public_result(&state, &app, run_id, expected, rejection).await; + // Worker exit after the finish, without a platform terminal + // acknowledgement, must preserve Petri's committed result. + runtime.end_worker(); + for _ in 0..500 { + if concluded_runs(&app).await == 1 { + break; + } + time::sleep(Duration::from_millis(10)).await; + } + assert_eq!(concluded_runs(&app).await, 1); + assert_public_result(&state, &app, run_id, expected, rejection).await; + // A host error arriving after cleanup must neither append a + // competing terminal record nor replace the useful failure. + let platform_count = fabro_petri::test_support::stored_platform_records( + &state.stores.run_summaries.pool(), + run_id, + ) + .await + .unwrap() + .len(); + crate::server::persist_run_failure( + &state, + run_id, + FailureReason::Terminated, + "worker wait failed during teardown".into(), + ) + .await; + assert_public_result(&state, &app, run_id, expected, rejection).await; + let after_count = fabro_petri::test_support::stored_platform_records( + &state.stores.run_summaries.pool(), + run_id, + ) + .await + .unwrap() + .len(); + assert_eq!(after_count, platform_count); + let (rebuilt, _, _) = fabro_petri::test_support::rebuild( + &state.db_pool, + &state.stores.run_summaries.pool(), + run_id, + ) + .await + .unwrap(); + let rebuilt = rebuilt.unwrap(); + assert_eq!(rebuilt.status, expected); + assert_eq!( + rebuilt + .conclusion + .unwrap() + .failure + .map(|failure| failure.detail.message), + rejection.map(str::to_string) + ); + } + } } diff --git a/lib/apps/fabro-server/src/server.rs b/lib/apps/fabro-server/src/server.rs index 2a2ae3079..c0069200d 100644 --- a/lib/apps/fabro-server/src/server.rs +++ b/lib/apps/fabro-server/src/server.rs @@ -282,6 +282,23 @@ struct ManagedRun { } impl ManagedRun { + /// Settle the run on a terminal `status`. The first terminal status + /// sticks: a later one that agrees only fills a missing error, and one + /// that disagrees is ignored. Returns whether the status was applied. + fn settle(&mut self, status: RunStatus, error: Option) -> bool { + if self.status.is_terminal() { + if self.status == status && self.error.is_none() { + self.error = error; + } + return false; + } + self.status = status; + self.error = error; + self.active_steerable_stages.clear(); + self.active_non_steerable_stages.clear(); + true + } + /// True if cancellation should still escalate to `worker_ref`; clears a /// stale escalation marker as a side effect. fn escalation_still_current(&mut self, worker_ref: &WorkerRef) -> bool { @@ -3505,24 +3522,90 @@ async fn reject_run_if_sandbox_provider_disabled( return false; }; tracing::warn!(run_id = %run_id, error = %error, "Sandbox provider disabled by server policy"); - fail_run_before_execution(state, run_id, FailureReason::LaunchFailed, error).await; + persist_run_failure(state, run_id, FailureReason::LaunchFailed, error).await; true } -async fn fail_run_before_execution( +/// Record a host failure only while no terminal result is committed. This +/// also handles a worker wait/launch error racing its durable Petri finish. +/// The managed run settles on whichever terminal result the store committed. +/// A failed commit is retried, so a brief storage fault does not leave the +/// run active. When nothing could be committed its live state is still +/// released, but its status is left alone: the API never reports an outcome +/// storage lacks, and the restart reconciliation fails the run. +pub(crate) async fn persist_run_failure( state: &Arc, run_id: RunId, reason: FailureReason, message: String, ) { - if let Err(err) = - run_records::lifecycle(state, run_id, run_records::failed(reason, message.clone())).await - { - error!(run_id = %run_id, error = %err, "Failed to persist run failure status"); + let mut retry_delays = HOST_FAILURE_RETRY_DELAYS.iter(); + loop { + match commit_host_failure(state, run_id, reason, message.clone()).await { + Ok(committed) => { + let failure = committed + .conclusion + .as_ref() + .and_then(|conclusion| conclusion.failure.as_ref()) + .map(|failure| failure.detail.message.clone()); + settle_managed_run_at_finish(state, run_id, committed.status, failure); + break; + } + Err(err) => { + let Some(delay) = retry_delays.next() else { + error!(run_id = %run_id, error = %err, "Failed to record a host failure"); + break; + }; + warn!( + run_id = %run_id, + error = %err, + retry_in_ms = delay.as_millis(), + "Failed to record a host failure; retrying" + ); + sleep(*delay).await; + } + } } + release_managed_run(state, run_id); +} - fail_managed_run(state, run_id, reason, message); - state.scheduler_notify.notify_one(); +/// How long [`persist_run_failure`] waits before each retry of a failed +/// commit. +const HOST_FAILURE_RETRY_DELAYS: [Duration; 3] = [ + Duration::from_millis(250), + Duration::from_secs(1), + Duration::from_secs(4), +]; + +/// Append the host failure unless the run already ended, and return the +/// terminal projection the store committed: the failure, or the finish that +/// won the race. +async fn commit_host_failure( + state: &AppState, + run_id: RunId, + reason: FailureReason, + message: String, +) -> anyhow::Result> { + let committed = run_records::projection(state, run_id) + .await? + .context("the run is missing")?; + if committed.status.is_terminal() { + return Ok(committed); + } + // The append settles the projector, so the read after it folds the + // failure, or the finish that won the race. + run_records::lifecycle(state, run_id, run_records::failed(reason, message)).await?; + let committed = state + .stores + .run_summaries + .load_petri_projection(&run_id) + .await? + .context("the run is missing")?; + anyhow::ensure!( + committed.status.is_terminal(), + "the stored host failure has no terminal projection" + ); + Ok(committed) } fn managed_run( @@ -3582,11 +3665,22 @@ async fn durable_run_status(state: &AppState, run_id: RunId) -> anyhow::Result, run_id: RunId, reason: FailureReason, message: String) { let mut runs = state.runs.lock().expect("runs lock poisoned"); if let Some(managed_run) = runs.get_mut(&run_id) { - managed_run.status = RunStatus::Failed { reason }; - managed_run.error = Some(message); + managed_run.settle(RunStatus::Failed { reason }, Some(message)); + } + drop(runs); + release_managed_run(state, run_id); +} + +/// Drop the run's live worker state and controls and free its scheduler +/// slot, leaving its status alone. +pub(in crate::server) fn release_managed_run(state: &AppState, run_id: RunId) { + let mut runs = state.runs.lock().expect("runs lock poisoned"); + if let Some(managed_run) = runs.get_mut(&run_id) { clear_live_run_state(managed_run); } - cleanup_worker_control_bus_for_run(state.as_ref(), run_id); + drop(runs); + cleanup_worker_control_bus_for_run(state, run_id); + state.scheduler_notify.notify_one(); } /// Fold one lifecycle record of the run's stream into the in-memory run: @@ -3601,6 +3695,7 @@ fn apply_lifecycle_to_managed_run(state: &AppState, run_id: RunId, record: &RunL // A settled run is immutable to the lifecycle: the follower still folds // the records before the terminal one after Petri's finish or the // worker's terminal record settled the run, and none may reopen it. + // Terminal records go through [`ManagedRun::settle`]. if managed_run.status.is_terminal() && !is_terminal_transition(record) { return; } @@ -3648,21 +3743,17 @@ fn apply_lifecycle_to_managed_run(state: &AppState, run_id: RunId, record: &RunL } RunLifecycleKind::Removing => managed_run.status = RunStatus::Removing, RunLifecycleKind::Succeeded => { - managed_run.status = record.status.unwrap_or(RunStatus::Succeeded { + let status = record.status.unwrap_or(RunStatus::Succeeded { reason: SuccessReason::Completed, }); - managed_run.error = None; - managed_run.active_steerable_stages.clear(); - managed_run.active_non_steerable_stages.clear(); + managed_run.settle(status, None); cleanup_worker_control_bus_for_run(state, run_id); } RunLifecycleKind::Failed | RunLifecycleKind::Dead => { - managed_run.status = record.status.unwrap_or(RunStatus::Failed { + let status = record.status.unwrap_or(RunStatus::Failed { reason: FailureReason::WorkflowError, }); - managed_run.error.clone_from(&record.reason); - managed_run.active_steerable_stages.clear(); - managed_run.active_non_steerable_stages.clear(); + managed_run.settle(status, record.reason.clone()); cleanup_worker_control_bus_for_run(state, run_id); } RunLifecycleKind::Runnable @@ -3683,34 +3774,23 @@ fn is_terminal_transition(record: &RunLifecycleRecord) -> bool { ) } -/// Settle the in-memory run at Petri's own finish, as its worker stores -/// the `run.finished` record: the view reports the run ended from the -/// moment that record is stored, so the managed run the delete precheck -/// prefers must not still say running while the worker tears down; a -/// delete in that window was refused as active. A run already settled -/// keeps its status. The worker's terminal lifecycle record, a moment -/// later, refines the status and its error and ends the worker's controls -/// ([`settle_managed_run_at_terminal_record`]); the worker's exit later -/// reaps the process and leaves the settled status alone. +/// Settle the in-memory run after Petri's authoritative finish is durable. +/// Required publication has completed before this record. Worker teardown +/// and subsequent lifecycle records cannot change its terminal outcome. pub(in crate::server) fn settle_managed_run_at_finish( state: &AppState, run_id: RunId, status: RunStatus, + failure: Option, ) { let mut runs = state.runs.lock().expect("runs lock poisoned"); - let Some(managed_run) = runs.get_mut(&run_id) else { - return; - }; - if managed_run.status.is_terminal() { - return; + if let Some(managed_run) = runs.get_mut(&run_id) { + managed_run.settle(status, failure); } - managed_run.status = status; - managed_run.active_steerable_stages.clear(); - managed_run.active_non_steerable_stages.clear(); } /// Settle the in-memory run at the terminal lifecycle record its worker -/// stores, ahead of the store: the same as [`settle_managed_run_at_finish`] +/// stores, after the store: the same as [`settle_managed_run_at_finish`] /// for a worker that ended the run without Petri's finish (it failed before /// the engine ran), and the record's status, error and control cleanup for /// one that did. A record that is not terminal is left to the stream @@ -3775,7 +3855,6 @@ async fn fail_worker_launch(state: &Arc, run_id: RunId, err: anyhow::E None } }; - let launch_message = format!("Failed to spawn worker: {err}"); let (error, reason) = failure_honoring_pending_cancel(pending_control, || { ( WorkflowError::engine_with_anyhow("Failed to spawn worker", err), @@ -3785,16 +3864,9 @@ async fn fail_worker_launch(state: &Arc, run_id: RunId, err: anyhow::E let message = if reason == FailureReason::Cancelled { "Run cancelled before worker launch completed".to_string() } else { - launch_message + collect_chain(&error).join(": ") }; - let _ = run_records::lifecycle( - state, - run_id, - run_records::failed(reason, error.to_string()), - ) - .await; - fail_managed_run(state, run_id, reason, message); - state.scheduler_notify.notify_one(); + persist_run_failure(state, run_id, reason, message).await; } /// A worker that exited without recording the run's end left it failed, @@ -4163,7 +4235,6 @@ async fn execute_run_subprocess(state: Arc, run_id: RunId) { FailureReason::WorkflowError, "Run not found at launch".to_string(), ); - state.scheduler_notify.notify_one(); return; } Err(err) => { @@ -4174,7 +4245,6 @@ async fn execute_run_subprocess(state: Arc, run_id: RunId) { FailureReason::WorkflowError, format!("Failed to load run state: {err}"), ); - state.scheduler_notify.notify_one(); return; } }; @@ -4195,7 +4265,7 @@ async fn execute_run_subprocess(state: Arc, run_id: RunId) { let github_app_private_key = match state.vault_secret(EnvVars::GITHUB_APP_PRIVATE_KEY).await { Ok(value) => value, Err(err) => { - fail_run_before_execution( + persist_run_failure( &state, run_id, FailureReason::WorkflowError, @@ -4255,15 +4325,17 @@ async fn execute_run_subprocess(state: Arc, run_id: RunId) { Err(err) => { tracing::error!(run_id = %run_id, error = %err, "Failed while waiting on worker"); let message = format!("Worker wait failed: {err}"); + let superseded = { + let runs = state.runs.lock().expect("runs lock poisoned"); + runs.get(&run_id) + .is_some_and(|run| run.worker_ref.as_ref() != Some(&worker_ref)) + }; + if superseded { + return; + } state.worker_runtime.force_stop(&worker_ref).await; - let _ = run_records::lifecycle( - &state, - run_id, - run_records::failed(FailureReason::Terminated, message.clone()), - ) - .await; - fail_managed_run(&state, run_id, FailureReason::Terminated, message); - state.scheduler_notify.notify_one(); + state.petri_runs.worker_exited(run_id); + persist_run_failure(&state, run_id, FailureReason::Terminated, message).await; return; } }; @@ -4321,7 +4393,6 @@ async fn execute_run_subprocess(state: Arc, run_id: RunId) { FailureReason::WorkflowError, "The run's final state is missing from the store".to_string(), ); - state.scheduler_notify.notify_one(); return; } Err(err) => { @@ -4332,7 +4403,6 @@ async fn execute_run_subprocess(state: Arc, run_id: RunId) { FailureReason::WorkflowError, format!("Failed to load final run state: {err}"), ); - state.scheduler_notify.notify_one(); return; } }; diff --git a/lib/apps/fabro-server/src/server/handler/petri.rs b/lib/apps/fabro-server/src/server/handler/petri.rs index 403a97cd1..57ee546f3 100644 --- a/lib/apps/fabro-server/src/server/handler/petri.rs +++ b/lib/apps/fabro-server/src/server/handler/petri.rs @@ -23,7 +23,7 @@ use fabro_api::types::{ PetriReleaseRequest, WriteBlobResponse, }; use fabro_petri::petri::{Access, Digest, LogId, OwnerId, Record, StoreError}; -use fabro_petri::projection::finished_run_status; +use fabro_petri::projection; use fabro_petri::run_store::{log_id_text, parse_log_id}; use fabro_store::{PlatformRecord, PlatformRecordKind, StagePosition, StoredPlatformRecord}; use fabro_types::BlobHash; @@ -157,16 +157,15 @@ async fn append_records( Ok(writer) => writer, Err(err) => return store_error_response(id, &err), }; - // The view ends the run at Petri's own finish, the moment the record - // is stored and a pass folds it: the managed run settles first, so a - // delete that lands while the worker still tears down is not refused. - if log == LogId::Coordinator { - if let Some(status) = records.iter().find_map(finished_run_status) { - settle_managed_run_at_finish(&state, id, status); - } - } match writer.append(&log, &records).await { Ok(()) => { + if log == LogId::Coordinator { + if let Some((status, failure)) = + records.iter().find_map(projection::finished_run_result) + { + settle_managed_run_at_finish(&state, id, status, failure); + } + } // The records are durable; the projection trails them from here. state.petri_projector.signal(id); StatusCode::NO_CONTENT.into_response() @@ -279,10 +278,6 @@ async fn append_platform_record( (Some(execution), Some(firing)) => Some(StagePosition { execution, firing }), _ => None, }; - // A terminal lifecycle record ends the run in the view as Petri's - // finish does, for a worker that ended the run without one: the - // managed run settles before the record is stored. - settle_managed_run_at_terminal_record(&state, id, &record); let summaries = &state.stores.run_summaries; match summaries .platform_records() @@ -290,6 +285,7 @@ async fn append_platform_record( .await { Ok(stored) => { + settle_managed_run_at_terminal_record(&state, id, &record); summaries.notify_platform_record(id); match wire_platform_record(&stored) { Ok(record) => Json(record).into_response(), diff --git a/lib/apps/fabro-server/src/server/petri_runs.rs b/lib/apps/fabro-server/src/server/petri_runs.rs index 356054d94..7a09f4743 100644 --- a/lib/apps/fabro-server/src/server/petri_runs.rs +++ b/lib/apps/fabro-server/src/server/petri_runs.rs @@ -22,7 +22,7 @@ //! of the server's vault, and its blobs go to the server's blob store. Its //! managed run settles at Petri's own finish, as a worker's does at the //! worker's records endpoint: the run store it executes over settles the -//! run before the `run.finished` record is stored ([`SettlingStore`]). No +//! run after the `run.finished` record is stored ([`SettlingStore`]). No //! stage or agent event is projected either way, which is the read-side //! item that follows. //! @@ -496,8 +496,8 @@ pub(crate) async fn execute(state: Arc, run_id: RunId) { run_id: run_id.to_string(), run_dir: run_dir.join("petri"), execution, - // The projector's signal follows each durable append; the managed - // run's settle at Petri's finish precedes it. + // The coordinator finish is stored before managed status settles; + // the projector also reads only durable records. store: Arc::new(SettlingStore { inner: state .petri_projector @@ -549,13 +549,7 @@ pub(crate) async fn execute(state: Arc, run_id: RunId) { failed(reason, message) } }; - if let Err(err) = run_records::lifecycle(&state, run_id, record).await { - error!(run_id = %run_id, error = %err, "Failed to persist run outcome"); - } - // The managed run settled at Petri's finish, ahead of the store; the - // terminal record refines its status and error and ends its live - // state, and is the settle of a run that ended without a finish. - finish(&state, run_id, status, error); + commit_and_finish(&state, run_id, record, status, error).await; // The view trails the terminal record; the aggregate reads the settled // projection, as the worker path reads the final state at worker exit. state.petri_projector.settle(run_id).await; @@ -814,38 +808,43 @@ fn failed( async fn fail_before_execution(state: &Arc, run_id: RunId, message: &str) { error!(run_id = %run_id, error = message, "Petri run cannot start"); let (status, error, record) = failed(FailureReason::WorkflowError, message.to_string()); - if let Err(err) = run_records::lifecycle(state, run_id, record).await { - error!(run_id = %run_id, error = %err, "Failed to persist run failure status"); + commit_and_finish(state, run_id, record, status, error).await; +} + +/// Append the run's terminal record, then finish the run. An append that +/// fails only releases the live state: it commits no terminal result. +async fn commit_and_finish( + state: &Arc, + run_id: RunId, + record: RunLifecycleRecord, + status: RunStatus, + error: Option, +) { + match run_records::lifecycle(state, run_id, record).await { + Ok(_) => finish(state, run_id, status, error), + Err(err) => { + error!(run_id = %run_id, error = %err, "Failed to persist run outcome"); + super::release_managed_run(state, run_id); + } } - finish(state, run_id, status, error); } /// Settle the managed run at its terminal record and release its /// scheduler slot. A run that Petri finished settled already, at the -/// `run.finished` record ([`SettlingStore`]); this refines its status and -/// error and ends its live state. A run deleted since is gone from the map -/// and stays gone. +/// `run.finished` record ([`SettlingStore`]); this preserves that status, +/// fills any missing failure detail and ends its live state. A run deleted +/// since is gone from the map and stays gone. fn finish(state: &Arc, run_id: RunId, status: RunStatus, error: Option) { let mut runs = state.runs.lock().expect("runs lock poisoned"); if let Some(managed_run) = runs.get_mut(&run_id) { - managed_run.status = status; - managed_run.error = error; - clear_live_run_state(managed_run); + managed_run.settle(status, error); } drop(runs); - state.scheduler_notify.notify_one(); + super::release_managed_run(state, run_id); } -/// The run store an in-process run executes over: the projector's -/// signalling store, whose coordinator appends settle the managed run at -/// Petri's own finish first. The view ends the run at the `run.finished` -/// record the moment it is stored and a pass folds it, so the managed run -/// the delete precheck prefers must not still say running while the engine -/// tears down: a delete in that window was refused as active. The worker's -/// records endpoint does the same for a worker-backed run, ahead of the -/// same store. The settle is in memory only; the terminal lifecycle record -/// [`execute`] stores once the engine returns refines the status -/// ([`finish`]), and stays the settle of a run that ends without a finish. +/// The run store an in-process run executes over. A coordinator finish +/// settles managed status only after the append succeeds, as on HTTP workers. struct SettlingStore { inner: Arc, state: Arc, @@ -864,7 +863,7 @@ impl RunStore for SettlingStore { } /// One run's logs, whose coordinator appends settle the managed run at -/// Petri's finish before the records reach the store. +/// Petri's finish after the records reach the store. struct SettlingLogs { inner: Arc, run_id: Option, @@ -878,12 +877,15 @@ impl RunLogs for SettlingLogs { } async fn append(&self, log: &LogId, records: &[Record]) -> Result<(), StoreError> { + self.inner.append(log, records).await?; if let Some(run_id) = self.run_id.filter(|_| *log == LogId::Coordinator) { - if let Some(status) = records.iter().find_map(projection::finished_run_status) { - super::settle_managed_run_at_finish(&self.state, run_id, status); + if let Some((status, failure)) = + records.iter().find_map(projection::finished_run_result) + { + super::settle_managed_run_at_finish(&self.state, run_id, status, failure); } } - self.inner.append(log, records).await + Ok(()) } async fn read(&self, log: &LogId) -> Result, StoreError> { @@ -922,6 +924,7 @@ mod tests { state: Arc, run_id: RunId, seen: Mutex>>, + reject: std::sync::atomic::AtomicBool, } #[async_trait::async_trait] @@ -933,6 +936,9 @@ mod tests { async fn append(&self, log: &LogId, records: &[Record]) -> Result<(), StoreError> { let status = self.state.test_managed_run_status(&self.run_id); self.seen.lock().expect("seen lock poisoned").push(status); + if self.reject.load(std::sync::atomic::Ordering::SeqCst) { + return Err(StoreError::StaleOwner); + } self.inner.append(log, records).await } @@ -992,6 +998,7 @@ mod tests { state: Arc::clone(&state), run_id, seen: Mutex::new(Vec::new()), + reject: std::sync::atomic::AtomicBool::new(false), }); let logs = SettlingLogs { inner: Arc::clone(&recorder) as Arc, @@ -1001,12 +1008,10 @@ mod tests { (state, run_id, logs, recorder) } - /// The in-process run settles at Petri's own finish, before the - /// `run.finished` record reaches the store: the view cannot report the - /// run ended while the managed run still says running. The records - /// before the finish leave the run in flight. + /// The in-process run settles only after its authoritative finish is + /// stored. Earlier records leave the run in flight. #[tokio::test] - async fn an_in_process_run_settles_before_its_finish_is_stored() { + async fn an_in_process_run_settles_after_its_finish_is_stored() { let (state, run_id, logs, recorder) = in_flight_run().await; logs.append(&LogId::Coordinator, &[coordinator_record( @@ -1032,12 +1037,33 @@ mod tests { }; assert_eq!( *recorder.seen.lock().expect("seen lock poisoned"), - vec![Some(RunStatus::Running), Some(succeeded)], - "the managed run settled before the finish reached the store" + vec![Some(RunStatus::Running), Some(RunStatus::Running)], + "the managed run stayed running until the finish reached the store" ); assert_eq!(state.test_managed_run_status(&run_id), Some(succeeded)); } + #[tokio::test] + async fn a_rejected_in_process_finish_does_not_settle_the_run() { + let (state, run_id, logs, recorder) = in_flight_run().await; + recorder + .reject + .store(true, std::sync::atomic::Ordering::SeqCst); + let result = logs + .append(&LogId::Coordinator, &[coordinator_record( + 0, + &json!({"event": "run.finished", "status": "failed", + "finalization_failure": {"code": "publish_failed", "message": "push rejected"}}), + )]) + .await; + assert!(matches!(result, Err(StoreError::StaleOwner))); + assert_eq!( + state.test_managed_run_status(&run_id), + Some(RunStatus::Running) + ); + assert!(logs.read(&LogId::Coordinator).await.unwrap().is_empty()); + } + /// A finish on another log than the coordinator's is not Petri's /// finish of the run: an execution's engine log ends an execution. #[tokio::test] diff --git a/lib/components/fabro-petri/Cargo.toml b/lib/components/fabro-petri/Cargo.toml index 4034e8a0f..628b5503a 100644 --- a/lib/components/fabro-petri/Cargo.toml +++ b/lib/components/fabro-petri/Cargo.toml @@ -16,7 +16,7 @@ workspace = true # Petri's test kit, re-exported for Fabro crates that check a store # implementation against Petri's contract from their own tests. Never on in # a normal build. -test-support = ["dep:petri_testkit"] +test-support = ["dep:petri_testkit", "dep:tempfile"] [dependencies] fabro-github = { path = "../fabro-github" } @@ -47,6 +47,7 @@ petri_frontend_attractor.workspace = true petri_frontend_fabro.workspace = true lithos-llm = { workspace = true, features = ["runtime"] } petri_testkit = { workspace = true, optional = true } +tempfile = { version = "3", optional = true } anyhow.workspace = true bytes.workspace = true async-trait.workspace = true diff --git a/lib/components/fabro-petri/src/engine.rs b/lib/components/fabro-petri/src/engine.rs index f9f2380ed..f6ed2d194 100644 --- a/lib/components/fabro-petri/src/engine.rs +++ b/lib/components/fabro-petri/src/engine.rs @@ -84,6 +84,7 @@ use crate::admission::AdmittedGraphs; use crate::blobs::{Blobs, RunBlobs}; use crate::controls::RunControls; use crate::hooks::{FabroHooks, HooksSpec}; +use crate::projection; use crate::runtime::RuntimeSpec; use crate::secrets::SharedSecrets; @@ -150,7 +151,7 @@ pub enum RunStatus { #[derive(Clone, Debug)] pub struct RunOutcome { pub status: RunStatus, - /// The root invocation's failure message, when it failed. + /// Required-finalization failure detail, or the root execution's failure. pub failure: Option, /// Whether the record is whole: the run recorded its finish and every /// log replays byte for byte. @@ -353,27 +354,7 @@ pub async fn run(request: RunRequest) -> Result { return Err(RunError::StoreFailed(message)); } let inspection = inspect(request.store.as_ref(), &key).await?; - let mut outcome = outcome(inspection, result.err())?; - // A failed checkpoint cancelled the run; what Fabro reports is the - // checkpoint failure, not a cancellation. - if let Some(failure) = fabro_hooks - .as_ref() - .and_then(|hooks| hooks.checkpoint_failure()) - { - outcome.status = RunStatus::Failed; - outcome.failure = Some(failure); - } - // A successful run whose publication failed is a failed run: its work - // did not reach where the settings sent it. - if let Some(failure) = fabro_hooks - .as_ref() - .and_then(|hooks| hooks.publish_failure()) - { - outcome.status = RunStatus::Failed; - outcome.failure = Some(failure); - outcome.publish_failed = true; - } - Ok(outcome) + outcome(inspection, result.err()) } /// When Petri keeps a run's workspaces after their scope is released. @@ -419,8 +400,9 @@ pub async fn outcome_of(store: &dyn RunStore, run_id: &str) -> Result) -> Conclusion { match result { @@ -543,31 +525,47 @@ fn outcome( inspection: RunInspection, host_error: Option, ) -> Result { - let status = match inspection.status.as_deref() { - Some("success") => RunStatus::Success, - Some("failed") => RunStatus::Failed, - Some("cancelled") => RunStatus::Cancelled, - _ => { - let mut reasons = inspection.incomplete.clone(); - if let Some(error) = host_error { - reasons.push(error.to_string()); - } - return Err(RunError::Unfinished(reasons)); + let Some(recorded) = inspection + .status + .as_deref() + .filter(|status| matches!(*status, "success" | "failed" | "cancelled")) + else { + let mut reasons = inspection.incomplete.clone(); + if let Some(error) = host_error { + reasons.push(error.to_string()); } + return Err(RunError::Unfinished(reasons)); }; - let failure = inspection + // The projection's reading of the finish, so the worker's return and + // the API agree: a failed checkpoint's cancellation is a failure. + let (status, publish_failed) = + match projection::finished_status(recorded, inspection.finalization_failure.as_ref()) { + fabro_types::RunStatus::Succeeded { .. } => (RunStatus::Success, false), + fabro_types::RunStatus::Failed { + reason: FailureReason::Cancelled, + } => (RunStatus::Cancelled, false), + fabro_types::RunStatus::Failed { reason } => { + (RunStatus::Failed, reason == FailureReason::PublishFailed) + } + _ => (RunStatus::Failed, false), + }; + let execution_failure = inspection .invocations .iter() .find(|invocation| invocation.invocation == inspection.root.invocation) .and_then(|root| root.result.as_ref()) .and_then(|result| result.failure.as_ref()) .map(|failure| failure.message.clone()); + let failure = inspection + .finalization_failure + .map(|failure| failure.message) + .or(execution_failure); Ok(RunOutcome { status, failure, complete: inspection.complete, incomplete: inspection.incomplete, - publish_failed: false, + publish_failed, }) } diff --git a/lib/components/fabro-petri/src/fork.rs b/lib/components/fabro-petri/src/fork.rs index d4145faa6..bf79af0eb 100644 --- a/lib/components/fabro-petri/src/fork.rs +++ b/lib/components/fabro-petri/src/fork.rs @@ -39,11 +39,12 @@ use fabro_workflow::operations::{StageLabel, StageLabels}; use petri_execution::host::{self, ForkOptions, ForkOrigin, ForkPosition, HostError}; use petri_execution::inspect::{self, InspectError}; use petri_execution::{ - Access, CoordinatorEvent, ExecutionId, InvocationId, RunKey, RunStore, + Access, CoordinatorEvent, ExecutionId, InvocationId, RunKey, RunLogs, RunStore, StoreError as CoordinatorStoreError, }; use petri_runtime::RunOptions; -use petri_runtime::ir::FiringId; +use petri_runtime::driver::lifecycle::{ExecutionHooks, HookContext, RunFinished}; +use petri_runtime::ir::{FinalizationFailure, FiringId}; use petri_store::StoreError; use tracing::{debug, info}; @@ -131,7 +132,12 @@ pub async fn check( .open(&RunKey::new(source.to_string()), Access::Read) .await .map_err(ForkError::Open)?; - let state = host::stored_state(&*logs).await.map_err(ForkError::Seed)?; + check_logs(&*logs, &position).await +} + +/// [`check`] over the source's opened logs. +async fn check_logs(logs: &dyn RunLogs, position: &ForkPosition) -> Result<(), ForkError> { + let state = host::stored_state(logs).await.map_err(ForkError::Seed)?; let Some(execution) = state.executions.get(&position.execution) else { return Err(ForkError::Refused(format!( "the source run has no execution {}", @@ -146,7 +152,7 @@ pub async fn check( position.execution ))); } - let inspection = inspect::inspect_run(&*logs) + let inspection = inspect::inspect_run(logs) .await .map_err(ForkError::Inspect)?; if inspection @@ -169,10 +175,33 @@ pub async fn check( Ok(()) } +/// Declaration-only hooks used while copying a fork's records: every Fabro +/// run requires finalization ([`crate::hooks::FabroHooks`]), so the fork's +/// records declare it too, and its worker's hooks match them on resume. +/// Never execute with these hooks. +struct ForkFinalizationRequirement; + +#[async_trait::async_trait] +impl ExecutionHooks for ForkFinalizationRequirement { + fn requires_run_finalization(&self) -> bool { + true + } + + async fn finalize_run( + &self, + _context: &HookContext, + _finished: RunFinished, + ) -> Result<(), FinalizationFailure> { + Err(FinalizationFailure::new( + "finalization_unavailable", + "the fork must install its worker publication hooks before execution", + )) + } +} + /// Seed the fork: Petri's records, the kept checkpoints and the run branch. The /// new run must not exist in the store yet. pub async fn fork(request: ForkRequest) -> Result { - check(request.store.as_ref(), request.source, request.position).await?; let source_key = RunKey::new(request.source.to_string()); let fork_key = RunKey::new(request.fork.to_string()); let source_logs = request @@ -181,12 +210,14 @@ pub async fn fork(request: ForkRequest) -> Result { .await .map_err(ForkError::Open)?; + check_logs(&*source_logs, &request.position).await?; let mut options = RunOptions::new(&request.fork_run_dir); options.run_key = Some(fork_key.clone()); // A fork only copies records and acquires no sandbox, so it needs no // provider configuration. let runtime = providers::standard_runtime(&SandboxProviderConfig::default()) .options(options) + .hooks(Arc::new(ForkFinalizationRequirement)) .store(Arc::clone(&request.store)); let forked = host::fork_from(&runtime, &*source_logs, request.position, ForkOptions { rerun_last: request.rerun_last, diff --git a/lib/components/fabro-petri/src/hooks.rs b/lib/components/fabro-petri/src/hooks.rs index 40f72d260..207e8975e 100644 --- a/lib/components/fabro-petri/src/hooks.rs +++ b/lib/components/fabro-petri/src/hooks.rs @@ -31,12 +31,14 @@ //! `artifact.collected` record, unless the same file with the same content //! was already collected earlier in the run. A failed write is a recorded //! problem on the transition, never a blocked route. -//! - `run_finished`: the run's diff, its run branch against its base commit, as -//! the `run.diff` platform record with the patch as a blob; for a successful -//! run, its publication ([`RunPublisher`]: the platform pushes the run branch -//! and opens a pull request), whose failure fails the run before its terminal -//! record; then the forwarded point, so the local service runs `run_complete` -//! and `run_failed` with the sandbox in place. +//! - `finalize_run`, required for every run: the run's diff, its run branch +//! against its base commit, as the `run.diff` platform record with the patch +//! as a blob; a failed checkpoint, which fails the run and skips publication; +//! for a successful run, its publication ([`RunPublisher`]: the platform +//! pushes the run branch and opens a pull request), whose failure fails the +//! run before its terminal record. +//! - `run_finished`: the forwarded point, so the local service runs +//! `run_complete` and `run_failed` with the sandbox in place. //! - `scope_acquired`: a fresh run's Git target checked out into the workspace //! from inside the scope ([`crate::source`]); a resumed run uses its //! surviving workspace, while an explicit fork fetches the source run's @@ -51,16 +53,18 @@ //! //! # Operation identities //! -//! Every external effect here is keyed on `(run key, execution, DecisionId, -//! effect kind)` from the hook context and deduplicated on retry: the -//! checkpoint's key is the attempt's decision in its execution, effect +//! Checkpoint and artifact effects are keyed on `(run key, execution, +//! DecisionId, effect kind)` from the hook context and deduplicated on retry: +//! the checkpoint's key is the attempt's decision in its execution, effect //! `checkpoint`; an artifact's is the same decision, effect `artifact`, with //! the file's path and content digest as the identity within it. A //! re-dispatched attempt whose commit already landed reuses it when the //! workspace still sits on it unchanged (see [`RunWorkspaces::commit`]); a //! reissued routing decision finds the record, or the commit by its //! trailers, and writes nothing twice; a file already collected under the -//! same path and digest is not collected again. +//! same path and digest is not collected again. Publication keeps the +//! publisher's reconciliation policy; required finalization adds no independent +//! effect ledger or guarantee of deduplication across every external crash. //! //! # Where the workspace is //! @@ -100,7 +104,9 @@ use petri_runtime::driver::lifecycle::{ ScopeReleased, Transition, TransitionError, TransitionReport, }; use petri_runtime::executor::{EnvError, ExecEnv}; -use petri_runtime::ir::{ExecutionId, FailureInfo, RunStatus, ScopeId, Status}; +use petri_runtime::ir::{ + ExecutionId, FailureInfo, FinalizationFailure, RunStatus, ScopeId, Status, +}; use serde_json::json; use tokio::sync::{Mutex as AsyncMutex, OnceCell}; use tokio::{fs, time}; @@ -114,6 +120,7 @@ use crate::checkpoint::{ }; use crate::fork::{self, ForkError}; use crate::platform_records::{PlatformRecordError, PlatformRecords}; +use crate::projection; use crate::recovery::{self, Plan, RecoveryError, RestoreTarget}; use crate::source::RunSource; use crate::workspace::{self, WorkspaceLookup, WorkspaceLookupError}; @@ -453,8 +460,6 @@ pub struct FabroHooks { /// The checkpoint failure that ended the run, when one did. failure: Mutex>, publisher: Option>, - /// Why the run's publication failed, when it did. - publish_failure: Mutex>, /// Whether the run continues from its records: a sandbox workspace is /// then brought to its snapshot when its scope is first acquired. resumed: bool, @@ -522,7 +527,6 @@ impl FabroHooks { scopes: ScopeEnvs::default(), failure: Mutex::default(), publisher: spec.publisher, - publish_failure: Mutex::default(), resumed, restore: OnceCell::new(), store, @@ -537,20 +541,13 @@ impl FabroHooks { } } - /// The checkpoint failure that ended the run, when one did: what the - /// engine reports the run failed with. + /// The checkpoint failure that ended the run, when one did: required + /// finalization commits it as the run's failure. #[must_use] pub fn checkpoint_failure(&self) -> Option { sync::lock(&self.failure).clone() } - /// Why the run's publication failed, when it did: the run then fails - /// with this message. - #[must_use] - pub fn publish_failure(&self) -> Option { - sync::lock(&self.publish_failure).clone() - } - /// The run's workspaces on this host, as the hooks reach them. #[must_use] pub fn workspaces(&self) -> &RunWorkspaces { @@ -1403,23 +1400,6 @@ impl FabroHooks { Ok(publication) } - /// Hand a successful run's work to the publisher, before the run's - /// terminal record. A failure fails the run with its reason. - async fn publish(&self, publisher: &dyn RunPublisher, publication: &Publication) { - match publisher.publish(publication).await { - Ok(()) => info!( - run_id = %self.run_id, - branch = publication.run_branch, - sha = publication.head_sha, - "run published" - ), - Err(message) => { - warn!(run_id = %self.run_id, error = %message, "the run's publication failed"); - *sync::lock(&self.publish_failure) = Some(message); - } - } - } - /// Hold at a test gate when one is set for this point and node. async fn gate(&self, point: &str, node: &str) { let Some(dir) = &self.test_gates else { @@ -1626,35 +1606,69 @@ impl ExecutionHooks for FabroHooks { Ok(report) } - async fn run_finished(&self, context: &HookContext, finished: RunFinished) -> Vec { - info!( - run_id = %self.run_id, - status = ?finished.status, - failure = finished.failure.as_deref().unwrap_or(""), - "Petri run finished; recording the run's diff and running the run-end hooks" - ); - let publication = match self.record_run_diff().await { + /// Every Fabro run requires finalization, whether or not it publishes + /// or checkpoints: one path for every run, and a declaration that a + /// resume or a fork always matches. + fn requires_run_finalization(&self) -> bool { + true + } + + async fn finalize_run( + &self, + context: &HookContext, + finished: RunFinished, + ) -> Result<(), FinalizationFailure> { + let diff = self.record_run_diff().await.map_err(|error| { + let message = error.render(); + warn!(run_id = %self.run_id, error = %message, "the run's diff was not recorded"); + message + }); + // A failed checkpoint fails the run whatever its execution status, + // and its work is never published. + if let Some(message) = self.checkpoint_failure() { + return Err(projection::checkpoint_failure(message)); + } + let publication = match diff { Ok(publication) => publication, - Err(error) => { - let message = error.render(); - warn!(run_id = %self.run_id, error = %message, "the run's diff was not recorded"); + Err(message) => { if self.publisher.is_some() && finished.status == RunStatus::Success { - *sync::lock(&self.publish_failure) = Some(message); + return Err(projection::publish_failure(message)); } None } }; if let Some(publisher) = &self.publisher { - if finished.status == RunStatus::Success && self.checkpoint_failure().is_none() { - if let Some(publication) = &publication { - self.publish(publisher.as_ref(), publication).await; - } else if self.publish_failure().is_none() { - *sync::lock(&self.publish_failure) = Some( - "the run has no recorded branch and checkpoint to publish".to_string(), + if finished.status == RunStatus::Success { + let publication = publication.ok_or_else(|| { + projection::publish_failure( + "the run has no recorded branch and checkpoint to publish", + ) + })?; + publisher.publish(&publication).await.map_err(|message| { + warn!( + run_id = %self.run_id, + error = %message, + "the run's publication failed" ); - } + projection::publish_failure(message) + })?; + info!( + run_id = %self.run_id, + branch = publication.run_branch, + sha = publication.head_sha, + "run published" + ); } } + if self.inner.requires_run_finalization() { + self.inner.finalize_run(context, finished).await?; + } + Ok(()) + } + + async fn run_finished(&self, context: &HookContext, finished: RunFinished) -> Vec { + // The diff and publication already ran in finalize_run, before Petri + // committed the outcome. self.inner.run_finished(context, finished).await } diff --git a/lib/components/fabro-petri/src/projection/coordinator.rs b/lib/components/fabro-petri/src/projection/coordinator.rs index f78ebc439..7225dd007 100644 --- a/lib/components/fabro-petri/src/projection/coordinator.rs +++ b/lib/components/fabro-petri/src/projection/coordinator.rs @@ -9,6 +9,7 @@ use fabro_types::{ }; use petri_execution::CoordinatorEvent; use petri_execution::events::RunEvent; +use petri_runtime::ir::FinalizationFailure; use super::{FiringKey, InvocationRef, RunView, apply_status, settle_control}; @@ -96,9 +97,16 @@ impl RunView { settle_control(projection, RunControlAction::Unpause); } } - CoordinatorEvent::RunFinished { status } => { + CoordinatorEvent::RunFinished { + status, + finalization_failure, + } => { self.state.finished = Some(status.to_string()); - self.conclude(status.to_string().as_str(), at); + self.conclude( + status.to_string().as_str(), + finalization_failure.as_ref(), + at, + ); } // ── Sandbox: the retention outcome (VIEWS.md "Sandbox") ───────── // The instance stays on `Run.sandbox`: it names what ran, and @@ -126,19 +134,28 @@ impl RunView { } } - /// The run's conclusion, from its recorded finish and what the stages /// The run's conclusion, from its recorded finish and what the stages /// summed to. - fn conclude(&mut self, status: &str, at: DateTime) { + fn conclude( + &mut self, + status: &str, + finalization_failure: Option<&FinalizationFailure>, + at: DateTime, + ) { let Some(projection) = self.projection.as_mut() else { return; }; + if projection.status.is_terminal() { + return; + } let root = self .state .root .and_then(|root| self.state.invocations.get(&root)); - let failure_message = root.and_then(|root| root.failure.clone()); - let run_status = finished_status(status); + let failure_message = finalization_failure + .map(|failure| failure.message.clone()) + .or_else(|| root.and_then(|root| root.failure.clone())); + let run_status = finished_status(status, finalization_failure); let (outcome, failure) = match run_status { RunStatus::Failed { reason: reason @ FailureReason::Cancelled, @@ -209,17 +226,28 @@ impl RunView { } /// The status Fabro gives a run at Petri's finish, by the status the finish -/// records (`success`, `cancelled`, or a failure). -pub(super) fn finished_status(status: &str) -> RunStatus { +/// records (`success`, `cancelled`, or a failure). A failed checkpoint +/// cancels the run, but the run failed: its finish says so with the +/// checkpoint's finalization failure. +pub(crate) fn finished_status( + status: &str, + finalization_failure: Option<&FinalizationFailure>, +) -> RunStatus { match status { "success" => RunStatus::Succeeded { reason: SuccessReason::Completed, }, - "cancelled" => RunStatus::Failed { - reason: FailureReason::Cancelled, - }, + "cancelled" if !finalization_failure.is_some_and(super::is_checkpoint_failure) => { + RunStatus::Failed { + reason: FailureReason::Cancelled, + } + } _ => RunStatus::Failed { - reason: FailureReason::WorkflowError, + reason: if finalization_failure.is_some_and(super::is_publish_failure) { + FailureReason::PublishFailed + } else { + FailureReason::WorkflowError + }, }, } } diff --git a/lib/components/fabro-petri/src/projection/mod.rs b/lib/components/fabro-petri/src/projection/mod.rs index 4410df7b9..d65b65341 100644 --- a/lib/components/fabro-petri/src/projection/mod.rs +++ b/lib/components/fabro-petri/src/projection/mod.rs @@ -37,18 +37,23 @@ use std::fmt; use std::str::FromStr; use chrono::{DateTime, TimeZone as _, Utc}; +pub(crate) use coordinator::finished_status; use fabro_store::StagePosition; use fabro_store::platform_records::StoredPlatformRecord; use fabro_types::{ - RunControlAction, RunDiff, RunId, RunProjection, RunStatus, StageId, StageProjection, + FailureReason, RunControlAction, RunDiff, RunId, RunProjection, RunStatus, StageId, + StageProjection, }; use petri_execution::events::{NodeRef, RunEvent, Subject}; use petri_execution::{CoordinatorEvent, CoordinatorRecord, ExecutionId}; +use petri_runtime::ir::FinalizationFailure; use petri_store::Record; use serde::{Deserialize, Deserializer, Serialize, Serializer, de}; use serde_json::Value; use tracing::debug; +use crate::checkpoint::CHECKPOINT_FAILED_CLASS; + /// One item the projector hands the fold, with its delivery sequence. pub enum Item<'a> { Petri(&'a RunEvent), @@ -369,17 +374,46 @@ pub fn run_id_of(key: &str) -> Option { key.parse().ok() } -/// The status the view gives the run at Petri's own finish, when the -/// stored record is the coordinator log's `run.finished`: the view reports -/// the run ended from the moment that record is stored, ahead of Fabro's -/// terminal lifecycle record. `None` for any other record. +/// The required-finalization failure Fabro's hooks record when the run's +/// publication fails; its code is [`FailureReason::PublishFailed`]'s. #[must_use] -pub fn finished_run_status(record: &Record) -> Option { +pub fn publish_failure(message: impl Into) -> FinalizationFailure { + FinalizationFailure::new(<&'static str>::from(FailureReason::PublishFailed), message) +} + +/// Whether a required-finalization failure is the run's failed publication. +#[must_use] +pub fn is_publish_failure(failure: &FinalizationFailure) -> bool { + failure.code == <&'static str>::from(FailureReason::PublishFailed) +} + +/// The required-finalization failure Fabro's hooks record when a checkpoint +/// failed during the run. The failed checkpoint cancelled the run, so Petri +/// may record it as cancelled; Fabro reports it as a workflow failure. +#[must_use] +pub fn checkpoint_failure(message: impl Into) -> FinalizationFailure { + FinalizationFailure::new(CHECKPOINT_FAILED_CLASS, message) +} + +/// Whether a required-finalization failure is a failed checkpoint. +#[must_use] +pub fn is_checkpoint_failure(failure: &FinalizationFailure) -> bool { + failure.code == CHECKPOINT_FAILED_CLASS +} + +/// The committed overall status and required-finalization failure message. +/// Execution failure details remain in the invocation records and projection. +#[must_use] +pub fn finished_run_result(record: &Record) -> Option<(RunStatus, Option)> { let record: CoordinatorRecord = serde_json::from_value(record.record.clone()).ok()?; match record.body { - CoordinatorEvent::RunFinished { status } => { - Some(coordinator::finished_status(&status.to_string())) - } + CoordinatorEvent::RunFinished { + status, + finalization_failure, + } => Some(( + coordinator::finished_status(&status.to_string(), finalization_failure.as_ref()), + finalization_failure.map(|failure| failure.message), + )), _ => None, } } @@ -412,10 +446,11 @@ mod tests { #[test] fn a_finish_record_names_the_status_the_view_ends_the_run_on() { let finished = |status: &str| { - finished_run_status(&coordinator_record(&serde_json::json!({ + finished_run_result(&coordinator_record(&serde_json::json!({ "event": "run.finished", "status": status, }))) + .map(|(status, _)| status) }; assert_eq!( finished("success"), @@ -436,13 +471,30 @@ mod tests { }) ); assert_eq!( - finished_run_status(&coordinator_record(&serde_json::json!({ + finished_run_result(&coordinator_record(&serde_json::json!({ + "event": "run.finished", + "status": "cancelled", + "finalization_failure": { + "code": "checkpoint_failed", + "message": "checkpoint commit of `wreck` failed", + }, + }))), + Some(( + RunStatus::Failed { + reason: fabro_types::FailureReason::WorkflowError, + }, + Some("checkpoint commit of `wreck` failed".to_string()), + )), + "a failed checkpoint's cancellation is the checkpoint's failure" + ); + assert_eq!( + finished_run_result(&coordinator_record(&serde_json::json!({ "event": "run.paused", }))), None ); assert_eq!( - finished_run_status(&Record { + finished_run_result(&Record { seq: 3, recorded_at: 1_000, record: serde_json::json!({"event": "run.finished", "status": "success"}), diff --git a/lib/components/fabro-petri/src/projection/platform.rs b/lib/components/fabro-petri/src/projection/platform.rs index ca8c2c8a9..eb0519d8a 100644 --- a/lib/components/fabro-petri/src/projection/platform.rs +++ b/lib/components/fabro-petri/src/projection/platform.rs @@ -286,6 +286,9 @@ fn fold_lifecycle(projection: &mut RunProjection, record: &RunLifecycleRecord, a settle_control(projection, RunControlAction::Pause); } Kind::Succeeded | Kind::Failed => { + if projection.status.is_terminal() { + return; + } if let Some(status) = record.status { apply_status(projection, status, at); } diff --git a/lib/components/fabro-petri/src/test_support.rs b/lib/components/fabro-petri/src/test_support.rs index 91793f60e..ff5b452f9 100644 --- a/lib/components/fabro-petri/src/test_support.rs +++ b/lib/components/fabro-petri/src/test_support.rs @@ -5,6 +5,8 @@ //! view with. Compiled only with the `test-support` feature, which a //! dev-dependency turns on. +pub mod finalization; + use std::collections::{BTreeSet, HashMap}; use std::sync::Mutex; use std::time::Duration; diff --git a/lib/components/fabro-petri/src/test_support/finalization.rs b/lib/components/fabro-petri/src/test_support/finalization.rs new file mode 100644 index 000000000..8775f59ab --- /dev/null +++ b/lib/components/fabro-petri/src/test_support/finalization.rs @@ -0,0 +1,151 @@ +//! Real command-run records for transport and projection tests. No providers +//! or model credentials are used. + +use std::sync::Arc; + +use petri_execution::host::{self, HostRun}; +use petri_execution::inspect; +use petri_runtime::driver::lifecycle::{ExecutionHooks, HookContext, RunFinished}; +use petri_runtime::frontend::CompileInputs; +use petri_runtime::ir::FinalizationFailure; +use petri_runtime::{RunOptions, Runtime}; +use petri_store::{Access, LogId, MemoryRunStore, Record, RunKey, RunStore as _}; +use tokio::fs; +use tokio::sync::{Notify, Semaphore}; + +use crate::projection; +use crate::providers::{self, SandboxProviderConfig}; + +/// A deterministic finalizer gate for testing the committed boundary. +pub struct TestFinalizer { + pub entered: Notify, + pub release: Semaphore, + rejection: Option, +} + +impl TestFinalizer { + pub fn new(rejection: Option<&str>) -> Self { + Self { + entered: Notify::new(), + release: Semaphore::new(0), + rejection: rejection.map(str::to_string), + } + } +} + +#[async_trait::async_trait] +impl ExecutionHooks for TestFinalizer { + fn requires_run_finalization(&self) -> bool { + true + } + + async fn finalize_run( + &self, + _: &HookContext, + _: RunFinished, + ) -> Result<(), FinalizationFailure> { + self.entered.notify_one(); + self.release + .acquire() + .await + .expect("the gate stays open") + .forget(); + self.rejection.as_ref().map_or(Ok(()), |message| { + Err(projection::publish_failure(message.clone())) + }) + } +} + +/// The command-only runtime installed with a test finalizer. +pub fn test_runtime(finalizer: Arc) -> Runtime { + petri_attractor_steps::register(providers::standard_runtime( + &SandboxProviderConfig::default(), + )) + .frontend(petri_frontend_fabro::Fabro::new()) + .hooks(finalizer) +} + +/// Logs and blobs of a real command run, for authenticated worker transport +/// tests. The caller can send the coordinator prefix before its final record. +pub struct TestRunRecords { + pub logs: Vec<(LogId, Vec)>, + pub blobs: Vec>, +} + +/// Run the command fixture to completion and return its logs and graph blobs, +/// with the finalizer rejecting the run when `rejection` names a failure. +pub async fn test_run_records( + run_id: fabro_types::RunId, + rejection: Option<&str>, +) -> TestRunRecords { + let root = tempfile::tempdir().expect("the fixture has an isolated directory"); + let workflow = root.path().join("workflow.fabro"); + fs::write( + &workflow, + r#"digraph Finalization { + graph [goal="Check required publication"] + start [shape=Mdiamond] + work [shape=parallelogram, script="echo completed"] + exit [shape=Msquare] + start -> work -> exit + }"#, + ) + .await + .expect("the fixture workflow writes"); + fs::write( + root.path().join("workflow.toml"), + "_version = 1\n[workflow]\ngraph = \"workflow.fabro\"\n", + ) + .await + .expect("the settings write"); + let finalizer = Arc::new(TestFinalizer::new(rejection)); + finalizer.release.add_permits(1); + let store = Arc::new(MemoryRunStore::new()); + let key = RunKey::new(run_id.to_string()); + let mut options = RunOptions::new(root.path().join("run")); + options.run_key = Some(key.clone()); + options.echo = false; + let runtime = test_runtime(finalizer) + .store(store.clone()) + .options(options); + let checked = runtime + .check(&workflow, None, None, &CompileInputs::new()) + .expect("the fixture compiles"); + host::run_configured( + &runtime, + HostRun::new(checked.graph.expect("the fixture is valid")).with_children(checked.children), + |_, _| {}, + ) + .await + .expect("the fixture completes"); + let logs = store + .open(&key, Access::Read) + .await + .expect("the fixture reads"); + let inspection = inspect::inspect_run(&*logs) + .await + .expect("the fixture inspects"); + let mut records = vec![( + LogId::Coordinator, + logs.read(&LogId::Coordinator) + .await + .expect("coordinator reads"), + )]; + for execution in inspection.executions { + let id = LogId::Execution(execution.execution); + records.push((id, logs.read(&id).await.expect("execution reads"))); + } + let mut blobs = Vec::new(); + for graph in inspection.graphs { + blobs.push( + logs.get_blob(graph) + .await + .expect("graph reads") + .expect("graph exists"), + ); + } + TestRunRecords { + logs: records, + blobs, + } +} diff --git a/lib/components/fabro-petri/tests/hooks.rs b/lib/components/fabro-petri/tests/hooks.rs index d548f3a79..50c758f6d 100644 --- a/lib/components/fabro-petri/tests/hooks.rs +++ b/lib/components/fabro-petri/tests/hooks.rs @@ -44,8 +44,9 @@ use fabro_types::{GitIdentitySource, RunId, SandboxProviderKind}; use object_store::local::LocalFileSystem; use petri_execution::inspect::{self, RunInspection}; use petri_store::{Access, MemoryRunStore, RunKey, RunStore as _}; -use tokio::fs; use tokio::process::Command; +use tokio::sync::{Notify, Semaphore}; +use tokio::{fs, time}; use tokio_util::sync::CancellationToken; mod support; @@ -688,8 +689,9 @@ async fn a_failed_stage_is_committed_and_its_route_sees_the_files() { } /// A checkpoint commit that fails is fatal: the stage's outcome is recorded -/// as `checkpoint_failed`, no route is taken, the run ends failed with the -/// checkpoint's error, and a restart reports it failed without resuming. +/// as `checkpoint_failed`, no route is taken, the run's committed finish +/// fails it with the checkpoint's error, and a restart reports it failed +/// without resuming. #[tokio::test] async fn a_failed_checkpoint_ends_the_run_with_no_route() { let harness = Harness::new().await; @@ -712,6 +714,16 @@ async fn a_failed_checkpoint_ends_the_run_with_no_route() { ); let inspection = harness.inspection().await; + let committed = inspection + .finalization_failure + .clone() + .expect("the finish commits the checkpoint failure"); + assert_eq!(committed.code, CHECKPOINT_FAILED_CLASS); + let stored = engine::outcome_of(&*harness.store, &harness.run_id.to_string()) + .await + .expect("the finished run reads back"); + assert_eq!(stored.status, RunStatus::Failed, "{stored:?}"); + assert_eq!(stored.failure, outcome.failure); let attempts: Vec<_> = inspection .executions .iter() @@ -1261,16 +1273,24 @@ async fn published_run( let (origin, _) = upstream(&harness.run_dir.with_file_name("upstream"), 2).await; harness.source = Some(file_source(&origin, "main", Some(1))); harness.publisher = Some(Arc::clone(publisher) as Arc); - let workflow = workflow( - &format!(" edit [shape=parallelogram, {attributes}]"), - " start -> edit -> exit", - ); let outcome = harness - .run_on(SandboxProviderKind::LOCAL, &workflow, SETTINGS) + .run_on( + SandboxProviderKind::LOCAL, + &published_workflow(attributes), + SETTINGS, + ) .await; (harness, outcome) } +/// The workflow [`published_run`] runs: one command stage with `attributes`. +fn published_workflow(attributes: &str) -> String { + workflow( + &format!(" edit [shape=parallelogram, {attributes}]"), + " start -> edit -> exit", + ) +} + /// A successful run hands its publisher the run branch, the commit it ends /// on (held by the snapshot repository) and its patch, before the run ends. #[tokio::test] @@ -1307,10 +1327,37 @@ async fn a_successful_run_is_published_with_its_branch_head_and_patch() { #[tokio::test] async fn a_failed_publication_fails_the_run() { let publisher = RecordingPublisher::new(Some("the push was rejected")); - let (_, outcome) = published_run("script=\"echo edited >> README.md\"", &publisher).await; + let attributes = "script=\"echo edited >> README.md\""; + let (harness, outcome) = published_run(attributes, &publisher).await; assert_eq!(outcome.status, RunStatus::Failed, "{outcome:?}"); assert!(outcome.publish_failed); assert_eq!(outcome.failure.as_deref(), Some("the push was rejected")); + let inspection = harness.inspection().await; + assert_eq!(inspection.status.as_deref(), Some("failed")); + let failure = inspection.finalization_failure.expect("durable failure"); + assert_eq!(failure.code, "publish_failed"); + assert_eq!(failure.message, "the push was rejected"); + let stored = engine::outcome_of(&*harness.store, &harness.run_id.to_string()) + .await + .unwrap(); + assert_eq!(stored.status, outcome.status); + assert_eq!(stored.failure, outcome.failure); + assert!(stored.publish_failed); + let resumed = harness + .execute_on( + SandboxProviderKind::LOCAL, + &published_workflow(attributes), + SETTINGS, + true, + ) + .await; + assert_eq!(resumed.status, outcome.status); + assert_eq!(resumed.failure, outcome.failure); + assert_eq!( + publisher.published.lock().unwrap().len(), + 1, + "committed resume does not publish twice" + ); } /// A run that fails (here, at a goal gate) is not published. @@ -1323,6 +1370,27 @@ async fn a_failed_run_is_not_published() { assert!(publisher.published.lock().unwrap().is_empty()); } +/// A run whose checkpoint failed is never published, and its committed +/// failure is the checkpoint's, not a publication's. +#[tokio::test] +async fn a_run_whose_checkpoint_failed_is_not_published() { + let publisher = RecordingPublisher::new(None); + let (harness, outcome) = published_run( + "script=\"rm -rf .git && echo garbage > .git && echo wrecked > out.txt\"", + &publisher, + ) + .await; + assert_eq!(outcome.status, RunStatus::Failed, "{outcome:?}"); + assert!(!outcome.publish_failed); + assert!(publisher.published.lock().unwrap().is_empty()); + let failure = harness + .inspection() + .await + .finalization_failure + .expect("the finish commits the checkpoint failure"); + assert_eq!(failure.code, CHECKPOINT_FAILED_CLASS); +} + struct OriginPublisher { origin: String, pushed: std::sync::Mutex>, @@ -1437,6 +1505,10 @@ async fn assert_shallow_fork(checkpoint_index: usize) { }) .await .expect("the fork is seeded"); + assert!( + forked.inspection().await.required_finalization, + "a publishing fork inherits its required finalization declaration" + ); assert_eq!( seeded.start.expect("the selected checkpoint").sha, *checkpoint_sha @@ -1570,6 +1642,117 @@ async fn a_run_diff_failure_cannot_silently_skip_publication() { let outcome = harness.run(&graph, SETTINGS).await; assert_eq!(outcome.status, RunStatus::Failed, "{outcome:?}"); assert!(outcome.publish_failed); - assert!(outcome.failure.unwrap().contains("run diff unavailable")); + assert!( + outcome + .failure + .as_ref() + .unwrap() + .contains("run diff unavailable") + ); + let stored = engine::outcome_of(&*harness.store, &harness.run_id.to_string()) + .await + .unwrap(); + assert_eq!(stored.status, outcome.status); + assert_eq!(stored.failure, outcome.failure); + assert!(stored.publish_failed); assert!(publisher.published.lock().unwrap().is_empty()); } + +struct GatedPublisher { + inner: Arc, + entered: Notify, + release: Semaphore, +} + +#[async_trait::async_trait] +impl RunPublisher for GatedPublisher { + async fn push(&self, site: &Site, branch: &str, sha: &str) -> Result<(), String> { + self.inner.push(site, branch, sha).await + } + + async fn publish(&self, publication: &Publication) -> Result<(), String> { + self.entered.notify_one(); + self.release + .acquire() + .await + .expect("publication gate stays open") + .forget(); + self.inner.publish(publication).await + } +} + +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn required_publication_blocks_the_terminal_result_and_cleanup() { + for rejection in [None, Some("the push was rejected")] { + let publisher = Arc::new(GatedPublisher { + inner: RecordingPublisher::new(rejection), + entered: Notify::new(), + release: Semaphore::new(0), + }); + let mut harness = Harness::new().await; + harness.publisher = Some(publisher.clone()); + let harness = Arc::new(harness); + let executing = harness.clone(); + let task = tokio::spawn(async move { + executing + .run( + &workflow( + "edit [shape=parallelogram, script=\"echo edited >> README.md\"]", + "start -> edit -> exit", + ), + SETTINGS, + ) + .await + }); + time::timeout( + std::time::Duration::from_secs(15), + publisher.entered.notified(), + ) + .await + .expect("publication reaches the gate"); + let pending = harness.inspection().await; + assert!(pending.required_finalization); + assert!( + pending.status.is_none(), + "no terminal result while publication waits" + ); + assert!(!task.is_finished()); + let logs = harness + .store + .open(&RunKey::new(harness.run_id.to_string()), Access::Read) + .await + .unwrap(); + let coordinator = logs.read(&petri_store::LogId::Coordinator).await.unwrap(); + assert!( + coordinator + .iter() + .all(|record| record.record["body"]["event"] != "scope.released"), + "scope cleanup waits for publication" + ); + assert!( + harness.workspace_path(&harness.workspace().await).exists(), + "workspace is available to publication" + ); + assert!(matches!( + engine::outcome_of(&*harness.store, &harness.run_id.to_string()).await, + Err(engine::RunError::Unfinished(_)) + )); + publisher.release.add_permits(1); + let outcome = task.await.unwrap(); + assert_eq!( + outcome.status, + if rejection.is_some() { + RunStatus::Failed + } else { + RunStatus::Success + } + ); + assert_eq!(outcome.publish_failed, rejection.is_some()); + let stored = engine::outcome_of(&*harness.store, &harness.run_id.to_string()) + .await + .unwrap(); + assert_eq!(stored.status, outcome.status); + assert_eq!(stored.failure, outcome.failure); + assert_eq!(publisher.inner.published.lock().unwrap().len(), 1); + } +} diff --git a/lib/components/fabro-petri/tests/projection.rs b/lib/components/fabro-petri/tests/projection.rs index 8e9a9f9b6..71546ed76 100644 --- a/lib/components/fabro-petri/tests/projection.rs +++ b/lib/components/fabro-petri/tests/projection.rs @@ -27,7 +27,8 @@ use fabro_petri::providers::SandboxProviderConfig; use fabro_petri::runtime::RuntimeSpec; use fabro_petri::{SqliteRunStore, providers, test_support as petri_support}; use fabro_store::platform_records::{ - PlatformRecord, PlatformRecordStore, RunCreatedRecord, RunLifecycleKind, RunLifecycleRecord, + CheckpointRecord, PlatformRecord, PlatformRecordStore, RunCreatedRecord, RunDiffRecord, + RunLifecycleKind, RunLifecycleRecord, }; use fabro_store::{BlobStore, test_support}; use fabro_types::{ @@ -42,7 +43,7 @@ use petri_runtime::frontend::CompileInputs; use petri_runtime::ir::RunStatus as PetriRunStatus; use petri_store::{RunKey, RunStore}; use tokio::fs; -use tokio::time::sleep; +use tokio::time::{self, sleep}; const COMMAND_WORKFLOW: &str = r#"digraph Command { graph [goal="Run one command"] @@ -1361,3 +1362,184 @@ async fn an_auto_approved_answer_closes_the_question_in_the_projection() { ); assert_view_equals_rebuild(&gate.scenario.pool, gate.scenario.run_id).await; } + +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn required_finalization_projects_only_the_committed_overall_result() { + use fabro_petri::test_support::finalization::{self, TestFinalizer}; + use fabro_types::{FailureReason, StageOutcome, SuccessReason}; + + for rejection in [None, Some("the push was rejected")] { + let scenario = command_scenario().await; + let projector = Projector::new(scenario.pool.clone(), scenario.pool.clone()); + let store = projector.observe_store(Arc::new(SqliteRunStore::new(scenario.pool.clone()))); + let finalizer = Arc::new(TestFinalizer::new(rejection)); + let runtime = finalization::test_runtime(finalizer.clone()) + .store(store) + .options(run_options(&scenario.run_dir, scenario.run_id)); + let checked = runtime + .check(&scenario.workflow, None, None, &CompileInputs::new()) + .unwrap(); + let task = tokio::spawn(async move { + host::run_configured( + &runtime, + HostRun::new(checked.graph.unwrap()).with_children(checked.children), + |_, _| {}, + ) + .await + .unwrap() + }); + time::timeout(Duration::from_secs(15), finalizer.entered.notified()) + .await + .unwrap(); + // The driver can still flush its execution journal after entering + // finalization. Wait for the exit-stage evidence, rather than racing + // that writer while comparing the live view with a full rebuild. + let pending = time::timeout(Duration::from_secs(5), async { + loop { + projector.signal(scenario.run_id); + projector.settle(scenario.run_id).await; + let pending = petri_support::stored_projection(&scenario.pool, scenario.run_id) + .await + .unwrap() + .unwrap(); + if pending.iter_stages().any(|(id, stage)| { + id.node_id() == "exit" && stage.state == StageState::Succeeded + }) { + break pending; + } + time::sleep(Duration::from_millis(5)).await; + } + }) + .await + .expect("execution evidence flushes while publication is held"); + assert_eq!(pending.status, RunStatus::Running); + assert!(pending.conclusion.is_none()); + assert!(!task.is_finished()); + assert_view_equals_rebuild(&scenario.pool, scenario.run_id).await; + let summary = fabro_types::DiffSummary { + files_changed: 1, + additions: 2, + deletions: 0, + }; + let head_sha = "0123456789012345678901234567890123456789"; + let patch_blob = BlobHash::new(b"final patch"); + let platform = PlatformRecordStore::new(scenario.pool.clone()); + for record in [ + PlatformRecord::Checkpoint(CheckpointRecord { + execution: 0, + firing: 0, + attempt: Some(1), + workspace: None, + git_commit_sha: Some(head_sha.to_string()), + diff_summary: Some(summary), + patch_blob: Some(patch_blob), + operation: None, + }), + PlatformRecord::RunDiff(RunDiffRecord { + base_sha: None, + head_sha: Some(head_sha.to_string()), + diff_summary: Some(summary), + patch_blob: Some(patch_blob), + }), + ] { + platform + .append(&scenario.run_id, &record, None) + .await + .unwrap(); + } + finalizer.release.add_permits(1); + let report = task.await.unwrap(); + assert_eq!(report.state.folded_status(), PetriRunStatus::Success); + assert_eq!( + report.status, + if rejection.is_some() { + PetriRunStatus::Failed + } else { + PetriRunStatus::Success + } + ); + projector.settle(scenario.run_id).await; + let stored = petri_support::stored_projection(&scenario.pool, scenario.run_id) + .await + .unwrap() + .unwrap(); + let expected = match rejection { + Some(_) => RunStatus::Failed { + reason: FailureReason::PublishFailed, + }, + None => RunStatus::Succeeded { + reason: SuccessReason::Completed, + }, + }; + assert_eq!(stored.status, expected); + let conclusion = stored.conclusion.as_ref().unwrap(); + assert_eq!( + conclusion + .failure + .as_ref() + .map(|failure| failure.detail.message.as_str()), + rejection + ); + assert_eq!( + conclusion.status, + if rejection.is_some() { + StageOutcome::Failed { + retry_requested: false, + } + } else { + StageOutcome::Succeeded + } + ); + assert!( + stored + .iter_stages() + .all(|(_, stage)| stage.state == StageState::Succeeded) + ); + assert_eq!(conclusion.final_git_commit_sha.as_deref(), Some(head_sha)); + assert_eq!(conclusion.diff.summary, Some(summary)); + assert_eq!( + conclusion.diff.patch.as_deref(), + Some(fabro_types::format_blob_ref(&patch_blob).as_str()) + ); + assert!(!conclusion.stages.is_empty()); + assert_view_equals_rebuild(&scenario.pool, scenario.run_id).await; + + // A later worker lifecycle cannot replace the committed terminal + // result or conclusion, even if its status disagrees. + let other = match rejection { + Some(_) => RunStatus::Succeeded { + reason: SuccessReason::Completed, + }, + None => RunStatus::Failed { + reason: FailureReason::PublishFailed, + }, + }; + PlatformRecordStore::new(scenario.pool.clone()) + .append( + &scenario.run_id, + &PlatformRecord::RunLifecycle( + RunLifecycleRecord::new(if rejection.is_some() { + RunLifecycleKind::Succeeded + } else { + RunLifecycleKind::Failed + }) + .with_status(other), + ), + None, + ) + .await + .unwrap(); + projector.signal(scenario.run_id); + projector.settle(scenario.run_id).await; + let after = petri_support::stored_projection(&scenario.pool, scenario.run_id) + .await + .unwrap() + .unwrap(); + assert_eq!(after.status, stored.status); + assert_eq!( + serde_json::to_value(after.conclusion).unwrap(), + serde_json::to_value(stored.conclusion).unwrap() + ); + assert_view_equals_rebuild(&scenario.pool, scenario.run_id).await; + } +}