diff --git a/AGENTS.md b/AGENTS.md index 9dd8aab71..382eb86d6 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -246,7 +246,7 @@ Fabro is an AI-powered workflow orchestration platform. Workflows are defined as - **lib/packages/fabro-api-client** — Auto-generated TypeScript Axios client from OpenAPI spec ### Key design patterns -- **Direct sandbox access** — Petri creates every run sandbox through the sandbox driver and records its provider, id and working directory on the run (`RunSandboxInstance`); every Docker and Daytona sandbox carries the `petri.run` label. The server reaches a run's sandbox (the sandbox tab, Run Files, terminal, SSH, preview URLs, VNC, `fabro cp`, Ask Fabro, deletion) through `fabro-server/src/sandbox_access.rs`: it connects the record's provider itself, keys ownership on `petri.run`, and works on the driver's `Arc` facets (exec, filesystem, search, git, pty). There is no fabro-side sandbox trait; tests use `fabro_pebble_sandbox::test_support::MockSandbox` over the driver's scripted doubles. +- **Direct sandbox access** — Petri creates every run sandbox through the sandbox driver and records its provider, id and working directory on the run (`RunSandboxInstance`); every Docker and Daytona sandbox carries the `petri.run` label. The server reaches a run's sandbox (the sandbox tab, Run Files, terminal, SSH, preview URLs, VNC, `fabro cp`, Ask Fabro) through `fabro-server/src/sandbox_access.rs`: it connects the record's provider itself, keys ownership on `petri.run`, and works on the driver's `Arc` facets (exec, filesystem, search, git, pty). Deleting a run deletes its sandboxes through Petri's lease ledger (`fabro_petri::prune`, what `petri sandbox prune` does), not through a provider call of Fabro's own. There is no fabro-side sandbox trait; tests use `fabro_pebble_sandbox::test_support::MockSandbox` over the driver's scripted doubles. - **Graphviz graph workflows** — Stages and transitions defined as Graphviz graph attributes - **OpenAPI-first** — `fabro-api.yaml` drives Rust type + client generation (progenitor) and TypeScript client generation (openapi-generator) - **Checkpoint/resume** — Workflows can be paused, checkpointed, and resumed diff --git a/lib/apps/fabro-server/src/petri_runs.rs b/lib/apps/fabro-server/src/petri_runs.rs index c640aa6f6..ba7efbdb5 100644 --- a/lib/apps/fabro-server/src/petri_runs.rs +++ b/lib/apps/fabro-server/src/petri_runs.rs @@ -23,7 +23,7 @@ use fabro_types::RunId; use tracing::debug; pub(crate) struct PetriRuns { - store: SqliteRunStore, + store: Arc, /// The writer handle each worker holds open, by run and owner. handles: Mutex>>, } @@ -31,11 +31,20 @@ pub(crate) struct PetriRuns { impl PetriRuns { pub(crate) fn new(pool: DbPool) -> Self { Self { - store: SqliteRunStore::new(pool), + store: Arc::new(SqliteRunStore::new(pool)), handles: Mutex::default(), } } + /// The store the workers' handles are open on, for the server's own + /// work on a run's record (the sandbox prune at deletion). The same + /// instance matters: a handle dropped here has its lease release + /// awaited by this store's next open, so a prune right after + /// [`worker_exited`](Self::worker_exited) finds the lease free. + pub(crate) fn shared_store(&self) -> Arc { + Arc::clone(&self.store) + } + /// The Petri run key of a Fabro run. pub(crate) fn key(run_id: &RunId) -> RunKey { RunKey::new(run_id.to_string()) diff --git a/lib/apps/fabro-server/src/sandbox_access.rs b/lib/apps/fabro-server/src/sandbox_access.rs index 2e770428f..657d1d043 100644 --- a/lib/apps/fabro-server/src/sandbox_access.rs +++ b/lib/apps/fabro-server/src/sandbox_access.rs @@ -184,7 +184,7 @@ pub(crate) enum ConnectError { /// Connects the provider behind `kind`, unscoped: every sandbox on the /// backend is visible to it. Callers that act on a persisted id narrow it -/// with [`run_provider`]. +/// with [`scope_to_run`]. /// /// Bundled kinds link the driver's provider crates in process. `local` is /// the driver's Host provider with a fresh registry: a run's directory is @@ -278,21 +278,12 @@ fn is_petri_sandbox(labels: &BTreeMap) -> bool { labels.contains_key(PETRI_RUN_LABEL) } -/// The provider for `kind`, narrowed to the sandboxes of `run_id`: an -/// attach to or a delete of an id whose sandbox does not carry the run's -/// `petri.run` label is refused. The `local` kind is returned unscoped: a -/// host directory carries no labels, and nothing else shares the host's -/// directories with Fabro. -pub(crate) async fn run_provider( - kind: &SandboxProviderKind, - access: &ProviderAccess, - run_id: RunId, -) -> Result, ConnectError> { - let provider = connect_provider(kind, access).await?; - Ok(scope_to_run(kind, provider, run_id)) -} - -/// `provider` narrowed to the sandboxes of `run_id`; see [`run_provider`]. +/// `provider` narrowed to the sandboxes of `run_id`: an attach to an id +/// whose sandbox does not carry the run's `petri.run` label is refused. +/// The `local` kind is returned unscoped: a host directory carries no +/// labels, and nothing else shares the host's directories with Fabro. +/// Deletion does not come through here: a run's sandboxes are deleted +/// through Petri's lease ledger (`fabro_petri::prune`). fn scope_to_run( kind: &SandboxProviderKind, provider: Arc, diff --git a/lib/apps/fabro-server/src/server.rs b/lib/apps/fabro-server/src/server.rs index 9e42dc1f1..f420551ab 100644 --- a/lib/apps/fabro-server/src/server.rs +++ b/lib/apps/fabro-server/src/server.rs @@ -63,6 +63,7 @@ use fabro_llm::{ClientOptions, FabroClient}; use fabro_mcp_store::McpServerStore; use fabro_petri::controls::{RunControls, SteerError}; use fabro_petri::projector::Projector; +use fabro_petri::prune::{self, PruneError, PruneRequest}; use fabro_redact::redact_jsonl_line; use fabro_slack::client::{PostedMessage as SlackPostedMessage, SlackClient}; use fabro_slack::config::{ @@ -107,7 +108,6 @@ use fabro_workflow::{Error as WorkflowError, operations, pull_request}; use futures_util::future::join_all; use lithos_llm::catalog::ProviderId; use lithos_llm::types::Usage; -use sandbox_driver::SandboxId; use tempfile::NamedTempFile; use tokio::fs; use tokio::io::{AsyncBufReadExt, AsyncRead, AsyncWriteExt, BufReader}; @@ -2739,6 +2739,11 @@ async fn delete_run_internal( .await; } + // Whatever Petri run handles the run's worker held open over the API + // drop here, before its sandboxes are pruned through the lease ledger: + // the worker is gone or was told to stop above, and a lease it still + // held would refuse the prune. + state.petri_runs.worker_exited(id); let delete_outcome = delete_run_sandbox_resource(state, id, force).await?; if let Some(mut managed_run) = managed_run { @@ -2759,7 +2764,6 @@ async fn delete_run_internal( .delete_run(&id) .await .map_err(|err| ApiError::new(StatusCode::INTERNAL_SERVER_ERROR, err.to_string()))?; - state.petri_runs.worker_exited(id); state .petri_projector .delete_run(id) @@ -2834,51 +2838,95 @@ async fn delete_run_sandbox_resource( else { return Ok(SandboxDeleteOutcome::Cleaned); }; - let runtime = &record.runtime; if preserve { return Ok(SandboxDeleteOutcome::Preserved(DeleteRunResponse { deleted: true, sandbox_preserved: true, sandbox: DeleteRunSandbox { provider: record.provider, - id: runtime.id.clone(), + id: record.runtime.id, }, })); } - let access = state - .provider_access() - .await - .map_err(|err| ApiError::new(StatusCode::INTERNAL_SERVER_ERROR, err.to_string()))?; - // Deleted by id through the provider scoped to the run: a sandbox that - // no longer carries the run's `petri.run` label is refused, an id the - // provider no longer knows is already gone, and a designated host - // directory is left in place. - let deleted = async { - let provider = sandbox_access::run_provider(&record.provider, &access, id) - .await - .with_context(|| format!("Failed to connect to the {} provider", record.provider))?; - let sandbox_id = SandboxId::try_new(&runtime.id) - .with_context(|| format!("Invalid {} sandbox id", record.provider))?; - provider.delete(&sandbox_id, None).await.with_context(|| { - format!( - "Failed to delete {} sandbox '{}'", - record.provider, runtime.id - ) - }) - } + // Deleted through Petri's lease ledger, as `petri sandbox prune` does: + // Petri owns the lease record, checks the provider's fingerprint, and + // writes the intent and the tombstone beside the run's other records. + // The run directory is the worker's Petri run dir, where the host + // registry and a host workspace live; it is removed after this. + let run_dir = Storage::new(state.server_storage_dir()) + .run_scratch(&id) + .root() + .join("petri"); + let report = prune::prune(PruneRequest { + run_id: id.to_string(), + run_dir, + store: state.petri_runs.shared_store(), + provider: record.provider.clone(), + }) .await; - match deleted { - Ok(()) => Ok(SandboxDeleteOutcome::Cleaned), - Err(err) if force || delete_started => { - tracing::warn!( + match report { + Ok(report) if report.is_clean() => { + tracing::debug!( run_id = %id, - error = %format!("{err:#}"), - "Skipping failed sandbox provider delete during run deletion" + provider = %record.provider, + deleted = report.deleted.len(), + "Run sandboxes pruned through Petri" ); Ok(SandboxDeleteOutcome::Cleaned) } - Err(err) => Err(ApiError::new(StatusCode::CONFLICT, format!("{err:#}"))), + // A lease Petri could not prune keeps its pending intent, so a + // later prune tries again; a forced or restarted delete goes on + // without it. + Ok(report) => { + let problems = report + .problems + .iter() + .map(|(lease, problem)| format!("lease {lease}: {problem}")) + .collect::>() + .join("; "); + if force || delete_started { + tracing::warn!( + run_id = %id, + provider = %record.provider, + problems = %problems, + "Skipping the sandboxes Petri could not prune during run deletion" + ); + Ok(SandboxDeleteOutcome::Cleaned) + } else { + Err(ApiError::new( + StatusCode::CONFLICT, + format!("Failed to delete the run's sandboxes: {problems}"), + )) + } + } + // A live process still holds the run: only a forced delete leaves + // its sandboxes behind. + Err(error @ PruneError::RunHeld { .. }) => { + if force { + tracing::warn!( + run_id = %id, + error = %error, + "Skipping the sandbox prune of a held run during forced deletion" + ); + Ok(SandboxDeleteOutcome::Cleaned) + } else { + Err(ApiError::new(StatusCode::CONFLICT, error.to_string())) + } + } + Err(error) => { + let message = fabro_util::error::collect_chain(&error).join(": "); + if force || delete_started { + tracing::warn!( + run_id = %id, + error = %message, + "Skipping the failed sandbox prune during run deletion" + ); + Ok(SandboxDeleteOutcome::Cleaned) + } else { + Err(ApiError::new(StatusCode::CONFLICT, message)) + } + } } } diff --git a/lib/apps/fabro-server/tests/it/scenario/petri.rs b/lib/apps/fabro-server/tests/it/scenario/petri.rs index d3018c3ee..9e53c3381 100644 --- a/lib/apps/fabro-server/tests/it/scenario/petri.rs +++ b/lib/apps/fabro-server/tests/it/scenario/petri.rs @@ -27,6 +27,7 @@ use std::sync::Arc; use axum::body::Body; use axum::http::{Request, StatusCode}; use fabro_petri::engine::{self, RunStatus}; +use fabro_petri::petri::{Access, OwnerId, RunKey, RunStore as _}; use fabro_petri::{SqliteRunStore, projector}; use fabro_server::server::AppState; use fabro_server::test_support::{ @@ -757,6 +758,124 @@ async fn a_runs_projection_carries_its_host_sandbox_instance() { assert_eq!(run["sandbox"]["instance"]["runtime"]["id"], id, "{run}"); } +/// Deleting a run deletes its sandboxes through Petri's lease ledger: a +/// run whose lease a live process holds is refused with a conflict and +/// keeps its workspace, and once the lease is free the delete removes the +/// host workspace Petri kept along with the run. +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn deleting_a_run_prunes_its_host_workspace_through_petri() { + 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 + ); + let projection = settled_state(&state, &app, &run_id).await; + let working_directory = PathBuf::from( + projection["sandbox"]["instance"]["runtime"]["working_directory"] + .as_str() + .expect("the working directory"), + ); + assert!( + working_directory.is_dir(), + "the workspace is retained after the run: {}", + working_directory.display() + ); + + // The run's lease is free once its execution let go of the record. + let store = state.test_petri_run_store(); + let key = RunKey::new(run_id.clone()); + wait_for_free_lease(store, &key).await; + + // A live handle on the run, as its worker holds one, refuses the + // delete: Petri will not prune under a lease someone holds. + let held = store + .open(&key, Access::Write { + owner: OwnerId::new("worker-1"), + }) + .await + .expect("the worker takes the run"); + let refused = response_json( + app.clone() + .oneshot(delete(&run_id)) + .await + .expect("delete route"), + StatusCode::CONFLICT, + format!("DELETE /api/v1/runs/{run_id}"), + ) + .await; + assert!( + refused["errors"][0]["detail"] + .as_str() + .is_some_and(|detail| detail.contains("held by a live process")), + "the conflict names the held lease: {refused}" + ); + assert!( + working_directory.is_dir(), + "the refused delete left the workspace" + ); + drop(held); + wait_for_free_lease(store, &key).await; + + crate::helpers::response_status( + app.clone() + .oneshot(delete(&run_id)) + .await + .expect("delete route"), + StatusCode::NO_CONTENT, + format!("DELETE /api/v1/runs/{run_id}"), + ) + .await; + assert!( + !working_directory.exists(), + "the host provider removed the workspace Petri kept" + ); + 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; +} + +/// Wait until no owner holds the run's lease. +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() { + return; + } + tokio::time::sleep(std::time::Duration::from_millis(10)).await; + } + panic!("the run's lease was not released"); +} + +fn delete(run_id: &str) -> Request { + Request::builder() + .method("DELETE") + .uri(api(&format!("/runs/{run_id}"))) + .body(Body::empty()) + .expect("delete request should build") +} + /// 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. diff --git a/lib/components/fabro-petri/src/engine.rs b/lib/components/fabro-petri/src/engine.rs index 589336ad4..b322cce1e 100644 --- a/lib/components/fabro-petri/src/engine.rs +++ b/lib/components/fabro-petri/src/engine.rs @@ -170,7 +170,9 @@ pub enum Conclusion { /// Execute the run to its end and report what the record says. pub async fn run(request: RunRequest) -> Result { - let backend = backend(&request.provider)?; + let backend = backend(&request.provider).ok_or_else(|| RunError::UnsupportedProvider { + provider: request.provider.clone(), + })?; let key = RunKey::new(request.run_id.as_str()); let mut options = RunOptions::new(&request.run_dir); options.run_key = Some(key.clone()); @@ -381,18 +383,17 @@ fn error_chain(error: &RunError) -> String { parts.join(": ") } -/// The sandbox backend for Fabro's provider kind. -fn backend(provider: &SandboxProviderKind) -> Result { +/// The sandbox backend for Fabro's provider kind; `None` for a kind Petri +/// does not serve. +pub(crate) fn backend(provider: &SandboxProviderKind) -> Option { if *provider == SandboxProviderKind::LOCAL { - Ok(SandboxBackend::Host) + Some(SandboxBackend::Host) } else if *provider == SandboxProviderKind::DOCKER { - Ok(SandboxBackend::Docker) + Some(SandboxBackend::Docker) } else if *provider == SandboxProviderKind::DAYTONA { - Ok(SandboxBackend::Daytona) + Some(SandboxBackend::Daytona) } else { - Err(RunError::UnsupportedProvider { - provider: provider.clone(), - }) + None } } diff --git a/lib/components/fabro-petri/src/lib.rs b/lib/components/fabro-petri/src/lib.rs index 3395031d3..80b73dce8 100644 --- a/lib/components/fabro-petri/src/lib.rs +++ b/lib/components/fabro-petri/src/lib.rs @@ -54,7 +54,9 @@ //! - [`fork`]: a run seeded from another's records up to a checkpoint's //! position, over Petri's `host::fork_from`, with the kept checkpoints, their //! snapshots and the run branch carried over: what rewind, fork and retry are -//! built on. +//! built on; +//! - [`prune`]: a run's sandboxes deleted through Petri's lease ledger, as +//! `petri sandbox prune` deletes them, when Fabro deletes the run. //! //! The Petri packages are pinned by revision in the workspace `Cargo.toml` //! under `petri_*` keys. @@ -74,6 +76,7 @@ pub mod petri; pub mod platform_records; pub mod projection; pub mod projector; +pub mod prune; pub mod recovery; pub mod run_graph; pub mod run_store; diff --git a/lib/components/fabro-petri/src/prune.rs b/lib/components/fabro-petri/src/prune.rs new file mode 100644 index 000000000..177ac7154 --- /dev/null +++ b/lib/components/fabro-petri/src/prune.rs @@ -0,0 +1,79 @@ +//! A run's sandboxes deleted through Petri's lease ledger. +//! +//! Petri records every sandbox a run creates as a lease: the provider, the +//! provider's id for the resource, and the fingerprint of the backend it +//! lives on. When Fabro deletes a run, its sandboxes go the way `petri +//! sandbox prune` deletes them, over the store the run's records live in, +//! rather than through a provider call of Fabro's own: Petri opens the run +//! for writing, so a live worker that still holds the lease refuses the +//! delete; it checks each lease's fingerprint against the plugin it +//! launches, so a changed daemon or account is a problem to report, never +//! a delete on another backend; it writes the delete intent before the +//! provider call and the tombstone after, beside the run's other records; +//! and each provider removes its sandbox's managed workspace, a host +//! workspace under the run directory included. +//! +//! The runtime a prune runs on is the run's as [`engine`](crate::engine) +//! assembles it, reduced to what a prune reads: the store, the run key, the +//! run directory (where Petri's host registry and action-host markers are) +//! and the sandbox backend. No step registry, frontend or model client +//! takes part. + +use std::path::PathBuf; +use std::sync::Arc; + +use fabro_types::SandboxProviderKind; +pub use petri_execution::prune::PruneReport; +use petri_execution::prune::{self as petri_prune}; +use petri_execution::{RunKey, RunStore}; +use petri_runtime::{RunOptions, Runtime}; + +use crate::engine; + +/// One run whose sandboxes are to be deleted. +pub struct PruneRequest { + /// The Fabro run id, which is Petri's run key. + pub run_id: String, + /// Where the run's worker ran Petri: its host registry and action-host + /// markers are under it, and so is a host workspace. + pub run_dir: PathBuf, + /// The run's durable record. + pub store: Arc, + /// The sandbox provider Fabro resolved for the run's environment. + pub provider: SandboxProviderKind, +} + +/// Why a run's sandboxes could not be pruned. +#[derive(Debug, thiserror::Error)] +pub enum PruneError { + #[error("the run's sandbox provider `{provider}` is not one Petri serves")] + UnsupportedProvider { provider: SandboxProviderKind }, + /// A live process holds the run's lease: pruning under it would delete + /// the sandboxes it is using. + #[error("run {locator} is held by a live process; stop it first")] + RunHeld { locator: String }, + #[error("the run's sandboxes could not be pruned")] + Petri(#[source] petri_prune::PruneError), +} + +/// Delete every sandbox the run still holds, through Petri's lease ledger. +/// The report says what was deleted, what needed nothing, and which leases +/// could not be pruned and why; their records keep the pending intent, so +/// the next prune tries again. +pub async fn prune(request: PruneRequest) -> Result { + let backend = + engine::backend(&request.provider).ok_or_else(|| PruneError::UnsupportedProvider { + provider: request.provider.clone(), + })?; + let mut options = RunOptions::new(&request.run_dir); + options.run_key = Some(RunKey::new(request.run_id.as_str())); + options.retention = engine::RETENTION; + options.sandbox.backend = backend; + let runtime = Runtime::bare().store(request.store).options(options); + petri_prune::prune(&runtime) + .await + .map_err(|error| match error { + petri_prune::PruneError::RunHeld(locator) => PruneError::RunHeld { locator }, + other => PruneError::Petri(other), + }) +} diff --git a/lib/components/fabro-petri/tests/prune.rs b/lib/components/fabro-petri/tests/prune.rs new file mode 100644 index 000000000..33e9ec18f --- /dev/null +++ b/lib/components/fabro-petri/tests/prune.rs @@ -0,0 +1,196 @@ +//! A finished run's sandboxes deleted through Petri's lease ledger: the +//! host workspace the run kept is removed by its provider and its lease +//! tombstoned in the run's record, a second prune has nothing to do, and a +//! run a live handle holds is refused. +//! +//! The run takes its scope through the sandbox-driver host plugin, so the +//! test skips, and says why, when the executable is not found. + +mod support; + +use std::path::PathBuf; +use std::sync::Arc; + +use fabro_petri::check::Launch; +use fabro_petri::engine::{self, RunStatus}; +use fabro_petri::prune::{PruneError, PruneRequest, prune}; +use fabro_petri::runtime::RuntimeSpec; +use fabro_petri::{SqliteRunStore, petri}; +use fabro_store::test_support; +use fabro_types::SandboxProviderKind; +use support::{Silent, admit, host_plugin, no_questions, run_request}; + +/// A command-only workflow whose one stage writes a file into its +/// workspace. +const WORKFLOW: &str = r#"digraph Command { + graph [goal="Leave a file behind"] + start [shape=Mdiamond] + exit [shape=Msquare] + write [shape=parallelogram, script="echo kept > kept.txt"] + start -> write -> exit +}"#; + +/// Every file under `root`, relative to it, in path order. +fn files_under(root: &std::path::Path) -> Vec { + fn walk(dir: &std::path::Path, root: &std::path::Path, out: &mut Vec) { + let Ok(entries) = std::fs::read_dir(dir) else { + return; + }; + for entry in entries { + let path = entry.expect("an entry reads").path(); + if path.is_dir() { + walk(&path, root, out); + } else { + out.push( + path.strip_prefix(root) + .expect("a path under the root") + .to_path_buf(), + ); + } + } + } + let mut out = Vec::new(); + walk(root, root, &mut out); + out.sort(); + out +} + +/// The one host workspace under the run directory: the host backend keeps a +/// scope's state under `scopes/` and its workspace under `work` there. +fn workspace(run_dir: &std::path::Path) -> Option { + let scopes = run_dir.join("scopes"); + let mut entries: Vec = std::fs::read_dir(&scopes) + .ok()? + .map(|entry| entry.expect("an entry reads").path()) + .collect(); + assert!(entries.len() <= 1, "one scope at most: {entries:?}"); + entries.pop().map(|scope| scope.join("work")) +} + +/// The state of each lease in the run's resource log, latest record per +/// lease. +async fn lease_states(store: &SqliteRunStore, run_id: &str) -> Vec { + let logs = petri::RunStore::open(store, &petri::RunKey::new(run_id), petri::Access::Read) + .await + .expect("the run opens for reading"); + let records = logs + .read(&petri::LogId::Resources) + .await + .expect("the resource log reads"); + let mut latest = std::collections::BTreeMap::new(); + for record in records { + let lease = record.record["body"]["lease"].clone(); + let state = record.record["body"]["state"] + .as_str() + .expect("a lease state") + .to_owned(); + latest.insert(lease.to_string(), state); + } + latest.into_values().collect() +} + +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn a_finished_runs_host_workspace_is_deleted_once_and_a_held_run_is_refused() { + if host_plugin().is_none() { + return; + } + let root = tempfile::tempdir().expect("a temp dir"); + let run_dir = root.path().join("run"); + let pool = test_support::in_memory_pool_with(&[ + fabro_db::BLOBS_MIGRATION_SQL, + fabro_db::PETRI_RECORDS_MIGRATION_SQL, + ]); + let store = Arc::new(SqliteRunStore::new(pool.clone())); + let runtime = RuntimeSpec::default(); + let graphs = admit( + &[ + ("workflow.fabro", WORKFLOW), + ("workflow.toml", support::SETTINGS), + ], + Launch::default(), + &runtime, + ); + let request = run_request( + "prune", + &run_dir, + graphs, + store.clone(), + runtime, + no_questions(Arc::new(Silent)), + ); + let outcome = engine::run(request).await.expect("the run ends"); + assert_eq!(outcome.status, RunStatus::Success, "{outcome:?}"); + + // Retention kept the workspace, and its lease is live in the record. + let kept = workspace(&run_dir).expect("the run's workspace is retained"); + let files = files_under(&kept); + assert!( + files + .iter() + .any(|file| file.file_name().is_some_and(|name| name == "kept.txt")), + "the stage's file is in the retained workspace: {files:?}" + ); + assert_eq!(lease_states(&store, "prune").await, ["stopped"]); + + let request = || PruneRequest { + run_id: "prune".to_string(), + run_dir: run_dir.clone(), + store: store.clone(), + provider: SandboxProviderKind::LOCAL, + }; + let report = prune(request()).await.expect("the run prunes"); + assert!(report.is_clean(), "{report:?}"); + assert_eq!(report.deleted.len(), 1, "{report:?}"); + assert!( + !kept.exists(), + "the host provider removed its managed workspace; left: {:?}", + files_under(&kept) + ); + assert!( + run_dir.is_dir(), + "the run directory itself is the caller's to remove" + ); + assert_eq!( + lease_states(&store, "prune").await, + ["deleted"], + "the tombstone is in the run's record" + ); + + // A second prune finds only the tombstone. + let again = prune(request()).await.expect("the run prunes again"); + assert!(again.is_clean() && again.deleted.is_empty(), "{again:?}"); + assert_eq!(again.clean.len(), 1, "{again:?}"); + + // A live handle on the run holds its lease; the prune is refused. + let held = petri::RunStore::open( + store.as_ref(), + &petri::RunKey::new("prune"), + petri::Access::Write { + owner: petri::OwnerId::new("worker-1"), + }, + ) + .await + .expect("the worker takes the run"); + let error = prune(request()) + .await + .expect_err("a held run is not pruned"); + assert!(matches!(error, PruneError::RunHeld { .. }), "{error}"); + drop(held); +} + +#[tokio::test] +async fn a_provider_petri_does_not_serve_is_refused_before_the_store_is_opened() { + let store = Arc::new(petri_store::MemoryRunStore::new()); + let error = prune(PruneRequest { + run_id: "e2b-run".to_string(), + run_dir: std::env::temp_dir().join("fabro-petri-prune-e2b"), + store, + provider: SandboxProviderKind::try_new("e2b").expect("a valid kind"), + }) + .await + .expect_err("an unknown backend cannot be pruned"); + assert!( + matches!(error, PruneError::UnsupportedProvider { ref provider } if provider.as_str() == "e2b"), + "{error}" + ); +}