diff --git a/docs/backends/wslc/wslc-state-aware.md b/docs/backends/wslc/wslc-state-aware.md index dde7f0632..ed975fb3b 100644 --- a/docs/backends/wslc/wslc-state-aware.md +++ b/docs/backends/wslc/wslc-state-aware.md @@ -43,9 +43,10 @@ Sandbox daemon pattern. The daemon owns the live SDK handles on a **single apartment-affine worker thread**, which services every lifecycle command. Any thread that has joined the MTA may use those handles, so an image pull and an `exec` each run on an MTA thread of their own and post their outcome back to the -worker, leaving it free to serve other sandboxes for the duration of a run. A command naming a -container with a run in flight waits for that run, because deleting the container would free a -handle the run is using. See [Known limitations](#known-limitations). +worker, leaving it free to serve other sandboxes for the duration of a run. A second `exec` on a +container with a run in flight is refused as `busy`; a lifecycle command naming that container +waits for the run, because deleting the container would free a handle the run is using. See +[Known limitations](#known-limitations). ## Components @@ -143,11 +144,39 @@ the streaming path reports it through its wait result. ### exec admission, cancellation, and failure containment -The daemon admits one exec at a time. A concurrent exec is rejected with -`backend_error` rather than queued behind an unknown-duration workload. Up to -eight additional control connections can be serviced while an exec owns the -stream slot; connections beyond the daemon's bounded client capacity are -refused. +The daemon admits up to eight exec streams at once, one per container. A second +exec naming a container that already has one in flight is refused with `busy` +before any admission reaches the client; the outward SDK error is +`backend_error`. The bound of eight comes from memory: a streaming exec's output +lives only in its bounded live-output queue, so the persistent per-user daemon +stays near 128 MB of live output even against clients that never drain. An +exec's slot is held until the run has reported back and its client has been +written to, so a client that disconnects mid-run keeps counting against that +bound while its process is still going, and one that drains slowly keeps +counting while its output is still queued. A run whose termination could not be +confirmed leaves its sandbox quarantined and keeps the slot until that sandbox +is deprovisioned, because the process may still be alive. + +Three conditions surface as `busy`, and all reach an SDK caller as +`backend_error`: + +| Condition | Message | Retry | +| --- | --- | --- | +| The daemon's eight exec slots are all occupied | `WSLc daemon exec capacity is exhausted` | Succeeds once any exec finishes | +| The named container already has an exec in flight | `sandbox already has an exec in flight` | Succeeds once that container's run finishes | +| Every client slot is occupied and the request is not a cancellation | `WSLc daemon client capacity is exhausted` | Succeeds once any client disconnects | + +Up to eight additional control connections can be serviced while every exec slot +is occupied. Beyond that, a connection is admitted only to cancel: cancellation +is the one request that never waits on the worker, and the only way to end a run +with no timeout, so it keeps capacity of its own that lifecycle work cannot +consume. A connection on that lane must send its request within two seconds and +in under 1 KB, so one that connects and stalls cannot hold a cancellation slot +for the general deadline. + +`start` / `stop` / `deprovision` naming a container with an exec in flight +**wait** for that run, because deleting the container would free a handle the +run is still using. Each exec carries an internal ID and per-run token. A duplicate live ID is rejected, and cancellation must match both values so a delayed cancellation @@ -277,9 +306,14 @@ it can observe idle-teardown within seconds. host that can reach a registry or already has the image cached, and `wxc-wslc-daemon.exe` staged next to `wxc-exec.exe`). It exercises core lifecycle, warm-reuse (a marker written by one `exec` is read back by a separate `exec` process — only possible if the container stayed warm), filesystem volumes, bridged networking + -proxy, validation rejections, and idle teardown. Fixtures live in +proxy, validation rejections, exec concurrency, and idle teardown. Fixtures live in `tests/configs/wslc_state_aware_*.json`. +The concurrency section launches a phase without waiting for it (`Start-StateAware` / +`Wait-StateAware`), which is what lets it observe two sandboxes running at once, a refused +same-container second exec, and a lifecycle command issued while a run is in flight. Every other +section drives one phase process at a time. + ### Running the fixtures (ordering + id substitution) The `wslc_state_aware_*.json` fixtures are **stateful** — unlike the one-shot configs, they cannot be @@ -302,15 +336,15 @@ fixtures **through the harness**, not by pointing `wxc-exec --config` at them di ## Known limitations -- **Multiple exec streams are deferred.** The daemon admits one exec stream at a time, so a - client's concurrent exec is refused rather than run alongside the first, and the per-container - single-flight slot is not reported as `Busy`. Raising that bound is tracked as follow-up work. - Lifecycle calls on another sandbox are a separate matter: they are admitted through the control - client reserve and proceed while a run is in flight. +- **Ordering is per-container, not global.** A lifecycle command naming a container with a run in + flight waits for that run; commands for other sandboxes proceed independently. A caller cannot + infer that work on one sandbox completed because work on another did. -- **Ordering is per-container, not global.** Commands naming a container with a run in flight wait - for that run; commands for other sandboxes proceed independently. A caller cannot infer that - work on one sandbox completed because work on another did. +- **`busy` collapses to `backend_error` (deferred).** A refused exec reaches an SDK caller as a + generic `backend_error` with no indication that retrying would succeed. A retryable wire code + needs a new `MxcErrorCode` variant, which is a closed set matching the SDK `ErrorCode` union + one-for-one, so it spans Rust, Node, .NET, the versioned references, and schema regeneration. + This is tracked as follow-up work. - **No typed SDK can set port mappings yet.** The Rust, Node, and .NET v1 SDKs all pin the published stable contract `1.0.0`, which does not declare the diff --git a/src/mxc-sdk/src/backends/wslc/common/container_steps.rs b/src/mxc-sdk/src/backends/wslc/common/container_steps.rs index 1c47f20d2..1d888795b 100644 --- a/src/mxc-sdk/src/backends/wslc/common/container_steps.rs +++ b/src/mxc-sdk/src/backends/wslc/common/container_steps.rs @@ -118,11 +118,8 @@ pub enum OutStream { Stderr, } -/// Optional live-output sink invoked from the SDK's stdout/stderr callbacks in -/// addition to the capped capture buffers. The daemon supplies one to stream a -/// container's output to the client as bytes arrive; paths that only need the -/// final captured blob (one-shot, detached init) leave it unset. It receives -/// the full callback bytes, independent of the capped buffers' truncation. +/// Optional live-output sink that replaces the capped capture buffers, so a +/// caller that sets one gets empty captured stdout/stderr. /// /// **Two distinct live-output architectures — why they don't share plumbing.** /// The one-shot runner streams via `OutputMode::Stream` (in `wsl_container_runner`) @@ -137,7 +134,7 @@ pub enum OutStream { /// block** — the same SDK thread also delivers the process-exit callback, so the /// daemon path uses a non-blocking `try_send` that drops on overflow rather than /// stalling teardown. Keep the two paths separate for those reasons; share only -/// the leaf primitives ([`OutStream`], the capped capture buffers). +/// the leaf primitives ([`OutStream`], [`IoContext`]). pub type OutputSink = Box; /// Shared buffer for capturing process I/O via SDK callbacks. Fields are @@ -147,10 +144,11 @@ pub struct IoContext { pub(crate) stdout: Arc>>, pub(crate) stderr: Arc>>, pub(crate) exited: Arc<(Mutex, Condvar)>, - /// Live sink for streaming output alongside the capped buffers; `None` when - /// only the final captured blob is needed. Its owning `Arc` is - /// released on the same schedule as the capture buffers, so on the - /// deliberate kill-path leak the sink (and its sender) is leaked too. + /// Live sink for streaming output; `None` when the caller wants the final + /// captured blob instead. A streaming caller reads its bytes from the sink + /// as they arrive, so the capture buffers stay empty. Its owning + /// `Arc` is released on the same schedule as those buffers, so + /// on the deliberate kill-path leak the sink (and its sender) is leaked too. sink: Option, } @@ -201,24 +199,20 @@ unsafe extern "C" fn io_callback( let ctx = &*(context as *const IoContext); let bytes = std::slice::from_raw_parts(data, data_size as usize); match io_handle { - WslcProcessIOHandle::WSLC_PROCESS_IO_HANDLE_STDOUT => { - { + WslcProcessIOHandle::WSLC_PROCESS_IO_HANDLE_STDOUT => match ctx.sink.as_ref() { + Some(sink) => sink(OutStream::Stdout, bytes), + None => { let mut buf = ctx.stdout.lock().unwrap_or_else(|e| e.into_inner()); append_capped(&mut buf, bytes); } - if let Some(sink) = ctx.sink.as_ref() { - sink(OutStream::Stdout, bytes); - } - } - WslcProcessIOHandle::WSLC_PROCESS_IO_HANDLE_STDERR => { - { + }, + WslcProcessIOHandle::WSLC_PROCESS_IO_HANDLE_STDERR => match ctx.sink.as_ref() { + Some(sink) => sink(OutStream::Stderr, bytes), + None => { let mut buf = ctx.stderr.lock().unwrap_or_else(|e| e.into_inner()); append_capped(&mut buf, bytes); } - if let Some(sink) = ctx.sink.as_ref() { - sink(OutStream::Stderr, bytes); - } - } + }, _ => {} } } @@ -1351,6 +1345,45 @@ mod tests { ); } + /// A streaming caller reads its bytes from the sink, so a second capped + /// copy would double the daemon's per-exec output memory. + #[test] + fn a_streaming_callback_captures_nothing() { + let streamed = Arc::new(Mutex::new(Vec::new())); + let seen = Arc::clone(&streamed); + let io_ctx = Arc::new(IoContext { + stdout: Arc::new(Mutex::new(Vec::new())), + stderr: Arc::new(Mutex::new(Vec::new())), + exited: Arc::new((Mutex::new(false), Condvar::new())), + sink: Some(Box::new(move |kind, bytes| { + seen.lock().unwrap().push((kind, bytes.to_vec())); + })), + }); + + let payload = b"streamed-not-captured"; + let raw = Arc::into_raw(Arc::clone(&io_ctx)); + // SAFETY: `raw` is a live `Arc` pointer and `payload` is + // valid for its full length, matching what the SDK passes. + unsafe { + io_callback( + WslcProcessIOHandle::WSLC_PROCESS_IO_HANDLE_STDOUT, + payload.as_ptr(), + payload.len() as u32, + raw as *mut c_void, + ); + drop(Arc::from_raw(raw)); + } + + assert_eq!( + streamed.lock().unwrap().as_slice(), + &[(OutStream::Stdout, payload.to_vec())] + ); + assert!( + io_ctx.stdout.lock().unwrap().is_empty(), + "a streamed chunk must not also be captured" + ); + } + #[test] fn null_exit_event_still_observes_cancellation() { let io_ctx = IoContext { diff --git a/src/mxc-sdk/src/bin/wslc_daemon/control_server.rs b/src/mxc-sdk/src/bin/wslc_daemon/control_server.rs index 6a4c6ecfd..ccd40d089 100644 --- a/src/mxc-sdk/src/bin/wslc_daemon/control_server.rs +++ b/src/mxc-sdk/src/bin/wslc_daemon/control_server.rs @@ -43,28 +43,59 @@ use windows::Win32::System::Threading::{GetCurrentProcess, OpenProcessToken}; use mxc_sdk::wslc_common::container_steps::OutStream; use mxc_sdk::wslc_common::daemon_protocol::{ encode_frame, DaemonRequest, DaemonResponse, DeprovisionConfig, ExecTerminal, StreamFrame, - MAX_FRAME_SIZE, + MAX_EXEC_ID_BYTES, MAX_FRAME_SIZE, }; -use crate::session_manager::{ExecStream, SessionHandle, WorkerError}; +use crate::session_manager::{ExecSlotGuard, ExecStream, SessionHandle, WorkerError}; -/// Exec streams admitted at once. A further request is refused rather than +/// Exec streams admitted at once; a further request is refused rather than /// queued behind a workload of unknown duration. -const MAX_CONCURRENT_EXECS: usize = 1; - -/// Capacity reserved for cancellation and lifecycle requests while all exec -/// stream slots are occupied. -const CONTROL_CLIENT_RESERVE: usize = 8; +/// +/// Each admitted exec owns a bounded live-output queue of +/// `LIVE_OUTPUT_CHANNEL_CAPACITY` x `LIVE_OUTPUT_MAX_CHUNK_BYTES`, and a +/// streaming exec captures nothing alongside it, so this bound holds the +/// persistent per-user daemon near 128 MB of live output against clients that +/// never drain. +const MAX_CONCURRENT_EXECS: usize = 8; + +/// Client capacity beyond the exec cap, so lifecycle work is still serviced +/// while every exec slot is occupied. +/// +/// A lifecycle command naming a container with a run in flight parks on the +/// worker and holds its slot for the whole wait, so this bounds how many such +/// waits can be outstanding. +const CONTROL_CLIENT_HEADROOM: usize = 8; /// Upper bound on concurrently-serviced client connections. Connections beyond /// this bound are refused without blocking the accept loop. -const MAX_CONCURRENT_CLIENTS: usize = MAX_CONCURRENT_EXECS + CONTROL_CLIENT_RESERVE; +const MAX_CONCURRENT_CLIENTS: usize = MAX_CONCURRENT_EXECS + CONTROL_CLIENT_HEADROOM; + +/// Connections admitted for cancellation alone once [`MAX_CONCURRENT_CLIENTS`] +/// is reached, one per exec that could need cancelling. +/// +/// Cancellation is the only way to end a run with no timeout, and the one +/// request that never waits on the worker, so it keeps capacity that lifecycle +/// work cannot consume. +const CANCEL_LANE_SLOTS: usize = MAX_CONCURRENT_EXECS; /// Deadline for a freshly-connected client to send its first (request) frame. A /// client that connects and then stalls must not pin a handler task — and a /// concurrency slot — indefinitely. const FIRST_FRAME_TIMEOUT: Duration = Duration::from_secs(30); +/// Deadline for a cancel-lane client to send its first frame, so a stalled +/// connection cannot hold a cancellation slot for the general one. +const LANE_FIRST_FRAME_TIMEOUT: Duration = Duration::from_secs(2); + +/// Upper bound on a cancel-lane request frame, large enough for every +/// `cancel_exec` the protocol admits: [`MAX_EXEC_ID_BYTES`] caps an +/// identifier's decoded length, not the escaped length it occupies on the wire. +const LANE_MAX_FRAME_BYTES: usize = 2 * MAX_EXEC_ID_BYTES * JSON_MAX_ESCAPE_BYTES + 256; + +/// Longest JSON escape a single byte of an identifier can produce, a control +/// character with no short form becoming `\u00XX`. +const JSON_MAX_ESCAPE_BYTES: usize = 6; + /// Bound on how long shutdown waits for in-flight handlers to finish before /// abandoning them, so a wedged handler cannot block daemon exit forever. const DRAIN_TIMEOUT: Duration = Duration::from_secs(30); @@ -112,6 +143,7 @@ pub async fn run( } = signals; let mut server = first_instance; let client_limiter = Arc::new(Semaphore::new(MAX_CONCURRENT_CLIENTS)); + let cancel_limiter = Arc::new(Semaphore::new(CANCEL_LANE_SLOTS)); let exec_limiter = Arc::new(Semaphore::new(MAX_CONCURRENT_EXECS)); let mut clients: JoinSet<()> = JoinSet::new(); @@ -165,6 +197,7 @@ pub async fn run( spawn_client_handler( &mut clients, &client_limiter, + &cancel_limiter, &exec_limiter, &session, &active_clients, @@ -204,18 +237,27 @@ pub async fn run( } /// Spawn a bounded task to service one accepted client connection. +/// +/// A connection that arrives once the general bound is reached is admitted into +/// the cancellation lane instead, where it is serviced only if it carries a +/// [`DaemonRequest::CancelExec`]. fn spawn_client_handler( clients: &mut JoinSet<()>, client_limiter: &Arc, + cancel_limiter: &Arc, exec_limiter: &Arc, session: &SessionHandle, active_clients: &Arc, connected: NamedPipeServer, ) { - let Ok(permit) = client_limiter.clone().try_acquire_owned() else { - // The connection is already accepted, so dropping it is the only - // bounded refusal path that cannot stall the accept loop. - return; + let (permit, cancel_only) = match client_limiter.clone().try_acquire_owned() { + Ok(permit) => (permit, false), + Err(_) => match cancel_limiter.clone().try_acquire_owned() { + Ok(permit) => (permit, true), + // The connection is already accepted, so dropping it is the only + // bounded refusal path that cannot stall the accept loop. + Err(_) => return, + }, }; let session = session.clone(); @@ -224,7 +266,7 @@ fn spawn_client_handler( active.fetch_add(1, Ordering::SeqCst); clients.spawn(async move { let _permit = permit; - if let Err(e) = handle_client(connected, session, exec_limiter).await { + if let Err(e) = handle_client(connected, session, exec_limiter, cancel_only).await { eprintln!("[wslc-daemon] client connection error: {e:#}"); } active.fetch_sub(1, Ordering::SeqCst); @@ -444,19 +486,41 @@ where } /// Service exactly one request on a freshly-connected pipe instance. +/// +/// `cancel_only` marks a connection admitted into the cancellation lane because +/// the general client bound was reached; anything other than a +/// [`DaemonRequest::CancelExec`] is refused there, and its request frame is held +/// to a lane-sized bound and deadline. async fn handle_client( mut pipe: S, session: SessionHandle, exec_limiter: Arc, + cancel_only: bool, ) -> Result<()> where S: AsyncRead + AsyncWrite + Unpin, { // Bound the wait for the request frame so a client that connects and then // stalls cannot pin this handler (and its concurrency slot) indefinitely. - let request: DaemonRequest = timeout(FIRST_FRAME_TIMEOUT, read_frame(&mut pipe)) + let (deadline, max_frame) = if cancel_only { + (LANE_FIRST_FRAME_TIMEOUT, LANE_MAX_FRAME_BYTES) + } else { + (FIRST_FRAME_TIMEOUT, MAX_FRAME_SIZE) + }; + let request: DaemonRequest = timeout(deadline, read_frame_capped(&mut pipe, max_frame)) .await .context("timed out waiting for the client's first frame")??; + if cancel_only && !matches!(request, DaemonRequest::CancelExec(_)) { + write_frame( + &mut pipe, + &DaemonResponse::Err { + kind: mxc_sdk::wslc_common::daemon_protocol::ErrKind::Busy, + message: "WSLc daemon client capacity is exhausted".to_string(), + }, + ) + .await?; + return Ok(()); + } match request { DaemonRequest::Ping => { write_frame(&mut pipe, &DaemonResponse::Pong).await?; @@ -513,11 +577,15 @@ where /// written, and — critically — admission is **atomic** with the claim the /// worker takes on the container (see [`SessionHandle::exec`]): the worker /// validates, claims the container and hands the run to a thread of its own -/// without yielding, and every later command naming that container parks behind -/// the claim, so no `Stop`/`Deprovision` can invalidate the checked state. An -/// unknown/not-started sandbox therefore comes back as a pre-admission typed +/// without yielding. A later `Stop`/`Deprovision` naming that container parks +/// behind the claim and a later `Exec` is refused with `Busy`, so neither can +/// invalidate the checked state. An unknown, not-started or already-busy +/// sandbox therefore comes back as a pre-admission typed /// [`DaemonResponse::Err`] rather than a post-admission stream `Error` frame. /// +/// The exec permit is shared with the run, so capacity frees only once the run +/// has reported back and this handler has finished writing. +/// /// Output streaming (process -> `Stdout`/`Stderr`) is live. Client `Stdin` /// frames are NOT forwarded: the WSLc SDK consumes all process IO handles once /// any `WslcSetProcessSettingsCallbacks` is registered (the callback path this @@ -529,25 +597,27 @@ async fn handle_exec( mut pipe: S, session: SessionHandle, config: mxc_sdk::wslc_common::daemon_protocol::ExecConfig, - _exec_permit: tokio::sync::OwnedSemaphorePermit, + exec_permit: tokio::sync::OwnedSemaphorePermit, ) -> Result<()> where S: AsyncWrite + Unpin, { let exec_id = config.exec_id.clone(); let run_token = config.run_token.clone(); + let slot: ExecSlotGuard = Arc::new(exec_permit); // Await the worker's admission decision before writing anything: a rejected // exec is a pre-admission typed error, never a post-admission stream frame. - let delivered = write_exec_result(&mut pipe, session.exec(config).await).await; + let admission = session.exec(config, Some(slot.clone())).await; + let delivered = write_exec_result(&mut pipe, admission).await; if delivered.is_err() { - // The run outlives this handler on a thread of its own, and the permit - // that bounds exec capacity is released as this returns. Without a kill - // a client could disconnect in a loop and leave runs going unbounded. + // The run outlives this handler on a thread of its own, so without a + // kill a client could disconnect in a loop and leave runs going. session.cancel_exec(&exec_id, &run_token); } + drop(slot); delivered } @@ -679,12 +749,22 @@ fn worker_err_response(e: WorkerError) -> DaemonResponse { } /// Read one length-prefixed frame and deserialise it. +#[cfg(test)] async fn read_frame(pipe: &mut S) -> Result { + read_frame_capped(pipe, MAX_FRAME_SIZE).await +} + +/// Read one length-prefixed frame, refusing a declared length above `max` +/// before allocating for it. +async fn read_frame_capped( + pipe: &mut S, + max: usize, +) -> Result { let mut len_buf = [0u8; 4]; pipe.read_exact(&mut len_buf).await?; let len = u32::from_le_bytes(len_buf) as usize; - if len > MAX_FRAME_SIZE { - bail!("incoming frame length {len} exceeds maximum {MAX_FRAME_SIZE}"); + if len > max { + bail!("incoming frame length {len} exceeds maximum {max}"); } let mut body = vec![0u8; len]; pipe.read_exact(&mut body).await?; @@ -702,7 +782,9 @@ async fn write_frame(pipe: &mut S, msg: &T) #[cfg(test)] mod tests { use super::*; - use crate::session_manager::{register_exec, spawn}; + use crate::session_manager::{ + register_exec, spawn, LIVE_OUTPUT_CHANNEL_CAPACITY, LIVE_OUTPUT_MAX_CHUNK_BYTES, + }; use mxc_sdk::wslc_common::daemon_protocol::{ CancelExecConfig, ErrKind, ExecConfig, ProvisionConfig, }; @@ -715,6 +797,100 @@ mod tests { Arc::new(register_exec(&active_execs, "test-exec", "test-run", &cancellation).unwrap()) } + /// A streaming exec's output memory is its live-output queue alone, so the + /// cap is what holds the daemon's worst case down. + #[test] + fn the_exec_cap_bounds_worst_case_live_output_memory() { + let per_exec = LIVE_OUTPUT_CHANNEL_CAPACITY * LIVE_OUTPUT_MAX_CHUNK_BYTES; + let worst_case = MAX_CONCURRENT_EXECS * per_exec; + + assert_eq!(per_exec, 16 * 1024 * 1024); + assert_eq!(worst_case, 128 * 1024 * 1024); + assert_eq!(MAX_CONCURRENT_CLIENTS, 16); + } + + /// Cancellation is the only way to end a run with no timeout, so it must + /// survive a client bound that parked lifecycle work has filled. + #[tokio::test] + async fn the_cancel_lane_outlives_exhausted_client_capacity() { + let session = spawn().unwrap(); + let client_limiter = Arc::new(Semaphore::new(MAX_CONCURRENT_CLIENTS)); + let cancel_limiter = Arc::new(Semaphore::new(CANCEL_LANE_SLOTS)); + + // Every general slot taken, as when each exec has a lifecycle command + // parked behind it. + let mut held = Vec::new(); + for _ in 0..MAX_CONCURRENT_CLIENTS { + held.push(client_limiter.clone().try_acquire_owned().unwrap()); + } + assert!(client_limiter.clone().try_acquire_owned().is_err()); + + // Every in-flight exec must stay cancellable, so the requirement is the + // exec cap. + let mut lane = Vec::new(); + for _ in 0..MAX_CONCURRENT_EXECS { + lane.push( + cancel_limiter + .clone() + .try_acquire_owned() + .expect("every in-flight exec must still be cancellable"), + ); + } + + for permit in lane { + let (mut client, server) = duplex(64 * 1024); + write_frame( + &mut client, + &DaemonRequest::CancelExec(CancelExecConfig { + exec_id: "stuck".to_string(), + run_token: "stuck-run".to_string(), + }), + ) + .await + .unwrap(); + + handle_client(server, session.clone(), Arc::new(Semaphore::new(0)), true) + .await + .unwrap(); + + let response: DaemonResponse = read_frame(&mut client).await.unwrap(); + assert_eq!(response, DaemonResponse::Ok); + drop(permit); + } + + session.shutdown().await.unwrap(); + } + + /// The cancel lane carries cancellations only; anything else is refused + /// with a typed error. + #[tokio::test] + async fn the_cancel_lane_refuses_non_cancel_requests() { + let session = spawn().unwrap(); + let (mut client, server) = duplex(64 * 1024); + write_frame( + &mut client, + &DaemonRequest::Stop(mxc_sdk::wslc_common::daemon_protocol::StopConfig { + sandbox_id: "wslc:test".to_string(), + }), + ) + .await + .unwrap(); + + handle_client(server, session.clone(), Arc::new(Semaphore::new(0)), true) + .await + .unwrap(); + + let response: DaemonResponse = read_frame(&mut client).await.unwrap(); + assert_eq!( + response, + DaemonResponse::Err { + kind: ErrKind::Busy, + message: "WSLc daemon client capacity is exhausted".to_string(), + } + ); + session.shutdown().await.unwrap(); + } + #[tokio::test] async fn exhausted_exec_capacity_returns_busy() { let session = spawn().unwrap(); @@ -737,7 +913,7 @@ mod tests { .await .unwrap(); - handle_client(server, session.clone(), exec_limiter) + handle_client(server, session.clone(), exec_limiter, false) .await .unwrap(); @@ -788,13 +964,339 @@ mod tests { assert!(delivered.is_err(), "the broken pipe must fail delivery"); assert!( cancellation.load(Ordering::Acquire), - "a run the client can no longer read must be cancelled, or it keeps \ - going after its capacity permit is released" + "a run the client can no longer read must be cancelled rather than \ + left going" ); drop(registration); session.shutdown().await.unwrap(); } + #[tokio::test] + async fn an_exec_slot_outlives_a_handler_that_is_still_writing() { + let session = spawn().unwrap(); + let limiter = Arc::new(Semaphore::new(1)); + let permit = limiter.clone().try_acquire_owned().unwrap(); + + // One byte of pipe, never drained, so the handler cannot finish writing. + let (server, mut client) = duplex(1); + let handler = tokio::spawn(handle_exec( + server, + session.clone(), + ExecConfig { + exec_id: "exec-writing".to_string(), + run_token: "run-writing".to_string(), + sandbox_id: "wslc:never-provisioned".to_string(), + script_code: "echo hi".to_string(), + working_directory: String::new(), + env: Vec::new(), + env_scope: mxc_sdk::wslc_common::process_env::EnvScope::Merge, + timeout_ms: 0, + }, + permit, + )); + + // One byte proves the handler is into its frame, and the rest of that + // frame cannot fit behind it. + let mut first = [0u8; 1]; + client.read_exact(&mut first).await.unwrap(); + assert_eq!( + limiter.available_permits(), + 0, + "a handler still writing to its client must not free its exec slot" + ); + + let mut rest = Vec::new(); + client.read_to_end(&mut rest).await.unwrap(); + handler.await.unwrap().unwrap(); + assert_eq!( + limiter.available_permits(), + 1, + "the slot must come back once the handler has finished writing" + ); + session.shutdown().await.unwrap(); + } + + /// A pipe name unique to this test run, so tests do not collide on one + /// instance. + fn test_pipe_name() -> String { + static NEXT: AtomicUsize = AtomicUsize::new(0); + format!( + r"\\.\pipe\mxc-wslc-lane-{}-{}", + std::process::id(), + NEXT.fetch_add(1, Ordering::SeqCst) + ) + } + + /// Route one accepted connection through the router with every general slot + /// taken, and hand back the client end of it. + async fn connection_past_general_capacity( + session: &SessionHandle, + clients: &mut JoinSet<()>, + cancel_limiter: &Arc, + ) -> tokio::net::windows::named_pipe::NamedPipeClient { + let name = test_pipe_name(); + let server = ServerOptions::new() + .first_pipe_instance(true) + .create(&name) + .unwrap(); + let client = tokio::net::windows::named_pipe::ClientOptions::new() + .open(&name) + .unwrap(); + server.connect().await.unwrap(); + + let client_limiter = Arc::new(Semaphore::new(1)); + let _held = client_limiter.clone().try_acquire_owned().unwrap(); + spawn_client_handler( + clients, + &client_limiter, + cancel_limiter, + &Arc::new(Semaphore::new(MAX_CONCURRENT_EXECS)), + session, + &Arc::new(AtomicUsize::new(0)), + server, + ); + client + } + + #[tokio::test] + async fn a_cancellation_survives_a_lane_saturated_by_stalled_connections() { + // Well past the lane's deadline and well inside the general one, so + // only a lane-sized deadline gets a cancellation through in time. + const CANCEL_BUDGET: Duration = Duration::from_secs(10); + + let session = spawn().unwrap(); + let cancel_limiter = Arc::new(Semaphore::new(CANCEL_LANE_SLOTS)); + let mut clients = JoinSet::new(); + + // Every lane slot taken by a connection that sends a length prefix and + // then nothing, so each one is mid-frame rather than idle. + let mut stalled = Vec::new(); + for _ in 0..CANCEL_LANE_SLOTS { + let mut client = + connection_past_general_capacity(&session, &mut clients, &cancel_limiter).await; + client.write_all(&300u32.to_le_bytes()).await.unwrap(); + stalled.push(client); + } + + let started = std::time::Instant::now(); + let mut cancelled = false; + while started.elapsed() < CANCEL_BUDGET { + let mut client = + connection_past_general_capacity(&session, &mut clients, &cancel_limiter).await; + if write_frame( + &mut client, + &DaemonRequest::CancelExec(CancelExecConfig { + exec_id: "stuck".to_string(), + run_token: "stuck-run".to_string(), + }), + ) + .await + .is_ok() + { + if let Ok(DaemonResponse::Ok) = read_frame::<_, DaemonResponse>(&mut client).await { + cancelled = true; + break; + } + } + } + + assert!( + cancelled, + "a cancellation behind stalled lane connections took longer than \ + {CANCEL_BUDGET:?}" + ); + clients.shutdown().await; + session.shutdown().await.unwrap(); + } + + #[tokio::test] + async fn an_escaped_cancellation_still_fits_the_lane() { + // Each byte escapes to its longest JSON form, at the longest identifier + // the protocol admits. + let escaped = "\u{1}".repeat(MAX_EXEC_ID_BYTES); + assert_eq!(escaped.len(), MAX_EXEC_ID_BYTES); + + let frame = encode_frame(&DaemonRequest::CancelExec(CancelExecConfig { + exec_id: escaped.clone(), + run_token: escaped.clone(), + })) + .unwrap(); + assert!( + frame.len() - 4 <= LANE_MAX_FRAME_BYTES, + "a {} byte cancellation does not fit the lane's {LANE_MAX_FRAME_BYTES} byte bound", + frame.len() - 4 + ); + + let session = spawn().unwrap(); + let cancel_limiter = Arc::new(Semaphore::new(1)); + let mut clients = JoinSet::new(); + let mut client = + connection_past_general_capacity(&session, &mut clients, &cancel_limiter).await; + + write_frame( + &mut client, + &DaemonRequest::CancelExec(CancelExecConfig { + exec_id: escaped.clone(), + run_token: escaped, + }), + ) + .await + .unwrap(); + + let response: DaemonResponse = read_frame(&mut client).await.unwrap(); + assert_eq!( + response, + DaemonResponse::Ok, + "the lane must carry every cancellation the protocol admits" + ); + clients.shutdown().await; + session.shutdown().await.unwrap(); + } + + #[tokio::test] + async fn a_connection_past_general_capacity_refuses_a_non_cancellation() { + let session = spawn().unwrap(); + let cancel_limiter = Arc::new(Semaphore::new(1)); + let mut clients = JoinSet::new(); + let mut client = + connection_past_general_capacity(&session, &mut clients, &cancel_limiter).await; + + write_frame( + &mut client, + &DaemonRequest::Exec(ExecConfig { + exec_id: "exec-lane".to_string(), + run_token: "run-lane".to_string(), + sandbox_id: "wslc:test".to_string(), + script_code: "echo hi".to_string(), + working_directory: String::new(), + env: Vec::new(), + env_scope: mxc_sdk::wslc_common::process_env::EnvScope::Merge, + timeout_ms: 0, + }), + ) + .await + .unwrap(); + + let response: DaemonResponse = read_frame(&mut client).await.unwrap(); + assert_eq!( + response, + DaemonResponse::Err { + kind: ErrKind::Busy, + message: "WSLc daemon client capacity is exhausted".to_string(), + } + ); + clients.shutdown().await; + session.shutdown().await.unwrap(); + } + + /// Cancellation is the only way to end a run with no timeout. + #[tokio::test] + async fn a_connection_past_general_capacity_still_serves_a_cancellation() { + let session = spawn().unwrap(); + let cancel_limiter = Arc::new(Semaphore::new(1)); + let mut clients = JoinSet::new(); + let mut client = + connection_past_general_capacity(&session, &mut clients, &cancel_limiter).await; + + write_frame( + &mut client, + &DaemonRequest::CancelExec(CancelExecConfig { + exec_id: "unknown-exec".to_string(), + run_token: "unknown-run".to_string(), + }), + ) + .await + .unwrap(); + + let response: DaemonResponse = read_frame(&mut client).await.unwrap(); + assert_eq!(response, DaemonResponse::Ok); + clients.shutdown().await; + session.shutdown().await.unwrap(); + } + + #[tokio::test] + async fn a_stalled_lane_connection_gives_up_its_slot_quickly() { + let session = spawn().unwrap(); + let (_client, server) = duplex(64 * 1024); + + let started = std::time::Instant::now(); + let outcome = + handle_client(server, session.clone(), Arc::new(Semaphore::new(0)), true).await; + let waited = started.elapsed(); + + assert!( + outcome.is_err(), + "a client that sends nothing must not be serviced" + ); + assert!( + waited < FIRST_FRAME_TIMEOUT, + "a stalled lane connection held its slot for {waited:?}, the general deadline" + ); + session.shutdown().await.unwrap(); + } + + #[tokio::test] + async fn an_oversized_lane_frame_is_refused_before_it_is_allocated() { + // Trivial for the general bound, so only a lane-sized bound refuses it. + const MODEST_FRAME_BYTES: usize = 64 * 1024; + + let session = spawn().unwrap(); + let (mut client, server) = duplex(64 * 1024); + + // A length prefix alone: a lane that accepted this would wait for a + // body that never comes. + client + .write_all(&(MODEST_FRAME_BYTES as u32).to_le_bytes()) + .await + .unwrap(); + + let outcome = + handle_client(server, session.clone(), Arc::new(Semaphore::new(0)), true).await; + + let message = outcome + .expect_err("an oversized lane frame must be refused") + .to_string(); + assert!( + message.contains("exceeds maximum"), + "a {MODEST_FRAME_BYTES}-byte frame must exceed the lane's bound of \ + {LANE_MAX_FRAME_BYTES}, got {message:?}" + ); + session.shutdown().await.unwrap(); + } + + #[tokio::test] + async fn a_general_connection_keeps_the_full_frame_bound() { + let session = spawn().unwrap(); + let (mut client, server) = duplex(64 * 1024); + + // Serviced concurrently so the body below cannot outgrow the pipe + // buffer with nothing draining it. + let handler = tokio::spawn(handle_client( + server, + session.clone(), + Arc::new(Semaphore::new(0)), + false, + )); + + let past_lane = LANE_MAX_FRAME_BYTES * 2; + client + .write_all(&(past_lane as u32).to_le_bytes()) + .await + .unwrap(); + client.write_all(&vec![b' '; past_lane]).await.unwrap(); + + let message = handler + .await + .unwrap() + .expect_err("whitespace is not a request") + .to_string(); + assert!( + !message.contains("exceeds maximum"), + "a general connection must not be held to the lane's bound, got {message:?}" + ); + session.shutdown().await.unwrap(); + } + /// A writer that always fails, standing in for a client that has gone. struct BrokenPipe; @@ -905,7 +1407,7 @@ mod tests { .unwrap(); drop(client); - let outcome = handle_client(server, session.clone(), exec_limiter).await; + let outcome = handle_client(server, session.clone(), exec_limiter, false).await; assert!( outcome.is_err(), "an undelivered reply must not be reported as a served request" @@ -932,7 +1434,7 @@ mod tests { .await .unwrap(); - handle_client(server, session.clone(), exec_limiter) + handle_client(server, session.clone(), exec_limiter, false) .await .unwrap(); diff --git a/src/mxc-sdk/src/bin/wslc_daemon/session_manager.rs b/src/mxc-sdk/src/bin/wslc_daemon/session_manager.rs index 3d330859a..6fc36d026 100644 --- a/src/mxc-sdk/src/bin/wslc_daemon/session_manager.rs +++ b/src/mxc-sdk/src/bin/wslc_daemon/session_manager.rs @@ -107,6 +107,9 @@ pub enum WorkerError { /// The sandbox exists but has not been started. NotStarted(String), + /// An exec already holds this container's single-flight slot. + Busy(String), + /// The host cannot run WSLc at all. Unavailable(anyhow::Error), @@ -123,6 +126,7 @@ impl WorkerError { match self { WorkerError::NotProvisioned(_) => ErrKind::NotProvisioned, WorkerError::NotStarted(_) => ErrKind::NotStarted, + WorkerError::Busy(_) => ErrKind::Busy, WorkerError::Unavailable(_) => ErrKind::Unavailable, WorkerError::Rejected(_) => ErrKind::Rejected, WorkerError::Backend(_) => ErrKind::Backend, @@ -135,6 +139,7 @@ impl std::fmt::Display for WorkerError { match self { WorkerError::NotProvisioned(id) => write!(f, "unknown sandbox {id}"), WorkerError::NotStarted(id) => write!(f, "sandbox {id} is not started"), + WorkerError::Busy(id) => write!(f, "sandbox {id} already has an exec in flight"), WorkerError::Unavailable(e) | WorkerError::Rejected(e) | WorkerError::Backend(e) => { write!(f, "{e:#}") } @@ -182,19 +187,29 @@ pub enum WorkerCommand { Shutdown { reply: oneshot::Sender<()> }, } +/// A resource an exec holds against the daemon's capacity, shared by the client +/// handler and the run so it is released only once both are done with it. +/// +/// Opaque so the worker never names the control server's permit type. +pub type ExecSlotGuard = Arc; + /// Everything an admitted exec needs, kept together so it travels as one value /// from the pipe handler to the run thread. pub struct ExecRequest { pub config: ExecConfig, - /// Live-output sink the worker hands to `exec_in_container`; the SDK's - /// stdout/stderr callbacks push chunks through it to the pipe handler as - /// bytes arrive, alongside the capped capture buffers. + /// Live-output sink the worker hands to `exec_in_container`, carrying the + /// SDK's stdout/stderr chunks to the pipe handler as bytes arrive. pub sink: OutputSink, pub cancellation: Arc, pub(crate) registration: Arc, pub admit: oneshot::Sender>, pub done: oneshot::Sender>, + + /// Counts this run against the daemon's exec capacity until it reports + /// back, so a disconnecting client does not free the slot while its run + /// thread and container process are still going. + pub slot: Option, } /// What an off-worker exec thread hands back, already classified. @@ -351,7 +366,7 @@ pub type OutputChunk = (OutStream, Vec); /// the SDK callback thread — which also delivers the process-exit callback — is /// never parked. Stalling that thread could otherwise block exit delivery and /// wedge teardown for every sandbox sharing the daemon. -const LIVE_OUTPUT_CHANNEL_CAPACITY: usize = 256; +pub(crate) const LIVE_OUTPUT_CHANNEL_CAPACITY: usize = 256; /// Max bytes per enqueued live-output chunk. A single SDK callback can deliver /// an arbitrarily large buffer; splitting it here bounds each queue entry's @@ -359,7 +374,7 @@ const LIVE_OUTPUT_CHANNEL_CAPACITY: usize = 256; /// protocol's `MAX_FRAME_SIZE` (a `Vec` serializes as a JSON number array, /// ~4x expansion), so a large callback can never overflow a frame and abort the /// stream before its terminal frame. -const LIVE_OUTPUT_MAX_CHUNK_BYTES: usize = 64 * 1024; +pub(crate) const LIVE_OUTPUT_MAX_CHUNK_BYTES: usize = 64 * 1024; /// Enqueue an SDK output callback, splitting it into `LIVE_OUTPUT_MAX_CHUNK_BYTES` /// pieces so each queue entry and its resulting frame stay bounded regardless of @@ -513,14 +528,19 @@ impl SessionHandle { } /// Admit and run a command in a started container. Awaits the worker's - /// **admission** decision first: on rejection (unknown/not-started sandbox) - /// this returns the typed error *before* the caller writes any admission to - /// the client. On admission it returns an [`ExecStream`] — the completion - /// receiver (the run's exit code) plus the live-output receiver, which the - /// caller drains into `Stdout`/`Stderr` frames as bytes arrive. Admission and + /// **admission** decision first: on rejection (unknown / not-started / + /// already-busy sandbox) this returns the typed error *before* the caller + /// writes any admission to the client. On admission it returns an + /// [`ExecStream`] — the completion receiver (the run's exit code) plus the + /// live-output receiver, which the caller drains into `Stdout`/`Stderr` + /// frames as bytes arrive. Admission and /// the claim on the container are atomic on the worker thread, so no /// lifecycle command can invalidate the checked state before the run starts. - pub async fn exec(&self, config: ExecConfig) -> Result { + pub async fn exec( + &self, + config: ExecConfig, + slot: Option, + ) -> Result { let (admit, admit_rx) = oneshot::channel(); let (done, done_rx) = oneshot::channel(); let (stream_tx, output) = mpsc::channel::(LIVE_OUTPUT_CHANNEL_CAPACITY); @@ -552,6 +572,7 @@ impl SessionHandle { registration: Arc::clone(®istration), admit, done, + slot, })))?; admit_rx.await.map_err(worker_gone)??; Ok(ExecStream { @@ -562,8 +583,8 @@ impl SessionHandle { }) } - /// Signal an admitted exec without waiting for the worker, which may have - /// parked the run behind an earlier exec on the same container. + /// Signal an admitted exec without waiting for the worker, which may not + /// have reached the run yet. pub fn cancel_exec(&self, exec_id: &str, run_token: &str) { if let Some(cancellation) = self .active_execs @@ -645,6 +666,10 @@ struct ContainerEntry { /// and deletes the container. retired: bool, container: WslcContainerGuard, + + /// Keeps a quarantined sandbox counted against exec capacity, because a run + /// whose termination was never confirmed may still hold a live process. + exec_slot: Option, } /// A provision waiting on its image to arrive from a registry. @@ -663,9 +688,10 @@ pub(crate) enum ContainerWork { /// a thread of its own. The two replies make admission **atomic** with the /// claim on the container: the worker validates, claims the container's /// in-flight slot, answers `admit` and starts the run thread without - /// yielding, and every later command naming that container parks behind the - /// claim, so none can interleave with the run. `admit` carries the pre-run - /// decision (so an unknown/not-started sandbox is a pre-admission typed + /// yielding. A later lifecycle command naming that container parks behind + /// the claim and a later `Exec` is refused with `Busy`, so none can + /// interleave with the run. `admit` carries the pre-run decision (so an + /// unknown, not-started or already-busy sandbox is a pre-admission typed /// error, never a post-admission stream `Error`); `done` carries the run's /// exit code once [`WorkerCommand::ExecFinished`] lands. Exec(ExecRequest), @@ -709,6 +735,50 @@ impl ContainerWork { } } +/// How [`Worker::begin_exec`] puts a claimed run on its own thread, injected so +/// a worker with no SDK loaded can still reach the claim. +type RunStarter = fn( + &Worker, + ExecConfig, + WslcContainer, + OutputSink, + Arc, + &mpsc::UnboundedSender, +) -> Result<(), WorkerError>; + +/// Hand a claimed run to its own MTA thread using the SDK this worker loaded. +fn start_run_on_sdk( + state: &Worker, + config: ExecConfig, + container: WslcContainer, + sink: OutputSink, + cancellation: Arc, + worker: &mpsc::UnboundedSender, +) -> Result<(), WorkerError> { + let Some(sdk) = state.sdk.as_ref() else { + return Err(WorkerError::Backend(anyhow::anyhow!( + "no active WSLc session" + ))); + }; + + let sandbox_id = config.sandbox_id.clone(); + let job = ExecJob { + // The box's contents, not the field: teardown may take the field while + // this run is still dereferencing the SDK. + sdk: &**sdk as *const WslcSdk, + container, + config, + sink, + cancellation, + }; + + start_exec(job, worker.clone(), &state.execs_in_flight).map_err(|e| { + WorkerError::Backend(anyhow::anyhow!( + "could not start a thread to run exec on sandbox {sandbox_id}: {e}" + )) + }) +} + /// An exec running on its own thread, and the work that parked behind it. /// /// `container` is the handle that run is using, so nothing may delete or @@ -719,6 +789,9 @@ struct InFlightExec { /// Held only so the exec id stays reserved for as long as the run lasts. _registration: Arc, + + /// Held only so the run counts against exec capacity until it reports back. + slot: Option, parked: Vec, } @@ -730,6 +803,8 @@ struct InFlightExec { /// so. struct Worker { logger: Logger, + start_run: RunStarter, + // Field order is load-bearing on implicit drop: `containers` and `session` // hold handles whose Drop calls into the SDK, so they must drop before `sdk` // unloads `wslcsdk.dll`. @@ -753,6 +828,7 @@ impl Worker { fn new() -> Self { Self { logger: Logger::new(Mode::Console), + start_run: start_run_on_sdk, sdk: None, session: None, containers: HashMap::new(), @@ -1003,6 +1079,7 @@ impl Worker { quarantined: false, retired: false, container, + exec_slot: None, }, ); Ok(sandbox_id) @@ -1042,7 +1119,8 @@ impl Worker { /// needed to run. Sole owner of the exists+started invariant: [`begin_exec`] /// trusts the handle it is given and never re-checks, because the worker /// validates and claims the container's in-flight slot without yielding, and - /// every later command naming that container parks behind the claim. + /// no later command naming that container can release the handle while the + /// claim stands. /// /// [`begin_exec`]: Worker::begin_exec fn validate_exec(&self, sandbox_id: &str) -> Result { @@ -1057,14 +1135,17 @@ impl Worker { } } - /// Run one container command, or park it behind the exec still using that - /// container's handle. - /// - /// A parked command waits for the run and then takes effect, so a caller - /// sees a delay rather than a refusal. + /// Run one container command, refusing a second exec on a container that + /// already has one and parking lifecycle work behind the run instead. fn dispatch(&mut self, work: ContainerWork, worker: &mpsc::UnboundedSender) { if let Some(in_flight) = self.exec_in_flight.get_mut(work.sandbox_id()) { - in_flight.parked.push(work); + match work { + ContainerWork::Exec(request) => { + let ExecRequest { config, admit, .. } = request; + let _ = admit.send(Err(WorkerError::Busy(config.sandbox_id))); + } + lifecycle => in_flight.parked.push(lifecycle), + } return; } @@ -1095,6 +1176,7 @@ impl Worker { registration, admit, done, + slot, } = request; let container = match self.validate_exec(&config.sandbox_id) { @@ -1118,28 +1200,10 @@ impl Worker { return; } - let Some(sdk) = self.sdk.as_ref() else { - let _ = done.send(Err(WorkerError::Backend(anyhow::anyhow!( - "no active WSLc session" - )))); - return; - }; - let sandbox_id = config.sandbox_id.clone(); - let job = ExecJob { - // The box's contents, not the field: teardown may take the field - // while this run is still dereferencing the SDK. - sdk: &**sdk as *const WslcSdk, - container, - config, - sink, - cancellation, - }; - - if let Err(e) = start_exec(job, worker.clone(), &self.execs_in_flight) { - let _ = done.send(Err(WorkerError::Backend(anyhow::anyhow!( - "could not start a thread to run exec on sandbox {sandbox_id}: {e}" - )))); + let start_run = self.start_run; + if let Err(e) = start_run(self, config, container, sink, cancellation, worker) { + let _ = done.send(Err(e)); return; } @@ -1149,6 +1213,7 @@ impl Worker { container, done, _registration: registration, + slot, parked: Vec::new(), }, ); @@ -1179,7 +1244,7 @@ impl Worker { let outcome = match report { ExecReport::Finished(outcome) => outcome, ExecReport::Unconfirmed(detail) => { - Err(self.quarantine(sandbox_id, in_flight.container, &detail)) + Err(self.quarantine(sandbox_id, in_flight.container, &detail, in_flight.slot)) } }; @@ -1201,6 +1266,7 @@ impl Worker { sandbox_id: &str, container: WslcContainer, detail: &str, + slot: Option, ) -> WorkerError { let delete_result = self.sdk.as_ref().map(|sdk| { // SAFETY: `sdk` is valid and `container` is the live handle stored @@ -1210,6 +1276,9 @@ impl Worker { match delete_result { Some(Ok(())) => { + // Deleting the container ends anything still running inside it, + // so this exec stops counting against capacity. + drop(slot); self.containers.remove(sandbox_id); WorkerError::Backend(anyhow::anyhow!( "exec on sandbox {sandbox_id} could not be confirmed terminated ({detail}); \ @@ -1219,6 +1288,7 @@ impl Worker { Some(Err(delete_error)) => { if let Some(entry) = self.containers.get_mut(sandbox_id) { entry.quarantined = true; + entry.exec_slot = slot; } WorkerError::Backend(anyhow::anyhow!( "exec on sandbox {sandbox_id} could not be confirmed terminated ({detail}); \ @@ -1230,6 +1300,7 @@ impl Worker { None => { if let Some(entry) = self.containers.get_mut(sandbox_id) { entry.quarantined = true; + entry.exec_slot = slot; } WorkerError::Backend(anyhow::anyhow!( "exec on sandbox {sandbox_id} could not be confirmed terminated ({detail}); \ @@ -1301,12 +1372,12 @@ impl Worker { Ok(()) } - /// Give up the SDK handles rather than release them, when a run thread may - /// still be using one. + /// Give up the SDK handles rather than release them, when an off-worker + /// thread may still be using one. /// /// Reached only on an unwind, where [`Worker::shutdown`]'s drain never ran. - fn abandon_if_execs_running(mut self, in_flight: usize) { - if in_flight > 0 { + fn abandon_if_borrowed(mut self, execs_in_flight: usize, pull_outstanding: bool) { + if execs_in_flight > 0 || pull_outstanding { std::mem::forget(std::mem::take(&mut self.containers)); std::mem::forget(self.session.take()); std::mem::forget(self.sdk.take()); @@ -1519,7 +1590,10 @@ pub fn spawn() -> Result { if served.is_err() { let in_flight = worker.execs_in_flight.load(Ordering::SeqCst); - worker.abandon_if_execs_running(in_flight); + + // A pull borrows the SDK and session the same way a run does. + let pull_outstanding = !image::wait_for_pulls_in_flight(Duration::ZERO); + worker.abandon_if_borrowed(in_flight, pull_outstanding); } }) .map_err(|e| anyhow::anyhow!("spawn WSLc worker thread: {e}"))?; @@ -1634,6 +1708,7 @@ mod tests { // SAFETY: `release_noop` never dereferences the handle, so the guard // owns a value it can release without touching memory. container: unsafe { WslcContainerGuard::from_raw(sentinel, release_noop) }, + exec_slot: None, } } @@ -1811,16 +1886,19 @@ mod tests { async fn exec_unknown_sandbox_errors() { let handle = spawn().unwrap(); let err = handle - .exec(ExecConfig { - exec_id: "unknown-1".to_string(), - run_token: "run-unknown-1".to_string(), - sandbox_id: "wslc:does-not-exist".to_string(), - script_code: "echo hi".to_string(), - working_directory: String::new(), - env: Vec::new(), - env_scope: EnvScope::Merge, - timeout_ms: 0, - }) + .exec( + ExecConfig { + exec_id: "unknown-1".to_string(), + run_token: "run-unknown-1".to_string(), + sandbox_id: "wslc:does-not-exist".to_string(), + script_code: "echo hi".to_string(), + working_directory: String::new(), + env: Vec::new(), + env_scope: EnvScope::Merge, + timeout_ms: 0, + }, + None, + ) .await .unwrap_err(); assert!(err.to_string().contains("unknown sandbox")); @@ -1871,16 +1949,19 @@ mod tests { async fn exec_unknown_sandbox_admission_is_not_provisioned() { let handle = spawn().unwrap(); let err = handle - .exec(ExecConfig { - exec_id: "unknown-2".to_string(), - run_token: "run-unknown-2".to_string(), - sandbox_id: "wslc:does-not-exist".to_string(), - script_code: "echo hi".to_string(), - working_directory: String::new(), - env: Vec::new(), - env_scope: EnvScope::Merge, - timeout_ms: 0, - }) + .exec( + ExecConfig { + exec_id: "unknown-2".to_string(), + run_token: "run-unknown-2".to_string(), + sandbox_id: "wslc:does-not-exist".to_string(), + script_code: "echo hi".to_string(), + working_directory: String::new(), + env: Vec::new(), + env_scope: EnvScope::Merge, + timeout_ms: 0, + }, + None, + ) .await .unwrap_err(); assert_eq!(err.kind(), ErrKind::NotProvisioned); @@ -2068,6 +2149,7 @@ mod tests { // SAFETY: `release` never dereferences the handle, so the guard owns // a value it can release without touching memory. container: unsafe { WslcContainerGuard::from_raw(sentinel, release) }, + exec_slot: None, } } @@ -2089,6 +2171,7 @@ mod tests { container: std::ptr::dangling_mut(), done, _registration: registration, + slot: None, parked: Vec::new(), }, done_rx, @@ -2104,6 +2187,14 @@ mod tests { } fn exec_work(sandbox_id: &str, exec_id: &str) -> TestExec { + exec_work_with_slot(sandbox_id, exec_id, None) + } + + fn exec_work_with_slot( + sandbox_id: &str, + exec_id: &str, + slot: Option, + ) -> TestExec { let (admit, admit_rx) = oneshot::channel(); let (done, done_rx) = oneshot::channel(); let active_execs: ActiveExecs = Arc::new(Mutex::new(HashMap::new())); @@ -2128,6 +2219,7 @@ mod tests { registration, admit, done, + slot, }), admit: admit_rx, done: done_rx, @@ -2149,6 +2241,41 @@ mod tests { (worker, done) } + /// A worker holding one started container it can start runs on without an + /// SDK loaded. + fn worker_that_can_run() -> Worker { + fn start_nothing( + _: &Worker, + _: ExecConfig, + _: WslcContainer, + _: OutputSink, + _: Arc, + _: &mpsc::UnboundedSender, + ) -> Result<(), WorkerError> { + Ok(()) + } + + let mut worker = Worker::new(); + worker.start_run = start_nothing; + worker + .containers + .insert("wslc:ready".to_string(), started_entry()); + worker + } + + fn start_work(sandbox_id: &str) -> (ContainerWork, oneshot::Receiver>) { + let (reply, reply_rx) = oneshot::channel(); + ( + ContainerWork::Start { + config: StartConfig { + sandbox_id: sandbox_id.to_string(), + }, + reply, + }, + reply_rx, + ) + } + fn stop_work(sandbox_id: &str) -> (ContainerWork, oneshot::Receiver>) { let (reply, reply_rx) = oneshot::channel(); ( @@ -2177,7 +2304,26 @@ mod tests { ) } - /// Deleting the container would free the handle the run thread is holding. + /// Deleting the container would free the handle the run thread is using, so + /// lifecycle work waits instead of being refused. + #[test] + fn lifecycle_work_parks_behind_an_exec_rather_than_being_refused() { + let (tx, _rx) = mpsc::unbounded_channel(); + let (mut worker, _done) = worker_with_exec_in_flight(); + let (start, mut start_reply) = start_work("wslc:busy"); + let (stop, mut stop_reply) = stop_work("wslc:busy"); + let (deprovision, mut deprovision_reply) = deprovision_work("wslc:busy"); + + worker.dispatch(start, &tx); + worker.dispatch(stop, &tx); + worker.dispatch(deprovision, &tx); + + assert!(start_reply.try_recv().is_err()); + assert!(stop_reply.try_recv().is_err()); + assert!(deprovision_reply.try_recv().is_err()); + assert_eq!(worker.exec_in_flight["wslc:busy"].parked.len(), 3); + } + #[test] fn deprovision_parks_behind_an_exec_using_the_same_container() { let (tx, _rx) = mpsc::unbounded_channel(); @@ -2236,40 +2382,275 @@ mod tests { assert!(!worker.containers.contains_key("wslc:busy")); } + /// A disconnected client leaves its run going, so the slot it claimed must + /// travel to the worker and stay claimed until that run reports back. + #[tokio::test] + async fn an_exec_slot_outlives_the_client_that_admitted_it() { + let (tx, mut rx) = mpsc::unbounded_channel(); + let handle = SessionHandle { + tx, + active_execs: Arc::new(Mutex::new(HashMap::new())), + }; + let limiter = Arc::new(tokio::sync::Semaphore::new(1)); + let permit = Arc::clone(&limiter).try_acquire_owned().unwrap(); + + let exec = tokio::spawn(async move { + handle + .exec( + ExecConfig { + exec_id: "slotted".to_string(), + run_token: "slotted-run".to_string(), + sandbox_id: "wslc:busy".to_string(), + script_code: "echo hi".to_string(), + working_directory: String::new(), + env: Vec::new(), + env_scope: EnvScope::Merge, + timeout_ms: 0, + }, + Some(Arc::new(permit)), + ) + .await + }); + + let Some(WorkerCommand::Container(ContainerWork::Exec(request))) = rx.recv().await else { + panic!("the exec never reached the worker"); + }; + assert!( + request.slot.is_some(), + "the admitting client's exec slot must reach the worker" + ); + + // Stand in for the worker: claim the slot for the run, then drop the + // client's end as a disconnect would. + let (mut worker, _done) = worker_with_exec_in_flight(); + worker.exec_in_flight.get_mut("wslc:busy").unwrap().slot = request.slot; + let _ = request.admit.send(Ok(())); + exec.await.unwrap().unwrap(); + assert_eq!( + limiter.available_permits(), + 0, + "a client that disconnected mid-run must not free its exec slot" + ); + + let (worker_tx, _worker_rx) = mpsc::unbounded_channel(); + worker.finish_exec( + "wslc:busy", + ExecReport::Finished(Ok(ExecTerminal::Exited(0))), + &worker_tx, + ); + + assert_eq!( + limiter.available_permits(), + 1, + "the slot must be released once the run reports back" + ); + } + + #[test] + fn the_worker_claims_the_slot_when_it_starts_the_run() { + let (tx, _rx) = mpsc::unbounded_channel(); + let mut worker = worker_that_can_run(); + let limiter = Arc::new(tokio::sync::Semaphore::new(1)); + let permit = Arc::clone(&limiter).try_acquire_owned().unwrap(); + let exec = exec_work_with_slot("wslc:ready", "claimed", Some(Arc::new(permit))); + + worker.dispatch(exec.work, &tx); + + assert_eq!( + limiter.available_permits(), + 0, + "the started run must hold the slot its request carried" + ); + + worker.finish_exec( + "wslc:ready", + ExecReport::Finished(Ok(ExecTerminal::Exited(0))), + &tx, + ); + assert_eq!( + limiter.available_permits(), + 1, + "retiring the run must release the slot" + ); + } + #[test] - fn a_second_exec_on_the_same_container_waits_for_the_first() { + fn parked_lifecycle_work_replays_in_the_order_it_arrived() { + let (tx, _rx) = mpsc::unbounded_channel(); + let (mut worker, _done) = worker_with_exec_in_flight(); + let (stop, mut stop_reply) = stop_work("wslc:busy"); + let (deprovision, mut deprovision_reply) = deprovision_work("wslc:busy"); + worker.dispatch(stop, &tx); + worker.dispatch(deprovision, &tx); + + worker.finish_exec( + "wslc:busy", + ExecReport::Finished(Ok(ExecTerminal::Exited(0))), + &tx, + ); + + let stopped = stop_reply + .try_recv() + .expect("the parked stop must be answered"); + assert!( + !matches!(stopped, Err(WorkerError::NotProvisioned(_))), + "the stop must run while its container still exists, got {stopped:?}" + ); + assert!(deprovision_reply + .try_recv() + .expect("the parked deprovision must be answered") + .is_ok()); + } + + #[test] + fn a_second_exec_on_the_same_container_is_refused_as_busy() { let (tx, _rx) = mpsc::unbounded_channel(); let (mut worker, _done) = worker_with_exec_in_flight(); let mut second = exec_work("wslc:busy", "second"); worker.dispatch(second.work, &tx); + let err = second + .admit + .try_recv() + .expect("a refused exec must be answered through admit") + .unwrap_err(); + assert_eq!(err.kind(), ErrKind::Busy); assert!( - second.admit.try_recv().is_err(), - "a second exec must not be admitted while the first holds the container" + worker.exec_in_flight["wslc:busy"].parked.is_empty(), + "a refused exec must not also queue behind the first" + ); + } + + #[test] + fn execs_on_different_containers_are_both_admitted() { + let (tx, _rx) = mpsc::unbounded_channel(); + let (mut worker, _done) = worker_with_exec_in_flight(); + worker + .containers + .insert("wslc:idle".to_string(), started_entry()); + let mut other = exec_work("wslc:idle", "other"); + + worker.dispatch(other.work, &tx); + + assert!( + other + .admit + .try_recv() + .expect("idle container answered") + .is_ok(), + "a container of its own must not inherit another container's slot" ); } /// The guarantee a cancellation must keep: observed before the run starts, /// it reports cancelled rather than creating the process. #[test] - fn a_parked_exec_cancelled_before_it_starts_never_runs() { + fn an_exec_cancelled_before_it_starts_never_runs() { let (tx, _rx) = mpsc::unbounded_channel(); - let (mut worker, _done) = worker_with_exec_in_flight(); - let mut queued = exec_work("wslc:busy", "queued"); + let mut worker = Worker::new(); + worker + .containers + .insert("wslc:idle".to_string(), started_entry()); + let mut queued = exec_work("wslc:idle", "queued"); + queued.cancellation.store(true, Ordering::Release); + worker.dispatch(queued.work, &tx); + assert!(queued.admit.try_recv().unwrap().is_ok()); + assert_eq!( + queued.done.try_recv().unwrap().unwrap(), + ExecTerminal::Cancelled + ); + assert!( + !worker.exec_in_flight.contains_key("wslc:idle"), + "a cancelled exec must not claim the container's slot" + ); + } + + #[test] + fn a_cancelled_exec_never_reaches_the_runner() { + static STARTED: AtomicBool = AtomicBool::new(false); + + fn record_start( + _: &Worker, + _: ExecConfig, + _: WslcContainer, + _: OutputSink, + _: Arc, + _: &mpsc::UnboundedSender, + ) -> Result<(), WorkerError> { + STARTED.store(true, Ordering::SeqCst); + Ok(()) + } + + let (tx, _rx) = mpsc::unbounded_channel(); + let mut worker = worker_that_can_run(); + worker.start_run = record_start; + let mut queued = exec_work("wslc:ready", "cancelled-before-start"); queued.cancellation.store(true, Ordering::Release); + + worker.dispatch(queued.work, &tx); + + assert!( + !STARTED.load(Ordering::SeqCst), + "a cancelled exec must not reach the runner that creates the process" + ); + assert_eq!( + queued.done.try_recv().unwrap().unwrap(), + ExecTerminal::Cancelled + ); + } + + #[test] + fn an_unconfirmed_exec_keeps_its_slot_claimed() { + let (tx, _rx) = mpsc::unbounded_channel(); + let (mut worker, _done) = worker_with_exec_in_flight(); + let limiter = Arc::new(tokio::sync::Semaphore::new(1)); + let permit = Arc::clone(&limiter).try_acquire_owned().unwrap(); + worker.exec_in_flight.get_mut("wslc:busy").unwrap().slot = Some(Arc::new(permit)); + worker.finish_exec( "wslc:busy", - ExecReport::Finished(Ok(ExecTerminal::Exited(0))), + ExecReport::Unconfirmed("the exit callback never fired".to_string()), &tx, ); - assert!(queued.admit.try_recv().unwrap().is_ok()); assert_eq!( - queued.done.try_recv().unwrap().unwrap(), - ExecTerminal::Cancelled + limiter.available_permits(), + 0, + "a quarantined sandbox whose process may still be running must keep its slot" + ); + } + + #[test] + fn deprovisioning_a_quarantined_sandbox_frees_its_slot() { + let (tx, _rx) = mpsc::unbounded_channel(); + let (mut worker, _done) = worker_with_exec_in_flight(); + let limiter = Arc::new(tokio::sync::Semaphore::new(1)); + let permit = Arc::clone(&limiter).try_acquire_owned().unwrap(); + worker.exec_in_flight.get_mut("wslc:busy").unwrap().slot = Some(Arc::new(permit)); + worker.finish_exec( + "wslc:busy", + ExecReport::Unconfirmed("the exit callback never fired".to_string()), + &tx, + ); + assert_eq!( + limiter.available_permits(), + 0, + "the quarantine holds a slot" + ); + + worker + .deprovision(DeprovisionConfig { + sandbox_id: "wslc:busy".to_string(), + }) + .unwrap(); + + assert_eq!( + limiter.available_permits(), + 1, + "deprovisioning the quarantined sandbox must return its slot" ); } @@ -2486,7 +2867,7 @@ mod tests { }, ); - worker.abandon_if_execs_running(1); + worker.abandon_if_borrowed(1, false); assert!( matches!(done.try_recv(), Err(oneshot::error::TryRecvError::Closed)), @@ -2524,7 +2905,7 @@ mod tests { .containers .insert("wslc:busy".to_string(), flagged_entry(release)); - worker.abandon_if_execs_running(1); + worker.abandon_if_borrowed(1, false); assert!( !RELEASED.load(Ordering::SeqCst), @@ -2532,6 +2913,28 @@ mod tests { ); } + #[test] + fn an_unwind_with_a_pull_outstanding_keeps_container_handles() { + static RELEASED: AtomicBool = AtomicBool::new(false); + + unsafe extern "C" fn release(_: mxc_sdk::wslc_common::wslc_bindings::WslcContainer) -> i32 { + RELEASED.store(true, Ordering::SeqCst); + 0 + } + + let mut worker = Worker::new(); + worker + .containers + .insert("wslc:pulling".to_string(), flagged_entry(release)); + + worker.abandon_if_borrowed(0, true); + + assert!( + !RELEASED.load(Ordering::SeqCst), + "a pull still borrows the SDK and session, so the handles must be kept" + ); + } + #[test] fn an_unwind_with_no_run_counted_releases_container_handles() { static RELEASED: AtomicBool = AtomicBool::new(false); @@ -2546,7 +2949,7 @@ mod tests { .containers .insert("wslc:idle".to_string(), flagged_entry(release)); - worker.abandon_if_execs_running(0); + worker.abandon_if_borrowed(0, false); assert!( RELEASED.load(Ordering::SeqCst), @@ -2588,6 +2991,82 @@ mod tests { )); } + /// Provision and start a sandbox on the live host, returning its id. + async fn provisioned_and_started(handle: &SessionHandle) -> String { + let id = handle + .provision(ProvisionConfig { + image: "alpine:latest".to_string(), + image_tar_path: None, + volumes: Vec::new(), + network: Default::default(), + port_mappings: Vec::new(), + }) + .await + .unwrap(); + handle + .start(StartConfig { + sandbox_id: id.clone(), + }) + .await + .unwrap(); + id + } + + async fn stop_and_deprovision(handle: &SessionHandle, sandbox_id: String) { + handle + .stop(StopConfig { + sandbox_id: sandbox_id.clone(), + }) + .await + .unwrap(); + handle + .deprovision(DeprovisionConfig { sandbox_id }) + .await + .unwrap(); + } + + fn sleep_exec(exec_id: &str, sandbox_id: &str, seconds: u32) -> ExecConfig { + ExecConfig { + exec_id: exec_id.to_string(), + run_token: format!("{exec_id}-run"), + sandbox_id: sandbox_id.to_string(), + script_code: format!("sleep {seconds}"), + working_directory: String::new(), + env: Vec::new(), + env_scope: EnvScope::Merge, + timeout_ms: 30_000, + } + } + + /// A run that prints an epoch second either side of its sleep. + fn stamped_sleep_exec(exec_id: &str, sandbox_id: &str, seconds: u32) -> ExecConfig { + ExecConfig { + script_code: format!("date +%s; sleep {seconds}; date +%s"), + ..sleep_exec(exec_id, sandbox_id, seconds) + } + } + + /// Drain a stamped run's output and return the epoch seconds it reported + /// either side of its sleep. + async fn stamped_interval(exec: &mut ExecStream) -> (i64, i64) { + let mut stdout = Vec::new(); + while let Some((stream, data)) = exec.output.recv().await { + if stream == OutStream::Stdout { + stdout.extend_from_slice(&data); + } + } + let text = String::from_utf8_lossy(&stdout); + let stamps: Vec = text + .lines() + .filter_map(|line| line.trim().parse::().ok()) + .collect(); + assert!( + stamps.len() >= 2, + "a stamped run must report a start and an end, got {text:?}" + ); + (stamps[0], stamps[stamps.len() - 1]) + } + // Exercises the real SDK path end to end: provision (boot VM + create // container) → start → exec → stop → deprovision → refcount back to 0. It // provisions with the default isolated posture, which refuses a registry @@ -2620,16 +3099,19 @@ mod tests { .unwrap(); let mut exec = handle - .exec(ExecConfig { - exec_id: "full-lifecycle".to_string(), - run_token: "full-lifecycle-run".to_string(), - sandbox_id: id.clone(), - script_code: "echo hi".to_string(), - working_directory: String::new(), - env: Vec::new(), - env_scope: EnvScope::Merge, - timeout_ms: 30_000, - }) + .exec( + ExecConfig { + exec_id: "full-lifecycle".to_string(), + run_token: "full-lifecycle-run".to_string(), + sandbox_id: id.clone(), + script_code: "echo hi".to_string(), + working_directory: String::new(), + env: Vec::new(), + env_scope: EnvScope::Merge, + timeout_ms: 30_000, + }, + None, + ) .await .unwrap(); // Drain the live output stream, then await the exit code. @@ -2659,55 +3141,47 @@ mod tests { handle.shutdown().await.unwrap(); } + /// The guarantee a cancellation must keep against the live SDK: observed + /// before the worker reaches the run, it creates no process. #[tokio::test] #[ignore = "requires a WSL2 host with alpine:latest already in the daemon session cache"] async fn cancelled_queued_exec_never_starts_process() { let handle = spawn().unwrap(); - let id = handle - .provision(ProvisionConfig { - image: "alpine:latest".to_string(), - image_tar_path: None, - volumes: Vec::new(), - network: Default::default(), - port_mappings: Vec::new(), - }) - .await - .unwrap(); - handle - .start(StartConfig { - sandbox_id: id.clone(), - }) - .await - .unwrap(); - - let blocker = handle - .exec(ExecConfig { - exec_id: "queue-blocker".to_string(), - run_token: "queue-blocker-run".to_string(), - sandbox_id: id.clone(), - script_code: "sleep 2".to_string(), - working_directory: String::new(), - env: Vec::new(), - env_scope: EnvScope::Merge, - timeout_ms: 30_000, - }) - .await - .unwrap(); + let id = provisioned_and_started(&handle).await; + + // Occupy the worker with a container creation, so the exec below is + // still queued when the cancellation lands. Spawned tasks are polled in + // order, so this provision reaches the worker first. + let blocker_handle = handle.clone(); + let blocker = tokio::spawn(async move { + blocker_handle + .provision(ProvisionConfig { + image: "alpine:latest".to_string(), + image_tar_path: None, + volumes: Vec::new(), + network: Default::default(), + port_mappings: Vec::new(), + }) + .await + }); let queued_handle = handle.clone(); let queued_id = id.clone(); let queued = tokio::spawn(async move { queued_handle - .exec(ExecConfig { - exec_id: "cancelled-queued".to_string(), - run_token: "cancelled-queued-run".to_string(), - sandbox_id: queued_id, - script_code: "touch /tmp/mxc-cancelled-queued-marker".to_string(), - working_directory: String::new(), - env: Vec::new(), - env_scope: EnvScope::Merge, - timeout_ms: 30_000, - }) + .exec( + ExecConfig { + exec_id: "cancelled-queued".to_string(), + run_token: "cancelled-queued-run".to_string(), + sandbox_id: queued_id, + script_code: "touch /tmp/mxc-cancelled-queued-marker".to_string(), + working_directory: String::new(), + env: Vec::new(), + env_scope: EnvScope::Merge, + timeout_ms: 30_000, + }, + None, + ) .await }); @@ -2731,25 +3205,25 @@ mod tests { "queued exec was not registered" ); handle.cancel_exec("cancelled-queued", "cancelled-queued-run"); - assert_eq!( - blocker.done.await.unwrap().unwrap(), - ExecTerminal::Exited(0) - ); + let blocker_id = blocker.await.unwrap().unwrap(); let queued = queued.await.unwrap().unwrap(); assert_eq!(queued.done.await.unwrap().unwrap(), ExecTerminal::Cancelled); let marker_check = handle - .exec(ExecConfig { - exec_id: "marker-check".to_string(), - run_token: "marker-check-run".to_string(), - sandbox_id: id.clone(), - script_code: "test ! -e /tmp/mxc-cancelled-queued-marker".to_string(), - working_directory: String::new(), - env: Vec::new(), - env_scope: EnvScope::Merge, - timeout_ms: 30_000, - }) + .exec( + ExecConfig { + exec_id: "marker-check".to_string(), + run_token: "marker-check-run".to_string(), + sandbox_id: id.clone(), + script_code: "test ! -e /tmp/mxc-cancelled-queued-marker".to_string(), + working_directory: String::new(), + env: Vec::new(), + env_scope: EnvScope::Merge, + timeout_ms: 30_000, + }, + None, + ) .await .unwrap(); assert_eq!( @@ -2757,16 +3231,85 @@ mod tests { ExecTerminal::Exited(0) ); + stop_and_deprovision(&handle, id).await; handle - .stop(StopConfig { - sandbox_id: id.clone(), + .deprovision(DeprovisionConfig { + sandbox_id: blocker_id, }) .await .unwrap(); - handle - .deprovision(DeprovisionConfig { sandbox_id: id }) + handle.shutdown().await.unwrap(); + } + + /// A second exec on a container that already has one is refused before any + /// admission reaches the client. + #[tokio::test] + #[ignore = "requires a WSL2 host with alpine:latest already in the daemon session cache"] + async fn a_second_exec_on_a_busy_container_is_refused() { + let handle = spawn().unwrap(); + let id = provisioned_and_started(&handle).await; + + let first = handle + .exec(sleep_exec("busy-first", &id, 3), None) .await .unwrap(); + + let refused = handle + .exec( + ExecConfig { + exec_id: "busy-second".to_string(), + run_token: "busy-second-run".to_string(), + sandbox_id: id.clone(), + script_code: "echo hi".to_string(), + working_directory: String::new(), + env: Vec::new(), + env_scope: EnvScope::Merge, + timeout_ms: 30_000, + }, + None, + ) + .await + .unwrap_err(); + assert_eq!(refused.kind(), ErrKind::Busy); + + assert_eq!(first.done.await.unwrap().unwrap(), ExecTerminal::Exited(0)); + + stop_and_deprovision(&handle, id).await; + handle.shutdown().await.unwrap(); + } + + /// Two sandboxes must run at the same time rather than one after the other. + #[tokio::test] + #[ignore = "requires a WSL2 host with alpine:latest already in the daemon session cache"] + async fn execs_on_two_sandboxes_overlap() { + let handle = spawn().unwrap(); + let first_id = provisioned_and_started(&handle).await; + let second_id = provisioned_and_started(&handle).await; + + let mut first = handle + .exec(stamped_sleep_exec("overlap-first", &first_id, 6), None) + .await + .unwrap(); + let mut second = handle + .exec(stamped_sleep_exec("overlap-second", &second_id, 6), None) + .await + .unwrap(); + + // Both sandboxes share one utility VM, so their clocks agree. + let (first_start, first_end) = stamped_interval(&mut first).await; + let (second_start, second_end) = stamped_interval(&mut second).await; + assert_eq!(first.done.await.unwrap().unwrap(), ExecTerminal::Exited(0)); + assert_eq!(second.done.await.unwrap().unwrap(), ExecTerminal::Exited(0)); + + let overlap = first_end.min(second_end) - first_start.max(second_start); + assert!( + overlap >= 2, + "the two 6s runs overlapped by {overlap}s, so they were serialized: \ + first {first_start}..{first_end}, second {second_start}..{second_end}" + ); + + stop_and_deprovision(&handle, first_id).await; + stop_and_deprovision(&handle, second_id).await; handle.shutdown().await.unwrap(); } @@ -2794,16 +3337,19 @@ mod tests { .unwrap(); let mut exec = handle - .exec(ExecConfig { - exec_id: "deprovision-overlap".to_string(), - run_token: "deprovision-overlap-run".to_string(), - sandbox_id: id.clone(), - script_code: "sleep 5; echo survived".to_string(), - working_directory: String::new(), - env: Vec::new(), - env_scope: EnvScope::Merge, - timeout_ms: 30_000, - }) + .exec( + ExecConfig { + exec_id: "deprovision-overlap".to_string(), + run_token: "deprovision-overlap-run".to_string(), + sandbox_id: id.clone(), + script_code: "sleep 5; echo survived".to_string(), + working_directory: String::new(), + env: Vec::new(), + env_scope: EnvScope::Merge, + timeout_ms: 30_000, + }, + None, + ) .await .unwrap(); diff --git a/src/mxc-sdk/tests/wslc_daemon_ipc.rs b/src/mxc-sdk/tests/wslc_daemon_ipc.rs index 220146e39..e85c2bdad 100644 --- a/src/mxc-sdk/tests/wslc_daemon_ipc.rs +++ b/src/mxc-sdk/tests/wslc_daemon_ipc.rs @@ -21,7 +21,9 @@ use std::process::{Child, Command, Stdio}; use std::time::{Duration, Instant}; -use mxc_sdk::wslc_common::daemon_client::DaemonClient; +use mxc_sdk::wslc_common::container_steps::OutStream; +use mxc_sdk::wslc_common::daemon_client::{DaemonClient, DaemonError}; +use mxc_sdk::wslc_common::daemon_protocol::{ErrKind, ExecConfig, ExecTerminal}; use mxc_sdk::wslc_common::daemon_record::{read_daemon_record, STATE_ROOT_ENV_VAR}; use mxc_sdk::wslc_common::process_env::EnvScope; @@ -81,6 +83,186 @@ impl Drop for DaemonProcess { } } +/// Provision and start a sandbox over the pipe, returning its id. +fn provisioned_and_started(client: &DaemonClient) -> String { + use mxc_sdk::wslc_common::daemon_protocol::{ProvisionConfig, StartConfig}; + + let sandbox_id = client + .provision(ProvisionConfig { + image: "alpine:latest".to_string(), + image_tar_path: None, + volumes: Vec::new(), + network: Default::default(), + port_mappings: Vec::new(), + }) + .expect("provision"); + client + .start(StartConfig { + sandbox_id: sandbox_id.clone(), + }) + .expect("start"); + sandbox_id +} + +/// A run that prints a marker, then an epoch second either side of a sleep. +fn stamped_exec(exec_id: &str, sandbox_id: &str, seconds: u32) -> ExecConfig { + ExecConfig { + exec_id: exec_id.to_string(), + run_token: format!("{exec_id}-run"), + sandbox_id: sandbox_id.to_string(), + script_code: format!("echo marker-{exec_id}; date +%s; sleep {seconds}; date +%s"), + working_directory: String::new(), + env: Vec::new(), + env_scope: EnvScope::Merge, + timeout_ms: 60_000, + } +} + +/// The epoch seconds a stamped run reported either side of its sleep. +fn stamped_interval(stdout: &str) -> (i64, i64) { + let stamps: Vec = stdout + .lines() + .filter_map(|line| line.trim().parse::().ok()) + .collect(); + assert!( + stamps.len() >= 2, + "a stamped run must report a start and an end, got {stdout:?}" + ); + (stamps[0], stamps[stamps.len() - 1]) +} + +fn deprovision_all(client: &DaemonClient, sandbox_ids: Vec) { + use mxc_sdk::wslc_common::daemon_protocol::DeprovisionConfig; + + for sandbox_id in sandbox_ids { + let _ = client.deprovision(DeprovisionConfig { sandbox_id }); + } +} + +/// Two clients running against their own sandboxes share the daemon: both are +/// admitted, each sees only its own output, and the runs overlap in the guest. +#[test] +#[ignore = "requires a WSL2 host with alpine:latest already in the daemon session cache"] +fn two_clients_exec_concurrently_over_the_pipe() { + let _daemon = DaemonProcess::spawn_ready(); + let client = DaemonClient::connect().expect("connect to daemon"); + + let sandboxes = vec![ + provisioned_and_started(&client), + provisioned_and_started(&client), + ]; + + let runs: Vec<_> = ["alpha", "beta"] + .iter() + .zip(sandboxes.iter()) + .map(|(tag, sandbox_id)| { + let client = client.clone(); + let config = stamped_exec(tag, sandbox_id, 6); + std::thread::spawn(move || { + let mut stdout = Vec::new(); + let completion = client + .exec_streaming(config, |stream, data| { + if stream == OutStream::Stdout { + stdout.extend_from_slice(data); + } + }) + .expect("both clients must be admitted"); + (completion, String::from_utf8_lossy(&stdout).into_owned()) + }) + }) + .collect(); + + let results: Vec<_> = runs + .into_iter() + .map(|handle| handle.join().expect("run thread")) + .collect(); + + for ((completion, stdout), tag) in results.iter().zip(["alpha", "beta"]) { + assert_eq!(completion.outcome, ExecTerminal::Exited(0)); + assert!(!completion.truncated, "{tag} reported dropped output"); + assert!( + stdout.contains(&format!("marker-{tag}")), + "{tag} did not receive its own output: {stdout:?}" + ); + } + assert!( + !results[0].1.contains("marker-beta") && !results[1].1.contains("marker-alpha"), + "each client must receive only its own stream" + ); + + // Both sandboxes share one utility VM, so their clocks agree. + let (alpha_start, alpha_end) = stamped_interval(&results[0].1); + let (beta_start, beta_end) = stamped_interval(&results[1].1); + let overlap = alpha_end.min(beta_end) - alpha_start.max(beta_start); + assert!( + overlap >= 2, + "the two 6s runs overlapped by {overlap}s, so they were serialized" + ); + + deprovision_all(&client, sandboxes); +} + +/// The daemon's exec capacity is finite, so enough simultaneous clients are +/// refused before admission. +#[test] +#[ignore = "requires a WSL2 host with alpine:latest already in the daemon session cache"] +fn the_global_exec_cap_refuses_the_excess() { + // The control server's own cap, which a client can observe only as the + // number of simultaneous runs the daemon admits. + const EXEC_CAP: usize = 8; + const CLIENTS: usize = EXEC_CAP + 2; + + let _daemon = DaemonProcess::spawn_ready(); + let client = DaemonClient::connect().expect("connect to daemon"); + + let sandboxes: Vec = (0..CLIENTS) + .map(|_| provisioned_and_started(&client)) + .collect(); + + // Long enough that every client has had its admission answered before the + // first run frees a slot. + let runs: Vec<_> = sandboxes + .iter() + .enumerate() + .map(|(n, sandbox_id)| { + let client = client.clone(); + let config = stamped_exec(&format!("capacity-{n}"), sandbox_id, 6); + std::thread::spawn(move || client.exec_streaming(config, |_, _| {})) + }) + .collect(); + + let results: Vec<_> = runs + .into_iter() + .map(|handle| handle.join().expect("run thread")) + .collect(); + + let admitted = results.iter().filter(|result| result.is_ok()).count(); + let refused = results + .iter() + .filter(|result| { + matches!( + result, + Err(DaemonError::Daemon { + kind: ErrKind::Busy, + .. + }) + ) + }) + .count(); + + assert_eq!( + admitted, EXEC_CAP, + "the daemon must admit exactly its cap, got {admitted} of {CLIENTS}" + ); + assert_eq!( + refused, + CLIENTS - EXEC_CAP, + "every client past the cap must be refused as busy, got {refused} of {CLIENTS}" + ); + + deprovision_all(&client, sandboxes); +} + #[test] fn ping_round_trip_over_pipe() { let _daemon = DaemonProcess::spawn_ready(); diff --git a/tests/scripts/run_wslc_state_aware_tests.ps1 b/tests/scripts/run_wslc_state_aware_tests.ps1 index a0a3dd455..c2c7e03b2 100644 --- a/tests/scripts/run_wslc_state_aware_tests.ps1 +++ b/tests/scripts/run_wslc_state_aware_tests.ps1 @@ -237,12 +237,19 @@ function ConvertTo-StateAwareInvocation { } } -# Parse or clone a state-aware request, move lifecycle routing to executor -# arguments, encode the remaining phase-specific payload, and run wxc-exec. -# A `wslc:{{SANDBOX_ID}}` fixture placeholder retains backend identity during -# static corpus parsing and is replaced as one unit by the full real ID. -# Returns @{ ExitCode; Stdout; Stderr }. -function Invoke-StateAware { +# Launch one state-aware phase and return without waiting, so a caller can start +# another phase while this one is still running. Complete it with Wait-StateAware. +# +# Drive wxc-exec via System.Diagnostics.Process rather than Start-Process +# -Wait: the state-aware phase process exits as soon as it has driven the +# daemon, but the daemon lives on. Start-Process -Wait does a process-tree +# (job-object) wait and would block until the surviving daemon also exits. +# WaitForExit()/ReadToEnd() here wait only on the direct phase process (the +# daemon does not inherit its stdio), matching how real SDK callers behave; +# capturing stderr this way also keeps wxc-exec's warnings off PowerShell's +# error stream (a bare `&` + `2>` throws under $ErrorActionPreference=Stop). +function Start-StateAware { + [CmdletBinding()] param( [hashtable]$Request, [string]$ConfigFile, @@ -262,14 +269,6 @@ function Invoke-StateAware { if ($Debug) { $argList += '--debug' } $argList += @('--config-base64', $invocation.ConfigBase64) - # Drive wxc-exec via System.Diagnostics.Process rather than Start-Process - # -Wait: the state-aware phase process exits as soon as it has driven the - # daemon, but the daemon lives on. Start-Process -Wait does a process-tree - # (job-object) wait and would block until the surviving daemon also exits. - # WaitForExit()/ReadToEnd() here wait only on the direct phase process (the - # daemon does not inherit its stdio), matching how real SDK callers behave; - # capturing stderr this way also keeps wxc-exec's warnings off PowerShell's - # error stream (a bare `&` + `2>` throws under $ErrorActionPreference=Stop). $psi = [System.Diagnostics.ProcessStartInfo]::new() $psi.FileName = $WxcExec foreach ($a in $argList) { $psi.ArgumentList.Add($a) } @@ -281,19 +280,155 @@ function Invoke-StateAware { $proc = [System.Diagnostics.Process]::new() $proc.StartInfo = $psi $null = $proc.Start() - # Read stdout to EOF (fires when the phase process exits) while draining - # stderr async so a large stderr can never deadlock the stdout read. - $stderrTask = $proc.StandardError.ReadToEndAsync() - $stdoutText = $proc.StandardOutput.ReadToEnd() - $proc.WaitForExit() - $stderrText = $stderrTask.GetAwaiter().GetResult() - $exitCode = $proc.ExitCode - $proc.Dispose() + + # Read both pipes async so neither can fill and stall the child while this + # caller is off starting another phase. @{ - ExitCode = $exitCode - Stdout = if ($null -eq $stdoutText) { "" } else { [string]$stdoutText } - Stderr = if ($null -eq $stderrText) { "" } else { [string]$stderrText } + Process = $proc + StdoutTask = $proc.StandardOutput.ReadToEndAsync() + StderrTask = $proc.StandardError.ReadToEndAsync() + StartedAt = [DateTime]::UtcNow + } +} + +# Complete a Start-StateAware handle. Returns +# @{ ExitCode; Stdout; Stderr; StartedAt; EndedAt } with UTC timestamps. +function Wait-StateAware { + param([hashtable]$Handle) + + $stdoutText = $Handle.StdoutTask.GetAwaiter().GetResult() + $stderrText = $Handle.StderrTask.GetAwaiter().GetResult() + $Handle.Process.WaitForExit() + $endedAt = [DateTime]::UtcNow + $exitCode = $Handle.Process.ExitCode + $Handle.Process.Dispose() + @{ + ExitCode = $exitCode + Stdout = if ($null -eq $stdoutText) { "" } else { [string]$stdoutText } + Stderr = if ($null -eq $stderrText) { "" } else { [string]$stderrText } + StartedAt = $Handle.StartedAt + EndedAt = $endedAt + } +} + +# Start a phase like Start-StateAware but leave stdout unread, so a caller can +# watch for a marker line while the process is still running. Pair with +# Wait-StateAwareLines, which returns the same shape as Wait-StateAware. +function Start-StateAwareLines { + [CmdletBinding()] + param( + [hashtable]$Request, + [string]$ConfigFile, + [string]$SandboxId, + [int]$WindowsPort + ) + + $invocation = ConvertTo-StateAwareInvocation ` + -Request $Request -ConfigFile $ConfigFile -SandboxId $SandboxId -WindowsPort $WindowsPort + + $argList = @('--operation', $invocation.Operation) + if ($invocation.Operation -ne 'provision') { + $argList += @('--container-id', $invocation.SandboxId) + } + if ($Debug) { $argList += '--debug' } + $argList += @('--config-base64', $invocation.ConfigBase64) + + $psi = [System.Diagnostics.ProcessStartInfo]::new() + $psi.FileName = $WxcExec + foreach ($a in $argList) { $psi.ArgumentList.Add($a) } + $psi.UseShellExecute = $false + $psi.RedirectStandardOutput = $true + $psi.RedirectStandardError = $true + $psi.CreateNoWindow = $true + + $proc = [System.Diagnostics.Process]::new() + $proc.StartInfo = $psi + $null = $proc.Start() + + # Stderr drains async so a large stderr cannot stall the child while this + # caller is reading stdout a line at a time. + @{ + Process = $proc + StderrTask = $proc.StandardError.ReadToEndAsync() + Buffered = New-Object System.Collections.Generic.List[string] + PendingLine = $null + StartedAt = [DateTime]::UtcNow + } +} + +# Read stdout lines until one matches $Marker, keeping them for the eventual +# Wait-StateAwareLines. $false if the stream ends or the wait runs out. +function Wait-StateAwareMarker { + param( + [hashtable]$Handle, + [string]$Marker, + [int]$TimeoutSec = 60 + ) + + $deadline = [DateTime]::UtcNow.AddSeconds($TimeoutSec) + while ([DateTime]::UtcNow -lt $deadline) { + if ($null -eq $Handle.PendingLine) { + $Handle.PendingLine = $Handle.Process.StandardOutput.ReadLineAsync() + } + if (-not $Handle.PendingLine.Wait(250)) { continue } + $line = $Handle.PendingLine.GetAwaiter().GetResult() + $Handle.PendingLine = $null + if ($null -eq $line) { return $false } + $null = $Handle.Buffered.Add($line) + if ($line -match $Marker) { return $true } + } + $false +} + +# Complete a Start-StateAwareLines handle, rejoining the lines already read with +# the rest of the stream. +function Wait-StateAwareLines { + param([hashtable]$Handle) + + if ($null -ne $Handle.PendingLine) { + $pending = $Handle.PendingLine.GetAwaiter().GetResult() + $Handle.PendingLine = $null + if ($null -ne $pending) { $null = $Handle.Buffered.Add($pending) } + } + $rest = $Handle.Process.StandardOutput.ReadToEndAsync().GetAwaiter().GetResult() + $stderrText = $Handle.StderrTask.GetAwaiter().GetResult() + $Handle.Process.WaitForExit() + $endedAt = [DateTime]::UtcNow + $exitCode = $Handle.Process.ExitCode + $Handle.Process.Dispose() + + $stdout = [string]::Join([Environment]::NewLine, $Handle.Buffered) + if (-not [string]::IsNullOrEmpty($rest)) { + if ($stdout.Length -gt 0) { $stdout += [Environment]::NewLine } + $stdout += $rest } + @{ + ExitCode = $exitCode + Stdout = $stdout + Stderr = if ($null -eq $stderrText) { "" } else { [string]$stderrText } + StartedAt = $Handle.StartedAt + EndedAt = $endedAt + } +} + +# Parse or clone a state-aware request, move lifecycle routing to executor +# arguments, encode the remaining phase-specific payload, and run wxc-exec. +# A `wslc:{{SANDBOX_ID}}` fixture placeholder retains backend identity during +# static corpus parsing and is replaced as one unit by the full real ID. +# Returns @{ ExitCode; Stdout; Stderr }. +function Invoke-StateAware { + [CmdletBinding()] + param( + [hashtable]$Request, + [string]$ConfigFile, + [string]$SandboxId, + [int]$WindowsPort, + [switch]$DryRun + ) + + Wait-StateAware (Start-StateAware ` + -Request $Request -ConfigFile $ConfigFile -SandboxId $SandboxId ` + -WindowsPort $WindowsPort -DryRun:$DryRun) } # Like Invoke-StateAware but records the wall-clock arrival time of each stdout @@ -376,6 +511,15 @@ function Envelope-Arm { '' } +# The epoch seconds a stamped exec printed before and after its sleep; $null if +# the run did not report both. +function Get-ExecInterval { + param([string]$Stdout) + $stamps = @($Stdout -split '\r?\n' | Where-Object { $_.Trim() -match '^\d+$' } | ForEach-Object { [long]$_.Trim() }) + if ($stamps.Count -lt 2) { return $null } + @{ Start = $stamps[0]; End = $stamps[-1] } +} + # Is the daemon process currently running? function Test-DaemonRunning { $null -ne (Get-Process -Name $DaemonProcName -ErrorAction SilentlyContinue) @@ -1532,6 +1676,189 @@ try { } } +# ---------------- Lifecycle J: exec concurrency ---------------- + +# Placed ahead of H because H drives the daemon to exit. Each scenario needs two +# phases in flight at once -- overlapping runs on separate sandboxes, a refused +# same-container second exec, and a lifecycle command issued mid-run -- so all +# three use Start-StateAware rather than the sequential Invoke-StateAware. +$script:ccSandboxA = $null +$script:ccSandboxB = $null +$ccADeprovisionedOk = $false +$ccBDeprovisionedOk = $false +try { + $ccAReady = $false + $ccBReady = $false + $ccAProvOk = Run-StateAwareTest "J: provision sandbox A" { + $r = Invoke-StateAware -ConfigFile 'wslc_state_aware_provision.json' + $envObj = Assert-ResultEnvelope $r "concurrency A provision" + if ($envObj) { $script:ccSandboxA = [string]$envObj.result.sandboxId } + } + $ccBProvOk = Run-StateAwareTest "J: provision sandbox B" { + $r = Invoke-StateAware -ConfigFile 'wslc_state_aware_provision.json' + $envObj = Assert-ResultEnvelope $r "concurrency B provision" + if ($envObj) { $script:ccSandboxB = [string]$envObj.result.sandboxId } + } + if ($ccAProvOk) { + $ccAReady = Run-StateAwareTest "J: start sandbox A" { + $r = Invoke-StateAware -ConfigFile 'wslc_state_aware_start.json' -SandboxId $script:ccSandboxA + $null = Assert-ResultEnvelope $r "concurrency A start" + } + } + if ($ccBProvOk) { + $ccBReady = Run-StateAwareTest "J: start sandbox B" { + $r = Invoke-StateAware -ConfigFile 'wslc_state_aware_start.json' -SandboxId $script:ccSandboxB + $null = Assert-ResultEnvelope $r "concurrency B start" + } + } + + # J1: both sandboxes share one utility VM, so their clocks agree and the two + # runs can be compared directly. Serialized runs cannot overlap at all. + if ($ccAReady -and $ccBReady) { + Run-StateAwareTest "J: execs on two sandboxes overlap in wall clock" { + $stamped = "sh -c 'date +%s; sleep 6; date +%s'" + $sleepA = @{ phase = 'exec'; sandboxId = $script:ccSandboxA; process = @{ commandLine = $stamped; timeout = 60000 } } + $sleepB = @{ phase = 'exec'; sandboxId = $script:ccSandboxB; process = @{ commandLine = $stamped; timeout = 60000 } } + + $handleA = Start-StateAware -Request $sleepA + $handleB = Start-StateAware -Request $sleepB + $ra = Wait-StateAware $handleA + $rb = Wait-StateAware $handleB + + Assert-True ($ra.ExitCode -eq 0) "A exec exit 0" + Assert-True ($rb.ExitCode -eq 0) "B exec exit 0" + + $ia = Get-ExecInterval -Stdout $ra.Stdout + $ib = Get-ExecInterval -Stdout $rb.Stdout + Assert-True ($null -ne $ia -and $null -ne $ib) "both runs reported an in-container start and end" + if ($null -ne $ia -and $null -ne $ib) { + $overlapSec = [math]::Min($ia.End, $ib.End) - [math]::Max($ia.Start, $ib.Start) + Assert-True ($overlapSec -ge 2) ` + "the two 6s runs overlapped by ${overlapSec}s inside their containers" + } + } | Out-Null + } + + # J2: the refusal is a pre-admission error, so it reaches the client as an + # envelope on stderr rather than a terminal stream frame. + if ($ccAReady) { + Run-StateAwareTest "J: second exec on the same sandbox is refused as busy" { + $blocker = @{ phase = 'exec'; sandboxId = $script:ccSandboxA; process = @{ commandLine = "sh -c 'echo blocker-ready; sleep 6; echo blocker-done'"; timeout = 30000 } } + $second = @{ phase = 'exec'; sandboxId = $script:ccSandboxA; process = @{ commandLine = 'echo should-not-run'; timeout = 30000 } } + + $blockerHandle = Start-StateAwareLines -Request $blocker + $ready = Wait-StateAwareMarker -Handle $blockerHandle -Marker 'blocker-ready' + Assert-True $ready "the first exec reported running inside its container" + $rs = Invoke-StateAware -Request $second + + Assert-True ($rs.ExitCode -ne 0) "the refused exec exits non-zero" + Assert-True ($rs.Stdout -notmatch 'should-not-run') "the refused command never ran" + $envObj = Parse-StderrEnvelope -Stderr $rs.Stderr + $code = if ($envObj) { $envObj.error.code } else { '' } + $message = if ($envObj) { [string]$envObj.error.message } else { '' } + Assert-True ($code -eq 'backend_error') "busy surfaces as 'backend_error' (got '$code')" + Assert-True ($message -match 'already has an exec in flight') ` + "the message names the single-flight slot (got '$message')" + + $rb = Wait-StateAwareLines $blockerHandle + Assert-True ($rb.ExitCode -eq 0) "the first exec is unaffected by the refusal" + Assert-True ($rb.Stdout -match 'blocker-done') "the first exec ran to completion" + } | Out-Null + } + + # J3: stop must wait for the run instead of deleting the container under it, + # so its wall clock covers the rest of the 4s exec. + if ($ccBReady) { + $ccStopWaited = Run-StateAwareTest "J: stop during an in-flight exec waits for the run" { $longRun = @{ phase = 'exec'; sandboxId = $script:ccSandboxB; process = @{ commandLine = "sh -c 'echo run-ready; sleep 6; echo run-survived'"; timeout = 30000 } } + + $execHandle = Start-StateAwareLines -Request $longRun + $ready = Wait-StateAwareMarker -Handle $execHandle -Marker 'run-ready' + Assert-True $ready "the exec reported running inside its container" + $stopHandle = Start-StateAware -ConfigFile 'wslc_state_aware_stop.json' -SandboxId $script:ccSandboxB + $rstop = Wait-StateAware $stopHandle + $rexec = Wait-StateAwareLines $execHandle + + Assert-True ($rexec.ExitCode -eq 0) "the in-flight exec still exits 0" + Assert-True ($rexec.Stdout -match 'run-survived') "the run was not cut short by the stop" + $null = Assert-ResultEnvelope $rstop "stop issued during an exec" + $stopSec = [math]::Round(($rstop.EndedAt - $rstop.StartedAt).TotalSeconds, 2) + Assert-True ($stopSec -ge 1.5) ` + "stop took ${stopSec}s, so it waited for the run rather than racing it" + } | Out-Null + } + + # J4: a long exec on A must not delay a full lifecycle on another sandbox, + # so C's own run has to fall inside A's. Both sandboxes share one utility VM, + # so their in-container stamps are on the same clock. + if ($ccAReady) { + Run-StateAwareTest "J: a long exec on A does not block lifecycle work on another sandbox" { + $blocker = @{ phase = 'exec'; sandboxId = $script:ccSandboxA; process = @{ commandLine = "sh -c 'echo blocker-ready; date +%s; sleep 25; date +%s; echo blocker-done'"; timeout = 60000 } } + $blockerHandle = Start-StateAwareLines -Request $blocker + $ready = Wait-StateAwareMarker -Handle $blockerHandle -Marker 'blocker-ready' + Assert-True $ready "A's exec reported running inside its container" + + $sandboxC = $null + $icC = $null + try { + $rp = Invoke-StateAware -ConfigFile 'wslc_state_aware_provision.json' + $envObj = Assert-ResultEnvelope $rp "C provision during A's exec" + if ($envObj) { $sandboxC = [string]$envObj.result.sandboxId } + + if ($sandboxC) { + $rs = Invoke-StateAware -ConfigFile 'wslc_state_aware_start.json' -SandboxId $sandboxC + $null = Assert-ResultEnvelope $rs "C start during A's exec" + + $req = @{ phase = 'exec'; sandboxId = $sandboxC; process = @{ commandLine = "sh -c 'date +%s; echo C-ran-during-A; date +%s'"; timeout = 30000 } } + $re = Invoke-StateAware -Request $req + Assert-True ($re.ExitCode -eq 0) "C exec during A's exec exits 0" + Assert-True ($re.Stdout -match 'C-ran-during-A') "C produced its own output" + $icC = Get-ExecInterval -Stdout $re.Stdout + + $rstop = Invoke-StateAware -ConfigFile 'wslc_state_aware_stop.json' -SandboxId $sandboxC + $null = Assert-ResultEnvelope $rstop "C stop during A's exec" + } + } finally { + if ($sandboxC) { + try { $null = Invoke-StateAware -ConfigFile 'wslc_state_aware_deprovision.json' -SandboxId $sandboxC } catch { } + } + } + + $rb = Wait-StateAwareLines $blockerHandle + Assert-True ($rb.ExitCode -eq 0) "A's exec still exits 0" + Assert-True ($rb.Stdout -match 'blocker-done') "A's exec ran to completion" + + $iaA = Get-ExecInterval -Stdout $rb.Stdout + Assert-True ($null -ne $iaA -and $null -ne $icC) "both runs reported in-container stamps" + if ($null -ne $iaA -and $null -ne $icC) { + Assert-True ($icC.Start -ge $iaA.Start -and $icC.End -le $iaA.End) ` + "C ran inside A's run: A $($iaA.Start)..$($iaA.End), C $($icC.Start)..$($icC.End)" + } + } | Out-Null + } + + if ($ccAProvOk) { + $ccADeprovPassed = Run-StateAwareTest "J: deprovision A" { + $r = Invoke-StateAware -ConfigFile 'wslc_state_aware_deprovision.json' -SandboxId $script:ccSandboxA + $null = Assert-ResultEnvelope $r "concurrency A deprovision" + } + if ($ccADeprovPassed) { $ccADeprovisionedOk = $true } + } + if ($ccBProvOk) { + $ccBDeprovPassed = Run-StateAwareTest "J: deprovision B" { + $r = Invoke-StateAware -ConfigFile 'wslc_state_aware_deprovision.json' -SandboxId $script:ccSandboxB + $null = Assert-ResultEnvelope $r "concurrency B deprovision" + } + if ($ccBDeprovPassed) { $ccBDeprovisionedOk = $true } + } +} finally { + if ($null -ne $script:ccSandboxA -and -not $ccADeprovisionedOk) { + try { $null = Invoke-StateAware -ConfigFile 'wslc_state_aware_deprovision.json' -SandboxId $script:ccSandboxA } catch { } + } + if ($null -ne $script:ccSandboxB -and -not $ccBDeprovisionedOk) { + try { $null = Invoke-StateAware -ConfigFile 'wslc_state_aware_deprovision.json' -SandboxId $script:ccSandboxB } catch { } + } +} + # ---------------- Lifecycle H: idle teardown ---------------- # All sandboxes above have been deprovisioned, so the daemon's live-container