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();