From 968c46d69baf2fedd9fddbd4d5d87bc60a1fce2a Mon Sep 17 00:00:00 2001 From: Bryan Helmkamp Date: Mon, 21 Sep 2026 12:08:48 -0400 Subject: [PATCH] Settle an in-process run at Petri's finish, not at its terminal record The in-process Petri path settled the run's in-memory managed run once Fabro's terminal lifecycle record was stored, after the engine returned. GET /runs/{id} reads the stored summary, which the projector ends at Petri's own `run.finished` record, a moment earlier, so a delete issued the moment the run read as ended could reach the delete precheck, which prefers the managed run, while it still said running, and was refused with 409 "cannot remove active run". Against the real engine the window hit eight times in thirty. The worker path settles its run at the worker's records endpoint, ahead of the store (#888). The in-process run now settles at the same record through the run store it executes over: a store whose coordinator appends settle the managed run at the `run.finished` record before the record reaches the store and the projector's signal, with the same finish mapping and settle the worker path uses. The settle is in memory only. The terminal lifecycle record stored once the engine returns refines the status and error and ends the run's live state as before, and stays the settle of a run that ends without an engine finish. The prune scenario no longer waits for the managed run to settle before its delete: the wait guarded only this window. Co-Authored-By: Claude Fable 5.1 --- .../fabro-server/src/server/petri_runs.rs | 252 +++++++++++++++++- .../fabro-server/tests/it/scenario/petri.rs | 92 +++++-- 2 files changed, 315 insertions(+), 29 deletions(-) diff --git a/lib/apps/fabro-server/src/server/petri_runs.rs b/lib/apps/fabro-server/src/server/petri_runs.rs index 91e54c5f7..5a54f187a 100644 --- a/lib/apps/fabro-server/src/server/petri_runs.rs +++ b/lib/apps/fabro-server/src/server/petri_runs.rs @@ -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, 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, 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, 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, run_id: RunId, status: RunStatus, error: Option) { 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, 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, + state: Arc, +} + +#[async_trait::async_trait] +impl RunStore for SettlingStore { + async fn open(&self, key: &RunKey, access: Access) -> Result, 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, + run_id: Option, + state: Arc, +} + +#[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, StoreError> { + self.inner.read(log).await + } + + async fn read_from(&self, log: &LogId, seq: u64) -> Result, StoreError> { + self.inner.read_from(log, seq).await + } + + async fn put_blob(&self, bytes: &[u8]) -> Result { + self.inner.put_blob(bytes).await + } + + async fn get_blob(&self, digest: Digest) -> Result>, 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, + state: Arc, + run_id: RunId, + seen: Mutex>>, + } + + #[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, StoreError> { + self.inner.read(log).await + } + + async fn put_blob(&self, bytes: &[u8]) -> Result { + self.inner.put_blob(bytes).await + } + + async fn get_blob(&self, digest: Digest) -> Result>, 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, RunId, SettlingLogs, Arc) { + 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, + 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 diff --git a/lib/apps/fabro-server/tests/it/scenario/petri.rs b/lib/apps/fabro-server/tests/it/scenario/petri.rs index 86300d554..f5da19a1f 100644 --- a/lib/apps/fabro-server/tests/it/scenario/petri.rs +++ b/lib/apps/fabro-server/tests/it/scenario/petri.rs @@ -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 { .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.