Merge pull request #939 from fabro-sh/required-run-finalization

Commit required publication before reporting run completion
This commit is contained in:
Scott Werner 2026-10-09 12:38:30 -04:00 • committed by GitHub
commit 5b3689c132
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
24 changed files with 1553 additions and 283 deletions

36
Cargo.lock generated
View file

@ -1642,7 +1642,7 @@ checksum = "d7a1e2f27636f116493b8b860f5546edb47c8d8f8ea73e1d2a20be88e28d1fea"
[[package]]
name = "daytona-api-client"
version = "0.1.0"
source = "git+https://github.com/brynary/daytona-sdk-rust?branch=main#0e69058c888a6562c70a5b16d707253914cff563"
source = "git+https://github.com/brynary/daytona-sdk-rust?branch=main#df596bad2093f83fc793d1890bd25ba0d2f93425"
dependencies = [
"reqwest 0.13.4",
"reqwest-middleware",
@ -1656,7 +1656,7 @@ dependencies = [
[[package]]
name = "daytona-sdk"
version = "0.1.0"
source = "git+https://github.com/brynary/daytona-sdk-rust?branch=main#0e69058c888a6562c70a5b16d707253914cff563"
source = "git+https://github.com/brynary/daytona-sdk-rust?branch=main#df596bad2093f83fc793d1890bd25ba0d2f93425"
dependencies = [
"daytona-api-client",
"daytona-toolbox-client",
@ -1676,7 +1676,7 @@ dependencies = [
[[package]]
name = "daytona-toolbox-client"
version = "0.1.0"
source = "git+https://github.com/brynary/daytona-sdk-rust?branch=main#0e69058c888a6562c70a5b16d707253914cff563"
source = "git+https://github.com/brynary/daytona-sdk-rust?branch=main#df596bad2093f83fc793d1890bd25ba0d2f93425"
dependencies = [
"reqwest 0.13.4",
"reqwest-middleware",
@ -5166,7 +5166,7 @@ checksum = "9b4f627cb1b25917193a259e49bdad08f671f8d9708acfd5fe0a8c1455d87220"
[[package]]
name = "petri-attractor-steps"
version = "0.1.0"
source = "git+https://github.com/lithoscomputer/petri.git?branch=main#e46845bd0139dd04e9e795d3201b5b9b4b8b1026"
source = "git+https://github.com/lithoscomputer/petri.git?branch=main#505a17839987f052ef0a47b799d30c565fd29475"
dependencies = [
"async-trait",
"globset",
@ -5197,7 +5197,7 @@ dependencies = [
[[package]]
name = "petri-driver"
version = "0.1.0"
source = "git+https://github.com/lithoscomputer/petri.git?branch=main#e46845bd0139dd04e9e795d3201b5b9b4b8b1026"
source = "git+https://github.com/lithoscomputer/petri.git?branch=main#505a17839987f052ef0a47b799d30c565fd29475"
dependencies = [
"async-trait",
"getrandom 0.3.4",
@ -5217,7 +5217,7 @@ dependencies = [
[[package]]
name = "petri-engine"
version = "0.1.0"
source = "git+https://github.com/lithoscomputer/petri.git?branch=main#e46845bd0139dd04e9e795d3201b5b9b4b8b1026"
source = "git+https://github.com/lithoscomputer/petri.git?branch=main#505a17839987f052ef0a47b799d30c565fd29475"
dependencies = [
"petri-ir",
"serde",
@ -5229,7 +5229,7 @@ dependencies = [
[[package]]
name = "petri-execution"
version = "0.1.0"
source = "git+https://github.com/lithoscomputer/petri.git?branch=main#e46845bd0139dd04e9e795d3201b5b9b4b8b1026"
source = "git+https://github.com/lithoscomputer/petri.git?branch=main#505a17839987f052ef0a47b799d30c565fd29475"
dependencies = [
"async-trait",
"petri-driver",
@ -5253,7 +5253,7 @@ dependencies = [
[[package]]
name = "petri-executor"
version = "0.1.0"
source = "git+https://github.com/lithoscomputer/petri.git?branch=main#e46845bd0139dd04e9e795d3201b5b9b4b8b1026"
source = "git+https://github.com/lithoscomputer/petri.git?branch=main#505a17839987f052ef0a47b799d30c565fd29475"
dependencies = [
"async-trait",
"libc",
@ -5268,7 +5268,7 @@ dependencies = [
[[package]]
name = "petri-executor-sandbox"
version = "0.1.0"
source = "git+https://github.com/lithoscomputer/petri.git?branch=main#e46845bd0139dd04e9e795d3201b5b9b4b8b1026"
source = "git+https://github.com/lithoscomputer/petri.git?branch=main#505a17839987f052ef0a47b799d30c565fd29475"
dependencies = [
"async-trait",
"petri-executor",
@ -5290,7 +5290,7 @@ dependencies = [
[[package]]
name = "petri-frontend"
version = "0.1.0"
source = "git+https://github.com/lithoscomputer/petri.git?branch=main#e46845bd0139dd04e9e795d3201b5b9b4b8b1026"
source = "git+https://github.com/lithoscomputer/petri.git?branch=main#505a17839987f052ef0a47b799d30c565fd29475"
dependencies = [
"marked-yaml",
"petri-ir",
@ -5304,7 +5304,7 @@ dependencies = [
[[package]]
name = "petri-frontend-attractor"
version = "0.1.0"
source = "git+https://github.com/lithoscomputer/petri.git?branch=main#e46845bd0139dd04e9e795d3201b5b9b4b8b1026"
source = "git+https://github.com/lithoscomputer/petri.git?branch=main#505a17839987f052ef0a47b799d30c565fd29475"
dependencies = [
"minijinja",
"petri-frontend",
@ -5321,7 +5321,7 @@ dependencies = [
[[package]]
name = "petri-frontend-fabro"
version = "0.1.0"
source = "git+https://github.com/lithoscomputer/petri.git?branch=main#e46845bd0139dd04e9e795d3201b5b9b4b8b1026"
source = "git+https://github.com/lithoscomputer/petri.git?branch=main#505a17839987f052ef0a47b799d30c565fd29475"
dependencies = [
"petri-frontend",
"petri-frontend-attractor",
@ -5337,7 +5337,7 @@ dependencies = [
[[package]]
name = "petri-frontend-native"
version = "0.1.0"
source = "git+https://github.com/lithoscomputer/petri.git?branch=main#e46845bd0139dd04e9e795d3201b5b9b4b8b1026"
source = "git+https://github.com/lithoscomputer/petri.git?branch=main#505a17839987f052ef0a47b799d30c565fd29475"
dependencies = [
"petri-frontend",
"petri-ir",
@ -5348,7 +5348,7 @@ dependencies = [
[[package]]
name = "petri-ir"
version = "0.1.0"
source = "git+https://github.com/lithoscomputer/petri.git?branch=main#e46845bd0139dd04e9e795d3201b5b9b4b8b1026"
source = "git+https://github.com/lithoscomputer/petri.git?branch=main#505a17839987f052ef0a47b799d30c565fd29475"
dependencies = [
"regex",
"serde",
@ -5361,7 +5361,7 @@ dependencies = [
[[package]]
name = "petri-runtime"
version = "0.1.0"
source = "git+https://github.com/lithoscomputer/petri.git?branch=main#e46845bd0139dd04e9e795d3201b5b9b4b8b1026"
source = "git+https://github.com/lithoscomputer/petri.git?branch=main#505a17839987f052ef0a47b799d30c565fd29475"
dependencies = [
"async-trait",
"petri-driver",
@ -5382,7 +5382,7 @@ dependencies = [
[[package]]
name = "petri-steps"
version = "0.1.0"
source = "git+https://github.com/lithoscomputer/petri.git?branch=main#e46845bd0139dd04e9e795d3201b5b9b4b8b1026"
source = "git+https://github.com/lithoscomputer/petri.git?branch=main#505a17839987f052ef0a47b799d30c565fd29475"
dependencies = [
"async-trait",
"petri-executor",
@ -5398,7 +5398,7 @@ dependencies = [
[[package]]
name = "petri-store"
version = "0.1.0"
source = "git+https://github.com/lithoscomputer/petri.git?branch=main#e46845bd0139dd04e9e795d3201b5b9b4b8b1026"
source = "git+https://github.com/lithoscomputer/petri.git?branch=main#505a17839987f052ef0a47b799d30c565fd29475"
dependencies = [
"async-trait",
"getrandom 0.3.4",
@ -5413,7 +5413,7 @@ dependencies = [
[[package]]
name = "petri-testkit"
version = "0.1.0"
source = "git+https://github.com/lithoscomputer/petri.git?branch=main#e46845bd0139dd04e9e795d3201b5b9b4b8b1026"
source = "git+https://github.com/lithoscomputer/petri.git?branch=main#505a17839987f052ef0a47b799d30c565fd29475"
dependencies = [
"async-trait",
"petri-driver",

View file

@ -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.

View file

@ -0,0 +1,59 @@
# Required run finalization
Petri owns the durable run result. Fabro declares required finalization for
every run and implements it in `FabroHooks::finalize_run`, so a resume or a
fork always matches the stored declaration. A failed checkpoint rejects with
`checkpoint_failed` and skips publication. The checkpoint failure cancelled the
run, so Petri may commit `cancelled`; Fabro reports that finish as a workflow
failure with the checkpoint's message.
For successful workflow execution this prepares the final diff, then uses the
existing publisher to retry the push and reconcile or create the pull request.
A preparation or publication error rejects with `publish_failed` and a rendered
message. Missing branch, checkpoint or workspace evidence is a rejection.
Failed or cancelled execution skips publication; its diff remains best effort.
Local workflow run-end hooks remain observational.
The coordinator awaits finalization and cleanup before committing `run.finished`.
Until that record, the run stays active without a conclusion, even when every
stage has succeeded. The run stream, API projection, managed server status,
worker return and CLI verdict use the committed overall status and structured
finalization failure. Stage outcomes, metrics and invocation results keep their
execution meaning. Successful stages stay successful when publication fails.
Both HTTP and in-process worker transports store the finish before settling
managed status. A rejected append cannot settle the run. The later platform
terminal lifecycle record acknowledges the same outcome; it cannot replace a
terminal result or conclusion. Early worker failures without a Petri finish
still end through their committed platform lifecycle record. Run scope,
authorization and lease ownership checks remain required for every worker write.
## Recovery and stored history
This integration requires a Petri with required run finalization (coordinator
format 9) and event contract 6. Petri retains its strict stored-format policy: version-8 coordinator
logs cannot resume, inspect or replay with this engine. Fabro does not rewrite
source history or relax that policy. Existing materialized views and stream
rows remain stored; a projector replay failure holds Petri positions and reports
incomplete record health. A historical run without a materialized view cannot
recover its old Petri stage evidence through the new engine. Operators requiring
inspection or recovery of old execution logs must keep the prior compatible
binary and its storage backup.
Existing incorrect historical publication-success projections are not repaired
by this change. Consumed cursors and immutable conclusions remain unchanged;
historical-view repair requires a separate, bounded process over preserved
source records. No publication is invoked by projection or replay.
A committed resume returns the same overall result without publishing again.
A fork declares required finalization while seeding its records, as every
Fabro run does; its worker restores the actual publication hooks before resume.
Unfinished recovery must restore the same finalization requirement. Petri may
call the restored finalizer again after interruption, including a crash after
the callback returns but before the terminal record commits. Existing GitHub
reconciliation is retained; this contract does not provide independent retry
orchestration or guarantee external-effect deduplication after every crash.
Lost terminal-execution workspaces are not recreated for finalization. Fabro
rejects publication when it cannot establish the required evidence. Push
attempts retain their existing five-minute timeout; this integration adds no
overall finalizer deadline. Cancellation arriving after
execution ends does not interrupt Petri's awaited finalizer.

View file

@ -228,7 +228,7 @@ async fn attach_petri_run_with_client(
for item in &items {
emit_stream_item(&mut progress_ui, item, opts.json_output)?;
if let Some(code) = petri_stream::exit_code_of(item) {
replayed_exit_code = Some(ExitCode::from(code));
replayed_exit_code.get_or_insert(ExitCode::from(code));
}
}
if let Some(exit_code) = replayed_exit_code.or_else(|| {
@ -1114,4 +1114,109 @@ mod tests {
assert_eq!(result.unwrap_err().to_string(), "questions unavailable");
assert_eq!(calls, 1);
}
#[tokio::test]
async fn required_publication_keeps_attach_open_until_the_committed_result() {
let run_id = fabro_types::fixtures::RUN_1;
let server = MockServer::start_async().await;
let mut state = terminal_run_state_response(run_id);
state["status"] = serde_json::json!({"kind": "running"});
state["conclusion"] = serde_json::Value::Null;
let state: server_client::RunProjection = serde_json::from_value(state).unwrap();
let executed = serde_json::json!({
"run_id": run_id, "stream_seq": 1, "kind": "petri", "id": "executed",
"recorded_at": 1000,
"item": {"origin": "external", "context": {}, "record": {"seq": 5,
"body": {"event": "invocation.finished", "invocation": 0,
"result": {"status": "success"}}}}
});
server
.mock_async(|when, then| {
when.method("GET")
.path(format!("/api/v1/runs/{run_id}/events"));
then.status(200).json_body(serde_json::json!({
"data": [executed], "meta": {"has_more": false}, "event_contract_version": 6
}));
})
.await;
server
.mock_async(|when, then| {
when.method("GET")
.path(format!("/api/v1/runs/{run_id}/state"));
then.status(200)
.json_body(serde_json::to_value(&state).unwrap());
})
.await;
server
.mock_async(|when, then| {
when.method("GET")
.path(format!("/api/v1/runs/{run_id}/questions"));
then.status(200)
.json_body(serde_json::json!({ "data": [], "meta": { "has_more": false } }));
})
.await;
let waiting = server
.mock_async(|when, then| {
when.method("GET")
.path(format!("/api/v1/runs/{run_id}/attach"));
then.status(200)
.header("Content-Type", "text/event-stream")
.body("");
})
.await;
let client = server_client::Client::new_no_proxy(&server.base_url()).unwrap();
let task = tokio::spawn(async move {
Box::pin(attach_petri_run_with_client(
&client,
&run_id,
&state,
no_color_styles(),
AttachOptions {
auto_approve: false,
verbose: false,
kill_on_detach: false,
json_output: true,
},
Printer::Default,
))
.await
.unwrap()
});
tokio::time::timeout(Duration::from_secs(5), async {
while waiting.calls_async().await == 0 {
sleep(Duration::from_millis(10)).await;
}
})
.await
.unwrap();
assert!(
!task.is_finished(),
"successful execution does not end attach while publication is pending"
);
let finished = serde_json::json!({
"run_id": run_id, "stream_seq": 2, "kind": "petri", "id": "finished",
"recorded_at": 2000,
"item": {"origin": "external", "context": {}, "record": {"seq": 6,
"body": {"event": "run.finished", "status": "failed",
"finalization_failure": {"code": "publish_failed", "message": "the push was rejected"}}}}
});
let final_stream = server
.mock_async(|when, then| {
when.method("GET")
.path(format!("/api/v1/runs/{run_id}/attach"));
then.status(200)
.header("Content-Type", "text/event-stream")
.body(format!("data: {finished}\n\n"));
})
.await;
waiting.delete_async().await;
assert_eq!(
tokio::time::timeout(Duration::from_secs(5), task)
.await
.unwrap()
.unwrap(),
ExitCode::from(1)
);
final_stream.assert_async().await;
}
}

View file

@ -523,7 +523,7 @@ pub(crate) fn format_pretty(
"run.finished" => {
let status = view.body()?.get("status").and_then(Value::as_str)?;
let duration = format_duration_ms(elapsed.unwrap_or(0));
Some(match status {
let verdict = match status {
"success" => format!(
"{ts} {} {}",
styles.bold_green.apply_to("\u{2713} SUCCEEDED"),
@ -539,6 +539,17 @@ pub(crate) fn format_pretty(
styles.bold_red.apply_to("\u{2717} FAILED"),
styles.bold.apply_to(&duration),
),
};
let failure = view
.body()?
.get("finalization_failure")
.and_then(|failure| failure.get("message"))
.and_then(Value::as_str);
Some(match failure {
Some(message) => {
format!("{verdict}\n{ts} {}", styles.bold_red.apply_to(message))
}
None => verdict,
})
}
_ => None,
@ -1276,6 +1287,32 @@ mod tests {
assert!(line.contains("\u{2717} FAILED"), "{line}");
}
#[test]
fn execution_success_does_not_finish_required_publication() {
let executed = petri(
8,
json!({ "origin": "external", "context": {},
"record": { "seq": 5, "body": { "event": "invocation.finished",
"invocation": 0, "result": { "status": "success" } } }
}),
);
assert_eq!(exit_code_of(&executed), None);
let finished = petri(
9,
json!({ "origin": "external", "context": {},
"record": { "seq": 6, "body": { "event": "run.finished", "status": "failed",
"finalization_failure": { "code": "publish_failed", "message": "the push was rejected" } } }
}),
);
assert_eq!(exit_code_of(&finished), Some(1));
let mut state = PrettyState::default();
let line = format_pretty(&finished, &Styles::new(false), &mut state).unwrap();
insta::assert_snapshot!(line, @"
04:43:08 ✗ FAILED 0ms
04:43:08 the push was rejected
");
}
#[test]
fn the_raw_line_is_the_envelope_as_json() {
let notice = platform(

View file

@ -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

View file

@ -13,9 +13,9 @@
//! long run's pushes into repeated failures. Re-minting near expiry keeps a
//! run of any length on a live token.
//!
//! Publication runs in Fabro's `run_finished` hook, after the last stage and
//! before the run's terminal record, as the legacy publish step did: the
//! final checkpoint is pushed from inside the sandbox to the run
//! Publication runs in Fabro's required `finalize_run` hook, after the last
//! stage and before the run's terminal record:
//! the final checkpoint is pushed from inside the sandbox to the run
//! branch on GitHub, and, when the run changed files and its settings ask
//! for one, a pull request is opened and recorded. A failure fails the run
//! with `publish_failed`.

View file

@ -1103,12 +1103,13 @@ fn attach_json_errors_without_prompting_for_human_input() {
"recorded_at": "[EPOCH_MS]",
"body": {
"event": "run.started",
"format_version": 8,
"format_version": 9,
"key": "[ULID]",
"root": 0,
"middleware_chain": [
"circuit-breaker"
]
],
"required_finalization": true
}
}
}

View file

@ -263,6 +263,12 @@ fn bare_fabro_with_unbound_inputs_in_template_partial_validates_structurally_wit
// The include error names the partial relative to the run's working
// directory, so the `../` run is as long as that directory is deep.
let mut filters = context.filters();
// A checkout outside the user home can be rendered through a relative
// path (including macOS's /tmp alias), rather than its canonical root.
filters.push((
r#"(?:\.\./)+[^"\s]*?/test/(templated_unbound_partial/)"#.to_string(),
"[UP][FIXTURES]/$1".to_string(),
));
filters.push((
r"(\.\./)*\.\.\[FIXTURES\]".to_string(),
"[UP][FIXTURES]".to_string(),
@ -333,6 +339,10 @@ fn validate_reports_missing_template_dependency() {
let mut cmd = context.validate();
cmd.arg(fixture("templates/missing_dependency/workflow.fabro"));
let mut filters = context.filters();
filters.push((
r"(?:\.\./)+[^`\s]*?/test/(templates/missing_dependency/)".to_string(),
"[FIXTURES]/$1".to_string(),
));
filters.push((
r"(?:\.\./)*\.\.\[FIXTURES\]/".to_string(),
"[FIXTURES]/".to_string(),

View file

@ -574,13 +574,25 @@ mod tests {
transition: RunLifecycleKind,
status: RunStatus,
) {
let response = append_lifecycle_as_worker(app, run_id, token, transition, status).await;
assert_eq!(response.status(), StatusCode::OK);
}
/// Append one lifecycle record as the run's worker, and return the
/// server's response whatever it is.
async fn append_lifecycle_as_worker(
app: &axum::Router,
run_id: RunId,
token: &str,
transition: RunLifecycleKind,
status: RunStatus,
) -> axum::response::Response {
let record =
PlatformRecord::RunLifecycle(RunLifecycleRecord::new(transition).with_status(status));
let body = json!({
"record": serde_json::to_value(&record).expect("the record encodes"),
});
let response = app
.clone()
app.clone()
.oneshot(
Request::builder()
.method("POST")
@ -591,8 +603,7 @@ mod tests {
.expect("the append request builds"),
)
.await
.expect("the append request completes");
assert_eq!(response.status(), StatusCode::OK);
.expect("the append request completes")
}
/// Petri's own finish of the run, stored the way its worker stores it:
@ -940,4 +951,308 @@ mod tests {
"the failure names the worker's exit"
);
}
async fn append_engine_records(
app: &axum::Router,
run_id: RunId,
token: &str,
owner: &str,
log: &fabro_petri::petri::LogId,
records: &[fabro_petri::petri::Record],
) -> axum::response::Response {
let body = json!({ "owner": owner, "records": records.iter().map(|record| json!({
"seq": record.seq, "recorded_at": record.recorded_at, "record": record.record,
})).collect::<Vec<_>>() });
app.clone()
.oneshot(
Request::builder()
.method("POST")
.uri(
format!("/api/v1/runs/{run_id}/petri/logs/{log}/records")
.replace(' ', "%20"),
)
.header(header::AUTHORIZATION, format!("Bearer {token}"))
.header(header::CONTENT_TYPE, "application/json")
.body(Body::from(body.to_string()))
.unwrap(),
)
.await
.unwrap()
}
async fn assert_public_result(
state: &AppState,
app: &axum::Router,
run_id: RunId,
expected: RunStatus,
message: Option<&str>,
) {
state.petri_projector.settle(run_id).await;
let mut values = Vec::new();
for suffix in ["", "/state"] {
let response = app
.clone()
.oneshot(
Request::builder()
.uri(format!("/api/v1/runs/{run_id}{suffix}"))
.body(Body::empty())
.unwrap(),
)
.await
.unwrap();
values.push(
fabro_test::expect_axum_json(
response,
StatusCode::OK,
"GET run and state at required finalization boundary",
)
.await,
);
}
assert_eq!(
values[0]["lifecycle"]["status"],
serde_json::to_value(expected).unwrap(),
"public state: {:#?}",
values[1]
);
assert_eq!(values[1]["status"], serde_json::to_value(expected).unwrap());
assert_eq!(state.test_managed_run_status(&run_id), Some(expected));
if expected.is_terminal() {
assert_eq!(
values[1]["conclusion"]["failure"]["detail"]["message"].as_str(),
message
);
assert!(
values[1]["conclusion"]["stages"]
.as_array()
.is_some_and(|stages| !stages.is_empty())
);
} else {
assert!(values[1]["conclusion"].is_null());
}
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn a_rejected_platform_failure_does_not_settle_the_run() {
let runtime = Arc::new(HeldWorkerRuntime::default());
let (state, app, run_id, token) = held_worker_run(&runtime).await;
run_to_running_as_worker(&app, run_id, &token).await;
let pool = state.stores.run_summaries.pool();
sqlx::query(
"CREATE TRIGGER reject_platform_record BEFORE INSERT ON platform_records \
BEGIN SELECT RAISE(FAIL, 'scripted append failure'); END",
)
.execute(&pool)
.await
.unwrap();
let response = append_lifecycle_as_worker(
&app,
run_id,
&token,
RunLifecycleKind::Failed,
RunStatus::Failed {
reason: FailureReason::LaunchFailed,
},
)
.await;
fabro_test::assert_axum_status(
response,
StatusCode::INTERNAL_SERVER_ERROR,
"rejected platform terminal append",
)
.await;
assert_public_result(&state, &app, run_id, RunStatus::Running, None).await;
crate::server::persist_run_failure(
&state,
run_id,
FailureReason::LaunchFailed,
"worker launch failed".into(),
)
.await;
assert_public_result(&state, &app, run_id, RunStatus::Running, None).await;
sqlx::query("DROP TRIGGER reject_platform_record")
.execute(&pool)
.await
.unwrap();
runtime.end_worker();
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn a_host_failure_is_retried_until_storage_accepts_it() {
let runtime = Arc::new(HeldWorkerRuntime::default());
let (state, app, run_id, token) = held_worker_run(&runtime).await;
run_to_running_as_worker(&app, run_id, &token).await;
let pool = state.stores.run_summaries.pool();
sqlx::query(
"CREATE TRIGGER reject_platform_record BEFORE INSERT ON platform_records \
BEGIN SELECT RAISE(FAIL, 'scripted append failure'); END",
)
.execute(&pool)
.await
.unwrap();
let failure = tokio::spawn({
let state = Arc::clone(&state);
async move {
crate::server::persist_run_failure(
&state,
run_id,
FailureReason::LaunchFailed,
"worker launch failed".into(),
)
.await;
}
});
tokio::time::sleep(std::time::Duration::from_millis(100)).await;
sqlx::query("DROP TRIGGER reject_platform_record")
.execute(&pool)
.await
.unwrap();
failure.await.unwrap();
let failed = RunStatus::Failed {
reason: FailureReason::LaunchFailed,
};
let stored = crate::server::run_records::projection(&state, run_id)
.await
.unwrap()
.expect("the run is stored");
assert_eq!(stored.status, failed, "the retried failure is durable");
assert_eq!(state.test_managed_run_status(&run_id), Some(failed));
runtime.end_worker();
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn required_publication_result_agrees_across_worker_api_projection_and_cleanup() {
use fabro_petri::petri::LogId;
use fabro_petri::test_support::finalization;
for rejection in [None, Some("the push was rejected")] {
let runtime = Arc::new(HeldWorkerRuntime::default());
let (state, app, run_id, token) = held_worker_run(&runtime).await;
run_to_running_as_worker(&app, run_id, &token).await;
let fixture = finalization::test_run_records(run_id, rejection).await;
for bytes in fixture.blobs {
let response = app
.clone()
.oneshot(
Request::builder()
.method("POST")
.uri(format!("/api/v1/runs/{run_id}/petri/blobs?owner=worker-1"))
.header(header::AUTHORIZATION, format!("Bearer {token}"))
.body(Body::from(bytes))
.unwrap(),
)
.await
.unwrap();
fabro_test::assert_axum_status(response, StatusCode::OK, "worker graph blob").await;
}
let mut terminal = None;
for (log, mut records) in fixture.logs {
if log == LogId::Coordinator {
let finish = records.pop().unwrap();
assert_eq!(finish.record["body"]["event"], "run.finished");
terminal = Some(finish);
}
let response =
append_engine_records(&app, run_id, &token, "worker-1", &log, &records).await;
fabro_test::assert_axum_status(
response,
StatusCode::NO_CONTENT,
"worker execution records",
)
.await;
}
assert_public_result(&state, &app, run_id, RunStatus::Running, None).await;
let terminal = terminal.unwrap();
// A rejected writer cannot settle the managed run or store its
// result. This exercises the same endpoint as the valid finish.
let response = append_engine_records(
&app,
run_id,
&token,
"superseded-worker",
&LogId::Coordinator,
std::slice::from_ref(&terminal),
)
.await;
fabro_test::assert_axum_status(
response,
StatusCode::CONFLICT,
"finish from a stale lease owner",
)
.await;
assert_public_result(&state, &app, run_id, RunStatus::Running, None).await;
let response =
append_engine_records(&app, run_id, &token, "worker-1", &LogId::Coordinator, &[
terminal,
])
.await;
fabro_test::assert_axum_status(
response,
StatusCode::NO_CONTENT,
"authoritative worker finish",
)
.await;
let expected = match rejection {
Some(_) => RunStatus::Failed {
reason: FailureReason::PublishFailed,
},
None => RunStatus::Succeeded {
reason: SuccessReason::Completed,
},
};
assert_public_result(&state, &app, run_id, expected, rejection).await;
// Worker exit after the finish, without a platform terminal
// acknowledgement, must preserve Petri's committed result.
runtime.end_worker();
for _ in 0..500 {
if concluded_runs(&app).await == 1 {
break;
}
time::sleep(Duration::from_millis(10)).await;
}
assert_eq!(concluded_runs(&app).await, 1);
assert_public_result(&state, &app, run_id, expected, rejection).await;
// A host error arriving after cleanup must neither append a
// competing terminal record nor replace the useful failure.
let platform_count = fabro_petri::test_support::stored_platform_records(
&state.stores.run_summaries.pool(),
run_id,
)
.await
.unwrap()
.len();
crate::server::persist_run_failure(
&state,
run_id,
FailureReason::Terminated,
"worker wait failed during teardown".into(),
)
.await;
assert_public_result(&state, &app, run_id, expected, rejection).await;
let after_count = fabro_petri::test_support::stored_platform_records(
&state.stores.run_summaries.pool(),
run_id,
)
.await
.unwrap()
.len();
assert_eq!(after_count, platform_count);
let (rebuilt, _, _) = fabro_petri::test_support::rebuild(
&state.db_pool,
&state.stores.run_summaries.pool(),
run_id,
)
.await
.unwrap();
let rebuilt = rebuilt.unwrap();
assert_eq!(rebuilt.status, expected);
assert_eq!(
rebuilt
.conclusion
.unwrap()
.failure
.map(|failure| failure.detail.message),
rejection.map(str::to_string)
);
}
}
}

View file

@ -282,6 +282,23 @@ struct ManagedRun {
}
impl ManagedRun {
/// Settle the run on a terminal `status`. The first terminal status
/// sticks: a later one that agrees only fills a missing error, and one
/// that disagrees is ignored. Returns whether the status was applied.
fn settle(&mut self, status: RunStatus, error: Option<String>) -> bool {
if self.status.is_terminal() {
if self.status == status && self.error.is_none() {
self.error = error;
}
return false;
}
self.status = status;
self.error = error;
self.active_steerable_stages.clear();
self.active_non_steerable_stages.clear();
true
}
/// True if cancellation should still escalate to `worker_ref`; clears a
/// stale escalation marker as a side effect.
fn escalation_still_current(&mut self, worker_ref: &WorkerRef) -> bool {
@ -3505,24 +3522,90 @@ async fn reject_run_if_sandbox_provider_disabled(
return false;
};
tracing::warn!(run_id = %run_id, error = %error, "Sandbox provider disabled by server policy");
fail_run_before_execution(state, run_id, FailureReason::LaunchFailed, error).await;
persist_run_failure(state, run_id, FailureReason::LaunchFailed, error).await;
true
}
async fn fail_run_before_execution(
/// Record a host failure only while no terminal result is committed. This
/// also handles a worker wait/launch error racing its durable Petri finish.
/// The managed run settles on whichever terminal result the store committed.
/// A failed commit is retried, so a brief storage fault does not leave the
/// run active. When nothing could be committed its live state is still
/// released, but its status is left alone: the API never reports an outcome
/// storage lacks, and the restart reconciliation fails the run.
pub(crate) async fn persist_run_failure(
state: &Arc<AppState>,
run_id: RunId,
reason: FailureReason,
message: String,
) {
if let Err(err) =
run_records::lifecycle(state, run_id, run_records::failed(reason, message.clone())).await
{
error!(run_id = %run_id, error = %err, "Failed to persist run failure status");
let mut retry_delays = HOST_FAILURE_RETRY_DELAYS.iter();
loop {
match commit_host_failure(state, run_id, reason, message.clone()).await {
Ok(committed) => {
let failure = committed
.conclusion
.as_ref()
.and_then(|conclusion| conclusion.failure.as_ref())
.map(|failure| failure.detail.message.clone());
settle_managed_run_at_finish(state, run_id, committed.status, failure);
break;
}
Err(err) => {
let Some(delay) = retry_delays.next() else {
error!(run_id = %run_id, error = %err, "Failed to record a host failure");
break;
};
warn!(
run_id = %run_id,
error = %err,
retry_in_ms = delay.as_millis(),
"Failed to record a host failure; retrying"
);
sleep(*delay).await;
}
}
}
release_managed_run(state, run_id);
}
fail_managed_run(state, run_id, reason, message);
state.scheduler_notify.notify_one();
/// How long [`persist_run_failure`] waits before each retry of a failed
/// commit.
const HOST_FAILURE_RETRY_DELAYS: [Duration; 3] = [
Duration::from_millis(250),
Duration::from_secs(1),
Duration::from_secs(4),
];
/// Append the host failure unless the run already ended, and return the
/// terminal projection the store committed: the failure, or the finish that
/// won the race.
async fn commit_host_failure(
state: &AppState,
run_id: RunId,
reason: FailureReason,
message: String,
) -> anyhow::Result<Arc<fabro_store::RunProjection>> {
let committed = run_records::projection(state, run_id)
.await?
.context("the run is missing")?;
if committed.status.is_terminal() {
return Ok(committed);
}
// The append settles the projector, so the read after it folds the
// failure, or the finish that won the race.
run_records::lifecycle(state, run_id, run_records::failed(reason, message)).await?;
let committed = state
.stores
.run_summaries
.load_petri_projection(&run_id)
.await?
.context("the run is missing")?;
anyhow::ensure!(
committed.status.is_terminal(),
"the stored host failure has no terminal projection"
);
Ok(committed)
}
fn managed_run(
@ -3582,11 +3665,22 @@ async fn durable_run_status(state: &AppState, run_id: RunId) -> anyhow::Result<O
fn fail_managed_run(state: &Arc<AppState>, run_id: RunId, reason: FailureReason, message: String) {
let mut runs = state.runs.lock().expect("runs lock poisoned");
if let Some(managed_run) = runs.get_mut(&run_id) {
managed_run.status = RunStatus::Failed { reason };
managed_run.error = Some(message);
managed_run.settle(RunStatus::Failed { reason }, Some(message));
}
drop(runs);
release_managed_run(state, run_id);
}
/// Drop the run's live worker state and controls and free its scheduler
/// slot, leaving its status alone.
pub(in crate::server) fn release_managed_run(state: &AppState, run_id: RunId) {
let mut runs = state.runs.lock().expect("runs lock poisoned");
if let Some(managed_run) = runs.get_mut(&run_id) {
clear_live_run_state(managed_run);
}
cleanup_worker_control_bus_for_run(state.as_ref(), run_id);
drop(runs);
cleanup_worker_control_bus_for_run(state, run_id);
state.scheduler_notify.notify_one();
}
/// Fold one lifecycle record of the run's stream into the in-memory run:
@ -3601,6 +3695,7 @@ fn apply_lifecycle_to_managed_run(state: &AppState, run_id: RunId, record: &RunL
// A settled run is immutable to the lifecycle: the follower still folds
// the records before the terminal one after Petri's finish or the
// worker's terminal record settled the run, and none may reopen it.
// Terminal records go through [`ManagedRun::settle`].
if managed_run.status.is_terminal() && !is_terminal_transition(record) {
return;
}
@ -3648,21 +3743,17 @@ fn apply_lifecycle_to_managed_run(state: &AppState, run_id: RunId, record: &RunL
}
RunLifecycleKind::Removing => managed_run.status = RunStatus::Removing,
RunLifecycleKind::Succeeded => {
managed_run.status = record.status.unwrap_or(RunStatus::Succeeded {
let status = record.status.unwrap_or(RunStatus::Succeeded {
reason: SuccessReason::Completed,
});
managed_run.error = None;
managed_run.active_steerable_stages.clear();
managed_run.active_non_steerable_stages.clear();
managed_run.settle(status, None);
cleanup_worker_control_bus_for_run(state, run_id);
}
RunLifecycleKind::Failed | RunLifecycleKind::Dead => {
managed_run.status = record.status.unwrap_or(RunStatus::Failed {
let status = record.status.unwrap_or(RunStatus::Failed {
reason: FailureReason::WorkflowError,
});
managed_run.error.clone_from(&record.reason);
managed_run.active_steerable_stages.clear();
managed_run.active_non_steerable_stages.clear();
managed_run.settle(status, record.reason.clone());
cleanup_worker_control_bus_for_run(state, run_id);
}
RunLifecycleKind::Runnable
@ -3683,34 +3774,23 @@ fn is_terminal_transition(record: &RunLifecycleRecord) -> bool {
)
}
/// Settle the in-memory run at Petri's own finish, as its worker stores
/// the `run.finished` record: the view reports the run ended from the
/// moment that record is stored, so the managed run the delete precheck
/// prefers must not still say running while the worker tears down; a
/// delete in that window was refused as active. A run already settled
/// keeps its status. The worker's terminal lifecycle record, a moment
/// later, refines the status and its error and ends the worker's controls
/// ([`settle_managed_run_at_terminal_record`]); the worker's exit later
/// reaps the process and leaves the settled status alone.
/// Settle the in-memory run after Petri's authoritative finish is durable.
/// Required publication has completed before this record. Worker teardown
/// and subsequent lifecycle records cannot change its terminal outcome.
pub(in crate::server) fn settle_managed_run_at_finish(
state: &AppState,
run_id: RunId,
status: RunStatus,
failure: Option<String>,
) {
let mut runs = state.runs.lock().expect("runs lock poisoned");
let Some(managed_run) = runs.get_mut(&run_id) else {
return;
};
if managed_run.status.is_terminal() {
return;
if let Some(managed_run) = runs.get_mut(&run_id) {
managed_run.settle(status, failure);
}
managed_run.status = status;
managed_run.active_steerable_stages.clear();
managed_run.active_non_steerable_stages.clear();
}
/// Settle the in-memory run at the terminal lifecycle record its worker
/// stores, ahead of the store: the same as [`settle_managed_run_at_finish`]
/// stores, after the store: the same as [`settle_managed_run_at_finish`]
/// for a worker that ended the run without Petri's finish (it failed before
/// the engine ran), and the record's status, error and control cleanup for
/// one that did. A record that is not terminal is left to the stream
@ -3775,7 +3855,6 @@ async fn fail_worker_launch(state: &Arc<AppState>, run_id: RunId, err: anyhow::E
None
}
};
let launch_message = format!("Failed to spawn worker: {err}");
let (error, reason) = failure_honoring_pending_cancel(pending_control, || {
(
WorkflowError::engine_with_anyhow("Failed to spawn worker", err),
@ -3785,16 +3864,9 @@ async fn fail_worker_launch(state: &Arc<AppState>, run_id: RunId, err: anyhow::E
let message = if reason == FailureReason::Cancelled {
"Run cancelled before worker launch completed".to_string()
} else {
launch_message
collect_chain(&error).join(": ")
};
let _ = run_records::lifecycle(
state,
run_id,
run_records::failed(reason, error.to_string()),
)
.await;
fail_managed_run(state, run_id, reason, message);
state.scheduler_notify.notify_one();
persist_run_failure(state, run_id, reason, message).await;
}
/// A worker that exited without recording the run's end left it failed,
@ -4163,7 +4235,6 @@ async fn execute_run_subprocess(state: Arc<AppState>, run_id: RunId) {
FailureReason::WorkflowError,
"Run not found at launch".to_string(),
);
state.scheduler_notify.notify_one();
return;
}
Err(err) => {
@ -4174,7 +4245,6 @@ async fn execute_run_subprocess(state: Arc<AppState>, run_id: RunId) {
FailureReason::WorkflowError,
format!("Failed to load run state: {err}"),
);
state.scheduler_notify.notify_one();
return;
}
};
@ -4195,7 +4265,7 @@ async fn execute_run_subprocess(state: Arc<AppState>, run_id: RunId) {
let github_app_private_key = match state.vault_secret(EnvVars::GITHUB_APP_PRIVATE_KEY).await {
Ok(value) => value,
Err(err) => {
fail_run_before_execution(
persist_run_failure(
&state,
run_id,
FailureReason::WorkflowError,
@ -4255,15 +4325,17 @@ async fn execute_run_subprocess(state: Arc<AppState>, run_id: RunId) {
Err(err) => {
tracing::error!(run_id = %run_id, error = %err, "Failed while waiting on worker");
let message = format!("Worker wait failed: {err}");
let superseded = {
let runs = state.runs.lock().expect("runs lock poisoned");
runs.get(&run_id)
.is_some_and(|run| run.worker_ref.as_ref() != Some(&worker_ref))
};
if superseded {
return;
}
state.worker_runtime.force_stop(&worker_ref).await;
let _ = run_records::lifecycle(
&state,
run_id,
run_records::failed(FailureReason::Terminated, message.clone()),
)
.await;
fail_managed_run(&state, run_id, FailureReason::Terminated, message);
state.scheduler_notify.notify_one();
state.petri_runs.worker_exited(run_id);
persist_run_failure(&state, run_id, FailureReason::Terminated, message).await;
return;
}
};
@ -4321,7 +4393,6 @@ async fn execute_run_subprocess(state: Arc<AppState>, run_id: RunId) {
FailureReason::WorkflowError,
"The run's final state is missing from the store".to_string(),
);
state.scheduler_notify.notify_one();
return;
}
Err(err) => {
@ -4332,7 +4403,6 @@ async fn execute_run_subprocess(state: Arc<AppState>, run_id: RunId) {
FailureReason::WorkflowError,
format!("Failed to load final run state: {err}"),
);
state.scheduler_notify.notify_one();
return;
}
};

View file

@ -23,7 +23,7 @@ use fabro_api::types::{
PetriReleaseRequest, WriteBlobResponse,
};
use fabro_petri::petri::{Access, Digest, LogId, OwnerId, Record, StoreError};
use fabro_petri::projection::finished_run_status;
use fabro_petri::projection;
use fabro_petri::run_store::{log_id_text, parse_log_id};
use fabro_store::{PlatformRecord, PlatformRecordKind, StagePosition, StoredPlatformRecord};
use fabro_types::BlobHash;
@ -157,16 +157,15 @@ async fn append_records(
Ok(writer) => writer,
Err(err) => return store_error_response(id, &err),
};
// The view ends the run at Petri's own finish, the moment the record
// is stored and a pass folds it: the managed run settles first, so a
// delete that lands while the worker still tears down is not refused.
if log == LogId::Coordinator {
if let Some(status) = records.iter().find_map(finished_run_status) {
settle_managed_run_at_finish(&state, id, status);
}
}
match writer.append(&log, &records).await {
Ok(()) => {
if log == LogId::Coordinator {
if let Some((status, failure)) =
records.iter().find_map(projection::finished_run_result)
{
settle_managed_run_at_finish(&state, id, status, failure);
}
}
// The records are durable; the projection trails them from here.
state.petri_projector.signal(id);
StatusCode::NO_CONTENT.into_response()
@ -279,10 +278,6 @@ async fn append_platform_record(
(Some(execution), Some(firing)) => Some(StagePosition { execution, firing }),
_ => None,
};
// A terminal lifecycle record ends the run in the view as Petri's
// finish does, for a worker that ended the run without one: the
// managed run settles before the record is stored.
settle_managed_run_at_terminal_record(&state, id, &record);
let summaries = &state.stores.run_summaries;
match summaries
.platform_records()
@ -290,6 +285,7 @@ async fn append_platform_record(
.await
{
Ok(stored) => {
settle_managed_run_at_terminal_record(&state, id, &record);
summaries.notify_platform_record(id);
match wire_platform_record(&stored) {
Ok(record) => Json(record).into_response(),

View file

@ -22,7 +22,7 @@
//! of the server's vault, and its blobs go to the server's blob store. Its
//! managed run settles at Petri's own finish, as a worker's does at the
//! worker's records endpoint: the run store it executes over settles the
//! run before the `run.finished` record is stored ([`SettlingStore`]). No
//! run after the `run.finished` record is stored ([`SettlingStore`]). No
//! stage or agent event is projected either way, which is the read-side
//! item that follows.
//!
@ -496,8 +496,8 @@ pub(crate) async fn execute(state: Arc<AppState>, run_id: RunId) {
run_id: run_id.to_string(),
run_dir: run_dir.join("petri"),
execution,
// The projector's signal follows each durable append; the managed
// run's settle at Petri's finish precedes it.
// The coordinator finish is stored before managed status settles;
// the projector also reads only durable records.
store: Arc::new(SettlingStore {
inner: state
.petri_projector
@ -549,13 +549,7 @@ pub(crate) async fn execute(state: Arc<AppState>, run_id: RunId) {
failed(reason, message)
}
};
if let Err(err) = run_records::lifecycle(&state, run_id, record).await {
error!(run_id = %run_id, error = %err, "Failed to persist run outcome");
}
// The managed run settled at Petri's finish, ahead of the store; the
// terminal record refines its status and error and ends its live
// state, and is the settle of a run that ended without a finish.
finish(&state, run_id, status, error);
commit_and_finish(&state, run_id, record, status, error).await;
// The view trails the terminal record; the aggregate reads the settled
// projection, as the worker path reads the final state at worker exit.
state.petri_projector.settle(run_id).await;
@ -814,38 +808,43 @@ fn failed(
async fn fail_before_execution(state: &Arc<AppState>, run_id: RunId, message: &str) {
error!(run_id = %run_id, error = message, "Petri run cannot start");
let (status, error, record) = failed(FailureReason::WorkflowError, message.to_string());
if let Err(err) = run_records::lifecycle(state, run_id, record).await {
error!(run_id = %run_id, error = %err, "Failed to persist run failure status");
commit_and_finish(state, run_id, record, status, error).await;
}
/// Append the run's terminal record, then finish the run. An append that
/// fails only releases the live state: it commits no terminal result.
async fn commit_and_finish(
state: &Arc<AppState>,
run_id: RunId,
record: RunLifecycleRecord,
status: RunStatus,
error: Option<String>,
) {
match run_records::lifecycle(state, run_id, record).await {
Ok(_) => finish(state, run_id, status, error),
Err(err) => {
error!(run_id = %run_id, error = %err, "Failed to persist run outcome");
super::release_managed_run(state, run_id);
}
}
finish(state, run_id, status, error);
}
/// Settle the managed run at its terminal record and release its
/// scheduler slot. A run that Petri finished settled already, at the
/// `run.finished` record ([`SettlingStore`]); this refines its status and
/// error and ends its live state. A run deleted since is gone from the map
/// and stays gone.
/// `run.finished` record ([`SettlingStore`]); this preserves that status,
/// fills any missing failure detail and ends its live state. A run deleted
/// since is gone from the map and stays gone.
fn finish(state: &Arc<AppState>, run_id: RunId, status: RunStatus, error: Option<String>) {
let mut runs = state.runs.lock().expect("runs lock poisoned");
if let Some(managed_run) = runs.get_mut(&run_id) {
managed_run.status = status;
managed_run.error = error;
clear_live_run_state(managed_run);
managed_run.settle(status, error);
}
drop(runs);
state.scheduler_notify.notify_one();
super::release_managed_run(state, run_id);
}
/// The run store an in-process run executes over: the projector's
/// signalling store, whose coordinator appends settle the managed run at
/// Petri's own finish first. The view ends the run at the `run.finished`
/// record the moment it is stored and a pass folds it, so the managed run
/// the delete precheck prefers must not still say running while the engine
/// tears down: a delete in that window was refused as active. The worker's
/// records endpoint does the same for a worker-backed run, ahead of the
/// same store. The settle is in memory only; the terminal lifecycle record
/// [`execute`] stores once the engine returns refines the status
/// ([`finish`]), and stays the settle of a run that ends without a finish.
/// The run store an in-process run executes over. A coordinator finish
/// settles managed status only after the append succeeds, as on HTTP workers.
struct SettlingStore {
inner: Arc<dyn RunStore>,
state: Arc<AppState>,
@ -864,7 +863,7 @@ impl RunStore for SettlingStore {
}
/// One run's logs, whose coordinator appends settle the managed run at
/// Petri's finish before the records reach the store.
/// Petri's finish after the records reach the store.
struct SettlingLogs {
inner: Arc<dyn RunLogs>,
run_id: Option<RunId>,
@ -878,12 +877,15 @@ impl RunLogs for SettlingLogs {
}
async fn append(&self, log: &LogId, records: &[Record]) -> Result<(), StoreError> {
self.inner.append(log, records).await?;
if let Some(run_id) = self.run_id.filter(|_| *log == LogId::Coordinator) {
if let Some(status) = records.iter().find_map(projection::finished_run_status) {
super::settle_managed_run_at_finish(&self.state, run_id, status);
if let Some((status, failure)) =
records.iter().find_map(projection::finished_run_result)
{
super::settle_managed_run_at_finish(&self.state, run_id, status, failure);
}
}
self.inner.append(log, records).await
Ok(())
}
async fn read(&self, log: &LogId) -> Result<Vec<Record>, StoreError> {
@ -922,6 +924,7 @@ mod tests {
state: Arc<AppState>,
run_id: RunId,
seen: Mutex<Vec<Option<RunStatus>>>,
reject: std::sync::atomic::AtomicBool,
}
#[async_trait::async_trait]
@ -933,6 +936,9 @@ mod tests {
async fn append(&self, log: &LogId, records: &[Record]) -> Result<(), StoreError> {
let status = self.state.test_managed_run_status(&self.run_id);
self.seen.lock().expect("seen lock poisoned").push(status);
if self.reject.load(std::sync::atomic::Ordering::SeqCst) {
return Err(StoreError::StaleOwner);
}
self.inner.append(log, records).await
}
@ -992,6 +998,7 @@ mod tests {
state: Arc::clone(&state),
run_id,
seen: Mutex::new(Vec::new()),
reject: std::sync::atomic::AtomicBool::new(false),
});
let logs = SettlingLogs {
inner: Arc::clone(&recorder) as Arc<dyn RunLogs>,
@ -1001,12 +1008,10 @@ mod tests {
(state, run_id, logs, recorder)
}
/// The in-process run settles at Petri's own finish, before the
/// `run.finished` record reaches the store: the view cannot report the
/// run ended while the managed run still says running. The records
/// before the finish leave the run in flight.
/// The in-process run settles only after its authoritative finish is
/// stored. Earlier records leave the run in flight.
#[tokio::test]
async fn an_in_process_run_settles_before_its_finish_is_stored() {
async fn an_in_process_run_settles_after_its_finish_is_stored() {
let (state, run_id, logs, recorder) = in_flight_run().await;
logs.append(&LogId::Coordinator, &[coordinator_record(
@ -1032,12 +1037,33 @@ mod tests {
};
assert_eq!(
*recorder.seen.lock().expect("seen lock poisoned"),
vec![Some(RunStatus::Running), Some(succeeded)],
"the managed run settled before the finish reached the store"
vec![Some(RunStatus::Running), Some(RunStatus::Running)],
"the managed run stayed running until the finish reached the store"
);
assert_eq!(state.test_managed_run_status(&run_id), Some(succeeded));
}
#[tokio::test]
async fn a_rejected_in_process_finish_does_not_settle_the_run() {
let (state, run_id, logs, recorder) = in_flight_run().await;
recorder
.reject
.store(true, std::sync::atomic::Ordering::SeqCst);
let result = logs
.append(&LogId::Coordinator, &[coordinator_record(
0,
&json!({"event": "run.finished", "status": "failed",
"finalization_failure": {"code": "publish_failed", "message": "push rejected"}}),
)])
.await;
assert!(matches!(result, Err(StoreError::StaleOwner)));
assert_eq!(
state.test_managed_run_status(&run_id),
Some(RunStatus::Running)
);
assert!(logs.read(&LogId::Coordinator).await.unwrap().is_empty());
}
/// A finish on another log than the coordinator's is not Petri's
/// finish of the run: an execution's engine log ends an execution.
#[tokio::test]

View file

@ -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

View file

@ -84,6 +84,7 @@ use crate::admission::AdmittedGraphs;
use crate::blobs::{Blobs, RunBlobs};
use crate::controls::RunControls;
use crate::hooks::{FabroHooks, HooksSpec};
use crate::projection;
use crate::runtime::RuntimeSpec;
use crate::secrets::SharedSecrets;
@ -150,7 +151,7 @@ pub enum RunStatus {
#[derive(Clone, Debug)]
pub struct RunOutcome {
pub status: RunStatus,
/// The root invocation's failure message, when it failed.
/// Required-finalization failure detail, or the root execution's failure.
pub failure: Option<String>,
/// Whether the record is whole: the run recorded its finish and every
/// log replays byte for byte.
@ -353,27 +354,7 @@ pub async fn run(request: RunRequest) -> Result<RunOutcome, RunError> {
return Err(RunError::StoreFailed(message));
}
let inspection = inspect(request.store.as_ref(), &key).await?;
let mut outcome = outcome(inspection, result.err())?;
// A failed checkpoint cancelled the run; what Fabro reports is the
// checkpoint failure, not a cancellation.
if let Some(failure) = fabro_hooks
.as_ref()
.and_then(|hooks| hooks.checkpoint_failure())
{
outcome.status = RunStatus::Failed;
outcome.failure = Some(failure);
}
// A successful run whose publication failed is a failed run: its work
// did not reach where the settings sent it.
if let Some(failure) = fabro_hooks
.as_ref()
.and_then(|hooks| hooks.publish_failure())
{
outcome.status = RunStatus::Failed;
outcome.failure = Some(failure);
outcome.publish_failed = true;
}
Ok(outcome)
outcome(inspection, result.err())
}
/// When Petri keeps a run's workspaces after their scope is released.
@ -419,8 +400,9 @@ pub async fn outcome_of(store: &dyn RunStore, run_id: &str) -> Result<RunOutcome
/// How Fabro reports what [`run`] returned. A cancelled run is a failure
/// with the cancelled reason, as the legacy executor reports one; a failed
/// store interrupted the run; every other shortfall is a workflow error
/// whose message says what the record, or the host, said.
/// store interrupted the run; a publication rejection keeps its
/// publish_failed reason. Other shortfalls are workflow errors whose message
/// says what the record, or the host, said.
#[must_use]
pub fn conclusion(result: &Result<RunOutcome, RunError>) -> Conclusion {
match result {
@ -543,31 +525,47 @@ fn outcome(
inspection: RunInspection,
host_error: Option<HostError>,
) -> Result<RunOutcome, RunError> {
let status = match inspection.status.as_deref() {
Some("success") => RunStatus::Success,
Some("failed") => RunStatus::Failed,
Some("cancelled") => RunStatus::Cancelled,
_ => {
let mut reasons = inspection.incomplete.clone();
if let Some(error) = host_error {
reasons.push(error.to_string());
}
return Err(RunError::Unfinished(reasons));
let Some(recorded) = inspection
.status
.as_deref()
.filter(|status| matches!(*status, "success" | "failed" | "cancelled"))
else {
let mut reasons = inspection.incomplete.clone();
if let Some(error) = host_error {
reasons.push(error.to_string());
}
return Err(RunError::Unfinished(reasons));
};
let failure = inspection
// The projection's reading of the finish, so the worker's return and
// the API agree: a failed checkpoint's cancellation is a failure.
let (status, publish_failed) =
match projection::finished_status(recorded, inspection.finalization_failure.as_ref()) {
fabro_types::RunStatus::Succeeded { .. } => (RunStatus::Success, false),
fabro_types::RunStatus::Failed {
reason: FailureReason::Cancelled,
} => (RunStatus::Cancelled, false),
fabro_types::RunStatus::Failed { reason } => {
(RunStatus::Failed, reason == FailureReason::PublishFailed)
}
_ => (RunStatus::Failed, false),
};
let execution_failure = inspection
.invocations
.iter()
.find(|invocation| invocation.invocation == inspection.root.invocation)
.and_then(|root| root.result.as_ref())
.and_then(|result| result.failure.as_ref())
.map(|failure| failure.message.clone());
let failure = inspection
.finalization_failure
.map(|failure| failure.message)
.or(execution_failure);
Ok(RunOutcome {
status,
failure,
complete: inspection.complete,
incomplete: inspection.incomplete,
publish_failed: false,
publish_failed,
})
}

View file

@ -39,11 +39,12 @@ use fabro_workflow::operations::{StageLabel, StageLabels};
use petri_execution::host::{self, ForkOptions, ForkOrigin, ForkPosition, HostError};
use petri_execution::inspect::{self, InspectError};
use petri_execution::{
Access, CoordinatorEvent, ExecutionId, InvocationId, RunKey, RunStore,
Access, CoordinatorEvent, ExecutionId, InvocationId, RunKey, RunLogs, RunStore,
StoreError as CoordinatorStoreError,
};
use petri_runtime::RunOptions;
use petri_runtime::ir::FiringId;
use petri_runtime::driver::lifecycle::{ExecutionHooks, HookContext, RunFinished};
use petri_runtime::ir::{FinalizationFailure, FiringId};
use petri_store::StoreError;
use tracing::{debug, info};
@ -131,7 +132,12 @@ pub async fn check(
.open(&RunKey::new(source.to_string()), Access::Read)
.await
.map_err(ForkError::Open)?;
let state = host::stored_state(&*logs).await.map_err(ForkError::Seed)?;
check_logs(&*logs, &position).await
}
/// [`check`] over the source's opened logs.
async fn check_logs(logs: &dyn RunLogs, position: &ForkPosition) -> Result<(), ForkError> {
let state = host::stored_state(logs).await.map_err(ForkError::Seed)?;
let Some(execution) = state.executions.get(&position.execution) else {
return Err(ForkError::Refused(format!(
"the source run has no execution {}",
@ -146,7 +152,7 @@ pub async fn check(
position.execution
)));
}
let inspection = inspect::inspect_run(&*logs)
let inspection = inspect::inspect_run(logs)
.await
.map_err(ForkError::Inspect)?;
if inspection
@ -169,10 +175,33 @@ pub async fn check(
Ok(())
}
/// Declaration-only hooks used while copying a fork's records: every Fabro
/// run requires finalization ([`crate::hooks::FabroHooks`]), so the fork's
/// records declare it too, and its worker's hooks match them on resume.
/// Never execute with these hooks.
struct ForkFinalizationRequirement;
#[async_trait::async_trait]
impl ExecutionHooks for ForkFinalizationRequirement {
fn requires_run_finalization(&self) -> bool {
true
}
async fn finalize_run(
&self,
_context: &HookContext,
_finished: RunFinished,
) -> Result<(), FinalizationFailure> {
Err(FinalizationFailure::new(
"finalization_unavailable",
"the fork must install its worker publication hooks before execution",
))
}
}
/// Seed the fork: Petri's records, the kept checkpoints and the run branch. The
/// new run must not exist in the store yet.
pub async fn fork(request: ForkRequest) -> Result<Forked, ForkError> {
check(request.store.as_ref(), request.source, request.position).await?;
let source_key = RunKey::new(request.source.to_string());
let fork_key = RunKey::new(request.fork.to_string());
let source_logs = request
@ -181,12 +210,14 @@ pub async fn fork(request: ForkRequest) -> Result<Forked, ForkError> {
.await
.map_err(ForkError::Open)?;
check_logs(&*source_logs, &request.position).await?;
let mut options = RunOptions::new(&request.fork_run_dir);
options.run_key = Some(fork_key.clone());
// A fork only copies records and acquires no sandbox, so it needs no
// provider configuration.
let runtime = providers::standard_runtime(&SandboxProviderConfig::default())
.options(options)
.hooks(Arc::new(ForkFinalizationRequirement))
.store(Arc::clone(&request.store));
let forked = host::fork_from(&runtime, &*source_logs, request.position, ForkOptions {
rerun_last: request.rerun_last,

View file

@ -31,12 +31,14 @@
//! `artifact.collected` record, unless the same file with the same content
//! was already collected earlier in the run. A failed write is a recorded
//! problem on the transition, never a blocked route.
//! - `run_finished`: the run's diff, its run branch against its base commit, as
//! the `run.diff` platform record with the patch as a blob; for a successful
//! run, its publication ([`RunPublisher`]: the platform pushes the run branch
//! and opens a pull request), whose failure fails the run before its terminal
//! record; then the forwarded point, so the local service runs `run_complete`
//! and `run_failed` with the sandbox in place.
//! - `finalize_run`, required for every run: the run's diff, its run branch
//! against its base commit, as the `run.diff` platform record with the patch
//! as a blob; a failed checkpoint, which fails the run and skips publication;
//! for a successful run, its publication ([`RunPublisher`]: the platform
//! pushes the run branch and opens a pull request), whose failure fails the
//! run before its terminal record.
//! - `run_finished`: the forwarded point, so the local service runs
//! `run_complete` and `run_failed` with the sandbox in place.
//! - `scope_acquired`: a fresh run's Git target checked out into the workspace
//! from inside the scope ([`crate::source`]); a resumed run uses its
//! surviving workspace, while an explicit fork fetches the source run's
@ -51,16 +53,18 @@
//!
//! # Operation identities
//!
//! Every external effect here is keyed on `(run key, execution, DecisionId,
//! effect kind)` from the hook context and deduplicated on retry: the
//! checkpoint's key is the attempt's decision in its execution, effect
//! Checkpoint and artifact effects are keyed on `(run key, execution,
//! DecisionId, effect kind)` from the hook context and deduplicated on retry:
//! the checkpoint's key is the attempt's decision in its execution, effect
//! `checkpoint`; an artifact's is the same decision, effect `artifact`, with
//! the file's path and content digest as the identity within it. A
//! re-dispatched attempt whose commit already landed reuses it when the
//! workspace still sits on it unchanged (see [`RunWorkspaces::commit`]); a
//! reissued routing decision finds the record, or the commit by its
//! trailers, and writes nothing twice; a file already collected under the
//! same path and digest is not collected again.
//! same path and digest is not collected again. Publication keeps the
//! publisher's reconciliation policy; required finalization adds no independent
//! effect ledger or guarantee of deduplication across every external crash.
//!
//! # Where the workspace is
//!
@ -100,7 +104,9 @@ use petri_runtime::driver::lifecycle::{
ScopeReleased, Transition, TransitionError, TransitionReport,
};
use petri_runtime::executor::{EnvError, ExecEnv};
use petri_runtime::ir::{ExecutionId, FailureInfo, RunStatus, ScopeId, Status};
use petri_runtime::ir::{
ExecutionId, FailureInfo, FinalizationFailure, RunStatus, ScopeId, Status,
};
use serde_json::json;
use tokio::sync::{Mutex as AsyncMutex, OnceCell};
use tokio::{fs, time};
@ -114,6 +120,7 @@ use crate::checkpoint::{
};
use crate::fork::{self, ForkError};
use crate::platform_records::{PlatformRecordError, PlatformRecords};
use crate::projection;
use crate::recovery::{self, Plan, RecoveryError, RestoreTarget};
use crate::source::RunSource;
use crate::workspace::{self, WorkspaceLookup, WorkspaceLookupError};
@ -453,8 +460,6 @@ pub struct FabroHooks {
/// The checkpoint failure that ended the run, when one did.
failure: Mutex<Option<String>>,
publisher: Option<Arc<dyn RunPublisher>>,
/// Why the run's publication failed, when it did.
publish_failure: Mutex<Option<String>>,
/// Whether the run continues from its records: a sandbox workspace is
/// then brought to its snapshot when its scope is first acquired.
resumed: bool,
@ -522,7 +527,6 @@ impl FabroHooks {
scopes: ScopeEnvs::default(),
failure: Mutex::default(),
publisher: spec.publisher,
publish_failure: Mutex::default(),
resumed,
restore: OnceCell::new(),
store,
@ -537,20 +541,13 @@ impl FabroHooks {
}
}
/// The checkpoint failure that ended the run, when one did: what the
/// engine reports the run failed with.
/// The checkpoint failure that ended the run, when one did: required
/// finalization commits it as the run's failure.
#[must_use]
pub fn checkpoint_failure(&self) -> Option<String> {
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<String> {
sync::lock(&self.publish_failure).clone()
}
/// The run's workspaces on this host, as the hooks reach them.
#[must_use]
pub fn workspaces(&self) -> &RunWorkspaces {
@ -1403,23 +1400,6 @@ impl FabroHooks {
Ok(publication)
}
/// Hand a successful run's work to the publisher, before the run's
/// terminal record. A failure fails the run with its reason.
async fn publish(&self, publisher: &dyn RunPublisher, publication: &Publication) {
match publisher.publish(publication).await {
Ok(()) => info!(
run_id = %self.run_id,
branch = publication.run_branch,
sha = publication.head_sha,
"run published"
),
Err(message) => {
warn!(run_id = %self.run_id, error = %message, "the run's publication failed");
*sync::lock(&self.publish_failure) = Some(message);
}
}
}
/// Hold at a test gate when one is set for this point and node.
async fn gate(&self, point: &str, node: &str) {
let Some(dir) = &self.test_gates else {
@ -1626,35 +1606,69 @@ impl ExecutionHooks for FabroHooks {
Ok(report)
}
async fn run_finished(&self, context: &HookContext, finished: RunFinished) -> Vec<Note> {
info!(
run_id = %self.run_id,
status = ?finished.status,
failure = finished.failure.as_deref().unwrap_or(""),
"Petri run finished; recording the run's diff and running the run-end hooks"
);
let publication = match self.record_run_diff().await {
/// Every Fabro run requires finalization, whether or not it publishes
/// or checkpoints: one path for every run, and a declaration that a
/// resume or a fork always matches.
fn requires_run_finalization(&self) -> bool {
true
}
async fn finalize_run(
&self,
context: &HookContext,
finished: RunFinished,
) -> Result<(), FinalizationFailure> {
let diff = self.record_run_diff().await.map_err(|error| {
let message = error.render();
warn!(run_id = %self.run_id, error = %message, "the run's diff was not recorded");
message
});
// A failed checkpoint fails the run whatever its execution status,
// and its work is never published.
if let Some(message) = self.checkpoint_failure() {
return Err(projection::checkpoint_failure(message));
}
let publication = match diff {
Ok(publication) => publication,
Err(error) => {
let message = error.render();
warn!(run_id = %self.run_id, error = %message, "the run's diff was not recorded");
Err(message) => {
if self.publisher.is_some() && finished.status == RunStatus::Success {
*sync::lock(&self.publish_failure) = Some(message);
return Err(projection::publish_failure(message));
}
None
}
};
if let Some(publisher) = &self.publisher {
if finished.status == RunStatus::Success && self.checkpoint_failure().is_none() {
if let Some(publication) = &publication {
self.publish(publisher.as_ref(), publication).await;
} else if self.publish_failure().is_none() {
*sync::lock(&self.publish_failure) = Some(
"the run has no recorded branch and checkpoint to publish".to_string(),
if finished.status == RunStatus::Success {
let publication = publication.ok_or_else(|| {
projection::publish_failure(
"the run has no recorded branch and checkpoint to publish",
)
})?;
publisher.publish(&publication).await.map_err(|message| {
warn!(
run_id = %self.run_id,
error = %message,
"the run's publication failed"
);
}
projection::publish_failure(message)
})?;
info!(
run_id = %self.run_id,
branch = publication.run_branch,
sha = publication.head_sha,
"run published"
);
}
}
if self.inner.requires_run_finalization() {
self.inner.finalize_run(context, finished).await?;
}
Ok(())
}
async fn run_finished(&self, context: &HookContext, finished: RunFinished) -> Vec<Note> {
// The diff and publication already ran in finalize_run, before Petri
// committed the outcome.
self.inner.run_finished(context, finished).await
}

View file

@ -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<Utc>) {
fn conclude(
&mut self,
status: &str,
finalization_failure: Option<&FinalizationFailure>,
at: DateTime<Utc>,
) {
let Some(projection) = self.projection.as_mut() else {
return;
};
if projection.status.is_terminal() {
return;
}
let root = self
.state
.root
.and_then(|root| self.state.invocations.get(&root));
let failure_message = root.and_then(|root| root.failure.clone());
let run_status = finished_status(status);
let failure_message = finalization_failure
.map(|failure| failure.message.clone())
.or_else(|| root.and_then(|root| root.failure.clone()));
let run_status = finished_status(status, finalization_failure);
let (outcome, failure) = match run_status {
RunStatus::Failed {
reason: reason @ FailureReason::Cancelled,
@ -209,17 +226,28 @@ impl RunView {
}
/// The status Fabro gives a run at Petri's finish, by the status the finish
/// records (`success`, `cancelled`, or a failure).
pub(super) fn finished_status(status: &str) -> RunStatus {
/// records (`success`, `cancelled`, or a failure). A failed checkpoint
/// cancels the run, but the run failed: its finish says so with the
/// checkpoint's finalization failure.
pub(crate) fn finished_status(
status: &str,
finalization_failure: Option<&FinalizationFailure>,
) -> RunStatus {
match status {
"success" => RunStatus::Succeeded {
reason: SuccessReason::Completed,
},
"cancelled" => RunStatus::Failed {
reason: FailureReason::Cancelled,
},
"cancelled" if !finalization_failure.is_some_and(super::is_checkpoint_failure) => {
RunStatus::Failed {
reason: FailureReason::Cancelled,
}
}
_ => RunStatus::Failed {
reason: FailureReason::WorkflowError,
reason: if finalization_failure.is_some_and(super::is_publish_failure) {
FailureReason::PublishFailed
} else {
FailureReason::WorkflowError
},
},
}
}

View file

@ -37,18 +37,23 @@ use std::fmt;
use std::str::FromStr;
use chrono::{DateTime, TimeZone as _, Utc};
pub(crate) use coordinator::finished_status;
use fabro_store::StagePosition;
use fabro_store::platform_records::StoredPlatformRecord;
use fabro_types::{
RunControlAction, RunDiff, RunId, RunProjection, RunStatus, StageId, StageProjection,
FailureReason, RunControlAction, RunDiff, RunId, RunProjection, RunStatus, StageId,
StageProjection,
};
use petri_execution::events::{NodeRef, RunEvent, Subject};
use petri_execution::{CoordinatorEvent, CoordinatorRecord, ExecutionId};
use petri_runtime::ir::FinalizationFailure;
use petri_store::Record;
use serde::{Deserialize, Deserializer, Serialize, Serializer, de};
use serde_json::Value;
use tracing::debug;
use crate::checkpoint::CHECKPOINT_FAILED_CLASS;
/// One item the projector hands the fold, with its delivery sequence.
pub enum Item<'a> {
Petri(&'a RunEvent),
@ -369,17 +374,46 @@ pub fn run_id_of(key: &str) -> Option<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.
/// The required-finalization failure Fabro's hooks record when the run's
/// publication fails; its code is [`FailureReason::PublishFailed`]'s.
#[must_use]
pub fn finished_run_status(record: &Record) -> Option<RunStatus> {
pub fn publish_failure(message: impl Into<String>) -> FinalizationFailure {
FinalizationFailure::new(<&'static str>::from(FailureReason::PublishFailed), message)
}
/// Whether a required-finalization failure is the run's failed publication.
#[must_use]
pub fn is_publish_failure(failure: &FinalizationFailure) -> bool {
failure.code == <&'static str>::from(FailureReason::PublishFailed)
}
/// The required-finalization failure Fabro's hooks record when a checkpoint
/// failed during the run. The failed checkpoint cancelled the run, so Petri
/// may record it as cancelled; Fabro reports it as a workflow failure.
#[must_use]
pub fn checkpoint_failure(message: impl Into<String>) -> FinalizationFailure {
FinalizationFailure::new(CHECKPOINT_FAILED_CLASS, message)
}
/// Whether a required-finalization failure is a failed checkpoint.
#[must_use]
pub fn is_checkpoint_failure(failure: &FinalizationFailure) -> bool {
failure.code == CHECKPOINT_FAILED_CLASS
}
/// The committed overall status and required-finalization failure message.
/// Execution failure details remain in the invocation records and projection.
#[must_use]
pub fn finished_run_result(record: &Record) -> Option<(RunStatus, Option<String>)> {
let record: CoordinatorRecord = serde_json::from_value(record.record.clone()).ok()?;
match record.body {
CoordinatorEvent::RunFinished { status } => {
Some(coordinator::finished_status(&status.to_string()))
}
CoordinatorEvent::RunFinished {
status,
finalization_failure,
} => Some((
coordinator::finished_status(&status.to_string(), finalization_failure.as_ref()),
finalization_failure.map(|failure| failure.message),
)),
_ => None,
}
}
@ -412,10 +446,11 @@ mod tests {
#[test]
fn a_finish_record_names_the_status_the_view_ends_the_run_on() {
let finished = |status: &str| {
finished_run_status(&coordinator_record(&serde_json::json!({
finished_run_result(&coordinator_record(&serde_json::json!({
"event": "run.finished",
"status": status,
})))
.map(|(status, _)| status)
};
assert_eq!(
finished("success"),
@ -436,13 +471,30 @@ mod tests {
})
);
assert_eq!(
finished_run_status(&coordinator_record(&serde_json::json!({
finished_run_result(&coordinator_record(&serde_json::json!({
"event": "run.finished",
"status": "cancelled",
"finalization_failure": {
"code": "checkpoint_failed",
"message": "checkpoint commit of `wreck` failed",
},
}))),
Some((
RunStatus::Failed {
reason: fabro_types::FailureReason::WorkflowError,
},
Some("checkpoint commit of `wreck` failed".to_string()),
)),
"a failed checkpoint's cancellation is the checkpoint's failure"
);
assert_eq!(
finished_run_result(&coordinator_record(&serde_json::json!({
"event": "run.paused",
}))),
None
);
assert_eq!(
finished_run_status(&Record {
finished_run_result(&Record {
seq: 3,
recorded_at: 1_000,
record: serde_json::json!({"event": "run.finished", "status": "success"}),

View file

@ -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);
}

View file

@ -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;

View file

@ -0,0 +1,151 @@
//! Real command-run records for transport and projection tests. No providers
//! or model credentials are used.
use std::sync::Arc;
use petri_execution::host::{self, HostRun};
use petri_execution::inspect;
use petri_runtime::driver::lifecycle::{ExecutionHooks, HookContext, RunFinished};
use petri_runtime::frontend::CompileInputs;
use petri_runtime::ir::FinalizationFailure;
use petri_runtime::{RunOptions, Runtime};
use petri_store::{Access, LogId, MemoryRunStore, Record, RunKey, RunStore as _};
use tokio::fs;
use tokio::sync::{Notify, Semaphore};
use crate::projection;
use crate::providers::{self, SandboxProviderConfig};
/// A deterministic finalizer gate for testing the committed boundary.
pub struct TestFinalizer {
pub entered: Notify,
pub release: Semaphore,
rejection: Option<String>,
}
impl TestFinalizer {
pub fn new(rejection: Option<&str>) -> Self {
Self {
entered: Notify::new(),
release: Semaphore::new(0),
rejection: rejection.map(str::to_string),
}
}
}
#[async_trait::async_trait]
impl ExecutionHooks for TestFinalizer {
fn requires_run_finalization(&self) -> bool {
true
}
async fn finalize_run(
&self,
_: &HookContext,
_: RunFinished,
) -> Result<(), FinalizationFailure> {
self.entered.notify_one();
self.release
.acquire()
.await
.expect("the gate stays open")
.forget();
self.rejection.as_ref().map_or(Ok(()), |message| {
Err(projection::publish_failure(message.clone()))
})
}
}
/// The command-only runtime installed with a test finalizer.
pub fn test_runtime(finalizer: Arc<TestFinalizer>) -> 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<Record>)>,
pub blobs: Vec<Vec<u8>>,
}
/// Run the command fixture to completion and return its logs and graph blobs,
/// with the finalizer rejecting the run when `rejection` names a failure.
pub async fn test_run_records(
run_id: fabro_types::RunId,
rejection: Option<&str>,
) -> TestRunRecords {
let root = tempfile::tempdir().expect("the fixture has an isolated directory");
let workflow = root.path().join("workflow.fabro");
fs::write(
&workflow,
r#"digraph Finalization {
graph [goal="Check required publication"]
start [shape=Mdiamond]
work [shape=parallelogram, script="echo completed"]
exit [shape=Msquare]
start -> work -> exit
}"#,
)
.await
.expect("the fixture workflow writes");
fs::write(
root.path().join("workflow.toml"),
"_version = 1\n[workflow]\ngraph = \"workflow.fabro\"\n",
)
.await
.expect("the settings write");
let finalizer = Arc::new(TestFinalizer::new(rejection));
finalizer.release.add_permits(1);
let store = Arc::new(MemoryRunStore::new());
let key = RunKey::new(run_id.to_string());
let mut options = RunOptions::new(root.path().join("run"));
options.run_key = Some(key.clone());
options.echo = false;
let runtime = test_runtime(finalizer)
.store(store.clone())
.options(options);
let checked = runtime
.check(&workflow, None, None, &CompileInputs::new())
.expect("the fixture compiles");
host::run_configured(
&runtime,
HostRun::new(checked.graph.expect("the fixture is valid")).with_children(checked.children),
|_, _| {},
)
.await
.expect("the fixture completes");
let logs = store
.open(&key, Access::Read)
.await
.expect("the fixture reads");
let inspection = inspect::inspect_run(&*logs)
.await
.expect("the fixture inspects");
let mut records = vec![(
LogId::Coordinator,
logs.read(&LogId::Coordinator)
.await
.expect("coordinator reads"),
)];
for execution in inspection.executions {
let id = LogId::Execution(execution.execution);
records.push((id, logs.read(&id).await.expect("execution reads")));
}
let mut blobs = Vec::new();
for graph in inspection.graphs {
blobs.push(
logs.get_blob(graph)
.await
.expect("graph reads")
.expect("graph exists"),
);
}
TestRunRecords {
logs: records,
blobs,
}
}

View file

@ -44,8 +44,9 @@ use fabro_types::{GitIdentitySource, RunId, SandboxProviderKind};
use object_store::local::LocalFileSystem;
use petri_execution::inspect::{self, RunInspection};
use petri_store::{Access, MemoryRunStore, RunKey, RunStore as _};
use tokio::fs;
use tokio::process::Command;
use tokio::sync::{Notify, Semaphore};
use tokio::{fs, time};
use tokio_util::sync::CancellationToken;
mod support;
@ -688,8 +689,9 @@ async fn a_failed_stage_is_committed_and_its_route_sees_the_files() {
}
/// A checkpoint commit that fails is fatal: the stage's outcome is recorded
/// as `checkpoint_failed`, no route is taken, the run ends failed with the
/// checkpoint's error, and a restart reports it failed without resuming.
/// as `checkpoint_failed`, no route is taken, the run's committed finish
/// fails it with the checkpoint's error, and a restart reports it failed
/// without resuming.
#[tokio::test]
async fn a_failed_checkpoint_ends_the_run_with_no_route() {
let harness = Harness::new().await;
@ -712,6 +714,16 @@ async fn a_failed_checkpoint_ends_the_run_with_no_route() {
);
let inspection = harness.inspection().await;
let committed = inspection
.finalization_failure
.clone()
.expect("the finish commits the checkpoint failure");
assert_eq!(committed.code, CHECKPOINT_FAILED_CLASS);
let stored = engine::outcome_of(&*harness.store, &harness.run_id.to_string())
.await
.expect("the finished run reads back");
assert_eq!(stored.status, RunStatus::Failed, "{stored:?}");
assert_eq!(stored.failure, outcome.failure);
let attempts: Vec<_> = inspection
.executions
.iter()
@ -1261,16 +1273,24 @@ async fn published_run(
let (origin, _) = upstream(&harness.run_dir.with_file_name("upstream"), 2).await;
harness.source = Some(file_source(&origin, "main", Some(1)));
harness.publisher = Some(Arc::clone(publisher) as Arc<dyn RunPublisher>);
let workflow = workflow(
&format!(" edit [shape=parallelogram, {attributes}]"),
" start -> edit -> exit",
);
let outcome = harness
.run_on(SandboxProviderKind::LOCAL, &workflow, SETTINGS)
.run_on(
SandboxProviderKind::LOCAL,
&published_workflow(attributes),
SETTINGS,
)
.await;
(harness, outcome)
}
/// The workflow [`published_run`] runs: one command stage with `attributes`.
fn published_workflow(attributes: &str) -> String {
workflow(
&format!(" edit [shape=parallelogram, {attributes}]"),
" start -> edit -> exit",
)
}
/// A successful run hands its publisher the run branch, the commit it ends
/// on (held by the snapshot repository) and its patch, before the run ends.
#[tokio::test]
@ -1307,10 +1327,37 @@ async fn a_successful_run_is_published_with_its_branch_head_and_patch() {
#[tokio::test]
async fn a_failed_publication_fails_the_run() {
let publisher = RecordingPublisher::new(Some("the push was rejected"));
let (_, outcome) = published_run("script=\"echo edited >> README.md\"", &publisher).await;
let attributes = "script=\"echo edited >> README.md\"";
let (harness, outcome) = published_run(attributes, &publisher).await;
assert_eq!(outcome.status, RunStatus::Failed, "{outcome:?}");
assert!(outcome.publish_failed);
assert_eq!(outcome.failure.as_deref(), Some("the push was rejected"));
let inspection = harness.inspection().await;
assert_eq!(inspection.status.as_deref(), Some("failed"));
let failure = inspection.finalization_failure.expect("durable failure");
assert_eq!(failure.code, "publish_failed");
assert_eq!(failure.message, "the push was rejected");
let stored = engine::outcome_of(&*harness.store, &harness.run_id.to_string())
.await
.unwrap();
assert_eq!(stored.status, outcome.status);
assert_eq!(stored.failure, outcome.failure);
assert!(stored.publish_failed);
let resumed = harness
.execute_on(
SandboxProviderKind::LOCAL,
&published_workflow(attributes),
SETTINGS,
true,
)
.await;
assert_eq!(resumed.status, outcome.status);
assert_eq!(resumed.failure, outcome.failure);
assert_eq!(
publisher.published.lock().unwrap().len(),
1,
"committed resume does not publish twice"
);
}
/// A run that fails (here, at a goal gate) is not published.
@ -1323,6 +1370,27 @@ async fn a_failed_run_is_not_published() {
assert!(publisher.published.lock().unwrap().is_empty());
}
/// A run whose checkpoint failed is never published, and its committed
/// failure is the checkpoint's, not a publication's.
#[tokio::test]
async fn a_run_whose_checkpoint_failed_is_not_published() {
let publisher = RecordingPublisher::new(None);
let (harness, outcome) = published_run(
"script=\"rm -rf .git && echo garbage > .git && echo wrecked > out.txt\"",
&publisher,
)
.await;
assert_eq!(outcome.status, RunStatus::Failed, "{outcome:?}");
assert!(!outcome.publish_failed);
assert!(publisher.published.lock().unwrap().is_empty());
let failure = harness
.inspection()
.await
.finalization_failure
.expect("the finish commits the checkpoint failure");
assert_eq!(failure.code, CHECKPOINT_FAILED_CLASS);
}
struct OriginPublisher {
origin: String,
pushed: std::sync::Mutex<Vec<String>>,
@ -1437,6 +1505,10 @@ async fn assert_shallow_fork(checkpoint_index: usize) {
})
.await
.expect("the fork is seeded");
assert!(
forked.inspection().await.required_finalization,
"a publishing fork inherits its required finalization declaration"
);
assert_eq!(
seeded.start.expect("the selected checkpoint").sha,
*checkpoint_sha
@ -1570,6 +1642,117 @@ async fn a_run_diff_failure_cannot_silently_skip_publication() {
let outcome = harness.run(&graph, SETTINGS).await;
assert_eq!(outcome.status, RunStatus::Failed, "{outcome:?}");
assert!(outcome.publish_failed);
assert!(outcome.failure.unwrap().contains("run diff unavailable"));
assert!(
outcome
.failure
.as_ref()
.unwrap()
.contains("run diff unavailable")
);
let stored = engine::outcome_of(&*harness.store, &harness.run_id.to_string())
.await
.unwrap();
assert_eq!(stored.status, outcome.status);
assert_eq!(stored.failure, outcome.failure);
assert!(stored.publish_failed);
assert!(publisher.published.lock().unwrap().is_empty());
}
struct GatedPublisher {
inner: Arc<RecordingPublisher>,
entered: Notify,
release: Semaphore,
}
#[async_trait::async_trait]
impl RunPublisher for GatedPublisher {
async fn push(&self, site: &Site, branch: &str, sha: &str) -> Result<(), String> {
self.inner.push(site, branch, sha).await
}
async fn publish(&self, publication: &Publication) -> Result<(), String> {
self.entered.notify_one();
self.release
.acquire()
.await
.expect("publication gate stays open")
.forget();
self.inner.publish(publication).await
}
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn required_publication_blocks_the_terminal_result_and_cleanup() {
for rejection in [None, Some("the push was rejected")] {
let publisher = Arc::new(GatedPublisher {
inner: RecordingPublisher::new(rejection),
entered: Notify::new(),
release: Semaphore::new(0),
});
let mut harness = Harness::new().await;
harness.publisher = Some(publisher.clone());
let harness = Arc::new(harness);
let executing = harness.clone();
let task = tokio::spawn(async move {
executing
.run(
&workflow(
"edit [shape=parallelogram, script=\"echo edited >> README.md\"]",
"start -> edit -> exit",
),
SETTINGS,
)
.await
});
time::timeout(
std::time::Duration::from_secs(15),
publisher.entered.notified(),
)
.await
.expect("publication reaches the gate");
let pending = harness.inspection().await;
assert!(pending.required_finalization);
assert!(
pending.status.is_none(),
"no terminal result while publication waits"
);
assert!(!task.is_finished());
let logs = harness
.store
.open(&RunKey::new(harness.run_id.to_string()), Access::Read)
.await
.unwrap();
let coordinator = logs.read(&petri_store::LogId::Coordinator).await.unwrap();
assert!(
coordinator
.iter()
.all(|record| record.record["body"]["event"] != "scope.released"),
"scope cleanup waits for publication"
);
assert!(
harness.workspace_path(&harness.workspace().await).exists(),
"workspace is available to publication"
);
assert!(matches!(
engine::outcome_of(&*harness.store, &harness.run_id.to_string()).await,
Err(engine::RunError::Unfinished(_))
));
publisher.release.add_permits(1);
let outcome = task.await.unwrap();
assert_eq!(
outcome.status,
if rejection.is_some() {
RunStatus::Failed
} else {
RunStatus::Success
}
);
assert_eq!(outcome.publish_failed, rejection.is_some());
let stored = engine::outcome_of(&*harness.store, &harness.run_id.to_string())
.await
.unwrap();
assert_eq!(stored.status, outcome.status);
assert_eq!(stored.failure, outcome.failure);
assert_eq!(publisher.inner.published.lock().unwrap().len(), 1);
}
}

View file

@ -27,7 +27,8 @@ use fabro_petri::providers::SandboxProviderConfig;
use fabro_petri::runtime::RuntimeSpec;
use fabro_petri::{SqliteRunStore, providers, test_support as petri_support};
use fabro_store::platform_records::{
PlatformRecord, PlatformRecordStore, RunCreatedRecord, RunLifecycleKind, RunLifecycleRecord,
CheckpointRecord, PlatformRecord, PlatformRecordStore, RunCreatedRecord, RunDiffRecord,
RunLifecycleKind, RunLifecycleRecord,
};
use fabro_store::{BlobStore, test_support};
use fabro_types::{
@ -42,7 +43,7 @@ use petri_runtime::frontend::CompileInputs;
use petri_runtime::ir::RunStatus as PetriRunStatus;
use petri_store::{RunKey, RunStore};
use tokio::fs;
use tokio::time::sleep;
use tokio::time::{self, sleep};
const COMMAND_WORKFLOW: &str = r#"digraph Command {
graph [goal="Run one command"]
@ -1361,3 +1362,184 @@ async fn an_auto_approved_answer_closes_the_question_in_the_projection() {
);
assert_view_equals_rebuild(&gate.scenario.pool, gate.scenario.run_id).await;
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn required_finalization_projects_only_the_committed_overall_result() {
use fabro_petri::test_support::finalization::{self, TestFinalizer};
use fabro_types::{FailureReason, StageOutcome, SuccessReason};
for rejection in [None, Some("the push was rejected")] {
let scenario = command_scenario().await;
let projector = Projector::new(scenario.pool.clone(), scenario.pool.clone());
let store = projector.observe_store(Arc::new(SqliteRunStore::new(scenario.pool.clone())));
let finalizer = Arc::new(TestFinalizer::new(rejection));
let runtime = finalization::test_runtime(finalizer.clone())
.store(store)
.options(run_options(&scenario.run_dir, scenario.run_id));
let checked = runtime
.check(&scenario.workflow, None, None, &CompileInputs::new())
.unwrap();
let task = tokio::spawn(async move {
host::run_configured(
&runtime,
HostRun::new(checked.graph.unwrap()).with_children(checked.children),
|_, _| {},
)
.await
.unwrap()
});
time::timeout(Duration::from_secs(15), finalizer.entered.notified())
.await
.unwrap();
// The driver can still flush its execution journal after entering
// finalization. Wait for the exit-stage evidence, rather than racing
// that writer while comparing the live view with a full rebuild.
let pending = time::timeout(Duration::from_secs(5), async {
loop {
projector.signal(scenario.run_id);
projector.settle(scenario.run_id).await;
let pending = petri_support::stored_projection(&scenario.pool, scenario.run_id)
.await
.unwrap()
.unwrap();
if pending.iter_stages().any(|(id, stage)| {
id.node_id() == "exit" && stage.state == StageState::Succeeded
}) {
break pending;
}
time::sleep(Duration::from_millis(5)).await;
}
})
.await
.expect("execution evidence flushes while publication is held");
assert_eq!(pending.status, RunStatus::Running);
assert!(pending.conclusion.is_none());
assert!(!task.is_finished());
assert_view_equals_rebuild(&scenario.pool, scenario.run_id).await;
let summary = fabro_types::DiffSummary {
files_changed: 1,
additions: 2,
deletions: 0,
};
let head_sha = "0123456789012345678901234567890123456789";
let patch_blob = BlobHash::new(b"final patch");
let platform = PlatformRecordStore::new(scenario.pool.clone());
for record in [
PlatformRecord::Checkpoint(CheckpointRecord {
execution: 0,
firing: 0,
attempt: Some(1),
workspace: None,
git_commit_sha: Some(head_sha.to_string()),
diff_summary: Some(summary),
patch_blob: Some(patch_blob),
operation: None,
}),
PlatformRecord::RunDiff(RunDiffRecord {
base_sha: None,
head_sha: Some(head_sha.to_string()),
diff_summary: Some(summary),
patch_blob: Some(patch_blob),
}),
] {
platform
.append(&scenario.run_id, &record, None)
.await
.unwrap();
}
finalizer.release.add_permits(1);
let report = task.await.unwrap();
assert_eq!(report.state.folded_status(), PetriRunStatus::Success);
assert_eq!(
report.status,
if rejection.is_some() {
PetriRunStatus::Failed
} else {
PetriRunStatus::Success
}
);
projector.settle(scenario.run_id).await;
let stored = petri_support::stored_projection(&scenario.pool, scenario.run_id)
.await
.unwrap()
.unwrap();
let expected = match rejection {
Some(_) => RunStatus::Failed {
reason: FailureReason::PublishFailed,
},
None => RunStatus::Succeeded {
reason: SuccessReason::Completed,
},
};
assert_eq!(stored.status, expected);
let conclusion = stored.conclusion.as_ref().unwrap();
assert_eq!(
conclusion
.failure
.as_ref()
.map(|failure| failure.detail.message.as_str()),
rejection
);
assert_eq!(
conclusion.status,
if rejection.is_some() {
StageOutcome::Failed {
retry_requested: false,
}
} else {
StageOutcome::Succeeded
}
);
assert!(
stored
.iter_stages()
.all(|(_, stage)| stage.state == StageState::Succeeded)
);
assert_eq!(conclusion.final_git_commit_sha.as_deref(), Some(head_sha));
assert_eq!(conclusion.diff.summary, Some(summary));
assert_eq!(
conclusion.diff.patch.as_deref(),
Some(fabro_types::format_blob_ref(&patch_blob).as_str())
);
assert!(!conclusion.stages.is_empty());
assert_view_equals_rebuild(&scenario.pool, scenario.run_id).await;
// A later worker lifecycle cannot replace the committed terminal
// result or conclusion, even if its status disagrees.
let other = match rejection {
Some(_) => RunStatus::Succeeded {
reason: SuccessReason::Completed,
},
None => RunStatus::Failed {
reason: FailureReason::PublishFailed,
},
};
PlatformRecordStore::new(scenario.pool.clone())
.append(
&scenario.run_id,
&PlatformRecord::RunLifecycle(
RunLifecycleRecord::new(if rejection.is_some() {
RunLifecycleKind::Succeeded
} else {
RunLifecycleKind::Failed
})
.with_status(other),
),
None,
)
.await
.unwrap();
projector.signal(scenario.run_id);
projector.settle(scenario.run_id).await;
let after = petri_support::stored_projection(&scenario.pool, scenario.run_id)
.await
.unwrap()
.unwrap();
assert_eq!(after.status, stored.status);
assert_eq!(
serde_json::to_value(after.conclusion).unwrap(),
serde_json::to_value(stored.conclusion).unwrap()
);
assert_view_equals_rebuild(&scenario.pool, scenario.run_id).await;
}
}