mirror of
https://github.com/fabro-sh/fabro.git
synced 2026-10-01 02:04:24 +00:00
Delete a run's sandboxes through Petri's lease ledger
Run deletion called the driver's `provider.delete(id)` under the run's `petri.run` scope, a delete of Fabro's own over a sandbox whose lease record Petri owns. It now goes the way `petri sandbox prune` goes: `fabro_petri::prune` builds the run's Petri runtime over the server's store (the run key, the run directory, the sandbox backend) and calls Petri's prune, which opens the run for writing, checks each lease's provider fingerprint, writes the delete intent and the tombstone beside the run's other records, and lets each provider remove its managed workspace, a host workspace included. A run a live process holds answers 409 unless the delete is forced; a lease Petri could not prune answers 409 with the problem text, or is warned and skipped under force or a delete that already started. The server drops the worker's handles before the prune, on the store instance the prune opens, so the lease a stopped worker held is released first. The projection reads only the coordinator and execution logs, so the resource records change nothing it reports. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
This commit is contained in:
parent
809b3891b5
commit
c7aa50c943
9 changed files with 507 additions and 61 deletions
|
|
@ -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<dyn Sandbox>` 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<dyn Sandbox>` 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
|
||||
|
|
|
|||
|
|
@ -23,7 +23,7 @@ use fabro_types::RunId;
|
|||
use tracing::debug;
|
||||
|
||||
pub(crate) struct PetriRuns {
|
||||
store: SqliteRunStore,
|
||||
store: Arc<SqliteRunStore>,
|
||||
/// The writer handle each worker holds open, by run and owner.
|
||||
handles: Mutex<HashMap<(RunId, OwnerId), Arc<dyn RunLogs>>>,
|
||||
}
|
||||
|
|
@ -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<SqliteRunStore> {
|
||||
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())
|
||||
|
|
|
|||
|
|
@ -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<String, String>) -> 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<Arc<dyn SandboxProvider>, 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<dyn SandboxProvider>,
|
||||
|
|
|
|||
|
|
@ -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::<Vec<_>>()
|
||||
.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))
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -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<Body> {
|
||||
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.
|
||||
|
|
|
|||
|
|
@ -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<RunOutcome, RunError> {
|
||||
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<SandboxBackend, RunError> {
|
||||
/// The sandbox backend for Fabro's provider kind; `None` for a kind Petri
|
||||
/// does not serve.
|
||||
pub(crate) fn backend(provider: &SandboxProviderKind) -> Option<SandboxBackend> {
|
||||
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
|
||||
}
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -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;
|
||||
|
|
|
|||
79
lib/components/fabro-petri/src/prune.rs
Normal file
79
lib/components/fabro-petri/src/prune.rs
Normal file
|
|
@ -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<dyn RunStore>,
|
||||
/// 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<PruneReport, PruneError> {
|
||||
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),
|
||||
})
|
||||
}
|
||||
196
lib/components/fabro-petri/tests/prune.rs
Normal file
196
lib/components/fabro-petri/tests/prune.rs
Normal file
|
|
@ -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<PathBuf> {
|
||||
fn walk(dir: &std::path::Path, root: &std::path::Path, out: &mut Vec<PathBuf>) {
|
||||
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/<id>` and its workspace under `work` there.
|
||||
fn workspace(run_dir: &std::path::Path) -> Option<PathBuf> {
|
||||
let scopes = run_dir.join("scopes");
|
||||
let mut entries: Vec<PathBuf> = 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<String> {
|
||||
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}"
|
||||
);
|
||||
}
|
||||
Loading…
Add table
Reference in a new issue