Commit publication failures through required run finalization

This commit is contained in:
Scott Werner 2026-10-05 16:59:50 -04:00
parent 36f8b61b60
commit 21d4db6766
21 changed files with 1040 additions and 148 deletions

View file

@ -108,13 +108,13 @@ pebble-cli-core = { git = "https://github.com/lithoscomputer/pebble", branch = "
# petri: the workflow engine Fabro runs its workflows on. Only `fabro-petri`
# and `fabro-dot` (the DOT parser alone) may depend on these packages; the
# keys carry the `petri_` prefix so the crate names say where they come from.
petri_runtime = { git = "https://github.com/lithoscomputer/petri.git", branch = "main", package = "petri-runtime" }
petri_execution = { git = "https://github.com/lithoscomputer/petri.git", branch = "main", package = "petri-execution" }
petri_store = { git = "https://github.com/lithoscomputer/petri.git", branch = "main", package = "petri-store" }
petri_attractor_steps = { git = "https://github.com/lithoscomputer/petri.git", branch = "main", package = "petri-attractor-steps" }
petri_frontend_attractor = { git = "https://github.com/lithoscomputer/petri.git", branch = "main", package = "petri-frontend-attractor" }
petri_frontend_fabro = { git = "https://github.com/lithoscomputer/petri.git", branch = "main", package = "petri-frontend-fabro" }
petri_testkit = { git = "https://github.com/lithoscomputer/petri.git", branch = "main", package = "petri-testkit" }
petri_runtime = { git = "https://github.com/lithoscomputer/petri.git", rev = "0198e66f6619e147e461c52c138497863d851969", package = "petri-runtime" }
petri_execution = { git = "https://github.com/lithoscomputer/petri.git", rev = "0198e66f6619e147e461c52c138497863d851969", package = "petri-execution" }
petri_store = { git = "https://github.com/lithoscomputer/petri.git", rev = "0198e66f6619e147e461c52c138497863d851969", package = "petri-store" }
petri_attractor_steps = { git = "https://github.com/lithoscomputer/petri.git", rev = "0198e66f6619e147e461c52c138497863d851969", package = "petri-attractor-steps" }
petri_frontend_attractor = { git = "https://github.com/lithoscomputer/petri.git", rev = "0198e66f6619e147e461c52c138497863d851969", package = "petri-frontend-attractor" }
petri_frontend_fabro = { git = "https://github.com/lithoscomputer/petri.git", rev = "0198e66f6619e147e461c52c138497863d851969", package = "petri-frontend-fabro" }
petri_testkit = { git = "https://github.com/lithoscomputer/petri.git", rev = "0198e66f6619e147e461c52c138497863d851969", package = "petri-testkit" }
sentry = { version = "0.35", default-features = false, features = ["backtrace", "contexts", "ureq", "rustls"] }
fork = "0.2"
exec = "0.3"

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

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

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, @r###"
12:43:08 ✗ FAILED 0s
12:43:08 the push was rejected
"###);
}
#[test]
fn the_raw_line_is_the_envelope_as_json() {
let notice = platform(

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, as the legacy publish step did:
//! the final checkpoint is pushed from inside the sandbox to the run
//! branch on GitHub, and, when the run changed files and its settings ask
//! for one, a pull request is opened and recorded. A failure fails the run
//! with `publish_failed`.

View file

@ -940,4 +940,189 @@ mod tests {
"the failure names the worker's exit"
);
}
async fn append_engine_records(
app: &axum::Router,
run_id: RunId,
token: &str,
owner: &str,
log: &fabro_petri::petri::LogId,
records: &[fabro_petri::petri::Record],
) -> axum::response::Response {
let body = json!({ "owner": owner, "records": records.iter().map(|record| json!({
"seq": record.seq, "recorded_at": record.recorded_at, "record": record.record,
})).collect::<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>,
) {
let mut values = Vec::new();
for suffix in ["", "/state"] {
let response = app
.clone()
.oneshot(
Request::builder()
.uri(format!("/api/v1/runs/{run_id}{suffix}"))
.body(Body::empty())
.unwrap(),
)
.await
.unwrap();
values.push(
fabro_test::expect_axum_json(
response,
StatusCode::OK,
"GET run and state at required finalization boundary",
)
.await,
);
}
assert_eq!(
values[0]["lifecycle"]["status"],
serde_json::to_value(expected).unwrap()
);
assert_eq!(values[1]["status"], serde_json::to_value(expected).unwrap());
assert_eq!(state.test_managed_run_status(&run_id), Some(expected));
if expected.is_terminal() {
assert_eq!(
values[1]["conclusion"]["failure"]["detail"]["message"].as_str(),
message
);
assert!(
values[1]["conclusion"]["stages"]
.as_array()
.is_some_and(|stages| !stages.is_empty())
);
} else {
assert!(values[1]["conclusion"].is_null());
}
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn required_publication_result_agrees_across_worker_api_projection_and_cleanup() {
use fabro_petri::petri::LogId;
use fabro_petri::test_support::finalization;
for rejection in [None, Some("the push was rejected")] {
let runtime = Arc::new(HeldWorkerRuntime::default());
let (state, app, run_id, token) = held_worker_run(&runtime).await;
run_to_running_as_worker(&app, run_id, &token).await;
let fixture = finalization::test_run_records(run_id, rejection).await;
for bytes in fixture.blobs {
let response = app
.clone()
.oneshot(
Request::builder()
.method("POST")
.uri(format!("/api/v1/runs/{run_id}/petri/blobs?owner=worker-1"))
.header(header::AUTHORIZATION, format!("Bearer {token}"))
.body(Body::from(bytes))
.unwrap(),
)
.await
.unwrap();
fabro_test::assert_axum_status(response, StatusCode::OK, "worker graph blob").await;
}
let mut terminal = None;
for (log, mut records) in fixture.logs {
if log == LogId::Coordinator {
let finish = records.pop().unwrap();
assert_eq!(finish.record["body"]["event"], "run.finished");
terminal = Some(finish);
}
let response =
append_engine_records(&app, run_id, &token, "worker-1", &log, &records).await;
fabro_test::assert_axum_status(
response,
StatusCode::NO_CONTENT,
"worker execution records",
)
.await;
}
assert_public_result(&state, &app, run_id, RunStatus::Running, None).await;
let terminal = terminal.unwrap();
// A rejected writer cannot settle the managed run or store its
// result. This exercises the same endpoint as the valid finish.
let response = append_engine_records(
&app,
run_id,
&token,
"superseded-worker",
&LogId::Coordinator,
std::slice::from_ref(&terminal),
)
.await;
fabro_test::assert_axum_status(
response,
StatusCode::CONFLICT,
"finish from a stale lease owner",
)
.await;
assert_public_result(&state, &app, run_id, RunStatus::Running, None).await;
let response =
append_engine_records(&app, run_id, &token, "worker-1", &LogId::Coordinator, &[
terminal,
])
.await;
fabro_test::assert_axum_status(
response,
StatusCode::NO_CONTENT,
"authoritative worker finish",
)
.await;
let expected = match rejection {
Some(_) => RunStatus::Failed {
reason: FailureReason::PublishFailed,
},
None => RunStatus::Succeeded {
reason: SuccessReason::Completed,
},
};
assert_public_result(&state, &app, run_id, expected, rejection).await;
// Worker exit after the finish, without a platform terminal
// acknowledgement, must preserve Petri's committed result.
runtime.end_worker();
for _ in 0..500 {
if concluded_runs(&app).await == 1 {
break;
}
time::sleep(Duration::from_millis(10)).await;
}
assert_eq!(concluded_runs(&app).await, 1);
assert_public_result(&state, &app, run_id, expected, rejection).await;
let (rebuilt, _, _) =
fabro_petri::test_support::rebuild(&state.db_pool, &state.db_pool, run_id)
.await
.unwrap();
let rebuilt = rebuilt.unwrap();
assert_eq!(rebuilt.status, expected);
assert_eq!(
rebuilt
.conclusion
.unwrap()
.failure
.map(|failure| failure.detail.message),
rejection.map(str::to_string)
);
}
}
}

View file

@ -3582,8 +3582,10 @@ 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);
if !managed_run.status.is_terminal() {
managed_run.status = RunStatus::Failed { reason };
managed_run.error = Some(message);
}
clear_live_run_state(managed_run);
}
cleanup_worker_control_bus_for_run(state.as_ref(), run_id);
@ -3601,7 +3603,9 @@ fn apply_lifecycle_to_managed_run(state: &AppState, run_id: RunId, record: &RunL
// A settled run is immutable to the lifecycle: the follower still folds
// the records before the terminal one after Petri's finish or the
// worker's terminal record settled the run, and none may reopen it.
if managed_run.status.is_terminal() && !is_terminal_transition(record) {
if managed_run.status.is_terminal()
&& (!is_terminal_transition(record) || record.status != Some(managed_run.status))
{
return;
}
match record.transition {
@ -3660,7 +3664,9 @@ fn apply_lifecycle_to_managed_run(state: &AppState, run_id: RunId, record: &RunL
managed_run.status = record.status.unwrap_or(RunStatus::Failed {
reason: FailureReason::WorkflowError,
});
managed_run.error.clone_from(&record.reason);
if managed_run.error.is_none() {
managed_run.error.clone_from(&record.reason);
}
managed_run.active_steerable_stages.clear();
managed_run.active_non_steerable_stages.clear();
cleanup_worker_control_bus_for_run(state, run_id);
@ -3683,19 +3689,14 @@ fn is_terminal_transition(record: &RunLifecycleRecord) -> bool {
)
}
/// Settle the in-memory run at Petri's own finish, as its worker stores
/// the `run.finished` record: the view reports the run ended from the
/// moment that record is stored, so the managed run the delete precheck
/// prefers must not still say running while the worker tears down; a
/// delete in that window was refused as active. A run already settled
/// keeps its status. The worker's terminal lifecycle record, a moment
/// later, refines the status and its error and ends the worker's controls
/// ([`settle_managed_run_at_terminal_record`]); the worker's exit later
/// reaps the process and leaves the settled status alone.
/// Settle the in-memory run after Petri's authoritative finish is durable.
/// Required publication has completed before this record. Worker teardown
/// and subsequent lifecycle records cannot change its terminal outcome.
pub(in crate::server) fn settle_managed_run_at_finish(
state: &AppState,
run_id: RunId,
status: RunStatus,
failure: Option<String>,
) {
let mut runs = state.runs.lock().expect("runs lock poisoned");
let Some(managed_run) = runs.get_mut(&run_id) else {
@ -3705,12 +3706,13 @@ pub(in crate::server) fn settle_managed_run_at_finish(
return;
}
managed_run.status = status;
managed_run.error = failure;
managed_run.active_steerable_stages.clear();
managed_run.active_non_steerable_stages.clear();
}
/// Settle the in-memory run at the terminal lifecycle record its worker
/// stores, ahead of the store: the same as [`settle_managed_run_at_finish`]
/// stores, after the store: the same as [`settle_managed_run_at_finish`]
/// for a worker that ended the run without Petri's finish (it failed before
/// the engine ran), and the record's status, error and control cleanup for
/// one that did. A record that is not terminal is left to the stream

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

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.
//!
@ -549,13 +549,13 @@ 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");
match run_records::lifecycle(&state, run_id, record).await {
Ok(()) => finish(&state, run_id, status, error),
Err(err) => {
error!(run_id = %run_id, error = %err, "Failed to persist run outcome");
release_live_state(&state, run_id);
}
}
// The managed run settled at Petri's finish, ahead of the store; the
// terminal record refines its status and error and ends its live
// state, and is the settle of a run that ended without a finish.
finish(&state, run_id, status, error);
// The view trails the terminal record; the aggregate reads the settled
// projection, as the worker path reads the final state at worker exit.
state.petri_projector.settle(run_id).await;
@ -814,10 +814,13 @@ fn failed(
async fn fail_before_execution(state: &Arc<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");
match run_records::lifecycle(state, run_id, record).await {
Ok(()) => finish(state, run_id, status, error),
Err(err) => {
error!(run_id = %run_id, error = %err, "Failed to persist run failure status");
release_live_state(state, run_id);
}
}
finish(state, run_id, status, error);
}
/// Settle the managed run at its terminal record and release its
@ -828,24 +831,31 @@ async fn fail_before_execution(state: &Arc<AppState>, run_id: RunId, message: &s
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;
if !managed_run.status.is_terminal() || managed_run.status == status {
managed_run.status = status;
if managed_run.error.is_none() {
managed_run.error = error;
}
}
clear_live_run_state(managed_run);
}
drop(runs);
state.scheduler_notify.notify_one();
}
/// The run store an in-process run executes over: the projector's
/// signalling store, whose coordinator appends settle the managed run at
/// Petri's own finish first. The view ends the run at the `run.finished`
/// record the moment it is stored and a pass folds it, so the managed run
/// the delete precheck prefers must not still say running while the engine
/// tears down: a delete in that window was refused as active. The worker's
/// records endpoint does the same for a worker-backed run, ahead of the
/// same store. The settle is in memory only; the terminal lifecycle record
/// [`execute`] stores once the engine returns refines the status
/// ([`finish`]), and stays the settle of a run that ends without a finish.
/// Release controls after an append failure without claiming a new terminal
/// result. A Petri finish already committed remains authoritative.
fn release_live_state(state: &Arc<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);
}
drop(runs);
state.scheduler_notify.notify_one();
}
/// The run store an in-process run executes over. A coordinator finish
/// settles managed status only after the append succeeds, as on HTTP workers.
struct SettlingStore {
inner: Arc<dyn RunStore>,
state: Arc<AppState>,
@ -864,7 +874,7 @@ impl RunStore for SettlingStore {
}
/// One run's logs, whose coordinator appends settle the managed run at
/// Petri's finish before the records reach the store.
/// Petri's finish after the records reach the store.
struct SettlingLogs {
inner: Arc<dyn RunLogs>,
run_id: Option<RunId>,
@ -878,12 +888,15 @@ impl RunLogs for SettlingLogs {
}
async fn append(&self, log: &LogId, records: &[Record]) -> Result<(), StoreError> {
self.inner.append(log, records).await?;
if let Some(run_id) = self.run_id.filter(|_| *log == LogId::Coordinator) {
if let Some(status) = records.iter().find_map(projection::finished_run_status) {
super::settle_managed_run_at_finish(&self.state, run_id, status);
if let Some((status, failure)) =
records.iter().find_map(projection::finished_run_result)
{
super::settle_managed_run_at_finish(&self.state, run_id, status, failure);
}
}
self.inner.append(log, records).await
Ok(())
}
async fn read(&self, log: &LogId) -> Result<Vec<Record>, StoreError> {
@ -922,6 +935,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 +947,9 @@ mod tests {
async fn append(&self, log: &LogId, records: &[Record]) -> Result<(), StoreError> {
let status = self.state.test_managed_run_status(&self.run_id);
self.seen.lock().expect("seen lock poisoned").push(status);
if self.reject.load(std::sync::atomic::Ordering::SeqCst) {
return Err(StoreError::StaleOwner);
}
self.inner.append(log, records).await
}
@ -992,6 +1009,7 @@ mod tests {
state: Arc::clone(&state),
run_id,
seen: Mutex::new(Vec::new()),
reject: std::sync::atomic::AtomicBool::new(false),
});
let logs = SettlingLogs {
inner: Arc::clone(&recorder) as Arc<dyn RunLogs>,
@ -1001,12 +1019,10 @@ mod tests {
(state, run_id, logs, recorder)
}
/// The in-process run settles at Petri's own finish, before the
/// `run.finished` record reaches the store: the view cannot report the
/// run ended while the managed run still says running. The records
/// before the finish leave the run in flight.
/// The in-process run settles only after its authoritative finish is
/// stored. Earlier records leave the run in flight.
#[tokio::test]
async fn an_in_process_run_settles_before_its_finish_is_stored() {
async fn an_in_process_run_settles_after_its_finish_is_stored() {
let (state, run_id, logs, recorder) = in_flight_run().await;
logs.append(&LogId::Coordinator, &[coordinator_record(
@ -1032,12 +1048,33 @@ mod tests {
};
assert_eq!(
*recorder.seen.lock().expect("seen lock poisoned"),
vec![Some(RunStatus::Running), Some(succeeded)],
"the managed run settled before the finish reached the store"
vec![Some(RunStatus::Running), Some(RunStatus::Running)],
"the managed run stayed running until the finish reached the store"
);
assert_eq!(state.test_managed_run_status(&run_id), Some(succeeded));
}
#[tokio::test]
async fn a_rejected_in_process_finish_does_not_settle_the_run() {
let (state, run_id, logs, recorder) = in_flight_run().await;
recorder
.reject
.store(true, std::sync::atomic::Ordering::SeqCst);
let result = logs
.append(&LogId::Coordinator, &[coordinator_record(
0,
&json!({"event": "run.finished", "status": "failed",
"finalization_failure": {"code": "publish_failed", "message": "push rejected"}}),
)])
.await;
assert!(matches!(result, Err(StoreError::StaleOwner)));
assert_eq!(
state.test_managed_run_status(&run_id),
Some(RunStatus::Running)
);
assert!(logs.read(&LogId::Coordinator).await.unwrap().is_empty());
}
/// A finish on another log than the coordinator's is not Petri's
/// finish of the run: an execution's engine log ends an execution.
#[tokio::test]

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

@ -363,16 +363,6 @@ pub async fn run(request: RunRequest) -> Result<RunOutcome, RunError> {
outcome.status = RunStatus::Failed;
outcome.failure = Some(failure);
}
// A successful run whose publication failed is a failed run: its work
// did not reach where the settings sent it.
if let Some(failure) = fabro_hooks
.as_ref()
.and_then(|hooks| hooks.publish_failure())
{
outcome.status = RunStatus::Failed;
outcome.failure = Some(failure);
outcome.publish_failed = true;
}
Ok(outcome)
}
@ -555,19 +545,27 @@ fn outcome(
return Err(RunError::Unfinished(reasons));
}
};
let failure = inspection
let execution_failure = inspection
.invocations
.iter()
.find(|invocation| invocation.invocation == inspection.root.invocation)
.and_then(|root| root.result.as_ref())
.and_then(|result| result.failure.as_ref())
.map(|failure| failure.message.clone());
let publish_failed = inspection
.finalization_failure
.as_ref()
.is_some_and(|failure| failure.code == "publish_failed");
let failure = inspection
.finalization_failure
.map(|failure| failure.message)
.or(execution_failure);
Ok(RunOutcome {
status,
failure,
complete: inspection.complete,
incomplete: inspection.incomplete,
publish_failed: false,
publish_failed,
})
}

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
//! - `finalize_run`: the run's diff, its run branch against its base commit, as
//! the `run.diff` platform record with the patch as a blob; for a successful
//! run, its publication ([`RunPublisher`]: the platform pushes the run branch
//! and opens a pull request), whose failure fails the run before its terminal
//! record; then the forwarded point, so the local service runs `run_complete`
//! and `run_failed` with the sandbox in place.
//! record.
//! - `run_finished`: best-effort diff preparation for nonpublishing runs, then
//! the forwarded point, so the local service runs `run_complete` and
//! `run_failed` with the sandbox in place.
//! - `scope_acquired`: a fresh run's Git target checked out into the workspace
//! from inside the scope ([`crate::source`]); a resumed run uses its
//! surviving workspace, while an explicit fork fetches the source run's
@ -100,7 +102,9 @@ use petri_runtime::driver::lifecycle::{
ScopeReleased, Transition, TransitionError, TransitionReport,
};
use petri_runtime::executor::{EnvError, ExecEnv};
use petri_runtime::ir::{ExecutionId, FailureInfo, RunStatus, ScopeId, Status};
use petri_runtime::ir::{
ExecutionId, FailureInfo, FinalizationFailure, RunStatus, ScopeId, Status,
};
use serde_json::json;
use tokio::sync::{Mutex as AsyncMutex, OnceCell};
use tokio::{fs, time};
@ -453,8 +457,6 @@ pub struct FabroHooks {
/// The checkpoint failure that ended the run, when one did.
failure: Mutex<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 +524,6 @@ impl FabroHooks {
scopes: ScopeEnvs::default(),
failure: Mutex::default(),
publisher: spec.publisher,
publish_failure: Mutex::default(),
resumed,
restore: OnceCell::new(),
store,
@ -544,13 +545,6 @@ impl FabroHooks {
sync::lock(&self.failure).clone()
}
/// Why the run's publication failed, when it did: the run then fails
/// with this message.
#[must_use]
pub fn publish_failure(&self) -> Option<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 +1397,6 @@ impl FabroHooks {
Ok(publication)
}
/// Hand a successful run's work to the publisher, before the run's
/// terminal record. A failure fails the run with its reason.
async fn publish(&self, publisher: &dyn RunPublisher, publication: &Publication) {
match publisher.publish(publication).await {
Ok(()) => info!(
run_id = %self.run_id,
branch = publication.run_branch,
sha = publication.head_sha,
"run published"
),
Err(message) => {
warn!(run_id = %self.run_id, error = %message, "the run's publication failed");
*sync::lock(&self.publish_failure) = Some(message);
}
}
}
/// Hold at a test gate when one is set for this point and node.
async fn gate(&self, point: &str, node: &str) {
let Some(dir) = &self.test_gates else {
@ -1626,33 +1603,53 @@ impl ExecutionHooks for FabroHooks {
Ok(report)
}
async fn run_finished(&self, context: &HookContext, finished: RunFinished) -> Vec<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"
);
fn requires_run_finalization(&self) -> bool {
self.publisher.is_some() || self.inner.requires_run_finalization()
}
async fn finalize_run(
&self,
context: &HookContext,
finished: RunFinished,
) -> Result<(), FinalizationFailure> {
let publication = match self.record_run_diff().await {
Ok(publication) => publication,
Err(error) => {
let message = error.render();
warn!(run_id = %self.run_id, error = %message, "the run's diff was not recorded");
if self.publisher.is_some() && finished.status == RunStatus::Success {
*sync::lock(&self.publish_failure) = Some(message);
return Err(FinalizationFailure::new("publish_failed", message));
}
None
}
};
if let Some(publisher) = &self.publisher {
if finished.status == RunStatus::Success && self.checkpoint_failure().is_none() {
if let Some(publication) = &publication {
self.publish(publisher.as_ref(), publication).await;
} else if self.publish_failure().is_none() {
*sync::lock(&self.publish_failure) = Some(
"the run has no recorded branch and checkpoint to publish".to_string(),
);
}
if finished.status == RunStatus::Success {
let publication = publication.ok_or_else(|| {
FinalizationFailure::new(
"publish_failed",
"the run has no recorded branch and checkpoint to publish",
)
})?;
publisher.publish(&publication).await.map_err(|message| {
warn!(run_id = %self.run_id, error = %message, "the run's publication failed");
FinalizationFailure::new("publish_failed", message)
})?;
info!(run_id = %self.run_id, branch = publication.run_branch, sha = publication.head_sha, "run published");
}
}
if self.inner.requires_run_finalization() {
self.inner.finalize_run(context, finished).await?;
}
Ok(())
}
async fn run_finished(&self, context: &HookContext, finished: RunFinished) -> Vec<Note> {
// Nonpublishing runs keep best-effort diff preparation. Required work
// has already run in finalize_run, before Petri commits its outcome.
if !self.requires_run_finalization() {
if let Err(error) = self.record_run_diff().await {
warn!(run_id = %self.run_id, error = %error.render(), "the run's diff was not recorded");
}
}
self.inner.run_finished(context, finished).await

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,
@ -210,7 +227,10 @@ impl RunView {
/// The status Fabro gives a run at Petri's finish, by the status the finish
/// records (`success`, `cancelled`, or a failure).
pub(super) fn finished_status(status: &str) -> RunStatus {
pub(super) fn finished_status(
status: &str,
finalization_failure: Option<&FinalizationFailure>,
) -> RunStatus {
match status {
"success" => RunStatus::Succeeded {
reason: SuccessReason::Completed,
@ -219,7 +239,12 @@ pub(super) fn finished_status(status: &str) -> RunStatus {
reason: FailureReason::Cancelled,
},
_ => RunStatus::Failed {
reason: FailureReason::WorkflowError,
reason: if finalization_failure.is_some_and(|failure| failure.code == "publish_failed")
{
FailureReason::PublishFailed
} else {
FailureReason::WorkflowError
},
},
}
}

View file

@ -375,11 +375,22 @@ pub fn run_id_of(key: &str) -> Option<RunId> {
/// terminal lifecycle record. `None` for any other record.
#[must_use]
pub fn finished_run_status(record: &Record) -> Option<RunStatus> {
finished_run_result(record).map(|(status, _)| status)
}
/// The committed overall status and required-finalization failure message.
/// Execution failure details remain in the invocation records and projection.
#[must_use]
pub fn finished_run_result(record: &Record) -> Option<(RunStatus, Option<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,
}
}

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

View file

@ -1307,10 +1307,31 @@ async fn a_successful_run_is_published_with_its_branch_head_and_patch() {
#[tokio::test]
async fn a_failed_publication_fails_the_run() {
let publisher = RecordingPublisher::new(Some("the push was rejected"));
let (_, outcome) = published_run("script=\"echo edited >> README.md\"", &publisher).await;
let (harness, outcome) = published_run("script=\"echo edited >> README.md\"", &publisher).await;
assert_eq!(outcome.status, RunStatus::Failed, "{outcome:?}");
assert!(outcome.publish_failed);
assert_eq!(outcome.failure.as_deref(), Some("the push was rejected"));
let inspection = harness.inspection().await;
assert_eq!(inspection.status.as_deref(), Some("failed"));
let failure = inspection.finalization_failure.expect("durable failure");
assert_eq!(failure.code, "publish_failed");
assert_eq!(failure.message, "the push was rejected");
let stored = engine::outcome_of(&*harness.store, &harness.run_id.to_string())
.await
.unwrap();
assert_eq!(stored.status, outcome.status);
assert_eq!(stored.failure, outcome.failure);
assert!(stored.publish_failed);
let resumed = harness
.execute_on(SandboxProviderKind::LOCAL, "", SETTINGS, true)
.await;
assert_eq!(resumed.status, outcome.status);
assert_eq!(resumed.failure, outcome.failure);
assert_eq!(
publisher.published.lock().unwrap().len(),
1,
"committed resume does not publish twice"
);
}
/// A run that fails (here, at a goal gate) is not published.
@ -1570,6 +1591,117 @@ async fn a_run_diff_failure_cannot_silently_skip_publication() {
let outcome = harness.run(&graph, SETTINGS).await;
assert_eq!(outcome.status, RunStatus::Failed, "{outcome:?}");
assert!(outcome.publish_failed);
assert!(outcome.failure.unwrap().contains("run diff unavailable"));
assert!(
outcome
.failure
.as_ref()
.unwrap()
.contains("run diff unavailable")
);
let stored = engine::outcome_of(&*harness.store, &harness.run_id.to_string())
.await
.unwrap();
assert_eq!(stored.status, outcome.status);
assert_eq!(stored.failure, outcome.failure);
assert!(stored.publish_failed);
assert!(publisher.published.lock().unwrap().is_empty());
}
struct GatedPublisher {
inner: Arc<RecordingPublisher>,
entered: tokio::sync::Notify,
release: tokio::sync::Semaphore,
}
#[async_trait::async_trait]
impl RunPublisher for GatedPublisher {
async fn push(&self, site: &Site, branch: &str, sha: &str) -> Result<(), String> {
self.inner.push(site, branch, sha).await
}
async fn publish(&self, publication: &Publication) -> Result<(), String> {
self.entered.notify_one();
self.release.acquire().await.unwrap().forget();
self.inner.publish(publication).await
}
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn required_publication_blocks_the_terminal_result_and_cleanup() {
for rejection in [None, Some("the push was rejected")] {
let publisher = Arc::new(GatedPublisher {
inner: RecordingPublisher::new(rejection),
entered: tokio::sync::Notify::new(),
release: tokio::sync::Semaphore::new(0),
});
let mut harness = Harness::new().await;
harness.publisher = Some(publisher.clone());
let harness = Arc::new(harness);
let executing = harness.clone();
let task = tokio::spawn(async move {
executing
.run(
&workflow(
"edit [shape=parallelogram, script=\"echo edited >> README.md\"]",
"start -> edit -> exit",
),
SETTINGS,
)
.await
});
tokio::time::timeout(
std::time::Duration::from_secs(15),
publisher.entered.notified(),
)
.await
.expect("publication reaches the gate");
let pending = harness.inspection().await;
assert!(pending.required_finalization);
assert!(
pending.status.is_none(),
"no terminal result while publication waits"
);
assert!(!task.is_finished());
let logs = harness
.store
.open(&RunKey::new(harness.run_id.to_string()), Access::Read)
.await
.unwrap();
let coordinator = logs.read(&petri_store::LogId::Coordinator).await.unwrap();
assert!(
coordinator
.iter()
.all(|record| record.record["body"]["event"] != "scope.released"),
"scope cleanup waits for publication"
);
assert!(
pending.invocations[0].result.is_some(),
"execution already ended"
);
assert!(
harness.workspace_path(&harness.workspace().await).exists(),
"workspace is available to publication"
);
assert!(matches!(
engine::outcome_of(&*harness.store, &harness.run_id.to_string()).await,
Err(engine::RunError::Unfinished(_))
));
publisher.release.add_permits(1);
let outcome = task.await.unwrap();
assert_eq!(
outcome.status,
if rejection.is_some() {
RunStatus::Failed
} else {
RunStatus::Success
}
);
assert_eq!(outcome.publish_failed, rejection.is_some());
let stored = engine::outcome_of(&*harness.store, &harness.run_id.to_string())
.await
.unwrap();
assert_eq!(stored.status, outcome.status);
assert_eq!(stored.failure, outcome.failure);
assert_eq!(publisher.inner.published.lock().unwrap().len(), 1);
}
}

View file

@ -1361,3 +1361,168 @@ async fn an_auto_approved_answer_closes_the_question_in_the_projection() {
);
assert_view_equals_rebuild(&gate.scenario.pool, gate.scenario.run_id).await;
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn required_finalization_projects_only_the_committed_overall_result() {
use fabro_petri::test_support::finalization::{self, TestFinalizer};
use fabro_types::{FailureReason, StageOutcome, SuccessReason};
for rejection in [None, Some("the push was rejected")] {
let scenario = command_scenario().await;
let projector = Projector::new(scenario.pool.clone(), scenario.pool.clone());
let store = projector.observe_store(Arc::new(SqliteRunStore::new(scenario.pool.clone())));
let finalizer = Arc::new(TestFinalizer::new(rejection));
let runtime = finalization::test_runtime(finalizer.clone())
.store(store)
.options(run_options(&scenario.run_dir, scenario.run_id));
let checked = runtime
.check(&scenario.workflow, None, None, &CompileInputs::new())
.unwrap();
let task = tokio::spawn(async move {
host::run_configured(
&runtime,
HostRun::new(checked.graph.unwrap()).with_children(checked.children),
|_, _| {},
)
.await
.unwrap()
});
tokio::time::timeout(Duration::from_secs(15), finalizer.entered.notified())
.await
.unwrap();
projector.settle(scenario.run_id).await;
let pending = petri_support::stored_projection(&scenario.pool, scenario.run_id)
.await
.unwrap()
.unwrap();
assert_eq!(pending.status, RunStatus::Running);
assert!(pending.conclusion.is_none());
assert!(!task.is_finished());
assert_view_equals_rebuild(&scenario.pool, scenario.run_id).await;
let summary = fabro_types::DiffSummary {
files_changed: 1,
additions: 2,
deletions: 0,
};
let head_sha = "0123456789012345678901234567890123456789";
let patch_blob = BlobHash::new(b"final patch");
let platform = PlatformRecordStore::new(scenario.pool.clone());
for record in [
PlatformRecord::Checkpoint(fabro_store::platform_records::CheckpointRecord {
execution: 0,
firing: 0,
attempt: Some(1),
workspace: None,
git_commit_sha: Some(head_sha.to_string()),
diff_summary: Some(summary),
patch_blob: Some(patch_blob),
operation: None,
}),
PlatformRecord::RunDiff(fabro_store::platform_records::RunDiffRecord {
base_sha: None,
head_sha: Some(head_sha.to_string()),
diff_summary: Some(summary),
patch_blob: Some(patch_blob),
}),
] {
platform
.append(&scenario.run_id, &record, None)
.await
.unwrap();
}
finalizer.release.add_permits(1);
let report = task.await.unwrap();
assert_eq!(report.execution_status, PetriRunStatus::Success);
assert_eq!(
report.status,
if rejection.is_some() {
PetriRunStatus::Failed
} else {
PetriRunStatus::Success
}
);
projector.settle(scenario.run_id).await;
let stored = petri_support::stored_projection(&scenario.pool, scenario.run_id)
.await
.unwrap()
.unwrap();
let expected = match rejection {
Some(_) => RunStatus::Failed {
reason: FailureReason::PublishFailed,
},
None => RunStatus::Succeeded {
reason: SuccessReason::Completed,
},
};
assert_eq!(stored.status, expected);
let conclusion = stored.conclusion.as_ref().unwrap();
assert_eq!(
conclusion
.failure
.as_ref()
.map(|failure| failure.detail.message.as_str()),
rejection
);
assert_eq!(
conclusion.status,
if rejection.is_some() {
StageOutcome::Failed {
retry_requested: false,
}
} else {
StageOutcome::Succeeded
}
);
assert!(
stored
.iter_stages()
.all(|(_, stage)| stage.state == StageState::Succeeded)
);
assert_eq!(conclusion.final_git_commit_sha.as_deref(), Some(head_sha));
assert_eq!(conclusion.diff.summary, Some(summary));
assert_eq!(
conclusion.diff.patch.as_deref(),
Some(fabro_types::format_blob_ref(&patch_blob).as_str())
);
assert!(!conclusion.stages.is_empty());
assert_view_equals_rebuild(&scenario.pool, scenario.run_id).await;
// A later worker lifecycle cannot replace the committed terminal
// result or conclusion, even if its status disagrees.
let other = match rejection {
Some(_) => RunStatus::Succeeded {
reason: SuccessReason::Completed,
},
None => RunStatus::Failed {
reason: FailureReason::PublishFailed,
},
};
PlatformRecordStore::new(scenario.pool.clone())
.append(
&scenario.run_id,
&PlatformRecord::RunLifecycle(
RunLifecycleRecord::new(if rejection.is_some() {
RunLifecycleKind::Succeeded
} else {
RunLifecycleKind::Failed
})
.with_status(other),
),
None,
)
.await
.unwrap();
projector.signal(scenario.run_id);
projector.settle(scenario.run_id).await;
let after = petri_support::stored_projection(&scenario.pool, scenario.run_id)
.await
.unwrap()
.unwrap();
assert_eq!(after.status, stored.status);
assert_eq!(
serde_json::to_value(after.conclusion).unwrap(),
serde_json::to_value(stored.conclusion).unwrap()
);
assert_view_equals_rebuild(&scenario.pool, scenario.run_id).await;
}
}