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 f361a3fdb0..f737b6e79e 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; @@ -2211,6 +2242,7 @@ impl VmDriver { state_dir: state_dir.clone(), process: None, provisioning_task: None, + preparation: None, gpu_bdf: None, deleting: false, }); @@ -2310,6 +2342,7 @@ impl VmDriver { state_dir: state_dir.clone(), process: None, provisioning_task: None, + preparation: None, gpu_bdf: None, deleting: false, }, @@ -2657,7 +2690,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; @@ -2672,10 +2705,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); } @@ -2693,6 +2732,235 @@ 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, + ) + }); + // Independent attempts can prepare images and private overlays at the + // same time. Workers serialize only shared cache publication. + 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), @@ -2820,28 +3088,53 @@ 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"); + let cache_root = image_cache_root_dir(&self.config.state_dir); tokio::task::spawn_blocking(move || { - ensure_sandbox_overlay_template_image(&template_path, overlay_size_bytes) + ensure_sandbox_overlay_template_image( + &cache_root, + &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}"))?; @@ -3252,7 +3545,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?; @@ -3365,7 +3658,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?; @@ -3506,7 +3799,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); @@ -3704,17 +3997,28 @@ impl VmDriver { return Err(Status::failed_precondition(message)); } - if tokio::fs::metadata(&image_path).await.is_ok() { - let _ = tokio::fs::remove_dir_all(staging_dir).await; - return Ok(()); - } - tokio::fs::rename(&prepared_image, &image_path) - .await - .map_err(|err| Status::internal(format!("store prepared image disk failed: {err}")))?; + self.publish_prepared_image(&prepared_image, &image_path) + .await?; let _ = tokio::fs::remove_dir_all(staging_dir).await; Ok(()) } + async fn publish_prepared_image( + &self, + staged: &Path, + destination: &Path, + ) -> Result<(), Status> { + let cache_root = image_cache_root_dir(&self.config.state_dir); + let staged = staged.to_path_buf(); + let destination = destination.to_path_buf(); + tokio::task::spawn_blocking(move || { + preparation::publish_cache_file(&cache_root, &staged, &destination, None) + }) + .await + .map_err(|error| Status::internal(format!("cache publication panicked: {error}")))? + .map_err(|error| Status::internal(format!("store cached rootfs image failed: {error}"))) + } + #[allow(clippy::similar_names)] async fn run_image_prep_vm( &self, @@ -3726,7 +4030,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); @@ -3821,7 +4126,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); @@ -3928,14 +4233,8 @@ impl VmDriver { return Err(Status::failed_precondition(err)); } - if tokio::fs::metadata(&image_path).await.is_ok() { - let _ = tokio::fs::remove_dir_all(&staging_dir).await; - return Ok(()); - } - - tokio::fs::rename(&prepared_image, &image_path) - .await - .map_err(|err| Status::internal(format!("store cached rootfs image failed: {err}")))?; + self.publish_prepared_image(&prepared_image, &image_path) + .await?; let _ = tokio::fs::remove_dir_all(&staging_dir).await; Ok(()) } @@ -3952,7 +4251,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); @@ -4079,23 +4378,8 @@ impl VmDriver { return Err(Status::failed_precondition(err)); } - if tokio::fs::metadata(&image_path).await.is_ok() { - info!( - image_identity = %image_identity, - "vm driver: another task wrote image while we were building, discarding ours" - ); - let _ = tokio::fs::remove_dir_all(&staging_dir).await; - return Ok(()); - } - - tokio::fs::rename(&prepared_image, &image_path) - .await - .map_err(|err| Status::internal(format!("store cached rootfs image failed: {err}")))?; - info!( - image_identity = %image_identity, - image_path = %image_path.display(), - "vm driver: root disk image committed to cache" - ); + self.publish_prepared_image(&prepared_image, &image_path) + .await?; let _ = tokio::fs::remove_dir_all(&staging_dir).await; Ok(()) } @@ -4305,6 +4589,89 @@ 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 (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; @@ -6277,8 +6644,10 @@ 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() @@ -6300,30 +6669,23 @@ 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)?; - fs::rename(&staging_image, template_path).map_err(|err| { - format!( - "move overlay template {} to {}: {err}", - staging_image.display(), - template_path.display() - ) - }) + preparation::publish_cache_file(cache_root, &staging_image, template_path, Some(size_bytes)) + .map_err(|err| { + format!( + "move overlay template {} to {}: {err}", + staging_image.display(), + template_path.display() + ) + }) })(); - if result.is_err() { - let _ = fs::remove_file(&staging_image); - } + // A competing publisher may have won; only our private staged file is disposable. + let _ = fs::remove_file(&staging_image); result } @@ -8467,6 +8829,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 { @@ -9131,6 +9548,7 @@ mod tests { state_dir: state_dir.clone(), process: None, provisioning_task: None, + preparation: None, gpu_bdf: None, deleting: false, }, @@ -9213,6 +9631,7 @@ mod tests { state_dir, process: None, provisioning_task: Some(provisioning_task), + preparation: None, gpu_bdf: None, deleting: false, }, @@ -9233,6 +9652,273 @@ mod tests { task.abort(); } + #[tokio::test] + async fn blocked_worker_does_not_delay_independent_image_and_overlay_requests() { + let root = tempfile::tempdir().unwrap(); + let mut driver = test_driver_with_extensions(LifecycleExtensionRegistry::new()); + driver.config.state_dir = root.path().to_path_buf(); + driver.config.bootstrap_image = "test-bootstrap".to_string(); + driver.launcher_bin = root.path().join("worker"); + let cached_root = root.path().join("images/warm/rootfs.ext4"); + fs::create_dir_all(cached_root.parent().unwrap()).unwrap(); + fs::write(&cached_root, b"cached bootstrap").unwrap(); + let images = + preparation::Message::Complete(Ok(preparation::Output::Images(RuntimeImagePlan { + root_disk: cached_root.clone(), + image_disk: None, + image_identity: "warm".to_string(), + bootstrap_image_identity: "warm".to_string(), + }))); + let overlay = preparation::Message::Complete(Ok(preparation::Output::Overlay( + SandboxOwnerIdentity { + uid: 1000, + gid: 1000, + }, + ))); + fs::write( + root.path().join("images.json"), + serde_json::to_vec(&images).unwrap(), + ) + .unwrap(); + fs::write( + root.path().join("overlay.json"), + serde_json::to_vec(&overlay).unwrap(), + ) + .unwrap(); + fs::write( + &driver.launcher_bin, + r#"#!/bin/sh + control=$(dirname "$0") + request=$(cat "$2") + case "$request" in + *'"sandbox_id":"slow-a"'*) + : > "$control/ready" + exec sleep 300 + ;; + *'"overlay":null'*) cat "$control/images.json" ;; + *) cat "$control/overlay.json" ;; + esac + printf '\n' + "#, + ) + .unwrap(); + fs::set_permissions(&driver.launcher_bin, fs::Permissions::from_mode(0o700)).unwrap(); + for id in ["slow-a", "warm-b"] { + let state_dir = sandboxes_root_dir(root.path()).join(id); + create_private_dir_all(&state_dir).await.unwrap(); + driver.registry.lock().await.insert( + id.to_string(), + SandboxRecord { + snapshot: Sandbox { + id: id.to_string(), + ..Default::default() + }, + state_dir, + process: None, + provisioning_task: None, + preparation: None, + gpu_bdf: None, + deleting: false, + }, + ); + } + + let slow_driver = driver.clone(); + let slow = tokio::spawn(async move { + slow_driver + .prepare_images_in_worker("slow-a", "test-bootstrap", None, true) + .await + }); + let warm_driver = driver.clone(); + let ready = root.path().join("ready"); + let mut warm = tokio::spawn(async move { + while !ready.exists() { + tokio::time::sleep(Duration::from_millis(10)).await; + } + let plan = warm_driver + .prepare_images_in_worker("warm-b", "test-bootstrap", None, true) + .await + .map_err(|error| error.to_string())?; + if plan.root_disk != cached_root { + return Err("warm image result was not preserved".to_string()); + } + let owner = warm_driver + .prepare_overlay_in_worker( + "warm-b", + &plan.root_disk, + OverlayPreparation::Fresh, + None, + ) + .await + .map_err(|error| error.to_string())?; + if owner + != (SandboxOwnerIdentity { + uid: 1000, + gid: 1000, + }) + { + return Err("warm overlay result was not preserved".to_string()); + } + Ok::<(), String>(()) + }); + let outcome = match tokio::time::timeout(Duration::from_secs(10), &mut warm).await { + Ok(Ok(result)) => result, + Ok(Err(error)) => Err(format!("warm task failed: {error}")), + Err(_) => { + warm.abort(); + let _ = warm.await; + Err("warm request waited for the unrelated slow worker".to_string()) + } + }; + let slow_still_running = !slow.is_finished(); + slow.abort(); + let _ = slow.await; + let warm_cleanup = driver.cleanup_image_preparation("warm-b").await; + let slow_cleanup = driver.cleanup_image_preparation("slow-a").await; + warm_cleanup.expect("warm worker cleanup"); + slow_cleanup.expect("slow worker cleanup"); + outcome.expect("independent image and overlay requests must complete"); + assert!( + slow_still_running, + "the slow worker must remain blocked during the proof" + ); + } + + 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, @@ -9356,6 +10042,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()), @@ -9395,6 +10082,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()), @@ -9429,6 +10117,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()), @@ -9457,6 +10146,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()), @@ -9486,6 +10176,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()), @@ -9510,6 +10201,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()), @@ -9531,6 +10223,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()), @@ -10177,6 +10870,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()), @@ -10242,6 +10936,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()), @@ -10262,6 +10957,7 @@ mod tests { state_dir: state_dir.clone(), process: None, provisioning_task: None, + preparation: None, gpu_bdf: None, deleting: false, }, @@ -10296,6 +10992,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()), @@ -10317,6 +11014,7 @@ mod tests { state_dir: state_dir.clone(), process: None, provisioning_task: None, + preparation: None, gpu_bdf: None, deleting: false, }, @@ -10671,6 +11369,7 @@ mod tests { state_dir, process: Some(process), provisioning_task: None, + preparation: None, gpu_bdf: None, deleting: false, }, @@ -10698,6 +11397,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()), @@ -11041,6 +11741,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), @@ -11649,6 +12350,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..51f087e5e5 --- /dev/null +++ b/crates/openshell-driver-vm/src/preparation.rs @@ -0,0 +1,1016 @@ +// 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) +} + +/// Publish a complete staged image without replacing a ready cache entry. +/// Preparation stays outside this lock; only the readiness recheck and atomic +/// rename are serialized across workers. Call from a blocking task so the +/// lock remains owned until publication finishes even if its waiter is dropped. +pub(super) fn publish_cache_file( + cache_root: &Path, + staged: &Path, + destination: &Path, + expected_size: Option, +) -> io::Result<()> { + let _lease = cache_lock(cache_root)?; + match fs::metadata(destination) { + Ok(metadata) + if expected_size.is_none_or(|size| metadata.is_file() && metadata.len() == size) => + { + return Ok(()); + } + Ok(_) => {} + Err(error) if error.kind() == io::ErrorKind::NotFound => {} + Err(error) => return Err(error), + } + fs::rename(staged, destination) +} + +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 cache_lock_child_process_helper() { + let Some(root) = std::env::var_os("OPENSHELL_TEST_PUBLICATION_LOCK_ROOT") else { + return; + }; + let root = PathBuf::from(root); + let name = std::env::var("OPENSHELL_TEST_PUBLICATION_LOCK_NAME").unwrap(); + fs::write(root.join(format!("{name}.waiting")), b"waiting").unwrap(); + if name == "second" { + let probe = OpenOptions::new() + .read(true) + .write(true) + .open(root.join("preparation.lock")) + .unwrap(); + assert!( + !lock(&probe, true).unwrap(), + "first process must own the kernel lock" + ); + fs::write(root.join("second.contended"), b"contended").unwrap(); + } + if name == "second" { + publish_cache_file( + &root, + &root.join("second.staged"), + &root.join("committed"), + Some(5), + ) + .unwrap(); + fs::write(root.join("second.acquired"), b"published").unwrap(); + return; + } + let _lock = cache_lock(&root).unwrap(); + fs::write(root.join(format!("{name}.acquired")), b"acquired").unwrap(); + let mut release = [0_u8; 1]; + io::stdin().read_exact(&mut release).unwrap(); + } + + struct LockHelper(std::process::Child); + + impl Drop for LockHelper { + fn drop(&mut self) { + let _ = self.0.kill(); + let _ = self.0.wait(); + } + } + + #[test] + fn publication_lock_is_exclusive_across_processes_and_released_on_death() { + let root = tempfile::tempdir().unwrap(); + let spawn = |name: &str| { + LockHelper( + std::process::Command::new(std::env::current_exe().unwrap()) + .args([ + "--exact", + "driver::preparation::tests::cache_lock_child_process_helper", + "--nocapture", + ]) + .env("OPENSHELL_TEST_PUBLICATION_LOCK_ROOT", root.path()) + .env("OPENSHELL_TEST_PUBLICATION_LOCK_NAME", name) + .stdin(Stdio::piped()) + .stdout(Stdio::null()) + .stderr(Stdio::inherit()) + .spawn() + .unwrap(), + ) + }; + let wait_for = |name: &str| { + let deadline = std::time::Instant::now() + Duration::from_secs(5); + while !root.path().join(name).exists() { + assert!( + std::time::Instant::now() < deadline, + "helper did not reach {name}" + ); + std::thread::sleep(Duration::from_millis(10)); + } + }; + fs::write(root.path().join("second.staged"), b"later").unwrap(); + let mut first = spawn("first"); + wait_for("first.acquired"); + let mut second = spawn("second"); + wait_for("second.contended"); + assert!( + !root.path().join("second.acquired").exists(), + "second process bypassed lock" + ); + // The lock holder commits before dying. The queued publisher must read + // readiness only after it owns the lock and preserve this complete image. + fs::write(root.path().join("committed"), b"first").unwrap(); + first.0.kill().unwrap(); + first.0.wait().unwrap(); + wait_for("second.acquired"); + assert!(second.0.wait().unwrap().success()); + assert_eq!(fs::read(root.path().join("committed")).unwrap(), b"first"); + assert!(root.path().join("second.staged").exists()); + assert!(root.path().join("preparation.lock").exists()); + // An invalid-sized template can be replaced; a valid rootfs cannot. + fs::write(root.path().join("second.staged"), b"new!").unwrap(); + publish_cache_file( + root.path(), + &root.path().join("second.staged"), + &root.path().join("committed"), + Some(4), + ) + .unwrap(); + assert_eq!(fs::read(root.path().join("committed")).unwrap(), b"new!"); + fs::write(root.path().join("third.staged"), b"third").unwrap(); + publish_cache_file( + root.path(), + &root.path().join("third.staged"), + &root.path().join("committed"), + None, + ) + .unwrap(); + assert_eq!(fs::read(root.path().join("committed")).unwrap(), b"new!"); + } + + #[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..ffa0992c9e --- /dev/null +++ b/crates/openshell-driver-vm/tests/image_preparation.rs @@ -0,0 +1,383 @@ +// 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" + ); +} + +struct OwnedPreparationChild { + child: std::process::Child, + reaped: bool, + _lease: File, +} + +impl OwnedPreparationChild { + fn exited(&self) -> std::io::Result { + // SAFETY: initialized output storage; WNOWAIT retains our child's PID. + let mut status: libc::siginfo_t = unsafe { std::mem::zeroed() }; + if unsafe { + libc::waitid( + libc::P_PID, + self.child.id(), + &raw mut status, + libc::WEXITED | libc::WNOHANG | libc::WNOWAIT, + ) + } != 0 + { + return Err(std::io::Error::last_os_error()); + } + Ok(matches!( + status.si_code, + libc::CLD_EXITED | libc::CLD_KILLED | libc::CLD_DUMPED + )) + } + + fn stop_and_reap(&mut self) -> std::io::Result { + let pid = i32::try_from(self.child.id()).map_err(std::io::Error::other)?; + let _ = nix::sys::signal::killpg( + nix::unistd::Pid::from_raw(pid), + nix::sys::signal::Signal::SIGKILL, + ); + let status = self.child.wait()?; + self.reaped = true; + Ok(status) + } +} + +impl Drop for OwnedPreparationChild { + fn drop(&mut self) { + if !self.reaped { + let _ = self.stop_and_reap(); + } + } +} + +fn spawn_overlay_fixture( + root: &std::path::Path, + id: &str, + number: u128, + controlled_copy: bool, +) -> OwnedPreparationChild { + let attempts = root.join("images/preparations"); + let name = format!("attempt-{number:032x}"); + 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: lease owns fd and stays alive in the returned process guard. + assert_eq!(unsafe { libc::flock(fd, libc::LOCK_EX) }, 0); + fs::create_dir_all(root.join("sandboxes").join(id)).unwrap(); + let config = VmDriverConfig { + state_dir: root.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": id, "image_ref": "", "rootfs_tar": null, + "bootstrap_only": false, + "overlay": { "source_disk": root.join("images/test-image/rootfs.ext4"), + "preparation": "Fresh", "requested_identity": null }, + "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(File::create(root.join(format!("{id}.stdout"))).unwrap()) + .stderr(File::create(root.join(format!("{id}.stderr"))).unwrap()) + .process_group(0); + if controlled_copy { + command + .env( + "PATH", + format!("{}:/usr/bin:/bin", root.join("tools").display()), + ) + .env("COPY_READY", root.join("copy-ready")); + } else { + command.env("PATH", "/usr/bin:/bin"); + } + // SAFETY: only fcntl runs between fork and exec; the guard keeps fd alive. + unsafe { + command.pre_exec(move || { + if libc::fcntl(fd, libc::F_SETFD, 0) == -1 { + return Err(std::io::Error::last_os_error()); + } + Ok(()) + }); + } + OwnedPreparationChild { + child: command.spawn().unwrap(), + reaped: false, + _lease: lease, + } +} + +#[test] +fn independent_overlay_finishes_while_another_copy_is_blocked() { + let root = tempfile::tempdir().unwrap(); + let template = root + .path() + .join("images/overlay-templates/sandbox-overlay-ext4-v1/1048576.ext4"); + fs::create_dir_all(template.parent().unwrap()).unwrap(); + let mut expected = vec![0_u8; 1024 * 1024]; + expected[..13].copy_from_slice(b"template-data"); + fs::write(&template, &expected).unwrap(); + let source = root.path().join("images/test-image/rootfs.ext4"); + fs::create_dir_all(source.parent().unwrap()).unwrap(); + fs::write(&source, b"configured owner fixture").unwrap(); + let tools = root.path().join("tools"); + fs::create_dir(&tools).unwrap(); + fs::write(tools.join("cp"), + "#!/bin/sh\nfor value in \"$@\"; do destination=\"$value\"; done\nprintf partial > \"$destination\"\n: > \"$COPY_READY\"\nexec sleep 300\n" + ).unwrap(); + fs::set_permissions(tools.join("cp"), fs::Permissions::from_mode(0o700)).unwrap(); + let mut slow = spawn_overlay_fixture(root.path(), "slow-a", 1, true); + let deadline = std::time::Instant::now() + std::time::Duration::from_secs(10); + while !root.path().join("copy-ready").exists() { + assert!( + !slow.exited().unwrap(), + "slow worker exited before controlled copy" + ); + assert!( + std::time::Instant::now() < deadline, + "controlled copy readiness timed out" + ); + std::thread::sleep(std::time::Duration::from_millis(10)); + } + let mut warm = spawn_overlay_fixture(root.path(), "warm-b", 2, false); + let deadline = std::time::Instant::now() + std::time::Duration::from_secs(10); + while !warm.exited().unwrap() { + assert!( + std::time::Instant::now() < deadline, + "warm overlay waited for unrelated copy" + ); + std::thread::sleep(std::time::Duration::from_millis(10)); + } + let status = warm.stop_and_reap().unwrap(); + assert!( + status.success(), + "{}", + fs::read_to_string(root.path().join("warm-b.stderr")).unwrap() + ); + assert!( + !slow.exited().unwrap(), + "slow copy must remain blocked during the proof" + ); + assert_eq!( + fs::read(root.path().join("sandboxes/warm-b/overlay.ext4")).unwrap(), + expected + ); + assert_eq!(fs::read(&template).unwrap(), expected); + assert!(!root.path().join("sandboxes/slow-a/overlay.ext4").exists()); + let output = fs::read_to_string(root.path().join("warm-b.stdout")).unwrap(); + let completed = output.lines().any(|line| { + serde_json::from_str::(line) + .unwrap() + .get("Complete") + == Some(&serde_json::json!({"Ok": {"Overlay": {"uid": 1000, "gid": 1000}}})) + }); + assert!(completed, "actual worker did not report overlay completion"); + slow.stop_and_reap().unwrap(); +}