From 16275912cb98dfb135a5cdbbde7c4b36e995e59a Mon Sep 17 00:00:00 2001 From: Shiju Date: Thu, 1 Oct 2026 15:03:39 +0530 Subject: [PATCH] fix(vm): stop image workers before cleaning staging files Run image preparation in an owned worker process, reserve its process identity until cleanup completes, and protect staging with leases so cancellation and recovery cannot race with another preparation attempt. Signed-off-by: Shiju --- crates/openshell-driver-vm/README.md | 4 + crates/openshell-driver-vm/src/driver.rs | 677 +++++++++++++- crates/openshell-driver-vm/src/main.rs | 9 + crates/openshell-driver-vm/src/preparation.rs | 876 ++++++++++++++++++ .../tests/image_preparation.rs | 197 ++++ 5 files changed, 1717 insertions(+), 46 deletions(-) create mode 100644 crates/openshell-driver-vm/src/preparation.rs create mode 100644 crates/openshell-driver-vm/tests/image_preparation.rs diff --git a/crates/openshell-driver-vm/README.md b/crates/openshell-driver-vm/README.md index a0d83e1e4d..334d968dee 100644 --- a/crates/openshell-driver-vm/README.md +++ b/crates/openshell-driver-vm/README.md @@ -219,6 +219,10 @@ The requested sandbox image is never selected as the bootstrap image. Operators must configure either `bootstrap_image` or `default_image`; when both are empty, the driver fails during startup. +Image and writable-overlay preparation run in owned worker processes. Stopping or deleting a sandbox cancels its worker and formatter processes, waits for them to exit, then removes the attempt's temporary files under `/images/preparations/`. New overlays, including retries after an interrupted first start, are built there and renamed into sandbox state only after completion. A file lock inherited by the worker's children keeps cleanup from deleting files that a process still owns. If cleanup cannot establish that the processes have stopped, the operation returns an error and retains both temporary files and sandbox state. + +At startup, the driver reclaims inactive attempts in this directory. It preserves active attempts, committed image caches, and unmarked staging from older releases. Concurrent preparations serialize image-cache publication; an interrupted attempt cannot publish a partial disk over an existing cache entry. Gateway upload slots for `rootfs_tar_path` keep their existing, separate lifecycle. + Each sandbox gets its own sparse writable `/sandboxes//overlay.ext4`. Guest init mounts overlayfs as `/` with the prepared image rootfs as lowerdir when present, otherwise the bootstrap diff --git a/crates/openshell-driver-vm/src/driver.rs b/crates/openshell-driver-vm/src/driver.rs index a03982038d..aed051b24d 100644 --- a/crates/openshell-driver-vm/src/driver.rs +++ b/crates/openshell-driver-vm/src/driver.rs @@ -3,6 +3,9 @@ #![allow(unsafe_code)] +#[path = "preparation.rs"] +mod preparation; + use crate::gpu::{GpuInventory, allocate_vsock_cid}; use crate::isolation::VmBoundarySpec; @@ -88,7 +91,7 @@ use std::process::Stdio; use std::sync::Arc; use std::sync::atomic::{AtomicU64, Ordering}; use std::time::Duration; -use tokio::io::AsyncWriteExt; +use tokio::io::{AsyncBufReadExt, AsyncWriteExt}; use tokio::process::{Child, Command}; use tokio::sync::{Mutex, broadcast, mpsc}; use tokio::task::JoinHandle; @@ -214,7 +217,7 @@ struct VmDriverTlsPaths { ca: PathBuf, } -#[derive(Debug, Clone)] +#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)] struct RuntimeImagePlan { root_disk: PathBuf, image_disk: Option, @@ -579,6 +582,7 @@ struct SandboxRecord { state_dir: PathBuf, process: Option>>, provisioning_task: Option>, + preparation: Option>>, gpu_bdf: Option, deleting: bool, } @@ -612,13 +616,13 @@ fn resolve_record_id( Ok(first) } -#[derive(Debug, Clone, Copy, PartialEq, Eq)] +#[derive(Debug, Clone, Copy, PartialEq, Eq, serde::Serialize, serde::Deserialize)] enum OverlayPreparation { Fresh, PreserveExisting, } -#[derive(Debug, Clone, Copy, PartialEq, Eq)] +#[derive(Debug, Clone, Copy, PartialEq, Eq, serde::Serialize, serde::Deserialize)] struct SandboxOwnerIdentity { uid: u32, gid: u32, @@ -666,6 +670,7 @@ pub struct VmDriver { launcher_bin: PathBuf, registry: Arc>>, image_cache_lock: Arc>, + preparation_root: Option, events: broadcast::Sender, gpu_inventory: Option>>, lifecycle_extensions: Arc, @@ -754,10 +759,13 @@ impl VmDriver { launcher_bin, registry: Arc::new(Mutex::new(HashMap::new())), image_cache_lock: Arc::new(Mutex::new(())), + preparation_root: None, events, gpu_inventory, lifecycle_extensions: Arc::new(lifecycle_extensions), }; + preparation::reconcile(&image_cache_root_dir(&driver.config.state_dir)) + .map_err(|error| format!("reconcile image preparation staging: {error}"))?; driver.restore_persisted_sandboxes().await; Ok(driver) } @@ -840,17 +848,21 @@ impl VmDriver { if validate_host_supervisor(&destination).is_ok() { return Ok(destination); } - let _cache_guard = self.image_cache_lock.lock().await; + let cache_guard = self.image_cache_lock.clone().lock_owned().await; if validate_host_supervisor(&destination).is_ok() { return Ok(destination); } let destination_for_extract = destination.clone(); - tokio::task::spawn_blocking(move || extract_host_supervisor(&destination_for_extract)) - .await - .map_err(|error| { - Status::internal(format!("host supervisor extraction panicked: {error}")) - })? - .map_err(Status::failed_precondition)?; + tokio::task::spawn_blocking(move || { + // This atomic shared-cache write does not target sandbox state. + // Retain its lock until the write finishes even if provisioning is + // cancelled, so a later caller cannot start a competing extractor. + let _cache_guard = cache_guard; + extract_host_supervisor(&destination_for_extract) + }) + .await + .map_err(|error| Status::internal(format!("host supervisor extraction panicked: {error}")))? + .map_err(Status::failed_precondition)?; validate_host_supervisor(&destination).map_err(Status::failed_precondition)?; Ok(destination) } @@ -1127,6 +1139,7 @@ impl VmDriver { state_dir: state_dir.clone(), process: None, provisioning_task: None, + preparation: None, gpu_bdf: None, deleting: false, }, @@ -1312,8 +1325,9 @@ impl VmDriver { })?; let bootstrap_image_ref = self.bootstrap_image_ref()?; let bootstrap_image_identity = self - .ensure_cached_bootstrap_rootfs_image(&sandbox.id, &bootstrap_image_ref) - .await?; + .prepare_images_in_worker(&sandbox.id, &bootstrap_image_ref, None, true) + .await? + .bootstrap_image_identity; let root_disk = image_cache_rootfs_image(&self.config.state_dir, &bootstrap_image_identity); let image_disk = image_cache_rootfs_image(&self.config.state_dir, &persisted_identity); @@ -1324,8 +1338,13 @@ impl VmDriver { bootstrap_image_identity, } } else { - self.prepare_runtime_images(&sandbox.id, &image_ref, rootfs_tar_path.as_deref()) - .await? + self.prepare_images_in_worker( + &sandbox.id, + &image_ref, + rootfs_tar_path.as_deref(), + false, + ) + .await? }; let image_identity = image_plan.image_identity.clone(); self.ensure_provisioning_active(&sandbox.id).await?; @@ -1387,9 +1406,8 @@ impl VmDriver { ), ); let sandbox_owner_state = self - .prepare_runtime_overlay( - &state_dir, - &overlay_disk, + .prepare_overlay_in_worker( + &sandbox.id, &owner_source_disk, overlay_preparation, sandbox @@ -1825,7 +1843,11 @@ impl VmDriver { if let Some(task) = provisioning_task { task.abort(); + // Wait for the future to drop its worker I/O before taking the + // registered attempt. Aborting alone does not stop its subprocess. + let _ = task.await; } + self.cleanup_image_preparation(&record_id).await?; if let Some(process) = process { let mut process = process.lock().await; process.deleting = true; @@ -1881,6 +1903,11 @@ impl VmDriver { let record = registry .get(&id) .ok_or_else(|| Status::not_found("sandbox not found"))?; + if record.preparation.is_some() && record.provisioning_task.is_none() { + return Err(Status::failed_precondition( + "image preparation cleanup is still in progress", + )); + } ( id, record.state_dir.clone(), @@ -2029,7 +2056,11 @@ impl VmDriver { if let Some(task) = provisioning_task { task.abort(); + // Wait for the future to drop its worker I/O before taking the + // registered attempt. Aborting alone does not stop its subprocess. + let _ = task.await; } + self.cleanup_image_preparation(&record_id).await?; if let Some(process) = process { let mut process = process.lock().await; @@ -2181,6 +2212,7 @@ impl VmDriver { state_dir: state_dir.clone(), process: None, provisioning_task: None, + preparation: None, gpu_bdf: None, deleting: false, }); @@ -2208,6 +2240,7 @@ impl VmDriver { state_dir: state_dir.clone(), process: None, provisioning_task: None, + preparation: None, gpu_bdf: None, deleting: false, }); @@ -2307,6 +2340,7 @@ impl VmDriver { state_dir: state_dir.clone(), process: None, provisioning_task: None, + preparation: None, gpu_bdf: None, deleting: false, }, @@ -2654,7 +2688,7 @@ impl VmDriver { remove_state: bool, ) { self.release_gpu(sandbox_id); - let snapshot = { + let (snapshot, may_remove_state) = { let mut registry = self.registry.lock().await; let Some(record) = registry.get_mut(sandbox_id) else { return; @@ -2669,10 +2703,16 @@ impl VmDriver { error_condition(reason, message), false, )); - Some(record.snapshot.clone()) + // A failed cleanup leaves the attempt registered. Preserve all + // sandbox state until its worker and descendants are confirmed + // stopped, even when provisioning itself has already failed. + ( + Some(record.snapshot.clone()), + remove_state && record.preparation.is_none(), + ) }; - if remove_state { + if may_remove_state { let _ = tokio::fs::remove_dir_all(state_dir).await; remove_sandbox_socket_dir(&self.socket_root_fd, sandbox_id); } @@ -2690,6 +2730,237 @@ impl VmDriver { } } + /// Keep every image write in a process owned by the sandbox record. The + /// registry retains the attempt when this future is cancelled so lifecycle + /// cleanup can kill and reap it before removing any staging or state. + #[tracing::instrument( + name = "vm.prepare_images", + skip(self), + fields(otel.name = "vm.prepare_images", otel.status_code = tracing::field::Empty, sandbox.id = %sandbox_id, image.ref = %image_ref) + )] + async fn prepare_images_in_worker( + &self, + sandbox_id: &str, + image_ref: &str, + rootfs_tar: Option<&Path>, + bootstrap_only: bool, + ) -> Result { + match self + .run_preparation_in_worker(sandbox_id, image_ref, rootfs_tar, bootstrap_only, None) + .await? + { + preparation::Output::Images(plan) => Ok(plan), + preparation::Output::Overlay(_) => { + Err(Status::internal("image worker returned an overlay result")) + } + } + } + + #[tracing::instrument( + name = "vm.prepare_overlay", + skip(self), + fields(otel.name = "vm.prepare_overlay", otel.status_code = tracing::field::Empty, sandbox.id = %sandbox_id) + )] + async fn prepare_overlay_in_worker( + &self, + sandbox_id: &str, + owner_source_disk: &Path, + preparation: OverlayPreparation, + requested_identity: Option<&WorkloadIdentityRequest>, + ) -> Result { + let overlay = preparation::OverlayRequest { + source_disk: owner_source_disk.to_path_buf(), + preparation, + requested_identity: requested_identity + .map(preparation::WorkloadIdentitySelectors::from), + }; + match self + .run_preparation_in_worker(sandbox_id, "", None, false, Some(overlay)) + .await? + { + preparation::Output::Overlay(owner) => Ok(owner), + preparation::Output::Images(_) => { + Err(Status::internal("overlay worker returned an image result")) + } + } + } + + /// The worker resolves the owner from the same state directory it will + /// modify, then rejects selectors that conflict with that owner before + /// preparing or repairing the writable disk. + async fn prepare_requested_overlay( + &self, + sandbox_id: &str, + overlay: preparation::OverlayRequest, + ) -> Result { + if !overlay + .source_disk + .starts_with(image_cache_root_dir(&self.config.state_dir)) + { + return Err(Status::invalid_argument( + "overlay source must be a cached image", + )); + } + let state_dir = sandbox_state_dir(&self.config.state_dir, sandbox_id)?; + let overlay_disk = sandbox_runtime_disk_paths(&state_dir).overlay_disk; + let requested_identity = overlay + .requested_identity + .map(WorkloadIdentityRequest::from); + self.prepare_runtime_overlay( + &state_dir, + &overlay_disk, + &overlay.source_disk, + overlay.preparation, + requested_identity.as_ref(), + ) + .await + .map_err(Status::internal) + } + + async fn run_preparation_in_worker( + &self, + sandbox_id: &str, + image_ref: &str, + rootfs_tar: Option<&Path>, + bootstrap_only: bool, + overlay: Option, + ) -> Result { + let span_status = openshell_otel::ErrorStatusGuard::current(); + let bootstrap_image_ref = self.bootstrap_image_ref()?; + // Keep the driver trace continuous across the worker boundary. The + // first Pulled event completes bootstrap resolution; later events + // belong to the requested workload image. + let mut bootstrap_span = overlay.is_none().then(|| { + tracing::info_span!( + "vm.resolve_bootstrap_image", + otel.name = "vm.resolve_bootstrap_image", + otel.status_code = tracing::field::Empty, + sandbox.id = %sandbox_id, + image.ref = %bootstrap_image_ref, + ) + }); + // Queue locally before spawning a worker. The inherited file lock in + // the worker additionally serializes publication across driver restarts + // and other driver processes sharing this cache. + let _cache_guard = self.image_cache_lock.lock().await; + let (attempt, stdout) = { + let mut registry = self.registry.lock().await; + let record = registry + .get_mut(sandbox_id) + .filter(|record| !record.deleting) + .ok_or_else(|| Status::cancelled("sandbox provisioning cancelled"))?; + if record.preparation.is_some() { + return Err(Status::failed_precondition( + "sandbox image preparation is already active", + )); + } + let mut attempt = + preparation::Attempt::create(&image_cache_root_dir(&self.config.state_dir)) + .map_err(|error| { + Status::internal(format!("create image preparation attempt: {error}")) + })?; + let stdout = attempt.spawn( + &self.launcher_bin, + preparation::Request { + config: self.config.clone(), + sandbox_id: sandbox_id.to_string(), + image_ref: image_ref.to_string(), + rootfs_tar: rootfs_tar.map(Path::to_path_buf), + bootstrap_only, + overlay, + lease_fd: -1, + parent_pid: 0, + }, + ); + let attempt = Arc::new(Mutex::new(attempt)); + record.preparation = Some(attempt.clone()); + (attempt, stdout) + }; + let _cleanup = preparation::CleanupOnDrop(attempt.clone()); + let outcome = async { + let stdout = stdout.map_err(|error| { + Status::internal(format!("start image preparation worker: {error}")) + })?; + let mut lines = tokio::io::BufReader::new(stdout).lines(); + let mut result = None; + while let Some(line) = lines.next_line().await.map_err(|error| { + Status::internal(format!("read image preparation progress: {error}")) + })? { + match serde_json::from_str::(&line).map_err(|error| { + Status::internal(format!("decode image preparation progress: {error}")) + })? { + preparation::Message::Event(bytes) => { + let event = + WatchSandboxesEvent::decode(bytes.as_slice()).map_err(|error| { + Status::internal(format!("decode image preparation event: {error}")) + })?; + if matches!(event.payload.as_ref(), Some(watch_sandboxes_event::Payload::PlatformEvent(value)) + if value.event.as_ref().is_some_and(|event| event.reason == "Pulled")) { + bootstrap_span.take(); + } + let _ = self.events.send(event); + } + preparation::Message::Complete(value) => result = Some(value), + } + } + let status = attempt.lock().await.wait().await.map_err(|error| { + Status::internal(format!("wait for image preparation: {error}")) + })?; + if !status.success() { + return Err(Status::failed_precondition(format!( + "image preparation worker exited with {status}" + ))); + } + result + .ok_or_else(|| { + Status::internal("image preparation worker exited without a result") + })? + .map_err(Status::failed_precondition) + } + .await; + if outcome.is_err() + && let Some(span) = bootstrap_span.as_ref() + { + span.record("otel.status_code", "ERROR"); + } + self.cleanup_image_preparation(sandbox_id).await?; + span_status.finish(outcome) + } + + async fn cleanup_image_preparation(&self, sandbox_id: &str) -> Result<(), Status> { + let attempt = self + .registry + .lock() + .await + .get(sandbox_id) + .and_then(|record| record.preparation.clone()); + if let Some(attempt) = attempt { + attempt.lock().await.cleanup().await?; + if let Some(record) = self.registry.lock().await.get_mut(sandbox_id) + && record + .preparation + .as_ref() + .is_some_and(|current| Arc::ptr_eq(current, &attempt)) + { + record.preparation = None; + } + } + Ok(()) + } + + fn image_staging_dir(&self, image_identity: &str) -> PathBuf { + self.preparation_root.as_ref().map_or_else( + || image_cache_staging_dir(&self.config.state_dir, image_identity), + |root| { + root.join(format!( + "{}.staging-{}", + sanitize_image_identity(image_identity), + unique_image_cache_suffix() + )) + }, + ) + } + #[tracing::instrument( name = "vm.prepare_images", skip(self), @@ -2817,28 +3088,51 @@ impl VmDriver { } let template_path = overlay_template_image(&self.config.state_dir, overlay_size_bytes); + let overlay_metadata = tokio::fs::metadata(&overlay_disk).await; let recover_preserved_overlay = preparation == OverlayPreparation::PreserveExisting - && tokio::fs::metadata(&overlay_disk) - .await - .is_ok_and(|metadata| metadata.is_file()); + && overlay_metadata.as_ref().is_ok_and(fs::Metadata::is_file); + let publish_new_overlay = preparation == OverlayPreparation::Fresh + || matches!(overlay_metadata, Err(error) if error.kind() == std::io::ErrorKind::NotFound); if !overlay_template_image_ready(&template_path, overlay_size_bytes).await? { let _cache_guard = self.image_cache_lock.lock().await; let template_path = template_path.clone(); + let staging_dir = self.image_staging_dir("overlay-template"); tokio::task::spawn_blocking(move || { - ensure_sandbox_overlay_template_image(&template_path, overlay_size_bytes) + ensure_sandbox_overlay_template_image( + &template_path, + overlay_size_bytes, + &staging_dir, + ) }) .await .map_err(|err| format!("overlay template preparation panicked: {err}"))??; } let overlay_to_recover = overlay_disk.clone(); + let staging_dir = self.image_staging_dir("writable-overlay"); let result = tokio::task::spawn_blocking(move || { - prepare_sandbox_overlay_image( - &template_path, - &overlay_disk, - preparation, - overlay_size_bytes, - ) + if publish_new_overlay { + // An interrupted first start can leave no persisted overlay. + // Stage that retry just like a fresh copy, so cancellation + // cannot publish a partial disk as existing VM state. Other + // metadata errors still go through the preservation checks. + fs::create_dir_all(&staging_dir).map_err(|error| error.to_string())?; + let staging_overlay = staging_dir.join(SANDBOX_OVERLAY_IMAGE); + prepare_sandbox_overlay_image( + &template_path, + &staging_overlay, + preparation, + overlay_size_bytes, + )?; + fs::rename(&staging_overlay, &overlay_disk).map_err(|error| error.to_string()) + } else { + prepare_sandbox_overlay_image( + &template_path, + &overlay_disk, + preparation, + overlay_size_bytes, + ) + } }) .await .map_err(|err| format!("overlay image preparation panicked: {err}"))?; @@ -3249,7 +3543,7 @@ impl VmDriver { }); } - let staging_dir = image_cache_staging_dir(&self.config.state_dir, &cache_identity); + let staging_dir = self.image_staging_dir(&cache_identity); let rootfs_archive = staging_dir.join(IMAGE_EXPORT_ROOTFS_ARCHIVE); self.reset_image_staging_dir(&staging_dir).await?; @@ -3362,7 +3656,7 @@ impl VmDriver { }); } - let staging_dir = image_cache_staging_dir(&self.config.state_dir, &cache_identity); + let staging_dir = self.image_staging_dir(&cache_identity); let rootfs_archive = staging_dir.join(IMAGE_EXPORT_ROOTFS_ARCHIVE); self.reset_image_staging_dir(&staging_dir).await?; @@ -3503,7 +3797,7 @@ impl VmDriver { }); } - let staging_dir = image_cache_staging_dir(&self.config.state_dir, &cache_identity); + let staging_dir = self.image_staging_dir(&cache_identity); self.reset_image_staging_dir(&staging_dir).await?; let layout_dir = staging_dir.join(GUEST_IMAGE_OCI_LAYOUT_DIR); @@ -3723,7 +4017,8 @@ impl VmDriver { let mut command = Command::new(&self.launcher_bin); command.kill_on_drop(true); command.stdin(Stdio::null()); - command.stdout(Stdio::inherit()); + // Worker stdout carries framed progress and completion messages. + command.stdout(Stdio::from(std::io::stderr())); command.stderr(Stdio::inherit()); command.arg("--internal-run-vm"); command.arg("--vm-root-disk").arg(bootstrap_root_disk); @@ -3818,7 +4113,7 @@ impl VmDriver { ) -> Result<(), Status> { let cache_dir = image_cache_dir(&self.config.state_dir, image_identity); let image_path = image_cache_rootfs_image(&self.config.state_dir, image_identity); - let staging_dir = image_cache_staging_dir(&self.config.state_dir, image_identity); + let staging_dir = self.image_staging_dir(image_identity); let exported_rootfs = staging_dir.join(IMAGE_EXPORT_ROOTFS_ARCHIVE); let prepared_rootfs = staging_dir.join("rootfs"); let prepared_image = staging_dir.join(IMAGE_CACHE_ROOTFS_IMAGE); @@ -3949,7 +4244,7 @@ impl VmDriver { ) -> Result<(), Status> { let cache_dir = image_cache_dir(&self.config.state_dir, image_identity); let image_path = image_cache_rootfs_image(&self.config.state_dir, image_identity); - let staging_dir = image_cache_staging_dir(&self.config.state_dir, image_identity); + let staging_dir = self.image_staging_dir(image_identity); let prepared_rootfs = staging_dir.join("rootfs"); let prepared_image = staging_dir.join(IMAGE_CACHE_ROOTFS_IMAGE); @@ -4302,6 +4597,94 @@ impl VmDriver { } } +/// Execute an internal image preparation request in its own process group. +/// The driver supplies a private request file and an inherited staging lease; +/// this entry point never restores sandboxes or launches a gateway listener. +#[doc(hidden)] +pub async fn run_image_preparation_worker(request_path: &Path) -> Result<(), String> { + let (request, directory) = preparation::read_request(request_path)?; + validate_sandbox_id(&request.sandbox_id).map_err(|error| error.message().to_string())?; + preparation::watch_parent(request.parent_pid)?; + tracing_subscriber::fmt() + .with_writer(std::io::stderr) + .with_env_filter(request.config.log_level.clone()) + .try_init() + .map_err(|error| format!("initialize image preparation logging: {error}"))?; + let cache_root = image_cache_root_dir(&request.config.state_dir); + // The lock spans every image publication, including bootstrap images. + // Workers from separate driver processes cannot replace one another's + // validated cache entries while readers continue using committed images. + let _cache_lease = preparation::cache_lock(&cache_root).map_err(|error| error.to_string())?; + let (events, mut receiver) = broadcast::channel(WATCH_BUFFER); + let socket_root_fd = fs::File::open(&directory) + .map_err(|error| error.to_string())? + .into(); + let launcher_bin = request + .config + .launcher_bin + .clone() + .map_or_else(std::env::current_exe, Ok) + .map_err(|error| error.to_string())?; + let driver = VmDriver { + config: request.config, + socket_root: directory.clone(), + socket_root_fd: Arc::new(socket_root_fd), + launcher_bin, + registry: Arc::new(Mutex::new(HashMap::new())), + image_cache_lock: Arc::new(Mutex::new(())), + preparation_root: Some(directory), + events, + gpu_inventory: None, + lifecycle_extensions: Arc::new(LifecycleExtensionRegistry::new()), + }; + let prepare = async { + if let Some(overlay) = request.overlay { + driver + .prepare_requested_overlay(&request.sandbox_id, overlay) + .await + .map(preparation::Output::Overlay) + } else if request.bootstrap_only { + let identity = driver + .ensure_cached_bootstrap_rootfs_image(&request.sandbox_id, &request.image_ref) + .await?; + Ok(preparation::Output::Images(RuntimeImagePlan { + root_disk: image_cache_rootfs_image(&driver.config.state_dir, &identity), + image_disk: None, + image_identity: identity.clone(), + bootstrap_image_identity: identity, + })) + } else { + driver + .prepare_runtime_images( + &request.sandbox_id, + &request.image_ref, + request.rootfs_tar.as_deref(), + ) + .await + .map(preparation::Output::Images) + } + }; + tokio::pin!(prepare); + let result = loop { + tokio::select! { + result = &mut prepare => break result, + event = receiver.recv() => { + if let Ok(event) = event { + preparation::send(&preparation::Message::Event(event.encode_to_vec())).map_err(|error| error.to_string())?; + } + } + } + }; + while let Ok(event) = receiver.try_recv() { + preparation::send(&preparation::Message::Event(event.encode_to_vec())) + .map_err(|error| error.to_string())?; + } + preparation::send(&preparation::Message::Complete( + result.map_err(|error| error.message().to_string()), + )) + .map_err(|error| error.to_string()) +} + fn read_vm_console_tail(path: &Path, limit: u64) -> Option { if limit == 0 { return None; @@ -6247,6 +6630,7 @@ async fn overlay_template_image_ready(path: &Path, size_bytes: u64) -> Result Result<(), String> { if let Ok(metadata) = fs::metadata(template_path) && metadata.is_file() @@ -6268,15 +6652,8 @@ fn ensure_sandbox_overlay_template_image( ) })?; - let staging_image = parent.join(format!( - ".{}.staging-{}-{}", - template_path - .file_name() - .and_then(|name| name.to_str()) - .unwrap_or("overlay-template.ext4"), - std::process::id(), - openshell_core::time::now_ms() - )); + fs::create_dir_all(staging_dir).map_err(|error| error.to_string())?; + let staging_image = staging_dir.join("overlay-template.ext4"); let result = (|| { create_empty_sandbox_overlay_image(&staging_image, size_bytes)?; @@ -8265,6 +8642,61 @@ mod tests { ); } + #[tokio::test] + async fn overlay_worker_rejects_identity_conflicts_before_disk_changes() { + let directory = tempfile::tempdir().unwrap(); + let mut driver = test_driver_with_extensions(LifecycleExtensionRegistry::new()); + driver.config.state_dir = directory.path().to_path_buf(); + // A changed host default must not replace the persisted owner's identity. + driver.config.sandbox_uid = Some(10000); + driver.config.sandbox_gid = Some(10001); + let sandbox_id = "identity-worker"; + let state_dir = sandbox_state_dir(directory.path(), sandbox_id).unwrap(); + create_private_dir_all(&state_dir).await.unwrap(); + let owner = SandboxOwnerIdentity { + uid: 1000, + gid: 1001, + }; + write_sandbox_owner_state(&state_dir, owner).await.unwrap(); + let overlay_disk = sandbox_runtime_disk_paths(&state_dir).overlay_disk; + std::fs::write(&overlay_disk, b"existing overlay must not be touched").unwrap(); + + for (user, group, field) in [ + ("10000", "", "run_as_user"), + ("sandbox", "10001", "run_as_group"), + ] { + let requested_identity = WorkloadIdentityRequest { + user: user.into(), + group: group.into(), + }; + let request = preparation::OverlayRequest { + source_disk: image_cache_root_dir(directory.path()).join("must-not-read-image"), + preparation: OverlayPreparation::PreserveExisting, + requested_identity: Some(preparation::WorkloadIdentitySelectors::from( + &requested_identity, + )), + }; + // Exercise the serialized request and the same entry point the + // worker uses, so dropping either selector cannot weaken the check. + let wire = serde_json::to_vec(&request).unwrap(); + let decoded = serde_json::from_slice(&wire).unwrap(); + let error = driver + .prepare_requested_overlay(sandbox_id, decoded) + .await + .unwrap_err(); + assert!(error.message().contains(field), "{error}"); + assert!(error.message().contains("1000:1001"), "{error}"); + assert_eq!( + std::fs::read(&overlay_disk).unwrap(), + b"existing overlay must not be touched" + ); + assert_eq!( + std::fs::read_to_string(state_dir.join(SANDBOX_OWNER_STATE_FILE)).unwrap(), + owner.marker_contents() + ); + } + } + #[test] fn vm_workload_identity_checks_independent_numeric_and_symbolic_selectors() { let owner = SandboxOwnerIdentity { @@ -8929,6 +9361,7 @@ mod tests { state_dir: state_dir.clone(), process: None, provisioning_task: None, + preparation: None, gpu_bdf: None, deleting: false, }, @@ -9011,6 +9444,7 @@ mod tests { state_dir, process: None, provisioning_task: Some(provisioning_task), + preparation: None, gpu_bdf: None, deleting: false, }, @@ -9031,6 +9465,141 @@ mod tests { task.abort(); } + async fn lifecycle_cancels_image_preparation(delete: bool) { + let temp = tempfile::tempdir().unwrap(); + let mut driver = test_driver_with_extensions(LifecycleExtensionRegistry::new()); + driver.config.state_dir = temp.path().to_path_buf(); + let id = "sandbox-cancel-image"; + let state_dir = sandboxes_root_dir(temp.path()).join(id); + create_private_dir_all(&state_dir).await.unwrap(); + let launcher = temp.path().join("worker"); + fs::write(&launcher, "#!/bin/sh\nroot=$(dirname \"$2\")\nsleep 300 &\nprintf ready > \"$root/ready\"\nwait\n").unwrap(); + fs::set_permissions(&launcher, fs::Permissions::from_mode(0o700)).unwrap(); + let cache = image_cache_root_dir(temp.path()); + let mut attempt = preparation::Attempt::create(&cache).unwrap(); + let directory = attempt.directory.clone(); + let _stdout = attempt + .spawn( + &launcher, + preparation::Request { + config: driver.config.clone(), + sandbox_id: id.to_string(), + image_ref: "test-image".to_string(), + rootfs_tar: None, + bootstrap_only: false, + overlay: None, + lease_fd: -1, + parent_pid: 0, + }, + ) + .unwrap(); + tokio::time::timeout(Duration::from_secs(5), async { + while !directory.join("ready").exists() { + tokio::time::sleep(Duration::from_millis(10)).await; + } + }) + .await + .expect("preparation worker readiness"); + let attempt = Arc::new(Mutex::new(attempt)); + driver.registry.lock().await.insert( + id.to_string(), + SandboxRecord { + snapshot: Sandbox { + id: id.to_string(), + name: "cancel-image".to_string(), + ..Default::default() + }, + state_dir: state_dir.clone(), + process: None, + provisioning_task: Some(tokio::spawn(std::future::pending())), + preparation: Some(attempt.clone()), + gpu_bdf: None, + deleting: false, + }, + ); + if delete { + let response = driver.delete_sandbox(id, "").await.unwrap(); + assert!(response.deleted); + assert!(!state_dir.exists()); + } else { + driver.stop_sandbox(id, "").await.unwrap(); + assert!(state_dir.join(SANDBOX_STOPPED_FILE).exists()); + } + assert!( + !directory.exists(), + "lifecycle completion must reclaim preparation staging" + ); + attempt + .lock() + .await + .cleanup() + .await + .expect("worker was reaped and cleanup is idempotent"); + } + + #[tokio::test] + async fn stop_waits_for_image_preparation_cleanup() { + lifecycle_cancels_image_preparation(false).await; + } + + #[tokio::test] + async fn delete_waits_for_image_preparation_cleanup() { + lifecycle_cancels_image_preparation(true).await; + } + + #[tokio::test] + async fn incomplete_preparation_cleanup_failure_preserves_sandbox_state() { + let temp = tempfile::tempdir().unwrap(); + let mut driver = test_driver_with_extensions(LifecycleExtensionRegistry::new()); + driver.config.state_dir = temp.path().to_path_buf(); + let id = "sandbox-incomplete-cleanup"; + let state_dir = sandboxes_root_dir(temp.path()).join(id); + create_private_dir_all(&state_dir).await.unwrap(); + let overlay = state_dir.join(SANDBOX_OVERLAY_IMAGE); + fs::write(&overlay, b"owned overlay").unwrap(); + let attempt = Arc::new(Mutex::new( + preparation::Attempt::create(&image_cache_root_dir(temp.path())).unwrap(), + )); + driver.registry.lock().await.insert( + id.to_string(), + SandboxRecord { + snapshot: Sandbox { + id: id.to_string(), + ..Default::default() + }, + state_dir: state_dir.clone(), + process: None, + provisioning_task: None, + preparation: Some(attempt), + gpu_bdf: None, + deleting: false, + }, + ); + + // Failed worker cleanup retains the registered attempt. Its failure + // epilogue must not remove state that a descendant could still write. + driver + .fail_provisioning( + id, + &state_dir, + "PreparationFailed", + "image preparation descendants still own staging; files retained", + true, + ) + .await; + let retained = fs::read(&overlay); + driver.cleanup_image_preparation(id).await.unwrap(); + assert_eq!(retained.unwrap(), b"owned overlay"); + + driver + .fail_provisioning(id, &state_dir, "PreparationFailed", "worker stopped", true) + .await; + assert!( + !state_dir.exists(), + "successful worker cleanup permits removing failed sandbox state" + ); + } + fn test_launch_authentication(label: &str) -> (Vec, openshell_core::SandboxSessionId) { use openshell_core::jwt::{ SandboxLaunchAuthentication, SecretJwt, SessionVerificationKey, SupervisorAuthBundle, @@ -9154,6 +9723,7 @@ mod tests { launcher_bin: PathBuf::from("/tmp/openshell-driver-vm"), registry: Arc::new(Mutex::new(HashMap::new())), image_cache_lock: Arc::new(Mutex::new(())), + preparation_root: None, events: broadcast::channel(WATCH_BUFFER).0, gpu_inventory: None, lifecycle_extensions: Arc::new(LifecycleExtensionRegistry::new()), @@ -9193,6 +9763,7 @@ mod tests { launcher_bin: PathBuf::from("/tmp/openshell-driver-vm"), registry: Arc::new(Mutex::new(HashMap::new())), image_cache_lock: Arc::new(Mutex::new(())), + preparation_root: None, events: broadcast::channel(WATCH_BUFFER).0, gpu_inventory: None, lifecycle_extensions: Arc::new(LifecycleExtensionRegistry::new()), @@ -9227,6 +9798,7 @@ mod tests { launcher_bin: PathBuf::from("/tmp/openshell-driver-vm"), registry: Arc::new(Mutex::new(HashMap::new())), image_cache_lock: Arc::new(Mutex::new(())), + preparation_root: None, events: broadcast::channel(WATCH_BUFFER).0, gpu_inventory: None, lifecycle_extensions: Arc::new(LifecycleExtensionRegistry::new()), @@ -9255,6 +9827,7 @@ mod tests { launcher_bin: PathBuf::from("/tmp/openshell-driver-vm"), registry: Arc::new(Mutex::new(HashMap::new())), image_cache_lock: Arc::new(Mutex::new(())), + preparation_root: None, events: broadcast::channel(WATCH_BUFFER).0, gpu_inventory: None, lifecycle_extensions: Arc::new(LifecycleExtensionRegistry::new()), @@ -9284,6 +9857,7 @@ mod tests { launcher_bin: PathBuf::from("/tmp/openshell-driver-vm"), registry: Arc::new(Mutex::new(HashMap::new())), image_cache_lock: Arc::new(Mutex::new(())), + preparation_root: None, events: broadcast::channel(WATCH_BUFFER).0, gpu_inventory: None, lifecycle_extensions: Arc::new(LifecycleExtensionRegistry::new()), @@ -9308,6 +9882,7 @@ mod tests { launcher_bin: PathBuf::from("/tmp/openshell-driver-vm"), registry: Arc::new(Mutex::new(HashMap::new())), image_cache_lock: Arc::new(Mutex::new(())), + preparation_root: None, events: broadcast::channel(WATCH_BUFFER).0, gpu_inventory: None, lifecycle_extensions: Arc::new(LifecycleExtensionRegistry::new()), @@ -9329,6 +9904,7 @@ mod tests { launcher_bin: PathBuf::from("/tmp/openshell-driver-vm"), registry: Arc::new(Mutex::new(HashMap::new())), image_cache_lock: Arc::new(Mutex::new(())), + preparation_root: None, events: broadcast::channel(WATCH_BUFFER).0, gpu_inventory: None, lifecycle_extensions: Arc::new(LifecycleExtensionRegistry::new()), @@ -9975,6 +10551,7 @@ mod tests { launcher_bin: PathBuf::from("openshell-driver-vm"), registry: Arc::new(Mutex::new(HashMap::new())), image_cache_lock: Arc::new(Mutex::new(())), + preparation_root: None, events, gpu_inventory: None, lifecycle_extensions: Arc::new(LifecycleExtensionRegistry::new()), @@ -10040,6 +10617,7 @@ mod tests { launcher_bin: PathBuf::from("openshell-driver-vm"), registry: Arc::new(Mutex::new(HashMap::new())), image_cache_lock: Arc::new(Mutex::new(())), + preparation_root: None, events, gpu_inventory: None, lifecycle_extensions: Arc::new(LifecycleExtensionRegistry::new()), @@ -10060,6 +10638,7 @@ mod tests { state_dir: state_dir.clone(), process: None, provisioning_task: None, + preparation: None, gpu_bdf: None, deleting: false, }, @@ -10094,6 +10673,7 @@ mod tests { launcher_bin: PathBuf::from("openshell-driver-vm"), registry: Arc::new(Mutex::new(HashMap::new())), image_cache_lock: Arc::new(Mutex::new(())), + preparation_root: None, events, gpu_inventory: None, lifecycle_extensions: Arc::new(LifecycleExtensionRegistry::new()), @@ -10115,6 +10695,7 @@ mod tests { state_dir: state_dir.clone(), process: None, provisioning_task: None, + preparation: None, gpu_bdf: None, deleting: false, }, @@ -10469,6 +11050,7 @@ mod tests { state_dir, process: Some(process), provisioning_task: None, + preparation: None, gpu_bdf: None, deleting: false, }, @@ -10496,6 +11078,7 @@ mod tests { launcher_bin: PathBuf::from("openshell-driver-vm"), registry: Arc::new(Mutex::new(HashMap::new())), image_cache_lock: Arc::new(Mutex::new(())), + preparation_root: None, events, gpu_inventory: None, lifecycle_extensions: Arc::new(LifecycleExtensionRegistry::new()), @@ -10839,6 +11422,7 @@ mod tests { launcher_bin: PathBuf::from("openshell-driver-vm"), registry: Arc::new(Mutex::new(HashMap::new())), image_cache_lock: Arc::new(Mutex::new(())), + preparation_root: None, events, gpu_inventory: None, lifecycle_extensions: Arc::new(extensions), @@ -11447,6 +12031,7 @@ mod tests { state_dir: state_dir.clone(), process: None, provisioning_task: None, + preparation: None, gpu_bdf: None, deleting: false, }, diff --git a/crates/openshell-driver-vm/src/main.rs b/crates/openshell-driver-vm/src/main.rs index d7195b4cc3..4b4bcf3f6a 100644 --- a/crates/openshell-driver-vm/src/main.rs +++ b/crates/openshell-driver-vm/src/main.rs @@ -35,6 +35,9 @@ struct Args { #[arg(long, hide = true, default_value_t = false)] internal_run_vm: bool, + #[arg(long, hide = true)] + internal_prepare_image: Option, + #[arg(long = "vm-root-disk", hide = true, alias = "vm-rootfs")] vm_root_disk: Option, @@ -249,6 +252,12 @@ struct Args { #[tokio::main] async fn main() -> Result<()> { let args = Args::parse(); + if let Some(request) = args.internal_prepare_image { + openshell_driver_vm::driver::run_image_preparation_worker(&request) + .await + .map_err(|error| miette::miette!("{error}"))?; + return Ok(()); + } if args.internal_run_vm { // The VM launcher arms procguard after resolving its runtime so its // libkrun worker cannot outlive the launcher. diff --git a/crates/openshell-driver-vm/src/preparation.rs b/crates/openshell-driver-vm/src/preparation.rs new file mode 100644 index 0000000000..44c9b60e7b --- /dev/null +++ b/crates/openshell-driver-vm/src/preparation.rs @@ -0,0 +1,876 @@ +// SPDX-FileCopyrightText: Copyright (c) 2025-2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved. +// SPDX-License-Identifier: Apache-2.0 + +//! Own image-preparation processes and the files they may still be writing. +//! +//! A lease is inherited by every preparation subprocess. Cancellation first +//! kills and reaps the worker; staging is removed only after the last lease +//! descriptor closes. A restarted driver uses the same proof of inactivity. + +#![allow(unsafe_code)] + +use std::fs::{self, File, OpenOptions}; +use std::io::{self, Read, Write}; +use std::os::fd::AsRawFd; +use std::os::unix::fs::OpenOptionsExt; +use std::path::{Path, PathBuf}; +use std::process::{ExitStatus, Stdio}; +use std::sync::Arc; +use std::time::Duration; + +use nix::sys::signal::{Signal, killpg}; +use nix::unistd::Pid; +use serde::{Deserialize, Serialize}; +use tokio::process::{Child, ChildStdout, Command}; +use tokio::sync::Mutex; +use tonic::Status; + +use super::{ + OverlayPreparation, RuntimeImagePlan, SandboxOwnerIdentity, VmDriverConfig, + WorkloadIdentityRequest, +}; + +const ATTEMPTS_DIR: &str = "preparations"; +const REQUEST_FILE: &str = "request.json"; +const LEASE_MARKER: &[u8] = b"openshell-image-preparation-v1\n"; +const CLEANUP_TIMEOUT: Duration = Duration::from_secs(10); + +#[derive(Serialize, Deserialize)] +#[serde(deny_unknown_fields)] +pub(super) struct Request { + pub config: VmDriverConfig, + pub sandbox_id: String, + pub image_ref: String, + pub rootfs_tar: Option, + pub bootstrap_only: bool, + pub overlay: Option, + pub lease_fd: i32, + pub parent_pid: u32, +} + +#[derive(Serialize, Deserialize)] +#[serde(deny_unknown_fields)] +pub(super) struct OverlayRequest { + pub source_disk: PathBuf, + pub preparation: OverlayPreparation, + pub requested_identity: Option, +} + +/// Keep the requested selectors across the worker boundary. The worker must +/// validate them against the persisted overlay owner before changing any files; +/// the parent's current default identity cannot substitute for that owner. +#[derive(Serialize, Deserialize)] +#[serde(deny_unknown_fields)] +pub(super) struct WorkloadIdentitySelectors { + user: String, + group: String, +} + +impl From<&WorkloadIdentityRequest> for WorkloadIdentitySelectors { + fn from(request: &WorkloadIdentityRequest) -> Self { + Self { + user: request.user.clone(), + group: request.group.clone(), + } + } +} + +impl From for WorkloadIdentityRequest { + fn from(selectors: WorkloadIdentitySelectors) -> Self { + Self { + user: selectors.user, + group: selectors.group, + } + } +} + +#[derive(Serialize, Deserialize)] +pub(super) enum Output { + Images(RuntimeImagePlan), + Overlay(SandboxOwnerIdentity), +} + +#[derive(Serialize, Deserialize)] +pub(super) enum Message { + Event(Vec), + Complete(Result), +} + +pub(super) struct Attempt { + pub directory: PathBuf, + lease_path: PathBuf, + lease: Option, + child: Option, +} + +/// Cleanup survives cancellation of the provisioning future itself. Lifecycle +/// calls also await the same mutex, so they cannot report completion early. +pub(super) struct CleanupOnDrop(pub Arc>); + +impl Drop for CleanupOnDrop { + fn drop(&mut self) { + let attempt = self.0.clone(); + tokio::spawn(async move { + if let Err(error) = attempt.lock().await.cleanup().await { + tracing::warn!(%error, "image preparation cleanup incomplete; staging retained"); + } + }); + } +} + +impl Attempt { + /// Create and lock the lease before publishing the directory. A concurrent + /// reconciler can never observe this attempt without its ownership lock. + pub fn create(cache_root: &Path) -> io::Result { + let root = cache_root.join(ATTEMPTS_DIR); + fs::create_dir_all(&root)?; + let name = format!("attempt-{:032x}", rand::random::()); + let directory = root.join(&name); + let lease_path = root.join(format!("{name}.lease")); + let mut lease = OpenOptions::new() + .read(true) + .write(true) + .create_new(true) + .mode(0o600) + .open(&lease_path)?; + lock(&lease, false)?; + lease.write_all(LEASE_MARKER)?; + lease.sync_data()?; + fs::create_dir(&directory)?; + Ok(Self { + directory, + lease_path, + lease: Some(lease), + child: None, + }) + } + + /// Spawn while the caller holds the sandbox registry lock, then register + /// this attempt before yielding. Stop/delete cannot miss a live worker. + pub fn spawn(&mut self, launcher: &Path, mut request: Request) -> io::Result { + let lease = self + .lease + .as_ref() + .ok_or_else(|| io::Error::other("preparation lease is closed"))?; + let fd = lease.as_raw_fd(); + request.lease_fd = fd; + request.parent_pid = std::process::id(); + let request_path = self.directory.join(REQUEST_FILE); + let mut file = OpenOptions::new() + .write(true) + .create_new(true) + .mode(0o600) + .open(&request_path)?; + serde_json::to_writer(&mut file, &request)?; + file.flush()?; + let mut command = Command::new(launcher); + command.arg("--internal-prepare-image").arg(request_path); + command + .stdin(Stdio::null()) + .stdout(Stdio::piped()) + .stderr(Stdio::inherit()); + command.process_group(0).kill_on_drop(true); + // SAFETY: the child only makes an async-signal-safe fcntl call. The + // parent keeps this descriptor open until the child has exited. + unsafe { + command.pre_exec(move || { + if libc::fcntl(fd, libc::F_SETFD, 0) == -1 { + return Err(io::Error::last_os_error()); + } + Ok(()) + }); + } + let mut child = command.spawn()?; + let stdout = child + .stdout + .take() + .ok_or_else(|| io::Error::other("preparation stdout is missing"))?; + self.child = Some(child); + Ok(stdout) + } + + pub async fn wait(&mut self) -> io::Result { + // EOF can precede observable process exit. Let normal worker teardown + // finish before signaling, preserving its real status and reserving + // its PID until any remaining descendants have been stopped. + wait_for_worker_exit(self.child.as_ref().and_then(Child::id)).await?; + self.kill_group()?; + self.child + .as_mut() + .ok_or_else(|| io::Error::other("preparation worker is missing"))? + .wait() + .await + } + + /// Do not report successful cleanup while a formatter or other descendant + /// still owns the lease. Preserve uncertain staging for the next reconcile. + pub async fn cleanup(&mut self) -> Result<(), Status> { + let signal_result = self.kill_group(); + #[cfg(target_os = "macos")] + let signal_result = match signal_result { + // Darwin can reject signaling an exiting process before waitid + // exposes its terminal status. Keep immediate cancellation first, + // then retry the same guarded signal only after observing exit. + Err(error) if error.raw_os_error() == Some(libc::EPERM) => { + wait_for_worker_exit(self.child.as_ref().and_then(Child::id)) + .await + .and_then(|()| self.kill_group()) + } + result => result, + }; + signal_result + .map_err(|error| Status::internal(format!("terminate image preparation: {error}")))?; + if let Some(child) = self.child.as_mut() { + tokio::time::timeout(CLEANUP_TIMEOUT, child.wait()) + .await + .map_err(|_| { + Status::deadline_exceeded("image preparation did not stop; staging retained") + })? + .map_err(|error| Status::internal(format!("reap image preparation: {error}")))?; + } + // Do not explicitly unlock: inherited descriptors must continue to + // protect files until the last worker or descendant closes its copy. + self.lease.take(); + let deadline = tokio::time::Instant::now() + CLEANUP_TIMEOUT; + loop { + match reclaim(&self.directory, &self.lease_path) { + Ok(true) => return Ok(()), + Ok(false) if tokio::time::Instant::now() < deadline => { + tokio::time::sleep(Duration::from_millis(20)).await; + } + Ok(false) => { + return Err(Status::deadline_exceeded( + "image preparation descendants still own staging; files retained", + )); + } + Err(error) => { + return Err(Status::internal(format!( + "reclaim image preparation staging: {error}" + ))); + } + } + } + } + + // macOS may replace the staging lease when proving an exited process group + // has no remaining owner; other platforms only read the worker state. + #[cfg_attr(not(target_os = "macos"), allow(clippy::needless_pass_by_ref_mut))] + fn kill_group(&mut self) -> io::Result<()> { + if let Some(id) = self.child.as_ref().and_then(Child::id) { + let pid = i32::try_from(id).map_err(io::Error::other)?; + match killpg(Pid::from_raw(pid), Signal::SIGKILL) { + Ok(()) | Err(nix::errno::Errno::ESRCH) => {} + #[cfg(target_os = "macos")] + Err(nix::errno::Errno::EPERM) + if self.exited_worker_has_no_descendant_owner(id)? => {} + Err(error) => return Err(io::Error::from_raw_os_error(error as i32)), + } + } + Ok(()) + } + + /// Darwin rejects signals to a group containing only an exited leader. + /// Accept that case only while its PID remains reserved and no descendant + /// owns the inherited staging lease. Live or uncertain ownership still fails. + #[cfg(target_os = "macos")] + fn exited_worker_has_no_descendant_owner(&mut self, id: u32) -> io::Result { + if !worker_has_exited(id)? { + return Ok(false); + } + + // Open the probe before dropping our copy so a concurrent reconciler + // cannot remove the lease pathname between release and open. Reacquire + // ownership on success and retain it until the caller reaps the worker. + let lease = OpenOptions::new() + .read(true) + .write(true) + .custom_flags(libc::O_NOFOLLOW) + .open(&self.lease_path)?; + self.lease.take(); + if !lock(&lease, true)? { + return Ok(false); + } + self.lease = Some(lease); + Ok(true) + } +} + +/// Observe termination without releasing the PID or process-group identity. +fn worker_has_exited(id: u32) -> io::Result { + // SAFETY: waitid writes initialized storage. WNOWAIT leaves the owned + // child unreaped, so its PID cannot be reused before descendant signaling. + let mut status: libc::siginfo_t = unsafe { std::mem::zeroed() }; + if unsafe { + libc::waitid( + libc::P_PID, + id, + &raw mut status, + libc::WEXITED | libc::WNOHANG | libc::WNOWAIT, + ) + } != 0 + { + return Err(io::Error::last_os_error()); + } + Ok(matches!( + status.si_code, + libc::CLD_EXITED | libc::CLD_KILLED | libc::CLD_DUMPED + )) +} + +/// Bound EOF-to-exit observation without blocking the runtime or reaping. +/// Dropping this future leaves the worker available to cancellation cleanup. +async fn wait_for_worker_exit(id: Option) -> io::Result<()> { + let Some(id) = id else { + return Ok(()); + }; + tokio::time::timeout(CLEANUP_TIMEOUT, async { + loop { + if worker_has_exited(id)? { + return Ok(()); + } + tokio::time::sleep(Duration::from_millis(10)).await; + } + }) + .await + .map_err(|_| { + io::Error::new( + io::ErrorKind::TimedOut, + "image preparation worker did not exit before cleanup deadline", + ) + })? +} + +impl Drop for Attempt { + fn drop(&mut self) { + // Runtime shutdown may prevent the asynchronous reaper from running. + // Signal only this worker's still-unreaped process group; the lease + // remains held by descendants until the kernel closes their files. + let _ = self.kill_group(); + } +} + +/// Remove only directories created by this protocol whose inherited lease has +/// no remaining owner. Legacy staging and shared committed images are outside +/// this namespace and are never inferred to be inactive from their age. +pub(super) fn reconcile(cache_root: &Path) -> io::Result<()> { + let root = cache_root.join(ATTEMPTS_DIR); + let entries = match fs::read_dir(&root) { + Ok(entries) => entries, + Err(error) if error.kind() == io::ErrorKind::NotFound => return Ok(()), + Err(error) => return Err(error), + }; + for entry in entries { + let entry = entry?; + let name = entry.file_name(); + let Some(name) = name.to_str().filter(|name| valid_attempt_name(name)) else { + continue; + }; + if !entry.file_type()?.is_dir() { + continue; + } + if let Err(error) = reclaim(&entry.path(), &root.join(format!("{name}.lease"))) { + tracing::warn!(path = %entry.path().display(), %error, "image preparation ownership uncertain; staging retained"); + } + } + Ok(()) +} + +fn valid_attempt_name(name: &str) -> bool { + name.strip_prefix("attempt-").is_some_and(|suffix| { + suffix.len() == 32 && suffix.bytes().all(|byte| byte.is_ascii_hexdigit()) + }) +} + +fn reclaim(directory: &Path, lease_path: &Path) -> io::Result { + if !directory.try_exists()? && !lease_path.try_exists()? { + return Ok(true); + } + let lease = match OpenOptions::new() + .read(true) + .write(true) + .custom_flags(libc::O_NOFOLLOW) + .open(lease_path) + { + Ok(lease) if lease.metadata()?.is_file() => lease, + Ok(_) => return Ok(false), + Err(error) if error.kind() == io::ErrorKind::NotFound => return Ok(false), + Err(error) => return Err(error), + }; + if !lock(&lease, true)? { + return Ok(false); + } + let mut marker = Vec::new(); + (&lease).take(64).read_to_end(&mut marker)?; + if marker != LEASE_MARKER { + return Ok(false); + } + match fs::remove_dir_all(directory) { + Ok(()) => {} + Err(error) if error.kind() == io::ErrorKind::NotFound => {} + Err(error) => return Err(error), + } + match fs::remove_file(lease_path) { + Ok(()) => {} + Err(error) if error.kind() == io::ErrorKind::NotFound => {} + Err(error) => return Err(error), + } + Ok(true) +} + +fn lock(file: &File, nonblocking: bool) -> io::Result { + let operation = libc::LOCK_EX | if nonblocking { libc::LOCK_NB } else { 0 }; + // SAFETY: flock receives a valid borrowed file descriptor and no pointers. + if unsafe { libc::flock(file.as_raw_fd(), operation) } == 0 { + return Ok(true); + } + let error = io::Error::last_os_error(); + if nonblocking && error.kind() == io::ErrorKind::WouldBlock { + Ok(false) + } else { + Err(error) + } +} + +pub(super) fn cache_lock(cache_root: &Path) -> io::Result { + let file = OpenOptions::new() + .read(true) + .write(true) + .create(true) + .truncate(false) + .mode(0o600) + .custom_flags(libc::O_NOFOLLOW) + .open(cache_root.join("preparation.lock"))?; + lock(&file, false)?; + Ok(file) +} + +pub(super) fn read_request(path: &Path) -> Result<(Request, PathBuf), String> { + use std::os::unix::fs::MetadataExt; + let file = OpenOptions::new() + .read(true) + .custom_flags(libc::O_NOFOLLOW | libc::O_NONBLOCK) + .open(path) + .map_err(|error| error.to_string())?; + let metadata = file.metadata().map_err(|error| error.to_string())?; + if !metadata.is_file() || metadata.len() > 1024 * 1024 { + return Err( + "image preparation request must be a regular file no larger than 1 MiB".to_string(), + ); + } + let request: Request = serde_json::from_reader(file).map_err(|error| error.to_string())?; + let directory = path + .parent() + .ok_or("preparation request has no directory")?; + let name = directory + .file_name() + .and_then(|name| name.to_str()) + .filter(|name| valid_attempt_name(name)) + .ok_or("invalid preparation directory name")?; + let expected_root = super::image_cache_root_dir(&request.config.state_dir).join(ATTEMPTS_DIR); + if path.file_name().and_then(|name| name.to_str()) != Some(REQUEST_FILE) + || directory.parent() != Some(expected_root.as_path()) + { + return Err("preparation request is outside the driver staging directory".to_string()); + } + let lease_file = OpenOptions::new() + .read(true) + .custom_flags(libc::O_NOFOLLOW | libc::O_NONBLOCK) + .open(expected_root.join(format!("{name}.lease"))) + .map_err(|error| error.to_string())?; + let lease = lease_file.metadata().map_err(|error| error.to_string())?; + if !lease.is_file() { + return Err("image preparation lease must be a regular file".to_string()); + } + let mut marker = Vec::new(); + lease_file + .take(64) + .read_to_end(&mut marker) + .map_err(|error| error.to_string())?; + if marker != LEASE_MARKER { + return Err("image preparation lease marker is invalid".to_string()); + } + // SAFETY: fstat writes only to the supplied initialized stat storage. An + // invalid inherited descriptor is rejected rather than taken into ownership. + let mut inherited: libc::stat = unsafe { std::mem::zeroed() }; + if request.lease_fd < 3 || unsafe { libc::fstat(request.lease_fd, &raw mut inherited) } != 0 { + return Err("image preparation lease was not inherited".to_string()); + } + // Match MetadataExt::dev's representation: Darwin dev_t is signed, while + // Linux dev_t is already u64. The cast intentionally preserves that API's + // conversion, including the sign extension of a signed device identifier. + #[allow(clippy::cast_sign_loss, trivial_numeric_casts)] + let inherited_device = inherited.st_dev as u64; + if inherited.st_ino != lease.ino() || inherited_device != lease.dev() { + return Err("image preparation lease was not inherited".to_string()); + } + Ok((request, directory.to_path_buf())) +} + +/// A process-group watcher also covers host utilities on Linux, where a +/// parent-death signal on the worker alone cannot kill its children. +pub(super) fn watch_parent(parent_pid: u32) -> Result<(), String> { + let expected = i32::try_from(parent_pid).map_err(|error| error.to_string())?; + if nix::unistd::getpgrp() != nix::unistd::getpid() { + return Err("preparation worker must own its process group".to_string()); + } + std::thread::Builder::new() + .name("image-prep-parent".to_string()) + .spawn(move || { + loop { + if nix::unistd::getppid().as_raw() != expected { + let _ = killpg(nix::unistd::getpgrp(), Signal::SIGKILL); + std::process::exit(1); + } + std::thread::sleep(Duration::from_millis(50)); + } + }) + .map_err(|error| error.to_string())?; + Ok(()) +} + +pub(super) fn send(message: &Message) -> io::Result<()> { + let mut output = io::stdout().lock(); + serde_json::to_writer(&mut output, message)?; + output.write_all(b"\n")?; + output.flush() +} + +#[cfg(test)] +mod tests { + use super::*; + use std::os::unix::fs::{PermissionsExt, symlink}; + + fn request(root: &Path) -> Request { + Request { + config: VmDriverConfig { + state_dir: root.to_path_buf(), + ..VmDriverConfig::default() + }, + sandbox_id: "test-sandbox".to_string(), + image_ref: "test-image".to_string(), + rootfs_tar: None, + bootstrap_only: false, + overlay: None, + lease_fd: -1, + parent_pid: 0, + } + } + + async fn wait_for_file(path: &Path) { + tokio::time::timeout(Duration::from_secs(5), async { + while !path.exists() { + tokio::time::sleep(Duration::from_millis(10)).await; + } + }) + .await + .expect("worker readiness"); + } + + async fn wait_for_worker_exit_without_reaping(attempt: &Attempt) { + let id = attempt.child.as_ref().unwrap().id().unwrap(); + tokio::time::timeout(Duration::from_secs(5), async { + loop { + let exited = { + // SAFETY: waitid writes initialized exit information and + // WNOWAIT leaves the owned child's PID reserved for cleanup. + let mut status: libc::siginfo_t = unsafe { std::mem::zeroed() }; + assert_eq!( + unsafe { + libc::waitid( + libc::P_PID, + id, + &raw mut status, + libc::WEXITED | libc::WNOHANG | libc::WNOWAIT, + ) + }, + 0 + ); + matches!( + status.si_code, + libc::CLD_EXITED | libc::CLD_KILLED | libc::CLD_DUMPED + ) + }; + if exited { + break; + } + tokio::time::sleep(Duration::from_millis(10)).await; + } + }) + .await + .expect("worker must exit without being reaped"); + } + + async fn worker_at_stdout_eof(temp: &tempfile::TempDir, script: &str) -> Attempt { + let launcher = temp.path().join("worker"); + fs::write(&launcher, script).unwrap(); + fs::set_permissions(&launcher, fs::Permissions::from_mode(0o700)).unwrap(); + let cache = super::super::image_cache_root_dir(temp.path()); + let mut attempt = Attempt::create(&cache).unwrap(); + let mut stdout = attempt.spawn(&launcher, request(temp.path())).unwrap(); + let mut output = Vec::new(); + tokio::time::timeout( + Duration::from_secs(5), + tokio::io::AsyncReadExt::read_to_end(&mut stdout, &mut output), + ) + .await + .expect("worker must close stdout") + .unwrap(); + attempt + } + + #[tokio::test] + async fn stdout_eof_keeps_successful_worker_exit_status() { + let temp = tempfile::tempdir().unwrap(); + // stdout can close before the kernel exposes the final exit status. + // Keep that interval deterministic without requiring scheduler timing. + let mut attempt = + worker_at_stdout_eof(&temp, "#!/bin/sh\nexec 1>&-\nsleep 0.1\nexit 0\n").await; + let status = attempt.wait().await; + attempt.cleanup().await.expect("cleanup EOF worker"); + + let status = status.expect("wait after stdout EOF"); + assert!(status.success(), "worker exit after stdout EOF: {status}"); + assert!(!attempt.directory.exists()); + } + + #[tokio::test] + async fn stdout_eof_with_live_descendant_cleans_group_after_worker_exit() { + let temp = tempfile::tempdir().unwrap(); + let mut attempt = worker_at_stdout_eof( + &temp, + "#!/bin/sh\nsleep 300 >/dev/null &\nexec 1>&-\nsleep 0.1\nexit 0\n", + ) + .await; + let status = attempt.wait().await; + attempt.cleanup().await.expect("stop inherited lease owner"); + + let status = status.expect("preserve leader status"); + assert!(status.success(), "worker exit with descendant: {status}"); + assert!(!attempt.directory.exists()); + assert!(attempt.child.as_ref().unwrap().id().is_none()); + } + + #[tokio::test] + async fn cancelling_stdout_eof_wait_leaves_worker_for_cleanup() { + let temp = tempfile::tempdir().unwrap(); + let mut attempt = worker_at_stdout_eof(&temp, "#!/bin/sh\nexec 1>&-\nsleep 300\n").await; + let result = tokio::time::timeout(Duration::from_millis(50), attempt.wait()).await; + attempt + .cleanup() + .await + .expect("cancelled wait cleans worker"); + + assert!(result.is_err(), "waiting for exit must remain cancellable"); + assert!(!attempt.directory.exists()); + assert!(attempt.child.as_ref().unwrap().id().is_none()); + } + + #[cfg(target_os = "macos")] + #[tokio::test] + async fn cleanup_after_stdout_eof_waits_for_exit_publication() { + // Repeated natural exits exercise Darwin's short interval between + // closing stdout and making terminal status observable to waitid. + for _ in 0..8 { + let temp = tempfile::tempdir().unwrap(); + let mut attempt = worker_at_stdout_eof(&temp, "#!/bin/sh\nexit 0\n").await; + let result = attempt.cleanup().await; + if result.is_err() { + // Reap this owned worker even when checking faulty behavior. + attempt.cleanup().await.expect("retry owned test cleanup"); + } + result.expect("natural EOF must not fail cancellation cleanup"); + assert!(!attempt.directory.exists()); + } + } + + #[tokio::test] + async fn completed_worker_without_descendants_keeps_successful_exit_status() { + let temp = tempfile::tempdir().unwrap(); + let cache = super::super::image_cache_root_dir(temp.path()); + let launcher = temp.path().join("worker"); + fs::write(&launcher, "#!/bin/sh\nexit 0\n").unwrap(); + fs::set_permissions(&launcher, fs::Permissions::from_mode(0o700)).unwrap(); + let mut attempt = Attempt::create(&cache).unwrap(); + let directory = attempt.directory.clone(); + let _stdout = attempt.spawn(&launcher, request(temp.path())).unwrap(); + wait_for_worker_exit_without_reaping(&attempt).await; + + // macOS rejects killpg for a group containing only a zombie. The + // completed worker must still return its original successful status. + assert!( + attempt + .wait() + .await + .expect("completed worker status") + .success() + ); + attempt.cleanup().await.expect("reclaim completed attempt"); + assert!(!directory.exists()); + assert!(attempt.child.as_ref().unwrap().id().is_none()); + } + + #[tokio::test] + async fn exited_worker_with_live_descendant_stops_writer_before_cleanup() { + let temp = tempfile::tempdir().unwrap(); + let cache = super::super::image_cache_root_dir(temp.path()); + let launcher = temp.path().join("worker"); + fs::write(&launcher, "#!/bin/sh\nsleep 300 &\nexit 0\n").unwrap(); + fs::set_permissions(&launcher, fs::Permissions::from_mode(0o700)).unwrap(); + let mut attempt = Attempt::create(&cache).unwrap(); + let directory = attempt.directory.clone(); + let _stdout = attempt.spawn(&launcher, request(temp.path())).unwrap(); + wait_for_worker_exit_without_reaping(&attempt).await; + + #[cfg(target_os = "macos")] + { + let id = attempt.child.as_ref().unwrap().id().unwrap(); + assert!( + !attempt.exited_worker_has_no_descendant_owner(id).unwrap(), + "an exited leader cannot prove its live descendant stopped" + ); + assert!(directory.exists()); + } + attempt + .cleanup() + .await + .expect("kill the live descendant before reclaiming its staging"); + assert!(!directory.exists()); + assert!(attempt.child.as_ref().unwrap().id().is_none()); + } + + #[tokio::test] + async fn cleanup_kills_and_reaps_worker_and_formatter_before_removing_staging() { + let temp = tempfile::tempdir().unwrap(); + let cache = super::super::image_cache_root_dir(temp.path()); + let launcher = temp.path().join("worker"); + // The shell and the simulated formatter both inherit the attempt lease. + // Killing only the shell leaves the lease locked and makes cleanup fail. + fs::write(&launcher, "#!/bin/sh\nroot=$(dirname \"$2\")\nsleep 300 &\nprintf '%s' \"$!\" > \"$root/ready\"\nwait\n").unwrap(); + fs::set_permissions(&launcher, fs::Permissions::from_mode(0o700)).unwrap(); + let mut attempt = Attempt::create(&cache).unwrap(); + let directory = attempt.directory.clone(); + let _stdout = attempt.spawn(&launcher, request(temp.path())).unwrap(); + wait_for_file(&directory.join("ready")).await; + #[cfg(target_os = "macos")] + { + let id = attempt.child.as_ref().unwrap().id().unwrap(); + assert!(!attempt.exited_worker_has_no_descendant_owner(id).unwrap()); + assert!(attempt.lease.is_some(), "a live worker retains its lease"); + } + reconcile(&cache).unwrap(); + assert!(directory.exists(), "a live formatter protects its staging"); + attempt + .cleanup() + .await + .expect("stop and reclaim preparation"); + assert!(!directory.exists()); + assert!(!attempt.lease_path.exists()); + assert!( + attempt + .child + .as_mut() + .unwrap() + .try_wait() + .unwrap() + .is_some() + ); + attempt.cleanup().await.expect("cleanup is idempotent"); + } + + #[tokio::test] + async fn aborting_provisioning_still_runs_cleanup() { + let temp = tempfile::tempdir().unwrap(); + let cache = super::super::image_cache_root_dir(temp.path()); + let attempt = Arc::new(Mutex::new(Attempt::create(&cache).unwrap())); + let directory = attempt.lock().await.directory.clone(); + fs::write(directory.join("partial.ext4"), b"unfinished image").unwrap(); + let (ready, started) = tokio::sync::oneshot::channel(); + let task = tokio::spawn({ + let attempt = attempt.clone(); + async move { + let _cleanup = CleanupOnDrop(attempt); + ready.send(()).unwrap(); + std::future::pending::<()>().await; + } + }); + started.await.unwrap(); + task.abort(); + assert!(task.await.unwrap_err().is_cancelled()); + tokio::time::timeout(Duration::from_secs(5), async { + while directory.exists() { + tokio::time::sleep(Duration::from_millis(10)).await; + } + }) + .await + .expect("cancelled future must reclaim its staging"); + } + + #[test] + fn restart_reclaims_inactive_attempts_and_preserves_active_shared_and_unknown_data() { + let temp = tempfile::tempdir().unwrap(); + let cache = super::super::image_cache_root_dir(temp.path()); + let inactive = Attempt::create(&cache).unwrap(); + let inactive_directory = inactive.directory.clone(); + fs::write(inactive_directory.join("partial.ext4"), b"unfinished").unwrap(); + drop(inactive); + let active = Attempt::create(&cache).unwrap(); + let committed = cache.join("shared-image"); + fs::create_dir(&committed).unwrap(); + fs::write(committed.join("rootfs.ext4"), b"valid shared image").unwrap(); + let legacy = cache.join("image.staging-legacy"); + fs::create_dir(&legacy).unwrap(); + let unmarked = cache.join(ATTEMPTS_DIR).join(format!("attempt-{:032x}", 1)); + fs::create_dir(&unmarked).unwrap(); + let link = cache.join(ATTEMPTS_DIR).join(format!("attempt-{:032x}", 2)); + symlink(&committed, &link).unwrap(); + fs::write(unmarked.with_extension("lease"), b"unrelated lock").unwrap(); + + reconcile(&cache).unwrap(); + + assert!(!inactive_directory.exists()); + assert!(active.directory.exists()); + assert_eq!( + fs::read(committed.join("rootfs.ext4")).unwrap(), + b"valid shared image" + ); + assert!(legacy.exists()); + assert!(unmarked.exists()); + assert!(link.is_symlink()); + } + + #[tokio::test] + async fn separate_workers_serialize_cache_publication() { + let temp = tempfile::tempdir().unwrap(); + let cache = temp.path().to_path_buf(); + let first = cache_lock(&cache).unwrap(); + let (entered, mut receiver) = tokio::sync::mpsc::channel(1); + let second = tokio::task::spawn_blocking(move || { + let _second = cache_lock(&cache).unwrap(); + entered.blocking_send(()).unwrap(); + }); + assert!( + tokio::time::timeout(Duration::from_millis(50), receiver.recv()) + .await + .is_err() + ); + drop(first); + tokio::time::timeout(Duration::from_secs(5), receiver.recv()) + .await + .unwrap() + .unwrap(); + second.await.unwrap(); + } + + #[test] + fn worker_rejects_request_without_inherited_lease() { + let temp = tempfile::tempdir().unwrap(); + let cache = super::super::image_cache_root_dir(temp.path()); + let attempt = Attempt::create(&cache).unwrap(); + let path = attempt.directory.join(REQUEST_FILE); + fs::write(&path, serde_json::to_vec(&request(temp.path())).unwrap()).unwrap(); + let error = read_request(&path) + .err() + .expect("missing inherited fd rejected"); + assert!(error.contains("lease was not inherited")); + } +} diff --git a/crates/openshell-driver-vm/tests/image_preparation.rs b/crates/openshell-driver-vm/tests/image_preparation.rs new file mode 100644 index 0000000000..39d1e27c38 --- /dev/null +++ b/crates/openshell-driver-vm/tests/image_preparation.rs @@ -0,0 +1,197 @@ +// SPDX-FileCopyrightText: Copyright (c) 2025-2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved. +// SPDX-License-Identifier: Apache-2.0 + +#![cfg(feature = "compute-driver")] +#![allow(unsafe_code)] + +use std::fs::{self, File}; +use std::os::fd::AsRawFd; +use std::os::unix::fs::PermissionsExt; +use std::os::unix::process::CommandExt; +use std::process::{Command, Stdio}; + +use openshell_driver_vm::VmDriverConfig; + +/// Exercise the actual worker entry point and a real sparse-file copy without +/// downloading an image, running a VM, or requiring host filesystem utilities. +#[test] +fn worker_publishes_complete_overlay_and_reports_resolved_owner() { + worker_overlay_case(false); +} + +/// Retrying an interrupted first start must never turn an incomplete copy +/// into persisted overlay state. The next retry must still be able to start. +#[test] +fn interrupted_missing_overlay_retry_does_not_publish_partial_state() { + worker_overlay_case(true); +} + +fn worker_overlay_case(interrupt_retry: bool) { + let root = tempfile::tempdir().unwrap(); + let cache = root.path().join("images"); + let attempts = cache.join("preparations"); + let name = "attempt-00000000000000000000000000000001"; + let attempt = attempts.join(name); + fs::create_dir_all(&attempt).unwrap(); + let lease_path = attempts.join(format!("{name}.lease")); + fs::write(&lease_path, b"openshell-image-preparation-v1\n").unwrap(); + let lease = File::open(&lease_path).unwrap(); + let fd = lease.as_raw_fd(); + // SAFETY: this descriptor belongs to the live fixture file. + assert_eq!(unsafe { libc::flock(fd, libc::LOCK_EX) }, 0); + + let size = 1024 * 1024; + let template = cache.join("overlay-templates/sandbox-overlay-ext4-v1/1048576.ext4"); + fs::create_dir_all(template.parent().unwrap()).unwrap(); + let mut expected = vec![0_u8; size]; + expected[..13].copy_from_slice(b"template-data"); + fs::write(&template, &expected).unwrap(); + let sandbox = root.path().join("sandboxes/worker-test"); + fs::create_dir_all(&sandbox).unwrap(); + if interrupt_retry { + // A cancelled Fresh attempt can persist the owner before publishing + // its first overlay. Start retries this state with PreserveExisting. + fs::write( + sandbox.join("sandbox-owner-state"), + b"sandbox-owner-v2:1000:1000\n", + ) + .unwrap(); + } + let source = cache.join("test-image/rootfs.ext4"); + fs::create_dir_all(source.parent().unwrap()).unwrap(); + fs::write(&source, b"identity comes from driver configuration").unwrap(); + let config = VmDriverConfig { + state_dir: root.path().to_path_buf(), + sandbox_uid: Some(1000), + sandbox_gid: Some(1000), + overlay_disk_mib: 1, + ..VmDriverConfig::default() + }; + let request = attempt.join("request.json"); + fs::write( + &request, + serde_json::to_vec(&serde_json::json!({ + "config": config, + "sandbox_id": "worker-test", + "image_ref": "", + "rootfs_tar": null, + "bootstrap_only": false, + "overlay": { + "source_disk": source, + "preparation": if interrupt_retry { "PreserveExisting" } else { "Fresh" }, + }, + "lease_fd": fd, + "parent_pid": std::process::id(), + })) + .unwrap(), + ) + .unwrap(); + let mut command = Command::new(env!("CARGO_BIN_EXE_openshell-driver-vm")); + command + .arg("--internal-prepare-image") + .arg(&request) + .stdin(Stdio::null()) + .stdout(Stdio::piped()) + .stderr(Stdio::piped()) + .process_group(0); + // SAFETY: only fcntl runs between fork and exec, and the fixture keeps its + // descriptor alive until the worker has exited. + unsafe { + command.pre_exec(move || { + if libc::fcntl(fd, libc::F_SETFD, 0) == -1 { + return Err(std::io::Error::last_os_error()); + } + Ok(()) + }); + } + if interrupt_retry { + let tools = root.path().join("tools"); + fs::create_dir(&tools).unwrap(); + let copy_ready = root.path().join("copy-ready"); + let cp = tools.join("cp"); + // Control only the external copy command. The actual worker selects + // the destination and executes its production preservation checks. + fs::write( + &cp, + "#!/bin/sh\nfor argument in \"$@\"; do destination=\"$argument\"; done\nprintf partial > \"$destination\"\nprintf '%s' \"$destination\" > \"$COPY_READY\"\nsleep 300\n", + ) + .unwrap(); + fs::set_permissions(&cp, fs::Permissions::from_mode(0o700)).unwrap(); + command + .env("PATH", format!("{}:/usr/bin:/bin", tools.display())) + .env("COPY_READY", ©_ready); + let mut interrupted = command.spawn().unwrap(); + let deadline = std::time::Instant::now() + std::time::Duration::from_secs(10); + while !copy_ready.exists() && std::time::Instant::now() < deadline { + if interrupted.try_wait().unwrap().is_some() { + break; + } + std::thread::sleep(std::time::Duration::from_millis(10)); + } + // Kill and reap even if readiness fails, so failed tests do not leave + // a blocked helper or worker behind. + if interrupted.try_wait().unwrap().is_none() { + let pid = i32::try_from(interrupted.id()).unwrap(); + let _ = nix::sys::signal::killpg( + nix::unistd::Pid::from_raw(pid), + nix::sys::signal::Signal::SIGKILL, + ); + } + interrupted.wait().unwrap(); + assert!(copy_ready.exists(), "worker must reach the controlled copy"); + assert!( + !sandbox.join("overlay.ext4").exists(), + "interrupted retry must leave no partial persisted overlay" + ); + let destination = fs::read_to_string(©_ready).unwrap(); + assert!(std::path::Path::new(&destination).starts_with(&attempt)); + // Retry using the real copy utility. The incomplete attempt image + // must not prevent publishing the complete overlay. + command + .env("PATH", "/usr/bin:/bin") + .env_remove("COPY_READY"); + } + let mut child = command.spawn().unwrap(); + let deadline = std::time::Instant::now() + std::time::Duration::from_secs(10); + while child.try_wait().unwrap().is_none() { + if std::time::Instant::now() >= deadline { + let pid = i32::try_from(child.id()).unwrap(); + let _ = nix::sys::signal::killpg( + nix::unistd::Pid::from_raw(pid), + nix::sys::signal::Signal::SIGKILL, + ); + let _ = child.wait(); + panic!("image preparation worker timed out"); + } + std::thread::sleep(std::time::Duration::from_millis(10)); + } + let output = child.wait_with_output().unwrap(); + assert!( + output.status.success(), + "{}", + String::from_utf8_lossy(&output.stderr) + ); + let messages = String::from_utf8(output.stdout).unwrap(); + let completion = messages + .lines() + .map(|line| { + serde_json::from_str::(line) + .expect("worker stdout must contain only protocol messages") + }) + .find_map(|message| message.get("Complete").cloned()) + .expect("worker completion"); + assert_eq!( + completion, + serde_json::json!({"Ok": {"Overlay": {"uid": 1000, "gid": 1000}}}) + ); + assert_eq!(fs::read(sandbox.join("overlay.ext4")).unwrap(), expected); + assert_eq!( + fs::read(&template).unwrap(), + expected, + "the shared template is immutable" + ); + assert_eq!( + fs::read_to_string(sandbox.join("sandbox-owner-state")).unwrap(), + "sandbox-owner-v2:1000:1000\n" + ); +}