From 3c14bb866494851a99d1d377fadfd6f77cf839c5 Mon Sep 17 00:00:00 2001 From: John Myers <9696606+johntmyers@users.noreply.github.com> Date: Wed, 30 Sep 2026 11:55:07 -0700 Subject: [PATCH] fix(gateway): delete finalized ephemeral sandboxes while connected Start driver cleanup after terminal finalization and retain disconnect fallback. Add detached success and failure e2e coverage across supervisor-based drivers. Closes #3938 Signed-off-by: John Myers <9696606+johntmyers@users.noreply.github.com> --- .../skills/launch-openshell-gator/SKILL.md | 2 +- crates/openshell-server/src/compute/mod.rs | 118 ++++++++++- .../src/supervisor_session.rs | 14 +- e2e/rust/Cargo.toml | 5 + e2e/rust/e2e-podman.sh | 1 + e2e/rust/e2e-vm.sh | 1 + e2e/rust/tests/ephemeral_cleanup.rs | 189 ++++++++++++++++++ 7 files changed, 322 insertions(+), 8 deletions(-) create mode 100644 e2e/rust/tests/ephemeral_cleanup.rs diff --git a/.agents/skills/launch-openshell-gator/SKILL.md b/.agents/skills/launch-openshell-gator/SKILL.md index 80a84ca744..dc29130659 100644 --- a/.agents/skills/launch-openshell-gator/SKILL.md +++ b/.agents/skills/launch-openshell-gator/SKILL.md @@ -158,7 +158,7 @@ sandbox_name="gator-pr-${pr_number}-supervised" "Review and monitor PR #${pr_number} through the gator-gate workflow. Scope this invocation only to PR #${pr_number}." ``` -The launcher queries the gateway's selected compute driver, builds the gator image in the matching Docker or Podman image store, stages the immutable payload, imports provider profiles, configures provider credentials and refresh, and starts the agent supervisor as the sandbox's canonical main process. The detached main process survives loss of the host CLI connection and reconnects to a restarted gateway. Unless `--keep` is set, the sandbox is marked ephemeral so the gateway deletes it after the supervisor exits. `CONTAINER_ENGINE`, when set, must match the gateway driver. +The launcher queries the gateway's selected compute driver, builds the gator image in the matching Docker or Podman image store, stages the immutable payload, imports provider profiles, configures provider credentials and refresh, and starts the agent supervisor as the sandbox's canonical main process. The detached main process survives loss of the host CLI connection and reconnects to a restarted gateway. Unless `--keep` is set, the sandbox is marked ephemeral so the gateway deletes it after the canonical main process exits and its terminal result is finalized. `CONTAINER_ENGINE`, when set, must match the gateway driver. The launcher streams image-build and provisioning output until the detached workload is ready, then exits. Use `openshell logs ` or the TUI for runtime output. diff --git a/crates/openshell-server/src/compute/mod.rs b/crates/openshell-server/src/compute/mod.rs index da0ebffef0..fb426ab3fc 100644 --- a/crates/openshell-server/src/compute/mod.rs +++ b/crates/openshell-server/src/compute/mod.rs @@ -2133,6 +2133,7 @@ impl ComputeRuntime { } } + #[cfg(test)] pub(crate) async fn delete_sandbox( &self, workspace: &str, @@ -4723,6 +4724,39 @@ impl ComputeRuntime { Ok(()) } + /// Start ephemeral cleanup only after the finalize RPC has recorded the + /// terminal result and marked its supervisor session finalized. + pub async fn cleanup_finalized_ephemeral_sandbox( + &self, + sandbox_id: &str, + instance_id: &str, + ) -> Result<(), String> { + let _guard = self.sync_lock.lock().await; + let Some(sandbox) = self + .store + .get_message::(sandbox_id) + .await + .map_err(|error| error.to_string())? + else { + return Ok(()); + }; + let phase = SandboxPhase::try_from(sandbox.phase()).unwrap_or(SandboxPhase::Unknown); + if phase != SandboxPhase::Completed && !is_failed_main_process_result(&sandbox) { + return Ok(()); + } + let Some(status) = sandbox.status.as_ref() else { + return Ok(()); + }; + if status.exit_code.is_none() + || (!status.main_process_instance_id.is_empty() + && status.main_process_instance_id != instance_id) + { + return Ok(()); + } + self.schedule_ephemeral_sandbox_delete(&sandbox); + Ok(()) + } + fn schedule_ephemeral_sandbox_delete(&self, sandbox: &Sandbox) { if provisioning_deadline::timed_out(sandbox) { return; @@ -4738,10 +4772,10 @@ impl ComputeRuntime { } let runtime = self.clone(); - let workspace = sandbox.object_workspace().to_string(); + let sandbox_id = sandbox.object_id().to_string(); let name = sandbox.object_name().to_string(); tokio::spawn(async move { - if let Err(error) = runtime.delete_sandbox(&workspace, &name).await { + if let Err(error) = runtime.delete_sandbox_by_id(&sandbox_id, &name).await { tracing::warn!( sandbox_name = %name, error = %error, @@ -9298,7 +9332,83 @@ mod tests { .unwrap(); assert_eq!(driver.delete_calls(), 0); runtime - .supervisor_session_disconnected("sb-1", true) + .cleanup_finalized_ephemeral_sandbox("sb-1", "instance-1") + .await + .unwrap(); + tokio::time::timeout(Duration::from_secs(1), async { + while driver.delete_calls() == 0 { + tokio::task::yield_now().await; + } + }) + .await + .expect("terminal finalization should delete before the supervisor disconnects"); + } + + #[tokio::test] + async fn finalized_ephemeral_cleanup_skips_retained_and_restarting_sandboxes() { + for (retention, restart_policy, exit_code) in [ + (None, SandboxRestartPolicy::Never, 0), + (Some("ephemeral"), SandboxRestartPolicy::OnFailure, 9), + ] { + let driver = ControlledDriver::new(); + let runtime = test_runtime(driver.clone()).await; + let mut sandbox = sandbox_record("sb-1", "sandbox-a", SandboxPhase::Provisioning); + if let Some(retention) = retention { + sandbox.metadata.as_mut().unwrap().annotations.insert( + "openshell.nvidia.com/retention".to_string(), + retention.to_string(), + ); + } + sandbox.spec = Some(SandboxSpec { + restart_policy: restart_policy as i32, + ..Default::default() + }); + runtime.store.put_message(&sandbox).await.unwrap(); + runtime + .supervisor_session_connected("sb-1", "instance-1") + .await + .unwrap(); + runtime + .report_main_process_exit("sb-1", "instance-1", exit_code) + .await + .unwrap(); + runtime + .finalize_main_process_exit("sb-1", "instance-1") + .await + .unwrap(); + runtime + .cleanup_finalized_ephemeral_sandbox("sb-1", "instance-1") + .await + .unwrap(); + tokio::task::yield_now().await; + assert_eq!(driver.delete_calls(), 0); + } + } + + #[tokio::test] + async fn finalized_failed_ephemeral_sandbox_deletes_while_connected() { + let driver = ControlledDriver::new(); + let runtime = test_runtime(driver.clone()).await; + let mut sandbox = sandbox_record("sb-1", "sandbox-a", SandboxPhase::Provisioning); + sandbox.metadata.as_mut().unwrap().annotations.insert( + "openshell.nvidia.com/retention".to_string(), + "ephemeral".to_string(), + ); + runtime.store.put_message(&sandbox).await.unwrap(); + runtime + .supervisor_session_connected("sb-1", "instance-1") + .await + .unwrap(); + runtime + .report_main_process_exit("sb-1", "instance-1", 9) + .await + .unwrap(); + runtime + .finalize_main_process_exit("sb-1", "instance-1") + .await + .unwrap(); + runtime + .cleanup_finalized_ephemeral_sandbox("sb-1", "instance-1") .await .unwrap(); tokio::time::timeout(Duration::from_secs(1), async { @@ -9307,7 +9417,7 @@ mod tests { } }) .await - .expect("terminal finalization should release ephemeral cleanup"); + .expect("failed canonical main should delete its ephemeral sandbox"); } #[tokio::test] diff --git a/crates/openshell-server/src/supervisor_session.rs b/crates/openshell-server/src/supervisor_session.rs index 710ecc1406..847d625a28 100644 --- a/crates/openshell-server/src/supervisor_session.rs +++ b/crates/openshell-server/src/supervisor_session.rs @@ -2072,10 +2072,18 @@ pub async fn handle_finalize_main_process_exit( .finalize_main_process_exit(&report.sandbox_id, &report.instance_id) .await .map_err(Status::failed_precondition)?; - if !state + let session_finalized = state .supervisor_sessions - .finalize_main_process_exit(&report.sandbox_id) - { + .finalize_main_process_exit(&report.sandbox_id); + // The session can close between durable result validation and this mark. + // Schedule cleanup in either case so a disconnect with an unfinalized + // in-memory session cannot strand the ephemeral sandbox. + state + .compute + .cleanup_finalized_ephemeral_sandbox(&report.sandbox_id, &report.instance_id) + .await + .map_err(Status::internal)?; + if !session_finalized { return Err(Status::failed_precondition( "supervisor session is not connected", )); diff --git a/e2e/rust/Cargo.toml b/e2e/rust/Cargo.toml index 4d0522ea86..e26ba1b107 100644 --- a/e2e/rust/Cargo.toml +++ b/e2e/rust/Cargo.toml @@ -93,6 +93,11 @@ name = "local_driver_token_restart" path = "tests/local_driver_token_restart.rs" required-features = ["e2e"] +[[test]] +name = "ephemeral_cleanup" +path = "tests/ephemeral_cleanup.rs" +required-features = ["e2e"] + [[test]] name = "podman_gateway_start" path = "tests/podman_gateway_start.rs" diff --git a/e2e/rust/e2e-podman.sh b/e2e/rust/e2e-podman.sh index bbbcfe6fa9..c0deb2c750 100755 --- a/e2e/rust/e2e-podman.sh +++ b/e2e/rust/e2e-podman.sh @@ -22,6 +22,7 @@ source "${ROOT}/e2e/support/conformance.sh" # stabilized and can be added here. PODMAN_CI_TESTS=( bypass_detection + ephemeral_cleanup core_dump_hardening credential_gating default_image diff --git a/e2e/rust/e2e-vm.sh b/e2e/rust/e2e-vm.sh index 6fd355170b..a30f5283d3 100755 --- a/e2e/rust/e2e-vm.sh +++ b/e2e/rust/e2e-vm.sh @@ -411,6 +411,7 @@ run_e2e_test() { if [ -n "${E2E_TEST_OVERRIDE}" ]; then run_e2e_test "${E2E_TEST_OVERRIDE}" else + run_e2e_test ephemeral_cleanup run_e2e_test host_gateway_alias run_e2e_test vm_overlay run_e2e_test vm_gateway_start diff --git a/e2e/rust/tests/ephemeral_cleanup.rs b/e2e/rust/tests/ephemeral_cleanup.rs new file mode 100644 index 0000000000..eca4bb1d68 --- /dev/null +++ b/e2e/rust/tests/ephemeral_cleanup.rs @@ -0,0 +1,189 @@ +// SPDX-FileCopyrightText: Copyright (c) 2025-2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved. +// SPDX-License-Identifier: Apache-2.0 + +//! Detached ephemeral lifecycle coverage shared by supervisor-based drivers. + +#![cfg(feature = "e2e")] + +use std::path::PathBuf; +use std::process::{Command, Stdio}; +use std::time::Duration; + +use openshell_e2e::harness::binary::{openshell_bin, openshell_cmd}; +use openshell_e2e::harness::cli::run_cli; +use openshell_e2e::harness::container::{ContainerEngine, e2e_driver}; +use openshell_e2e::harness::sandbox::{E2E_WORKLOAD_IMAGE, unique_sandbox_name}; +use serial_test::serial; +use tokio::time::{Instant, sleep}; + +const CLEANUP_TIMEOUT: Duration = Duration::from_secs(90); + +struct DeleteOnFailure { + name: String, + armed: bool, +} + +impl Drop for DeleteOnFailure { + fn drop(&mut self) { + if self.armed { + let _ = Command::new(openshell_bin()) + .args(["sandbox", "delete", &self.name]) + .stdout(Stdio::null()) + .stderr(Stdio::null()) + .status(); + } + } +} + +fn command_output(mut command: Command) -> Result { + let output = command + .output() + .map_err(|error| format!("run driver resource query: {error}"))?; + let stdout = String::from_utf8_lossy(&output.stdout); + let stderr = String::from_utf8_lossy(&output.stderr); + if !output.status.success() { + return Err(format!( + "driver resource query failed ({}): {stdout}{stderr}", + output.status + )); + } + Ok(stdout.trim().to_string()) +} + +fn driver_resources_present(sandbox_id: &str) -> Result { + match e2e_driver().as_deref() { + Some("docker" | "podman") => { + let engine = ContainerEngine::from_env()?; + let mut command = engine.command(); + command.args([ + "ps", + "-aq", + "--filter", + &format!("label=openshell.ai/sandbox-id={sandbox_id}"), + ]); + Ok(!command_output(command)?.is_empty()) + } + Some("kubernetes") => { + let mut command = Command::new("kubectl"); + command.args([ + "get", + "pods,sandboxes.agents.x-k8s.io", + "--all-namespaces", + "--selector", + &format!("openshell.ai/sandbox-id={sandbox_id}"), + "--output=name", + ]); + Ok(!command_output(command)?.is_empty()) + } + Some("vm") => { + let state_dir = std::env::var_os("OPENSHELL_E2E_VM_STATE_DIR") + .map(PathBuf::from) + .ok_or("OPENSHELL_E2E_VM_STATE_DIR must be set for VM resource checks")?; + Ok(state_dir.join("sandboxes").join(sandbox_id).exists()) + } + other => Err(format!( + "unsupported e2e driver for ephemeral cleanup: {other:?}" + )), + } +} + +async fn run_detached_ephemeral_cleanup(exit_code: i32) -> Result<(), String> { + let name = unique_sandbox_name(); + let mut cleanup = DeleteOnFailure { + name: name.clone(), + armed: true, + }; + let release_path = format!("/sandbox/.ephemeral-release-{name}"); + let script = + format!("while [ ! -e '{release_path}' ]; do sleep 0.1; done; sleep 1; exit {exit_code}"); + let mut create = openshell_cmd(); + create.args([ + "sandbox", + "create", + "--name", + &name, + "--from", + E2E_WORKLOAD_IMAGE, + "--no-keep", + "--detach", + "--", + "sh", + "-c", + &script, + ]); + let created = create + .output() + .await + .map_err(|error| format!("create detached ephemeral sandbox: {error}"))?; + if !created.status.success() { + return Err(format!( + "create detached ephemeral sandbox failed ({}): {}{}", + created.status, + String::from_utf8_lossy(&created.stdout), + String::from_utf8_lossy(&created.stderr) + )); + } + + let (details, get_code) = run_cli(&["sandbox", "get", &name, "--output", "json"]).await; + if get_code != 0 { + return Err(format!("get detached sandbox failed: {details}")); + } + let details: serde_json::Value = + serde_json::from_str(&details).map_err(|error| format!("parse sandbox JSON: {error}"))?; + let sandbox_id = details["id"] + .as_str() + .ok_or_else(|| format!("sandbox has no id: {details}"))?; + if !driver_resources_present(sandbox_id)? { + return Err(format!( + "driver resources were never observed for sandbox {name} ({sandbox_id})" + )); + } + + let (release, release_code) = run_cli(&[ + "sandbox", + "exec", + "--name", + &name, + "--no-tty", + "--no-login-shell", + "--", + "touch", + &release_path, + ]) + .await; + if release_code != 0 { + return Err(format!("release canonical process failed: {release}")); + } + + let deadline = Instant::now() + CLEANUP_TIMEOUT; + loop { + let (names, list_code) = run_cli(&["sandbox", "list", "--names"]).await; + if list_code != 0 { + return Err(format!("list sandboxes failed: {names}")); + } + let record_present = names.lines().any(|line| line.trim() == name); + let resources_present = driver_resources_present(sandbox_id)?; + if !record_present && !resources_present { + cleanup.armed = false; + return Ok(()); + } + if Instant::now() >= deadline { + return Err(format!( + "ephemeral sandbox {name} ({sandbox_id}) remained after {CLEANUP_TIMEOUT:?}: record_present={record_present}, resources_present={resources_present}" + )); + } + sleep(Duration::from_millis(500)).await; + } +} + +#[tokio::test] +#[serial(ephemeral_cleanup)] +async fn detached_ephemeral_success_removes_sandbox_and_driver_resources() { + run_detached_ephemeral_cleanup(0).await.unwrap(); +} + +#[tokio::test] +#[serial(ephemeral_cleanup)] +async fn detached_ephemeral_failure_removes_sandbox_and_driver_resources() { + run_detached_ephemeral_cleanup(17).await.unwrap(); +}