mirror of
https://github.com/fabro-sh/fabro.git
synced 2026-10-01 02:04:24 +00:00
Merge pull request #888 from fabro-sh/settle-worker-run-on-terminal-record
Settle a worker's run at Petri's finish, not at the worker's exit
This commit is contained in:
commit
81b7427add
7 changed files with 588 additions and 28 deletions
|
|
@ -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();
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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<HeldWorkerRuntime>,
|
||||
) -> (Arc<AppState>, 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<dyn WorkerRuntime>)
|
||||
.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);
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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<Box<dyn AsyncRead + Send + 'static>>,
|
||||
|
|
@ -4204,13 +4284,8 @@ async fn execute_run_subprocess(state: Arc<AppState>, 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()
|
||||
|
|
|
|||
|
|
@ -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()
|
||||
|
|
|
|||
|
|
@ -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
|
||||
);
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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:<fork>@<firing>:<index>:<target>`.
|
||||
fn branch_slot(slot: &str) -> Option<(u64, u32)> {
|
||||
|
|
|
|||
|
|
@ -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<RunId> {
|
|||
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<RunStatus> {
|
||||
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();
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue