Merge pull request #890 from fabro-sh/settle-in-process-run-at-finish
Some checks are pending
Rust / Format (push) Waiting to run
Rust / Clippy (push) Waiting to run
Rust / Rustdoc (push) Waiting to run
Rust / Generated Docs (push) Waiting to run
Rust / Test (Linux) (push) Waiting to run
Rust / Sandbox providers (Docker) (push) Waiting to run
Rust / Test (macOS) (push) Waiting to run

Settle an in-process run at Petri's finish, not at its terminal record
This commit is contained in:
Bryan Helmkamp 2026-09-21 12:51:01 -04:00 • committed by GitHub
commit 7518045203
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
2 changed files with 315 additions and 29 deletions

View file

@ -19,7 +19,10 @@
//! run store in the server's database, so the scenario tests need no
//! worker binary; its questions go to an in-process control interviewer
//! the answer endpoint reaches directly, its secrets come from a snapshot
//! of the server's vault, and its blobs go to the server's blob store. No
//! 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
//! stage or agent event is projected either way, which is the read-side
//! item that follows.
//!
@ -41,12 +44,12 @@ use fabro_petri::controls::RunControls;
use fabro_petri::engine::{self, Conclusion, Execution, RunRequest};
use fabro_petri::hooks::HooksSpec;
use fabro_petri::interview::{Approval, FabroInterviewer};
use fabro_petri::petri::StoreError;
use fabro_petri::petri::{Access, Digest, LogId, Record, RunKey, RunLogs, RunStore, StoreError};
use fabro_petri::platform_records::SqlitePlatformRecords;
use fabro_petri::recovery::{self, Recovery, RecoveryRequest};
use fabro_petri::runtime::{self, RuntimeSpec};
use fabro_petri::secrets::VaultSecrets;
use fabro_petri::{SqliteRunStore, admission, run_graph};
use fabro_petri::{SqliteRunStore, admission, projection, run_graph};
use fabro_store::platform_records::{RunLifecycleKind, RunLifecycleRecord};
use fabro_types::settings::McpTransport;
use fabro_types::settings::run::{ApprovalMode, McpServerSettings, RunMode};
@ -457,9 +460,14 @@ pub(crate) async fn execute(state: Arc<AppState>, run_id: RunId) {
run_id: run_id.to_string(),
run_dir: run_dir.join("petri"),
execution,
store: state
.petri_projector
.observe_store(Arc::new(SqliteRunStore::new(state.db_pool.clone()))),
// The projector's signal follows each durable append; the managed
// run's settle at Petri's finish precedes it.
store: Arc::new(SettlingStore {
inner: state
.petri_projector
.observe_store(Arc::new(SqliteRunStore::new(state.db_pool.clone()))),
state: Arc::clone(&state),
}),
runtime: runtime_spec(&state, &eligible, dry_run),
provider: run_state.spec.settings.run.environment.provider.clone(),
cancel,
@ -498,10 +506,9 @@ pub(crate) async fn execute(state: Arc<AppState>, run_id: RunId) {
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 run reads as ended from the moment its terminal record is stored,
// so the managed run settles here, before the view catches up: a delete
// that arrives between the record and the settle otherwise refuses the
// run as active while the API already reports it ended.
// 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.
@ -639,7 +646,11 @@ async fn fail_before_execution(state: &Arc<AppState>, run_id: RunId, message: &s
finish(state, run_id, status, error);
}
/// Settle the managed run and release its scheduler slot.
/// 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.
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) {
@ -651,11 +662,230 @@ fn finish(state: &Arc<AppState>, run_id: RunId, status: RunStatus, error: Option
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.
struct SettlingStore {
inner: Arc<dyn RunStore>,
state: Arc<AppState>,
}
#[async_trait::async_trait]
impl RunStore for SettlingStore {
async fn open(&self, key: &RunKey, access: Access) -> Result<Arc<dyn RunLogs>, StoreError> {
let logs = self.inner.open(key, access).await?;
Ok(Arc::new(SettlingLogs {
inner: logs,
run_id: projection::run_id_of(key.as_str()),
state: Arc::clone(&self.state),
}))
}
}
/// One run's logs, whose coordinator appends settle the managed run at
/// Petri's finish before the records reach the store.
struct SettlingLogs {
inner: Arc<dyn RunLogs>,
run_id: Option<RunId>,
state: Arc<AppState>,
}
#[async_trait::async_trait]
impl RunLogs for SettlingLogs {
fn locator(&self) -> String {
self.inner.locator()
}
async fn append(&self, log: &LogId, records: &[Record]) -> Result<(), StoreError> {
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);
}
}
self.inner.append(log, records).await
}
async fn read(&self, log: &LogId) -> Result<Vec<Record>, StoreError> {
self.inner.read(log).await
}
async fn read_from(&self, log: &LogId, seq: u64) -> Result<Vec<Record>, StoreError> {
self.inner.read_from(log, seq).await
}
async fn put_blob(&self, bytes: &[u8]) -> Result<Digest, StoreError> {
self.inner.put_blob(bytes).await
}
async fn get_blob(&self, digest: Digest) -> Result<Option<Vec<u8>>, StoreError> {
self.inner.get_blob(digest).await
}
}
#[cfg(test)]
mod tests {
use std::sync::Mutex;
use chrono::Utc;
use fabro_petri::petri::OwnerId;
use fabro_types::settings::run::McpHttpProtocol;
use serde_json::json;
use super::*;
use crate::test_support::TestAppStateBuilder;
/// The run's logs as the store keeps them, recording what the server
/// held for the managed run at the moment each append reached them.
struct RecordingLogs {
inner: Arc<dyn RunLogs>,
state: Arc<AppState>,
run_id: RunId,
seen: Mutex<Vec<Option<RunStatus>>>,
}
#[async_trait::async_trait]
impl RunLogs for RecordingLogs {
fn locator(&self) -> String {
self.inner.locator()
}
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);
self.inner.append(log, records).await
}
async fn read(&self, log: &LogId) -> Result<Vec<Record>, StoreError> {
self.inner.read(log).await
}
async fn put_blob(&self, bytes: &[u8]) -> Result<Digest, StoreError> {
self.inner.put_blob(bytes).await
}
async fn get_blob(&self, digest: Digest) -> Result<Option<Vec<u8>>, StoreError> {
self.inner.get_blob(digest).await
}
}
/// A coordinator record, as the engine appends one.
fn coordinator_record(seq: u64, body: &serde_json::Value) -> Record {
Record {
seq,
recorded_at: 1_000,
record: json!({
"seq": seq,
"origin": "external",
"recorded_at": 1_000,
"body": body,
}),
}
}
/// A server with a managed run in flight in the server process, and
/// the run's logs as the in-process engine writes them: the settling
/// logs over the store, with a recorder between them.
async fn in_flight_run() -> (Arc<AppState>, RunId, SettlingLogs, Arc<RecordingLogs>) {
let state = TestAppStateBuilder::new().in_process_execution().build();
let run_id = RunId::new();
let run_dir = Storage::new(state.server_storage_dir())
.run_scratch(&run_id)
.root()
.to_path_buf();
state.runs.lock().expect("runs lock poisoned").insert(
run_id,
super::super::managed_run(
String::new(),
RunStatus::Running,
Utc::now(),
run_dir,
RunExecutionMode::Start,
),
);
let store = SqliteRunStore::new(state.db_pool.clone());
let logs = store
.open(&RunKey::new(run_id.to_string()), Access::Create {
owner: OwnerId::new("in-process"),
})
.await
.expect("the run is created in the store");
let recorder = Arc::new(RecordingLogs {
inner: logs,
state: Arc::clone(&state),
run_id,
seen: Mutex::new(Vec::new()),
});
let logs = SettlingLogs {
inner: Arc::clone(&recorder) as Arc<dyn RunLogs>,
run_id: Some(run_id),
state: Arc::clone(&state),
};
(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.
#[tokio::test]
async fn an_in_process_run_settles_before_its_finish_is_stored() {
let (state, run_id, logs, recorder) = in_flight_run().await;
logs.append(&LogId::Coordinator, &[coordinator_record(
0,
&json!({ "event": "run.started" }),
)])
.await
.expect("the record appends");
assert_eq!(
state.test_managed_run_status(&run_id),
Some(RunStatus::Running),
"a record that is not the finish leaves the run in flight"
);
logs.append(&LogId::Coordinator, &[coordinator_record(
1,
&json!({ "event": "run.finished", "status": "success" }),
)])
.await
.expect("the finish appends");
let succeeded = RunStatus::Succeeded {
reason: SuccessReason::Completed,
};
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"
);
assert_eq!(state.test_managed_run_status(&run_id), Some(succeeded));
}
/// 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]
async fn an_execution_logs_finish_does_not_settle_the_run() {
let (state, run_id, logs, _recorder) = in_flight_run().await;
logs.append(
&LogId::Execution(fabro_petri::petri::ExecutionId::new(1)),
&[coordinator_record(
0,
&json!({ "event": "run.finished", "status": "success" }),
)],
)
.await
.expect("the record appends");
assert_eq!(
state.test_managed_run_status(&run_id),
Some(RunStatus::Running)
);
}
/// Every transport of the catalog serializes in the inline shape Petri's
/// Fabro frontend reads, keyed by catalog id, with the timeouts as

View file

@ -802,9 +802,6 @@ async fn deleting_a_run_prunes_its_host_workspace_through_petri() {
let store = state.test_petri_run_store();
let key = RunKey::new(run_id.clone());
wait_for_free_lease(store, &key).await;
// The view reports the run ended from Petri's own finish, before the
// server settles the managed run the delete precheck reads.
wait_for_managed_settle(&state, &run_id).await;
// A live handle on the run, as its worker holds one, refuses the
// delete: Petri will not prune under a lease someone holds.
@ -861,21 +858,6 @@ async fn deleting_a_run_prunes_its_host_workspace_through_petri() {
}
/// Wait until no owner holds the run's lease.
/// Wait until the server's own map holds the run as ended.
async fn wait_for_managed_settle(state: &AppState, run_id: &str) {
let run_id: RunId = run_id.parse().expect("a run id");
for _ in 0..500 {
if state
.test_managed_run_status(&run_id)
.is_none_or(fabro_types::RunStatus::is_terminal)
{
return;
}
tokio::time::sleep(std::time::Duration::from_millis(10)).await;
}
panic!("the managed run did not settle");
}
async fn wait_for_free_lease(store: &SqliteRunStore, key: &RunKey) {
for _ in 0..500 {
if store.owner(key).await.expect("reads the lease").is_none() {
@ -894,6 +876,80 @@ fn delete(run_id: &str) -> Request<Body> {
.expect("delete request should build")
}
/// A delete issued the moment the run reads as ended is accepted while its
/// execution still tears down in the server process: the server settles
/// the managed run at Petri's own finish, the record the view ends the run
/// on, not at the terminal record it stores after the engine returns, so
/// the delete precheck does not refuse the run as active. The execution's
/// end after the delete brings nothing back.
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn a_delete_right_after_the_run_reads_ended_is_accepted() {
if host_plugin().is_none() {
return;
}
let workspace = tempfile::tempdir().expect("workspace tempdir");
let settings = settings_from_toml("_version = 1\n\n[run.environment]\nid = \"local\"\n");
let state = test_app_state_with_options(settings, 5);
let app = test_app_with_scheduler(Arc::clone(&state));
let version_id = register_version(&app, &[
("workflow.fabro", COMMAND_DOT),
("workflow.toml", PLAIN_SETTINGS),
])
.await;
let run_id =
create_and_start_run_from_intent(&app, intent(&version_id, workspace.path())).await;
let status = wait_for_run_status(&app, &run_id, &["succeeded", "failed"]).await;
assert_eq!(
status,
"succeeded",
"run: {}",
run_json(&app, &run_id).await
);
// The run's lease is free once its execution let go of the record: a
// delete under a held lease is refused for the lease, which the prune
// scenario covers, not for the managed run's status.
let store = state.test_petri_run_store();
let key = RunKey::new(run_id.clone());
wait_for_free_lease(store, &key).await;
let response = app
.clone()
.oneshot(delete(&run_id))
.await
.expect("delete route");
let status = response.status();
let detail = axum::body::to_bytes(response.into_body(), usize::MAX)
.await
.map(|bytes| String::from_utf8_lossy(&bytes).into_owned())
.unwrap_or_default();
assert_eq!(
status,
StatusCode::NO_CONTENT,
"the delete was refused: {detail}"
);
// The execution ends in the background after the delete and must leave
// the run gone.
for _ in 0..20 {
crate::helpers::response_status(
app.clone()
.oneshot(get(&format!("/runs/{run_id}")))
.await
.expect("run route"),
StatusCode::NOT_FOUND,
format!("GET /api/v1/runs/{run_id}"),
)
.await;
assert_eq!(
state.test_managed_run_status(&run_id.parse().expect("a run id")),
None,
"the execution's end brought the managed run back"
);
tokio::time::sleep(std::time::Duration::from_millis(10)).await;
}
}
/// The same on the Docker provider: the instance is the run's container,
/// with the image it runs and the container's workspace, so a reconnect
/// attaches to it on the daemon.