From 2b30689e77a707a5b4bac3a5bbd6183c1ddc2e33 Mon Sep 17 00:00:00 2001 From: Bryan Helmkamp Date: Mon, 21 Sep 2026 03:49:44 -0400 Subject: [PATCH] Settle a worker's run at Petri's finish, not at the worker's exit A worker-backed run's in-memory managed run settled only when the worker process exited. GET /runs/{id} reads the stored summary, which the projector ends at Petri's own `run.finished` record, a moment before the worker stores Fabro's terminal lifecycle record and exits. The delete precheck prefers the managed run, so a delete issued in that window was refused with 409 "cannot remove active run". Against a real worker the window hit about six times in ten. The server sees both records before they are stored: `run.finished` on the coordinator log through the worker's records endpoint, and the terminal lifecycle record through the platform-records endpoint. The managed run now settles at either, ahead of the store, so the view never reports the run ended while the managed run still says running. The stream follower no longer reopens a settled run with the records that precede its terminal one, and the worker's exit keeps the settled status: it reaps the process, records a missing terminal record as it did, and takes the store's status only when the store ended the run differently. The mapping from Petri's finish to the run's status is the projection's own, shared with its fold. Co-Authored-By: Claude Fable 5.1 --- lib/apps/fabro-cli/tests/it/scenario/petri.rs | 53 ++++ lib/apps/fabro-server/src/petri_runs.rs | 294 +++++++++++++++++- lib/apps/fabro-server/src/server.rs | 89 +++++- .../fabro-server/src/server/handler/petri.rs | 16 +- lib/apps/fabro-server/src/server/tests.rs | 44 +++ .../fabro-petri/src/projection/coordinator.rs | 43 +-- .../fabro-petri/src/projection/mod.rs | 77 ++++- 7 files changed, 588 insertions(+), 28 deletions(-) diff --git a/lib/apps/fabro-cli/tests/it/scenario/petri.rs b/lib/apps/fabro-cli/tests/it/scenario/petri.rs index d92e6be65..6767c6bdf 100644 --- a/lib/apps/fabro-cli/tests/it/scenario/petri.rs +++ b/lib/apps/fabro-cli/tests/it/scenario/petri.rs @@ -1752,3 +1752,56 @@ async fn a_failed_checkpoint_fails_the_run_and_a_restart_leaves_it_failed() { ); server.shutdown(); } + +/// A delete issued the moment the run reads as ended is accepted while its +/// worker is still tearing down: the server settles the run at Petri's own +/// finish, the record the view ends the run on, not at the worker's exit, +/// so the delete precheck does not refuse the run as active. The worker's +/// exit after the delete brings nothing back. +#[tokio::test(flavor = "multi_thread")] +async fn a_delete_right_after_the_run_reads_ended_is_accepted() { + if host_plugin().is_none() { + return; + } + let context = test_context!(); + let server = RunningServer::start().await; + let workspace = write_petri_workspace(&context, "true"); + let run_id = run_detached(&context, &server, &workspace); + let status = wait_for_status(&server, &run_id, &["succeeded", "failed"]).await; + assert_eq!(status, "succeeded"); + + let response = fabro_test::test_http_client() + .delete(format!("{}/api/v1/runs/{run_id}", server.api_base_url)) + .bearer_auth(TEST_DEV_TOKEN) + .send() + .await + .expect("the delete sends"); + let status = response.status(); + let detail = response.text().await.unwrap_or_default(); + assert_eq!( + status, + fabro_http::StatusCode::NO_CONTENT, + "the delete was refused: {detail}" + ); + + let deadline = Instant::now() + RUN_TIMEOUT; + while worker_pid(&run_id).is_some() { + assert!( + Instant::now() < deadline, + "the worker of run {run_id} outlived the delete" + ); + tokio::time::sleep(POLL).await; + } + let response = fabro_test::test_http_client() + .get(format!("{}/api/v1/runs/{run_id}", server.api_base_url)) + .bearer_auth(TEST_DEV_TOKEN) + .send() + .await + .expect("the request sends"); + assert_eq!( + response.status(), + fabro_http::StatusCode::NOT_FOUND, + "the worker's exit brought the run back" + ); + server.shutdown(); +} diff --git a/lib/apps/fabro-server/src/petri_runs.rs b/lib/apps/fabro-server/src/petri_runs.rs index 16bef7e2b..1690fbe0e 100644 --- a/lib/apps/fabro-server/src/petri_runs.rs +++ b/lib/apps/fabro-server/src/petri_runs.rs @@ -171,7 +171,9 @@ mod tests { use fabro_petri::petri::RunStore as _; use fabro_static::EnvVars; use fabro_store::platform_records::{PlatformRecord, RunLifecycleKind, RunLifecycleRecord}; - use fabro_types::{RunId, RunStatus, WorkflowPath, WorkflowVersion}; + use fabro_types::{ + FailureReason, RunId, RunStatus, SuccessReason, WorkflowPath, WorkflowVersion, + }; use serde_json::json; use tokio::io::AsyncRead; use tokio::sync::Notify; @@ -512,4 +514,294 @@ mod tests { // crash releases nothing, and dropping them here would. drop(before); } + + /// A server whose one worker is held open by the test, with a run the + /// worker has taken over the API: the run's id and the worker's token. + async fn held_worker_run( + runtime: &Arc, + ) -> (Arc, axum::Router, RunId, String) { + let state = TestAppStateBuilder::new() + .vault_entries([(EnvVars::OPENAI_API_KEY, "test-openai-api-key")]) + .worker_runtime(Arc::clone(runtime) as Arc) + .build(); + write_test_server_record(&state); + let app = build_test_router(Arc::clone(&state)); + let run_id = create_and_start_petri_run(&app).await; + spawn_scheduler(Arc::clone(&state)); + runtime.wait_for_start().await; + let token = state.test_issue_worker_token(&run_id); + let response = app + .clone() + .oneshot( + Request::builder() + .method("POST") + .uri(format!("/api/v1/runs/{run_id}/petri/open")) + .header(header::AUTHORIZATION, format!("Bearer {token}")) + .header(header::CONTENT_TYPE, "application/json") + .body(Body::from( + json!({ "access": "create", "owner": "worker-1" }).to_string(), + )) + .expect("the open request builds"), + ) + .await + .expect("the open request completes"); + assert_eq!(response.status(), StatusCode::OK); + (state, app, run_id, token) + } + + /// One lifecycle record of the run, stored the way its worker stores + /// one: through the platform-records endpoint. + async fn record_lifecycle_as_worker( + app: &axum::Router, + run_id: RunId, + token: &str, + transition: RunLifecycleKind, + status: RunStatus, + ) { + 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() + .oneshot( + Request::builder() + .method("POST") + .uri(format!("/api/v1/runs/{run_id}/petri/platform-records")) + .header(header::AUTHORIZATION, format!("Bearer {token}")) + .header(header::CONTENT_TYPE, "application/json") + .body(Body::from(body.to_string())) + .expect("the append request builds"), + ) + .await + .expect("the append request completes"); + assert_eq!(response.status(), StatusCode::OK); + } + + /// Petri's own finish of the run, stored the way its worker stores it: + /// the `run.finished` record on the coordinator log, over the records + /// endpoint. + async fn record_finish_as_worker(app: &axum::Router, run_id: RunId, token: &str) { + let body = json!({ + "owner": "worker-1", + "records": [{ + "seq": 0, + "recorded_at": 1_000, + "record": { + "seq": 0, + "origin": "external", + "recorded_at": 1_000, + "body": { "event": "run.finished", "status": "success" }, + }, + }], + }); + let response = app + .clone() + .oneshot( + Request::builder() + .method("POST") + .uri(format!( + "/api/v1/runs/{run_id}/petri/logs/coordinator/records" + )) + .header(header::AUTHORIZATION, format!("Bearer {token}")) + .header(header::CONTENT_TYPE, "application/json") + .body(Body::from(body.to_string())) + .expect("the records request builds"), + ) + .await + .expect("the records request completes"); + assert_eq!(response.status(), StatusCode::NO_CONTENT); + } + + /// The worker takes the run as far as running. + async fn run_to_running_as_worker(app: &axum::Router, run_id: RunId, token: &str) { + for (transition, status) in [ + (RunLifecycleKind::Starting, RunStatus::Starting), + (RunLifecycleKind::Running, RunStatus::Running), + ] { + record_lifecycle_as_worker(app, run_id, token, transition, status).await; + } + } + + async fn status_of(app: &axum::Router, run_id: RunId) -> StatusCode { + app.clone() + .oneshot( + Request::builder() + .method("GET") + .uri(format!("/api/v1/runs/{run_id}")) + .body(Body::empty()) + .expect("the run request builds"), + ) + .await + .expect("the run request completes") + .status() + } + + /// The runs the server's usage aggregate has counted: it counts a run + /// at its worker's exit, once the exit is fully handled. + async fn concluded_runs(app: &axum::Router) -> i64 { + let response = app + .clone() + .oneshot( + Request::builder() + .method("GET") + .uri("/api/v1/usage") + .body(Body::empty()) + .expect("the usage request builds"), + ) + .await + .expect("the usage request completes"); + assert_eq!(response.status(), StatusCode::OK); + let body = to_bytes(response.into_body(), usize::MAX) + .await + .expect("the usage body reads"); + let body: serde_json::Value = + serde_json::from_slice(&body).expect("the usage body is JSON"); + body["totals"]["runs"].as_i64().expect("the run count") + } + + /// A delete of the run, as a client issues one. + async fn delete_run(app: &axum::Router, run_id: RunId) -> StatusCode { + app.clone() + .oneshot( + Request::builder() + .method("DELETE") + .uri(format!("/api/v1/runs/{run_id}")) + .body(Body::empty()) + .expect("the delete request builds"), + ) + .await + .expect("the delete request completes") + .status() + } + + /// The server settles a worker's run at Petri's own finish, the record + /// the view ends the run on, not at the worker's exit: a delete issued + /// the moment the finish lands, while the worker still tears down, is + /// accepted, and the worker's exit after it brings nothing back. + #[tokio::test(flavor = "multi_thread", worker_threads = 2)] + async fn a_delete_right_after_petris_finish_is_accepted() { + 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; + record_finish_as_worker(&app, run_id, &token).await; + assert_eq!( + state.test_managed_run_status(&run_id), + Some(RunStatus::Succeeded { + reason: SuccessReason::Completed, + }), + "the managed run settled at the finish" + ); + assert!( + runtime.running.load(Ordering::SeqCst), + "the worker is still up when the delete lands" + ); + assert_eq!(delete_run(&app, run_id).await, StatusCode::NO_CONTENT); + + // The delete ended the worker; its exit is handled in the + // background and must leave the run gone. + assert!(!runtime.running.load(Ordering::SeqCst)); + for _ in 0..20 { + assert_eq!(status_of(&app, run_id).await, StatusCode::NOT_FOUND); + assert_eq!(state.test_managed_run_status(&run_id), None); + time::sleep(Duration::from_millis(10)).await; + } + } + + /// A worker that ends the run without Petri's finish (it failed before + /// the engine ran) settles it at its terminal lifecycle record, so a + /// delete issued the moment that record lands is accepted too. + #[tokio::test(flavor = "multi_thread", worker_threads = 2)] + async fn a_delete_right_after_the_workers_terminal_record_is_accepted() { + let runtime = Arc::new(HeldWorkerRuntime::default()); + let (state, app, run_id, token) = held_worker_run(&runtime).await; + + let failed = RunStatus::Failed { + reason: FailureReason::WorkflowError, + }; + record_lifecycle_as_worker( + &app, + run_id, + &token, + RunLifecycleKind::Starting, + RunStatus::Starting, + ) + .await; + record_lifecycle_as_worker(&app, run_id, &token, RunLifecycleKind::Failed, failed).await; + assert_eq!( + state.test_managed_run_status(&run_id), + Some(failed), + "the managed run settled at the terminal record" + ); + assert!( + runtime.running.load(Ordering::SeqCst), + "the worker is still up when the delete lands" + ); + + let response = app + .clone() + .oneshot( + Request::builder() + .method("DELETE") + .uri(format!("/api/v1/runs/{run_id}")) + .body(Body::empty()) + .expect("the delete request builds"), + ) + .await + .expect("the delete request completes"); + assert_eq!(response.status(), StatusCode::NO_CONTENT); + + // The delete ended the worker; its exit is handled in the + // background and must leave the run gone. + assert!(!runtime.running.load(Ordering::SeqCst)); + for _ in 0..20 { + assert_eq!(status_of(&app, run_id).await, StatusCode::NOT_FOUND); + assert_eq!(state.test_managed_run_status(&run_id), None); + time::sleep(Duration::from_millis(10)).await; + } + } + + /// A worker that exits, even unsuccessfully, after its run settled at + /// its finish and its terminal record leaves the settled status alone: + /// the exit reaps the process and counts the run, nothing more. + #[tokio::test(flavor = "multi_thread", worker_threads = 2)] + async fn a_worker_exit_after_the_run_settled_keeps_the_settled_status() { + let runtime = Arc::new(HeldWorkerRuntime::default()); + let (state, app, run_id, token) = held_worker_run(&runtime).await; + let succeeded = RunStatus::Succeeded { + reason: SuccessReason::Completed, + }; + + run_to_running_as_worker(&app, run_id, &token).await; + record_finish_as_worker(&app, run_id, &token).await; + record_lifecycle_as_worker(&app, run_id, &token, RunLifecycleKind::Succeeded, succeeded) + .await; + assert_eq!(state.test_managed_run_status(&run_id), Some(succeeded)); + + // The held worker exits unsuccessfully, as a worker that crashed + // after its terminal record would. + runtime.end_worker(); + let mut counted = concluded_runs(&app).await; + for _ in 0..500 { + if counted == 1 { + break; + } + time::sleep(Duration::from_millis(10)).await; + counted = concluded_runs(&app).await; + } + assert_eq!(counted, 1, "the worker's exit was handled"); + + assert_eq!( + state.test_managed_run_status(&run_id), + Some(succeeded), + "the exit did not move the settled status" + ); + let run_state = run_records::projection(&state, run_id) + .await + .expect("the run state loads") + .expect("the run projects"); + assert_eq!(run_state.status, succeeded); + } } diff --git a/lib/apps/fabro-server/src/server.rs b/lib/apps/fabro-server/src/server.rs index 05502d0f1..100ee3285 100644 --- a/lib/apps/fabro-server/src/server.rs +++ b/lib/apps/fabro-server/src/server.rs @@ -3531,6 +3531,12 @@ fn apply_lifecycle_to_managed_run(state: &AppState, run_id: RunId, record: &RunL let Some(managed_run) = runs.get_mut(&run_id) else { return; }; + // 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) { + return; + } match record.transition { RunLifecycleKind::Submitted => managed_run.status = RunStatus::Submitted, RunLifecycleKind::Pending => { @@ -3602,6 +3608,80 @@ fn apply_lifecycle_to_managed_run(state: &AppState, run_id: RunId, record: &RunL } } +/// Whether the transition ends the run. +fn is_terminal_transition(record: &RunLifecycleRecord) -> bool { + matches!( + record.transition, + RunLifecycleKind::Succeeded | RunLifecycleKind::Failed | RunLifecycleKind::Dead + ) +} + +/// 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. +pub(in crate::server) fn settle_managed_run_at_finish( + state: &AppState, + run_id: RunId, + status: RunStatus, +) { + 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; + } + 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`] +/// 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 +/// follower, which folds the stream in order. The worker's exit later +/// leaves the settled status alone, unless the store ended the run +/// differently: the worker's append failed and the exit recorded the +/// failure. +pub(in crate::server) fn settle_managed_run_at_terminal_record( + state: &AppState, + run_id: RunId, + record: &PlatformRecord, +) { + if let PlatformRecord::RunLifecycle(record) = record { + if is_terminal_transition(record) { + apply_lifecycle_to_managed_run(state, run_id, record); + } + } +} + +/// The live status once the run's worker is gone. A run settled at its +/// terminal record keeps that status: the exit only reaps the process. The +/// store's final status stands when the run was not settled, or when the +/// store ended it differently; and a worker that died with no terminal +/// record and no failure recorded for it is a termination. +fn status_after_worker_exit(live: RunStatus, stored: RunStatus, exit_success: bool) -> RunStatus { + if live.is_terminal() { + if stored.is_terminal() { stored } else { live } + } else if stored != live { + stored + } else if exit_success { + live + } else { + RunStatus::Failed { + reason: FailureReason::Terminated, + } + } +} + async fn drain_worker_stderr( run_id: RunId, stderr: std::pin::Pin>, @@ -4204,13 +4284,8 @@ async fn execute_run_subprocess(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) { - if final_state.status != managed_run.status { - managed_run.status = final_state.status; - } else if !worker_exit.success { - managed_run.status = RunStatus::Failed { - reason: FailureReason::Terminated, - }; - } + managed_run.status = + status_after_worker_exit(managed_run.status, final_state.status, worker_exit.success); managed_run.error = final_state .conclusion .as_ref() diff --git a/lib/apps/fabro-server/src/server/handler/petri.rs b/lib/apps/fabro-server/src/server/handler/petri.rs index 8d255ff1c..403a97cd1 100644 --- a/lib/apps/fabro-server/src/server/handler/petri.rs +++ b/lib/apps/fabro-server/src/server/handler/petri.rs @@ -22,7 +22,8 @@ use fabro_api::types::{ PetriPlatformRecordAppendRequest, PetriPlatformRecordList, PetriRecord, PetriRecordList, PetriReleaseRequest, WriteBlobResponse, }; -use fabro_petri::petri::{Access, Digest, OwnerId, Record, StoreError}; +use fabro_petri::petri::{Access, Digest, LogId, OwnerId, Record, StoreError}; +use fabro_petri::projection::finished_run_status; use fabro_petri::run_store::{log_id_text, parse_log_id}; use fabro_store::{PlatformRecord, PlatformRecordKind, StagePosition, StoredPlatformRecord}; use fabro_types::BlobHash; @@ -32,6 +33,7 @@ use serde_json::{Map, Value, json}; use super::super::{ ApiError, AppState, Bytes, IntoResponse, Json, Query, RequireWorkerRunScoped, RequireWorkerRunSegment, Response, Router, RunId, State, StatusCode, octet_stream_response, + settle_managed_run_at_finish, settle_managed_run_at_terminal_record, }; /// The largest batch of records one append may carry. Petri batches an @@ -155,6 +157,14 @@ 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(()) => { // The records are durable; the projection trails them from here. @@ -269,6 +279,10 @@ 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() diff --git a/lib/apps/fabro-server/src/server/tests.rs b/lib/apps/fabro-server/src/server/tests.rs index 127e6667e..c6262bb39 100644 --- a/lib/apps/fabro-server/src/server/tests.rs +++ b/lib/apps/fabro-server/src/server/tests.rs @@ -10930,3 +10930,47 @@ async fn run_tools_worker_registers_contents_then_creates_by_version_id() { assert_eq!(projection.spec.target, Some(RunTarget::None {})); assert!(projection.start.is_none()); } + +/// A run settled at its terminal record keeps that status at its worker's +/// exit, whatever the exit status, unless the store ended the run +/// differently: then the store's terminal status stands. +#[test] +fn a_settled_run_keeps_its_status_at_worker_exit() { + let settled = RunStatus::Succeeded { + reason: SuccessReason::Completed, + }; + assert_eq!(status_after_worker_exit(settled, settled, false), settled); + assert_eq!(status_after_worker_exit(settled, settled, true), settled); + assert_eq!( + status_after_worker_exit(settled, RunStatus::Running, false), + settled + ); + let recorded = RunStatus::Failed { + reason: FailureReason::Terminated, + }; + assert_eq!(status_after_worker_exit(settled, recorded, false), recorded); +} + +/// A run its worker left unsettled takes the store's final status; with +/// none recorded, an unsuccessful exit is a termination and a successful +/// one changes nothing. +#[test] +fn an_unsettled_run_takes_the_stores_status_or_a_termination_at_worker_exit() { + let failed = RunStatus::Failed { + reason: FailureReason::WorkflowError, + }; + assert_eq!( + status_after_worker_exit(RunStatus::Running, failed, false), + failed + ); + assert_eq!( + status_after_worker_exit(RunStatus::Running, RunStatus::Running, false), + RunStatus::Failed { + reason: FailureReason::Terminated, + } + ); + assert_eq!( + status_after_worker_exit(RunStatus::Running, RunStatus::Running, true), + RunStatus::Running + ); +} diff --git a/lib/components/fabro-petri/src/projection/coordinator.rs b/lib/components/fabro-petri/src/projection/coordinator.rs index eb51c0529..f78ebc439 100644 --- a/lib/components/fabro-petri/src/projection/coordinator.rs +++ b/lib/components/fabro-petri/src/projection/coordinator.rs @@ -138,23 +138,16 @@ impl RunView { .root .and_then(|root| self.state.invocations.get(&root)); let failure_message = root.and_then(|root| root.failure.clone()); - let (run_status, outcome, failure) = match status { - "success" => ( - RunStatus::Succeeded { - reason: SuccessReason::Completed, - }, - StageOutcome::Succeeded, - None, - ), - "cancelled" => ( - RunStatus::Failed { - reason: FailureReason::Cancelled, - }, + let run_status = finished_status(status); + let (outcome, failure) = match run_status { + RunStatus::Failed { + reason: reason @ FailureReason::Cancelled, + } => ( StageOutcome::Failed { retry_requested: false, }, Some(RunFailure { - reason: FailureReason::Cancelled, + reason, detail: FailureDetail::new( failure_message .clone() @@ -163,15 +156,12 @@ impl RunView { ), }), ), - _ => ( - RunStatus::Failed { - reason: FailureReason::WorkflowError, - }, + RunStatus::Failed { reason } => ( StageOutcome::Failed { retry_requested: false, }, Some(RunFailure { - reason: FailureReason::WorkflowError, + reason, detail: FailureDetail::new( failure_message .clone() @@ -180,6 +170,7 @@ impl RunView { ), }), ), + _ => (StageOutcome::Succeeded, None), }; apply_status(projection, run_status, at); projection.pending_control = None; @@ -217,6 +208,22 @@ 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 { + match status { + "success" => RunStatus::Succeeded { + reason: SuccessReason::Completed, + }, + "cancelled" => RunStatus::Failed { + reason: FailureReason::Cancelled, + }, + _ => RunStatus::Failed { + reason: FailureReason::WorkflowError, + }, + } +} + /// The fork firing and branch index a branch child's call slot names: /// `branch:@::`. fn branch_slot(slot: &str) -> Option<(u64, u32)> { diff --git a/lib/components/fabro-petri/src/projection/mod.rs b/lib/components/fabro-petri/src/projection/mod.rs index e7a4b9cd8..4410df7b9 100644 --- a/lib/components/fabro-petri/src/projection/mod.rs +++ b/lib/components/fabro-petri/src/projection/mod.rs @@ -42,8 +42,9 @@ use fabro_store::platform_records::StoredPlatformRecord; use fabro_types::{ RunControlAction, RunDiff, RunId, RunProjection, RunStatus, StageId, StageProjection, }; -use petri_execution::ExecutionId; use petri_execution::events::{NodeRef, RunEvent, Subject}; +use petri_execution::{CoordinatorEvent, CoordinatorRecord, ExecutionId}; +use petri_store::Record; use serde::{Deserialize, Deserializer, Serialize, Serializer, de}; use serde_json::Value; use tracing::debug; @@ -368,6 +369,21 @@ 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. +#[must_use] +pub fn finished_run_status(record: &Record) -> 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())) + } + _ => None, + } +} + #[cfg(test)] mod tests { use fabro_store::platform_records::{PlatformRecord, RunCreatedRecord}; @@ -377,6 +393,65 @@ mod tests { use super::*; + /// A stored coordinator record, as the worker's append carries one. + fn coordinator_record(body: &serde_json::Value) -> Record { + Record { + seq: 3, + recorded_at: 1_000, + record: serde_json::json!({ + "seq": 3, + "origin": "external", + "recorded_at": 1_000, + "body": body, + }), + } + } + + /// Petri's finish gives the run the status the view folds it to; any + /// other record, or a line that is not a coordinator record, gives none. + #[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!({ + "event": "run.finished", + "status": status, + }))) + }; + assert_eq!( + finished("success"), + Some(RunStatus::Succeeded { + reason: fabro_types::SuccessReason::Completed, + }) + ); + assert_eq!( + finished("cancelled"), + Some(RunStatus::Failed { + reason: fabro_types::FailureReason::Cancelled, + }) + ); + assert_eq!( + finished("failed"), + Some(RunStatus::Failed { + reason: fabro_types::FailureReason::WorkflowError, + }) + ); + assert_eq!( + finished_run_status(&coordinator_record(&serde_json::json!({ + "event": "run.paused", + }))), + None + ); + assert_eq!( + finished_run_status(&Record { + seq: 3, + recorded_at: 1_000, + record: serde_json::json!({"event": "run.finished", "status": "success"}), + }), + None, + "an engine line is not a coordinator record" + ); + } + #[test] fn a_taken_label_is_made_unique_by_the_execution() { let mut view = RunView::new();