diff --git a/Cargo.toml b/Cargo.toml index cbad48387..3a958567f 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -108,13 +108,13 @@ pebble-cli-core = { git = "https://github.com/lithoscomputer/pebble", branch = " # petri: the workflow engine Fabro runs its workflows on. Only `fabro-petri` # and `fabro-dot` (the DOT parser alone) may depend on these packages; the # keys carry the `petri_` prefix so the crate names say where they come from. -petri_runtime = { git = "https://github.com/lithoscomputer/petri.git", branch = "main", package = "petri-runtime" } -petri_execution = { git = "https://github.com/lithoscomputer/petri.git", branch = "main", package = "petri-execution" } -petri_store = { git = "https://github.com/lithoscomputer/petri.git", branch = "main", package = "petri-store" } -petri_attractor_steps = { git = "https://github.com/lithoscomputer/petri.git", branch = "main", package = "petri-attractor-steps" } -petri_frontend_attractor = { git = "https://github.com/lithoscomputer/petri.git", branch = "main", package = "petri-frontend-attractor" } -petri_frontend_fabro = { git = "https://github.com/lithoscomputer/petri.git", branch = "main", package = "petri-frontend-fabro" } -petri_testkit = { git = "https://github.com/lithoscomputer/petri.git", branch = "main", package = "petri-testkit" } +petri_runtime = { git = "https://github.com/lithoscomputer/petri.git", rev = "0198e66f6619e147e461c52c138497863d851969", package = "petri-runtime" } +petri_execution = { git = "https://github.com/lithoscomputer/petri.git", rev = "0198e66f6619e147e461c52c138497863d851969", package = "petri-execution" } +petri_store = { git = "https://github.com/lithoscomputer/petri.git", rev = "0198e66f6619e147e461c52c138497863d851969", package = "petri-store" } +petri_attractor_steps = { git = "https://github.com/lithoscomputer/petri.git", rev = "0198e66f6619e147e461c52c138497863d851969", package = "petri-attractor-steps" } +petri_frontend_attractor = { git = "https://github.com/lithoscomputer/petri.git", rev = "0198e66f6619e147e461c52c138497863d851969", package = "petri-frontend-attractor" } +petri_frontend_fabro = { git = "https://github.com/lithoscomputer/petri.git", rev = "0198e66f6619e147e461c52c138497863d851969", package = "petri-frontend-fabro" } +petri_testkit = { git = "https://github.com/lithoscomputer/petri.git", rev = "0198e66f6619e147e461c52c138497863d851969", package = "petri-testkit" } sentry = { version = "0.35", default-features = false, features = ["backtrace", "contexts", "ureq", "rustls"] } fork = "0.2" exec = "0.3" 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..7faad090c --- /dev/null +++ b/docs/internal/run-finalization.md @@ -0,0 +1,54 @@ +# Required run finalization + +Petri owns the durable run result. Fabro declares required finalization when +its hooks have a publisher and implements it in `FabroHooks::finalize_run`. +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. Worker generation, +run scope and lease ownership checks remain required for every worker write. + +## Recovery and stored history + +This integration pins Petri revision +`0198e66f6619e147e461c52c138497863d851969`, 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. +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..d82251dd2 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,100 @@ 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"}); + 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; + 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 { + 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" + ); + waiting.delete_async().await; + 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; + 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..a5e3bd86b 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, @r###" + 12:43:08 ✗ FAILED 0s + 12: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..02cfb925f 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, as the legacy publish step did: +//! 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-server/src/petri_runs.rs b/lib/apps/fabro-server/src/petri_runs.rs index 28380e8f2..57462b0a3 100644 --- a/lib/apps/fabro-server/src/petri_runs.rs +++ b/lib/apps/fabro-server/src/petri_runs.rs @@ -940,4 +940,189 @@ 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>, + ) { + 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() + ); + 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 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; + let (rebuilt, _, _) = + fabro_petri::test_support::rebuild(&state.db_pool, &state.db_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..5d6308c03 100644 --- a/lib/apps/fabro-server/src/server.rs +++ b/lib/apps/fabro-server/src/server.rs @@ -3582,8 +3582,10 @@ 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); + if !managed_run.status.is_terminal() { + managed_run.status = RunStatus::Failed { reason }; + managed_run.error = Some(message); + } clear_live_run_state(managed_run); } cleanup_worker_control_bus_for_run(state.as_ref(), run_id); @@ -3601,7 +3603,9 @@ 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. - if managed_run.status.is_terminal() && !is_terminal_transition(record) { + if managed_run.status.is_terminal() + && (!is_terminal_transition(record) || record.status != Some(managed_run.status)) + { return; } match record.transition { @@ -3660,7 +3664,9 @@ fn apply_lifecycle_to_managed_run(state: &AppState, run_id: RunId, record: &RunL managed_run.status = record.status.unwrap_or(RunStatus::Failed { reason: FailureReason::WorkflowError, }); - managed_run.error.clone_from(&record.reason); + if managed_run.error.is_none() { + managed_run.error.clone_from(&record.reason); + } managed_run.active_steerable_stages.clear(); managed_run.active_non_steerable_stages.clear(); cleanup_worker_control_bus_for_run(state, run_id); @@ -3683,19 +3689,14 @@ 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 { @@ -3705,12 +3706,13 @@ pub(in crate::server) fn settle_managed_run_at_finish( return; } managed_run.status = status; + managed_run.error = failure; 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 diff --git a/lib/apps/fabro-server/src/server/handler/petri.rs b/lib/apps/fabro-server/src/server/handler/petri.rs index 403a97cd1..b50595df2 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::finished_run_result; 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,13 @@ 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(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 +276,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 +283,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..88880c642 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. //! @@ -549,13 +549,13 @@ 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"); + 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"); + release_live_state(&state, run_id); + } } - // 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); // 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,10 +814,13 @@ 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"); + 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 failure status"); + release_live_state(state, run_id); + } } - finish(state, run_id, status, error); } /// Settle the managed run at its terminal record and release its @@ -828,24 +831,31 @@ async fn fail_before_execution(state: &Arc, run_id: RunId, message: &s 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; + if !managed_run.status.is_terminal() || managed_run.status == status { + managed_run.status = status; + if managed_run.error.is_none() { + managed_run.error = error; + } + } clear_live_run_state(managed_run); } drop(runs); state.scheduler_notify.notify_one(); } -/// 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. +/// Release controls after an append failure without claiming a new terminal +/// result. A Petri finish already committed remains authoritative. +fn release_live_state(state: &Arc, 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); + } + drop(runs); + state.scheduler_notify.notify_one(); +} + +/// 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 +874,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 +888,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 +935,7 @@ mod tests { state: Arc, run_id: RunId, seen: Mutex>>, + reject: std::sync::atomic::AtomicBool, } #[async_trait::async_trait] @@ -933,6 +947,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 +1009,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 +1019,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 +1048,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..dc1ae0d9d 100644 --- a/lib/components/fabro-petri/src/engine.rs +++ b/lib/components/fabro-petri/src/engine.rs @@ -363,16 +363,6 @@ pub async fn run(request: RunRequest) -> Result { 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) } @@ -555,19 +545,27 @@ fn outcome( return Err(RunError::Unfinished(reasons)); } }; - let failure = inspection + 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 publish_failed = inspection + .finalization_failure + .as_ref() + .is_some_and(|failure| failure.code == "publish_failed"); + 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/hooks.rs b/lib/components/fabro-petri/src/hooks.rs index 40f72d260..8148ce7d9 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 +//! - `finalize_run`: 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. +//! record. +//! - `run_finished`: best-effort diff preparation for nonpublishing runs, then +//! 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 @@ -100,7 +102,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}; @@ -453,8 +457,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 +524,6 @@ impl FabroHooks { scopes: ScopeEnvs::default(), failure: Mutex::default(), publisher: spec.publisher, - publish_failure: Mutex::default(), resumed, restore: OnceCell::new(), store, @@ -544,13 +545,6 @@ impl FabroHooks { 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 +1397,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,33 +1603,53 @@ 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" - ); + fn requires_run_finalization(&self) -> bool { + self.publisher.is_some() || self.inner.requires_run_finalization() + } + + async fn finalize_run( + &self, + context: &HookContext, + finished: RunFinished, + ) -> Result<(), FinalizationFailure> { let publication = match self.record_run_diff().await { Ok(publication) => publication, Err(error) => { let message = error.render(); warn!(run_id = %self.run_id, error = %message, "the run's diff was not recorded"); if self.publisher.is_some() && finished.status == RunStatus::Success { - *sync::lock(&self.publish_failure) = Some(message); + return Err(FinalizationFailure::new("publish_failed", 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(|| { + FinalizationFailure::new( + "publish_failed", + "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"); + FinalizationFailure::new("publish_failed", 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 { + // Nonpublishing runs keep best-effort diff preparation. Required work + // has already run in finalize_run, before Petri commits its outcome. + if !self.requires_run_finalization() { + if let Err(error) = self.record_run_diff().await { + warn!(run_id = %self.run_id, error = %error.render(), "the run's diff was not recorded"); } } 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..888f77c35 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, @@ -210,7 +227,10 @@ 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 { +pub(super) fn finished_status( + status: &str, + finalization_failure: Option<&FinalizationFailure>, +) -> RunStatus { match status { "success" => RunStatus::Succeeded { reason: SuccessReason::Completed, @@ -219,7 +239,12 @@ pub(super) fn finished_status(status: &str) -> RunStatus { reason: FailureReason::Cancelled, }, _ => RunStatus::Failed { - reason: FailureReason::WorkflowError, + reason: if finalization_failure.is_some_and(|failure| failure.code == "publish_failed") + { + 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..d93a449b7 100644 --- a/lib/components/fabro-petri/src/projection/mod.rs +++ b/lib/components/fabro-petri/src/projection/mod.rs @@ -375,11 +375,22 @@ pub fn run_id_of(key: &str) -> Option { /// terminal lifecycle record. `None` for any other record. #[must_use] pub fn finished_run_status(record: &Record) -> Option { + finished_run_result(record).map(|(status, _)| status) +} + +/// 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, } } 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..d21a4332b --- /dev/null +++ b/lib/components/fabro-petri/src/test_support/finalization.rs @@ -0,0 +1,147 @@ +//! 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::sync::{Notify, Semaphore}; + +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(FinalizationFailure::new("publish_failed", 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>, +} + +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"); + tokio::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"); + tokio::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.clone(), 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..c18e70b77 100644 --- a/lib/components/fabro-petri/tests/hooks.rs +++ b/lib/components/fabro-petri/tests/hooks.rs @@ -1307,10 +1307,31 @@ 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 (harness, outcome) = published_run("script=\"echo edited >> README.md\"", &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, "", 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. @@ -1570,6 +1591,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: tokio::sync::Notify, + release: tokio::sync::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.unwrap().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: tokio::sync::Notify::new(), + release: tokio::sync::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 + }); + tokio::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!( + pending.invocations[0].result.is_some(), + "execution already ended" + ); + 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..4175b85ef 100644 --- a/lib/components/fabro-petri/tests/projection.rs +++ b/lib/components/fabro-petri/tests/projection.rs @@ -1361,3 +1361,168 @@ 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() + }); + tokio::time::timeout(Duration::from_secs(15), finalizer.entered.notified()) + .await + .unwrap(); + projector.settle(scenario.run_id).await; + let pending = petri_support::stored_projection(&scenario.pool, scenario.run_id) + .await + .unwrap() + .unwrap(); + 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(fabro_store::platform_records::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(fabro_store::platform_records::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.execution_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; + } +}