Tidy host failure handling and restore main's lockfile versions

Share the managed-run cleanup between persist_run_failure and
fail_managed_run, return the committed projection from
commit_host_failure, and keep the run's status when the failure cannot
be stored. Restore the cancelled-launch message and keep the spawn
error's cause in the launch failure.

Rebuild Cargo.lock from main with only the Petri bump so the Windows
crates stay on their newer versions. Import the projection module
rather than the function, wrap long tracing calls, and document the
finalization test fixture.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
This commit is contained in:
Scott Werner 2026-10-08 15:38:10 -04:00
parent af41c7c66c
commit b4aea6128d
6 changed files with 89 additions and 57 deletions

28
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",
@ -1820,7 +1820,7 @@ dependencies = [
"libc",
"option-ext",
"redox_users",
"windows-sys 0.59.0",
"windows-sys 0.61.2",
]
[[package]]
@ -1934,7 +1934,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "39cab71617ae0d63f51a36d69f866391735b51691dbda63cf6f96d042b63efeb"
dependencies = [
"libc",
"windows-sys 0.59.0",
"windows-sys 0.61.2",
]
[[package]]
@ -3774,7 +3774,7 @@ dependencies = [
"js-sys",
"log",
"wasm-bindgen",
"windows-core 0.61.2",
"windows-core 0.62.2",
]
[[package]]
@ -4584,7 +4584,7 @@ version = "0.50.3"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "7957b9740744892f114936ab4a57b3f487491bbeafaf8083688b16841a4240e5"
dependencies = [
"windows-sys 0.59.0",
"windows-sys 0.61.2",
]
[[package]]
@ -5752,7 +5752,7 @@ dependencies = [
"once_cell",
"socket2",
"tracing",
"windows-sys 0.59.0",
"windows-sys 0.60.2",
]
[[package]]
@ -6200,7 +6200,7 @@ dependencies = [
"errno 0.3.14",
"libc",
"linux-raw-sys",
"windows-sys 0.59.0",
"windows-sys 0.61.2",
]
[[package]]
@ -6259,7 +6259,7 @@ dependencies = [
"security-framework",
"security-framework-sys",
"webpki-root-certs",
"windows-sys 0.59.0",
"windows-sys 0.61.2",
]
[[package]]
@ -6907,7 +6907,7 @@ version = "1.4.8"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "c4db69cba1110affc0e9f7bcd48bbf87b3f4fc7c61fc9155afd4c469eb3d6c1b"
dependencies = [
"errno 0.3.14",
"errno 0.2.8",
"libc",
]
@ -7385,7 +7385,7 @@ dependencies = [
"getrandom 0.4.1",
"once_cell",
"rustix",
"windows-sys 0.59.0",
"windows-sys 0.61.2",
]
[[package]]
@ -7420,7 +7420,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "230a1b821ccbd75b185820a1f1ff7b14d21da1e442e22c0863ea5f08771a8874"
dependencies = [
"rustix",
"windows-sys 0.59.0",
"windows-sys 0.61.2",
]
[[package]]
@ -8467,7 +8467,7 @@ version = "0.1.11"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "c2a7b1c03c876122aa43f3020e6c3c3ee5c05081c9a00739faf7503aeba10d22"
dependencies = [
"windows-sys 0.59.0",
"windows-sys 0.61.2",
]
[[package]]

View file

@ -14,7 +14,7 @@
//! run of any length on a live token.
//!
//! 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:
//! 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

View file

@ -3511,59 +3511,59 @@ async fn reject_run_if_sandbox_provider_disabled(
/// 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.
/// 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.
pub(crate) async fn persist_run_failure(
state: &Arc<AppState>,
run_id: RunId,
reason: FailureReason,
message: String,
) {
if let Err(err) = commit_host_failure(state, run_id, reason, message).await {
error!(run_id = %run_id, error = %err, "Failed to record a host failure");
match commit_host_failure(state, run_id, reason, message).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);
}
Err(err) => error!(run_id = %run_id, error = %err, "Failed to record a host failure"),
}
// Resource cleanup is operational; it cannot substitute for a committed
// outcome if storage is unavailable.
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);
cleanup_worker_control_bus_for_run(state, run_id);
release_managed_run(state, run_id);
state.scheduler_notify.notify_one();
}
/// Append the host failure unless the run already ended, then settle the
/// managed run on whichever terminal result the store committed.
/// 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<()> {
let mut committed = run_records::projection(state, run_id)
) -> 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() {
run_records::lifecycle(state, run_id, run_records::failed(reason, message)).await?;
// The append waited for the projector, so the stored projection
// already folds it, or the finish that won the race.
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"
);
if committed.status.is_terminal() {
return 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);
Ok(())
// 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(
@ -3627,9 +3627,19 @@ fn fail_managed_run(state: &Arc<AppState>, run_id: RunId, reason: FailureReason,
managed_run.status = RunStatus::Failed { reason };
managed_run.error = Some(message);
}
}
drop(runs);
release_managed_run(state, run_id);
}
/// Drop the run's live worker state and controls, leaving its status alone.
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);
}
/// Fold one lifecycle record of the run's stream into the in-memory run:
@ -3824,7 +3834,12 @@ async fn fail_worker_launch(state: &Arc<AppState>, run_id: RunId, err: anyhow::E
FailureReason::LaunchFailed,
)
});
persist_run_failure(state, run_id, reason, error.to_string()).await;
let message = if reason == FailureReason::Cancelled {
"Run cancelled before worker launch completed".to_string()
} else {
collect_chain(&error).join(": ")
};
persist_run_failure(state, run_id, reason, message).await;
}
/// A worker that exited without recording the run's end left it failed,

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_result;
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;
@ -160,7 +160,9 @@ async fn append_records(
match writer.append(&log, &records).await {
Ok(()) => {
if log == LogId::Coordinator {
if let Some((status, failure)) = records.iter().find_map(finished_run_result) {
if let Some((status, failure)) =
records.iter().find_map(projection::finished_run_result)
{
settle_managed_run_at_finish(&state, id, status, failure);
}
}

View file

@ -1634,10 +1634,19 @@ impl ExecutionHooks for FabroHooks {
)
})?;
publisher.publish(&publication).await.map_err(|message| {
warn!(run_id = %self.run_id, error = %message, "the run's publication failed");
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");
info!(
run_id = %self.run_id,
branch = publication.run_branch,
sha = publication.head_sha,
"run published"
);
}
}
if self.inner.requires_run_finalization() {
@ -1651,7 +1660,11 @@ impl ExecutionHooks for FabroHooks {
// 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");
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

@ -72,6 +72,8 @@ pub struct TestRunRecords {
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>,