diff --git a/docs/backends/wslc/wslc-state-aware.md b/docs/backends/wslc/wslc-state-aware.md index 03ef2462e..dde7f0632 100644 --- a/docs/backends/wslc/wslc-state-aware.md +++ b/docs/backends/wslc/wslc-state-aware.md @@ -40,10 +40,12 @@ the live SDK handles. Each phase process is a thin client that contacts the daem pipe; the daemon performs the actual SDK calls and streams stdio back. This mirrors the Windows Sandbox daemon pattern. -The daemon runs all SDK calls on a **single apartment-affine worker thread** (the WSLc SDK handles -are not thread-agnostic). Today `exec` blocks that worker for the duration of the run, so commands -against different sandboxes are serialized — correct, just not concurrent. See -[Known limitations](#known-limitations). +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). ## Components @@ -141,8 +143,7 @@ the streaming path reports it through its wait result. ### exec admission, cancellation, and failure containment -The daemon runs one exec at a time because all WSLc SDK operations are confined -to one apartment-affine worker. A concurrent exec is rejected with +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 @@ -301,12 +302,15 @@ fixtures **through the harness**, not by pointing `wxc-exec --config` at them di ## Known limitations -- **Serialized exec (deferred).** Because the daemon's single worker thread blocks on - `WaitForSingleObject` for the whole `exec`, no other sandbox can provision or exec while one - command runs, and a per-container single-flight `Busy` guard is not yet meaningful. The intended - fix splits `exec` into an on-worker `ExecStart` (extract the thread-agnostic Win32 exit-event - handle) + an off-thread wait + an on-worker `ExecReap`, with a per-container `in_flight` slot. This - is tracked as follow-up work. +- **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.** 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. - **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/bin/wslc_daemon/control_server.rs b/src/mxc-sdk/src/bin/wslc_daemon/control_server.rs index 0831cf15e..6a4c6ecfd 100644 --- a/src/mxc-sdk/src/bin/wslc_daemon/control_server.rs +++ b/src/mxc-sdk/src/bin/wslc_daemon/control_server.rs @@ -48,8 +48,8 @@ use mxc_sdk::wslc_common::daemon_protocol::{ use crate::session_manager::{ExecStream, SessionHandle, WorkerError}; -/// The WSLc SDK worker is apartment-affine and executes one command at a time. -/// Refuse additional execs instead of admitting a queue that cannot run. +/// 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 @@ -510,10 +510,11 @@ where /// [`StreamFrame`]s, followed by a terminal frame. /// /// The sandbox is validated (exists + started) *before* the `Ok` admission is -/// written, and — critically — admission is **atomic** with the start of the -/// run on the worker thread (see [`SessionHandle::exec`]): the worker validates -/// and begins running within one command handler, so no `Stop`/`Deprovision` -/// can invalidate the checked state between the admission and the run. An +/// 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 /// [`DaemonResponse::Err`] rather than a post-admission stream `Error` frame. /// @@ -533,9 +534,21 @@ async fn handle_exec( where S: AsyncWrite + Unpin, { + let exec_id = config.exec_id.clone(); + let run_token = config.run_token.clone(); + // Await the worker's admission decision before writing anything: a rejected // exec is a pre-admission typed error, never a post-admission stream frame. - write_exec_result(&mut pipe, session.exec(config).await).await + let delivered = write_exec_result(&mut pipe, session.exec(config).await).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. + session.cancel_exec(&exec_id, &run_token); + } + + delivered } /// Turn an exec **admission** outcome into the client's frame sequence, generic @@ -739,6 +752,49 @@ mod tests { session.shutdown().await.unwrap(); } + #[tokio::test] + async fn an_undelivered_exec_is_cancelled_so_its_run_cannot_outlive_the_permit() { + use crate::session_manager::register_exec; + + let session = crate::session_manager::spawn().unwrap(); + let cancellation = Arc::new(AtomicBool::new(false)); + let registration = register_exec( + session.active_execs(), + "orphan-1", + "orphan-run-1", + &cancellation, + ) + .unwrap(); + + let limiter = Arc::new(Semaphore::new(MAX_CONCURRENT_EXECS)); + let permit = limiter.clone().try_acquire_owned().unwrap(); + let delivered = handle_exec( + BrokenPipe, + session.clone(), + mxc_sdk::wslc_common::daemon_protocol::ExecConfig { + exec_id: "orphan-1".to_string(), + run_token: "orphan-run-1".to_string(), + sandbox_id: "wslc:does-not-exist".to_string(), + script_code: "sleep 600".to_string(), + working_directory: String::new(), + env: Vec::new(), + env_scope: mxc_sdk::wslc_common::process_env::EnvScope::Merge, + timeout_ms: 0, + }, + permit, + ) + .await; + + 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" + ); + drop(registration); + session.shutdown().await.unwrap(); + } + /// A writer that always fails, standing in for a client that has gone. struct BrokenPipe; 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 0c4f9676f..3d330859a 100644 --- a/src/mxc-sdk/src/bin/wslc_daemon/session_manager.rs +++ b/src/mxc-sdk/src/bin/wslc_daemon/session_manager.rs @@ -5,10 +5,12 @@ //! thread. //! //! The WSLc SDK's `WslcSession` / `WslcContainer` / `WslcProcess` handles are -//! thread-affine and must all be created and used from one apartment. This -//! module confines every SDK call to a single long-lived worker thread that -//! joins the MTA for its whole lifetime; async pipe handlers dispatch typed -//! [`WorkerCommand`]s to it over a channel and await the reply. +//! apartment-affine: any thread that has joined the MTA may use them. A single +//! long-lived worker thread owns the container map and every short SDK call; +//! async pipe handlers dispatch typed [`WorkerCommand`]s to it over a channel +//! and await the reply. An image pull and an `exec` both run for an unbounded +//! time, so each takes a thread of its own and posts its outcome back to the +//! worker as another command. //! //! State-aware topology (decided): **one** shared `WslcSession` (the WSL2 //! utility VM, booted lazily on first provision and amortised across all @@ -20,15 +22,16 @@ //! [`mxc_sdk::wslc_common::container_steps`] and [`mxc_sdk::wslc_common::image`]: `provision` //! ensures the session + resolves the image + creates a container with a //! keepalive init process; `start` boots -//! it; `exec` runs a fresh `WslcCreateContainerProcess` to completion, streaming -//! its stdout/stderr live to the pipe handler via an [`OutputSink`]; `stop` / -//! `deprovision` stop + delete. The completion reply carries the exit code; -//! output flows over the sink. (Client `Stdin` forwarding is a later fill-in.) +//! it; `exec` runs a fresh `WslcCreateContainerProcess` to completion on a +//! thread of its own, streaming its stdout/stderr live to the pipe handler via +//! an [`OutputSink`]; `stop` / `deprovision` stop + delete. The completion reply +//! carries the exit code; output flows over the sink. (Client `Stdin` forwarding +//! is a later fill-in.) use std::collections::{hash_map::Entry, HashMap}; -use std::sync::atomic::{AtomicBool, Ordering}; +use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering}; use std::sync::{Arc, Mutex, Weak}; -use std::time::Duration; +use std::time::{Duration, Instant}; use anyhow::Result; use tokio::sync::mpsc::error::TrySendError; @@ -54,6 +57,25 @@ const SESSION_NAME: &str = "mxc-wslc-daemon"; /// How long teardown waits for an abandoned pull before giving up on it. const PULL_DRAIN_TIMEOUT: Duration = Duration::from_secs(30); +/// How long teardown waits for an off-worker exec before giving up on it. +const EXEC_DRAIN_TIMEOUT: Duration = Duration::from_secs(30); + +/// Block until `in_flight` reads zero, or `budget` expires. +/// +/// Reports whether the wait drained. A run that has not reported back is still +/// using the container handle it was given, so releasing that handle while one +/// is outstanding would free memory the SDK holds. +fn wait_for_execs_in_flight(in_flight: &AtomicUsize, budget: Duration) -> bool { + let deadline = Instant::now() + budget; + while in_flight.load(Ordering::SeqCst) > 0 { + if Instant::now() >= deadline { + return false; + } + std::thread::sleep(Duration::from_millis(50)); + } + true +} + /// Default on-disk WSLc image/session store. Matches the one-shot runner's /// default so both surfaces read and fill the same cache. fn default_storage_path() -> String { @@ -135,41 +157,20 @@ pub enum WorkerCommand { config: ProvisionConfig, reply: oneshot::Sender>, }, - Start { - config: StartConfig, - reply: oneshot::Sender>, - }, - /// Validate the sandbox (exists + started) and, if admitted, run the - /// command to completion. The two replies make admission **atomic** with the - /// run: because the worker services this whole command on its single thread - /// without yielding, no `Stop`/`Deprovision` can interleave between the - /// validation and the run. `admit` carries the pre-run decision (so an - /// unknown/not-started sandbox is a pre-admission typed error, never a - /// post-admission stream `Error`); `done` carries the run's exit code. - Exec { - 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. - sink: OutputSink, - cancellation: Arc, - registration: Arc, - admit: oneshot::Sender>, - done: oneshot::Sender>, - }, - Stop { - config: StopConfig, - reply: oneshot::Sender>, - }, - Deprovision { - config: DeprovisionConfig, - reply: oneshot::Sender>, - }, + /// Anything addressed to one container. Routing every such command through + /// a single variant is what forces it through [`Worker::dispatch`], so a new + /// one cannot reach a container whose exec is still using its handle. + Container(ContainerWork), /// Report a parked provision's pull, from the pull thread or its deadline. PullFinished { token: u64, outcome: Result<(), ScriptResponse>, }, + /// Report an off-worker exec, from its run thread. + ExecFinished { + sandbox_id: String, + report: ExecReport, + }, /// Give up on a sandbox whose provision reply never reached a client. Retire { sandbox_id: String, @@ -181,6 +182,161 @@ pub enum WorkerCommand { Shutdown { reply: oneshot::Sender<()> }, } +/// 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. + pub sink: OutputSink, + pub cancellation: Arc, + pub(crate) registration: Arc, + pub admit: oneshot::Sender>, + pub done: oneshot::Sender>, +} + +/// What an off-worker exec thread hands back, already classified. +/// +/// Quarantine stays on the worker because it mutates the container map and +/// issues its own SDK delete, so an unconfirmed run reports the detail instead +/// of acting on it. +#[derive(Debug)] +pub enum ExecReport { + Finished(Result), + Unconfirmed(String), +} + +/// Everything an off-worker exec needs, in a form that can cross a thread. +struct ExecJob { + sdk: *const WslcSdk, + container: WslcContainer, + config: ExecConfig, + sink: OutputSink, + cancellation: Arc, +} + +// SAFETY: the SDK is apartment-affine, not thread-affine. Both the worker and +// the spawned thread join the MTA, where a handle may be used from any member +// thread. Measured with two containers in one session running `sleep 5` from two +// MTA threads: both processes started within 0.5ms of each other and overlapped +// for 5.1100s of a 5.1177s wall clock, both exited 0, and each process's output +// arrived only on its own capture buffer and live sink. +unsafe impl Send for ExecJob {} + +/// Keeps the in-flight count accurate even if the run panics, so teardown +/// cannot be blocked forever by a thread that is already gone. +struct ExecCount(Arc); + +impl ExecCount { + fn enter(counter: &Arc) -> Self { + counter.fetch_add(1, Ordering::SeqCst); + Self(Arc::clone(counter)) + } +} + +impl Drop for ExecCount { + fn drop(&mut self) { + self.0.fetch_sub(1, Ordering::SeqCst); + } +} + +/// Run an exec on a thread of its own, reporting back to the worker. +/// +/// Returns as soon as the thread starts. The caller keeps `sdk` loaded and the +/// container handle alive until [`WorkerCommand::ExecFinished`] lands -- see +/// [`wait_for_execs_in_flight`]. +fn start_exec( + job: ExecJob, + worker: mpsc::UnboundedSender, + in_flight: &Arc, +) -> Result<(), std::io::Error> { + let counted = ExecCount::enter(in_flight); + let spawned = std::thread::Builder::new() + .name("wslc-exec".to_string()) + .spawn(move || { + let job = job; + + // Releasing the count only after the report is queued is what lets a + // teardown that sees zero rely on finding it. + let _counted = counted; + let sandbox_id = job.config.sandbox_id.clone(); + + let run = std::panic::AssertUnwindSafe(|| match ComApartment::enter() { + Err(e) => ExecReport::Finished(Err(WorkerError::Backend(e))), + Ok(_apartment) => { + let mut log = Logger::new(Mode::Console); + run_exec(job, &mut log) + } + }); + + // A panic leaves the process's fate unknown, which is what + // `Unconfirmed` already means. + let report = std::panic::catch_unwind(run).unwrap_or_else(|_| { + ExecReport::Unconfirmed("the exec thread panicked".to_string()) + }); + + let _ = worker.send(WorkerCommand::ExecFinished { sandbox_id, report }); + }); + + // A spawn failure drops the guard along with the closure, so there is no + // release to do here. + spawned.map(|_| ()) +} + +/// The off-worker half of an exec: run it to completion and classify it. +fn run_exec(job: ExecJob, logger: &mut Logger) -> ExecReport { + let ExecJob { + sdk, + container, + config, + sink, + cancellation, + } = job; + + // ProcessSettings::build expects `NAME=VALUE` env entries. + let env: Vec = config + .env + .iter() + .map(|(k, v)| format!("{}={}", k, v)) + .collect(); + + // SAFETY: `sdk` is valid and `container` is a live, started handle; the + // worker releases neither while this run is still counted in flight. + let outcome = unsafe { + container_steps::exec_in_container( + &*sdk, + container, + &config.script_code, + &env, + config.env_scope, + &config.working_directory, + config.timeout_ms, + &cancellation, + Some(sink), + logger, + ) + }; + + let outcome = match outcome { + Ok(outcome) => outcome, + Err(e) => return ExecReport::Finished(Err(sr_err(e))), + }; + + match outcome.completion { + container_steps::ProcessCompletion::TerminationUnconfirmed => ExecReport::Unconfirmed( + outcome + .post_launch_error + .map(|error| error.error_message) + .unwrap_or_else(|| "process termination could not be confirmed".to_string()), + ), + container_steps::ProcessCompletion::Confirmed(terminal) => { + ExecReport::Finished(Ok(terminal)) + } + } +} + /// A chunk of live process output streamed from the worker to the pipe handler: /// which stream it came from and the bytes (owned, so it can cross the channel). pub type OutputChunk = (OutStream, Vec); @@ -332,6 +488,13 @@ pub struct SessionHandle { } impl SessionHandle { + /// The live-exec registry, so a test can observe what [`Self::cancel_exec`] + /// reached. + #[cfg(test)] + pub(crate) fn active_execs(&self) -> &ActiveExecs { + &self.active_execs + } + /// Provision a container, returning its minted `sandbox_id`. pub async fn provision(&self, config: ProvisionConfig) -> Result { let (reply, rx) = oneshot::channel(); @@ -342,7 +505,10 @@ impl SessionHandle { /// Start a provisioned container. pub async fn start(&self, config: StartConfig) -> Result<(), WorkerError> { let (reply, rx) = oneshot::channel(); - self.send(WorkerCommand::Start { config, reply })?; + self.send(WorkerCommand::Container(ContainerWork::Start { + config, + reply, + }))?; rx.await.map_err(worker_gone)? } @@ -352,8 +518,8 @@ impl SessionHandle { /// 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 start of the run are atomic on the worker thread, so no lifecycle - /// command can invalidate the checked state between the two. + /// 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 { let (admit, admit_rx) = oneshot::channel(); let (done, done_rx) = oneshot::channel(); @@ -379,14 +545,14 @@ impl SessionHandle { &config.run_token, &cancellation, )?); - self.send(WorkerCommand::Exec { + self.send(WorkerCommand::Container(ContainerWork::Exec(ExecRequest { config, sink, cancellation, registration: Arc::clone(®istration), admit, done, - })?; + })))?; admit_rx.await.map_err(worker_gone)??; Ok(ExecStream { done: done_rx, @@ -396,8 +562,8 @@ impl SessionHandle { }) } - /// Signal an admitted exec without waiting for the apartment-affine worker, - /// which is blocked in that exec until the process exits. + /// Signal an admitted exec without waiting for the worker, which may have + /// parked the run behind an earlier exec on the same container. pub fn cancel_exec(&self, exec_id: &str, run_token: &str) { if let Some(cancellation) = self .active_execs @@ -414,14 +580,20 @@ impl SessionHandle { /// Stop a running container. pub async fn stop(&self, config: StopConfig) -> Result<(), WorkerError> { let (reply, rx) = oneshot::channel(); - self.send(WorkerCommand::Stop { config, reply })?; + self.send(WorkerCommand::Container(ContainerWork::Stop { + config, + reply, + }))?; rx.await.map_err(worker_gone)? } /// Deprovision (delete) a container. pub async fn deprovision(&self, config: DeprovisionConfig) -> Result<(), WorkerError> { let (reply, rx) = oneshot::channel(); - self.send(WorkerCommand::Deprovision { config, reply })?; + self.send(WorkerCommand::Container(ContainerWork::Deprovision { + config, + reply, + }))?; rx.await.map_err(worker_gone)? } @@ -484,11 +656,78 @@ struct PendingProvision { _retire_deadline: std::sync::mpsc::Sender<()>, } +/// A command addressed to one container, either about to run or parked behind +/// an exec that is still using that container's handle. +pub(crate) enum ContainerWork { + /// Validate the sandbox (exists + started) and, if admitted, hand the run to + /// 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 + /// error, never a post-admission stream `Error`); `done` carries the run's + /// exit code once [`WorkerCommand::ExecFinished`] lands. + Exec(ExecRequest), + Start { + config: StartConfig, + reply: oneshot::Sender>, + }, + Stop { + config: StopConfig, + reply: oneshot::Sender>, + }, + Deprovision { + config: DeprovisionConfig, + reply: oneshot::Sender>, + }, +} + +impl ContainerWork { + fn sandbox_id(&self) -> &str { + match self { + ContainerWork::Exec(request) => &request.config.sandbox_id, + ContainerWork::Start { config, .. } => &config.sandbox_id, + ContainerWork::Stop { config, .. } => &config.sandbox_id, + ContainerWork::Deprovision { config, .. } => &config.sandbox_id, + } + } + + /// Answer this work with `message`, for a daemon that is shutting down. + fn refuse(self, message: &str) { + let error = WorkerError::Backend(anyhow::anyhow!("{message}")); + match self { + ContainerWork::Exec(request) => { + let _ = request.admit.send(Err(error)); + } + ContainerWork::Start { reply, .. } + | ContainerWork::Stop { reply, .. } + | ContainerWork::Deprovision { reply, .. } => { + let _ = reply.send(Err(error)); + } + } + } +} + +/// 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 +/// release it until [`Worker::finish_exec`] removes this entry. +struct InFlightExec { + container: WslcContainer, + done: oneshot::Sender>, + + /// Held only so the exec id stays reserved for as long as the run lasts. + _registration: Arc, + parked: Vec, +} + /// The single-threaded WSLc session owner. Constructed and run entirely on the /// worker thread. Holds the lazily-loaded SDK and the one shared session (the /// WSL2 utility VM), plus the `sandbox_id -> container` map. The SDK/session/ -/// guard handles are `!Send` raw pointers, but they never leave this thread — -/// only [`WorkerCommand`]s cross the channel. +/// guard handles are `!Send` raw pointers; ownership stays here, and a pull or +/// exec thread borrows one only for as long as its own in-flight counter says +/// so. struct Worker { logger: Logger, // Field order is load-bearing on implicit drop: `containers` and `session` @@ -496,9 +735,18 @@ struct Worker { // unloads `wslcsdk.dll`. containers: HashMap, pending: HashMap, + exec_in_flight: HashMap, + + /// Runs this worker started that have not reported back. Shared with the run + /// threads, which outlive this struct whenever teardown abandons a handle. + execs_in_flight: Arc, next_pull_token: u64, session: Option, - sdk: Option, + + /// Boxed so the address an off-worker pull or exec was handed survives this + /// struct being moved or the field being taken, which teardown does to leak + /// the SDK rather than unload it under a running thread. + sdk: Option>, } impl Worker { @@ -509,6 +757,8 @@ impl Worker { session: None, containers: HashMap::new(), pending: HashMap::new(), + exec_in_flight: HashMap::new(), + execs_in_flight: Arc::new(AtomicUsize::new(0)), next_pull_token: 0, } } @@ -519,7 +769,7 @@ impl Worker { // SAFETY: the worker thread is already in the MTA (see `ComApartment`). let sdk = unsafe { container_steps::load_sdk_checked(&mut self.logger) }.map_err(sr_err)?; - self.sdk = Some(sdk); + self.sdk = Some(Box::new(sdk)); } if self.session.is_none() { let storage = default_storage_path(); @@ -789,9 +1039,12 @@ impl Worker { } /// Validate that a sandbox exists and is started, returning the live handle - /// needed to run. Sole owner of the exists+started invariant: [`exec`] trusts - /// the handle it is given and never re-checks, because the worker services - /// admission and the run on one thread without yielding between them. + /// 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. + /// + /// [`begin_exec`]: Worker::begin_exec fn validate_exec(&self, sandbox_id: &str) -> Result { match self.containers.get(sandbox_id) { None => Err(WorkerError::NotProvisioned(sandbox_id.to_string())), @@ -804,55 +1057,143 @@ impl Worker { } } - /// Run a command in a sandbox whose existence/started state was already - /// confirmed by [`validate_exec`]; `container` is that validated handle. - /// `sink` streams the run's stdout/stderr live to the pipe handler. - fn exec( - &mut self, - config: ExecConfig, - container: WslcContainer, - sink: OutputSink, - cancellation: &AtomicBool, - ) -> Result { - let sdk = self - .sdk - .as_ref() - .ok_or_else(|| anyhow::anyhow!("no active WSLc session"))?; + /// 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. + 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); + return; + } - // ProcessSettings::build expects `NAME=VALUE` env entries. - let env: Vec = config - .env - .iter() - .map(|(k, v)| format!("{}={}", k, v)) - .collect(); + match work { + ContainerWork::Exec(request) => self.begin_exec(request, worker), + ContainerWork::Start { config, reply } => { + let _ = reply.send(self.start(config)); + } + ContainerWork::Stop { config, reply } => { + let _ = reply.send(self.stop(config)); + } + ContainerWork::Deprovision { config, reply } => { + let _ = reply.send(self.deprovision(config)); + } + } + } - // SAFETY: `sdk` is valid and `container` is a live, started handle. - let outcome = unsafe { - container_steps::exec_in_container( - sdk, + /// Admit an exec and hand the run to a thread of its own. + /// + /// Validation, the admission reply, the claim on the container and the + /// spawn all happen here without yielding, so nothing can delete the + /// validated handle before the run thread takes it. + fn begin_exec(&mut self, request: ExecRequest, worker: &mpsc::UnboundedSender) { + let ExecRequest { + config, + sink, + cancellation, + registration, + admit, + done, + } = request; + + let container = match self.validate_exec(&config.sandbox_id) { + Ok(container) => container, + Err(e) => { + let _ = admit.send(Err(e)); + return; + } + }; + + if cancellation.load(Ordering::Acquire) { + if admit.send(Ok(())).is_ok() { + let _ = done.send(Ok(ExecTerminal::Cancelled)); + } + return; + } + + // Starting a run the client handler already abandoned would hold the + // container's slot for the full timeout with nobody left to read it. + if admit.send(Ok(())).is_err() { + 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}" + )))); + return; + } + + self.exec_in_flight.insert( + sandbox_id, + InFlightExec { container, - &config.script_code, - &env, - config.env_scope, - &config.working_directory, - config.timeout_ms, - cancellation, - Some(sink), - &mut self.logger, - ) + done, + _registration: registration, + parked: Vec::new(), + }, + ); + } + + /// Retire a finished exec: answer its client, then release the work that + /// parked behind it. + fn finish_exec( + &mut self, + sandbox_id: &str, + report: ExecReport, + worker: &mpsc::UnboundedSender, + ) { + for work in self.retire_exec(sandbox_id, report) { + self.dispatch(work, worker); } - .map_err(sr_err)?; + } - match outcome.completion { - container_steps::ProcessCompletion::TerminationUnconfirmed => { - let detail = outcome - .post_launch_error - .map(|error| error.error_message) - .unwrap_or_else(|| "process termination could not be confirmed".to_string()); - Err(self.quarantine(&config.sandbox_id, container, &detail)) + /// Answer a finished exec and hand back whatever parked behind it. + /// + /// Separate from [`Worker::finish_exec`] so teardown can answer the client + /// without releasing parked work into a daemon that is shutting down. + fn retire_exec(&mut self, sandbox_id: &str, report: ExecReport) -> Vec { + let Some(in_flight) = self.exec_in_flight.remove(sandbox_id) else { + return Vec::new(); + }; + + let outcome = match report { + ExecReport::Finished(outcome) => outcome, + ExecReport::Unconfirmed(detail) => { + Err(self.quarantine(sandbox_id, in_flight.container, &detail)) } - container_steps::ProcessCompletion::Confirmed(terminal) => Ok(terminal), + }; + + if let Err(orphaned) = in_flight.done.send(outcome) { + // The client handler is gone (e.g. its post-admission Ok write + // failed) but the run already happened. Record the result so a + // completed exec is never silently lost. + self.logger.log_line(&format!( + "exec on {sandbox_id} completed after the client disconnected; \ + orphaned result: {orphaned:?}" + )); } + + in_flight.parked } fn quarantine( @@ -960,11 +1301,60 @@ impl Worker { Ok(()) } + /// Give up the SDK handles rather than release them, when a run 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 { + std::mem::forget(std::mem::take(&mut self.containers)); + std::mem::forget(self.session.take()); + std::mem::forget(self.sdk.take()); + } + + // Only the handles are abandoned, so the reply channels close as the + // bookkeeping drops and a waiting client observes the worker is gone. + } + + /// Answer a command pulled off the queue during teardown. + /// + /// [`WorkerCommand::ExecFinished`] is the one the drain is looking for and + /// is handled by its caller. + fn refuse_while_shutting_down(&self, cmd: WorkerCommand) { + const SHUTTING_DOWN: &str = "the WSLc daemon is shutting down"; + + match cmd { + WorkerCommand::Provision { reply, .. } => { + let _ = reply.send(Err(WorkerError::Backend(anyhow::anyhow!( + "{SHUTTING_DOWN}" + )))); + } + WorkerCommand::Container(work) => work.refuse(SHUTTING_DOWN), + WorkerCommand::Retire { reply, .. } => { + let _ = reply.send(()); + } + WorkerCommand::ContainerCount { reply } => { + let _ = reply.send(self.live_container_count()); + } + WorkerCommand::Shutdown { reply } => { + let _ = reply.send(()); + } + + // The provision this would resume was already answered above. + WorkerCommand::PullFinished { .. } => {} + WorkerCommand::ExecFinished { .. } => {} + } + } + /// Release every container and the session. /// /// A parked provision is answered rather than dropped, since a dropped /// reply reaches its client as a bare "worker gone". - fn shutdown(&mut self) { + fn shutdown( + &mut self, + rx: &mut mpsc::UnboundedReceiver, + drain_budget: Duration, + ) { for (_, pending) in self.pending.drain() { let _ = pending.reply.send(Err(WorkerError::Backend(anyhow::anyhow!( "the WSLc daemon shut down while pulling image '{}'", @@ -972,6 +1362,58 @@ impl Worker { )))); } + // An off-worker exec is still using its container handle, so nothing + // below may stop, delete or release one until the runs report back. + let execs_drained = wait_for_execs_in_flight(&self.execs_in_flight, drain_budget); + + if execs_drained { + // Every drained run queued its report before releasing its count, so + // the reports are all in the channel now. Collecting them is what + // lets a finished exec keep its own exit code instead of the + // shutdown error below. Anything else found along the way is + // answered rather than dropped, since a dropped reply reaches its + // client as a bare "worker gone". + while !self.exec_in_flight.is_empty() { + let Ok(cmd) = rx.try_recv() else { + break; + }; + match cmd { + WorkerCommand::ExecFinished { sandbox_id, report } => { + for work in self.retire_exec(&sandbox_id, report) { + work.refuse(&format!( + "the WSLc daemon shut down before sandbox {sandbox_id} was free" + )); + } + } + other => self.refuse_while_shutting_down(other), + } + } + } + + for (sandbox_id, in_flight) in self.exec_in_flight.drain() { + let _ = in_flight + .done + .send(Err(WorkerError::Backend(anyhow::anyhow!( + "the WSLc daemon shut down while running an exec on sandbox {sandbox_id}" + )))); + for work in in_flight.parked { + work.refuse(&format!( + "the WSLc daemon shut down before sandbox {sandbox_id} was free" + )); + } + } + + if !execs_drained { + self.logger.log_line( + "a WSLC exec is still running; leaking the session rather than \ + releasing handles it is using", + ); + std::mem::forget(std::mem::take(&mut self.containers)); + std::mem::forget(self.session.take()); + std::mem::forget(self.sdk.take()); + return; + } + if let Some(sdk) = self.sdk.as_ref() { for (_, entry) in self.containers.drain() { // SAFETY: `sdk` is valid and `entry.container` is a live handle. @@ -1023,7 +1465,8 @@ pub fn spawn() -> Result { let (tx, mut rx) = mpsc::unbounded_channel::(); let active_execs = Arc::new(Mutex::new(HashMap::new())); - // A parked provision resumes by posting back into this queue. + // A parked provision or a finished exec resumes by posting back into this + // queue. let worker_tx = tx.clone(); std::thread::Builder::new() @@ -1039,77 +1482,44 @@ pub fn spawn() -> Result { }; let mut worker = Worker::new(); - while let Some(cmd) = rx.blocking_recv() { - match cmd { - WorkerCommand::Provision { config, reply } => { - worker.begin_provision(config, reply, &worker_tx); - } - WorkerCommand::PullFinished { token, outcome } => { - worker.finish_provision(token, outcome); - } - WorkerCommand::Start { config, reply } => { - let _ = reply.send(worker.start(config)); - } - WorkerCommand::Exec { - config, - sink, - cancellation, - registration: _registration, - admit, - done, - } => { - // Validate and run in one handler so admission is atomic - // with the start of the run: the worker never yields - // between the two, so no Stop/Deprovision can interleave. - match worker.validate_exec(&config.sandbox_id) { - Err(e) => { - let _ = admit.send(Err(e)); - } - Ok(_) if cancellation.load(Ordering::Acquire) => { - if admit.send(Ok(())).is_ok() { - let _ = done.send(Ok(ExecTerminal::Cancelled)); - } - } - // Only run if the admission receiver is still there: - // if the client handler was dropped before it read - // admission, the blocking exec would otherwise starve - // every other lifecycle command for its full timeout. - Ok(container) if admit.send(Ok(())).is_ok() => { - let sandbox_id = config.sandbox_id.clone(); - let outcome = worker.exec(config, container, sink, &cancellation); - if let Err(orphaned) = done.send(outcome) { - // The client handler is gone (e.g. its - // post-admission Ok write failed) but the run - // already happened. Record the result so a - // completed exec is never silently lost. - worker.logger.log_line(&format!( - "exec on {sandbox_id} completed after the client \ - disconnected; orphaned result: {orphaned:?}" - )); - } - } - Ok(_) => {} + + // A panic here would otherwise drop the container handles an exec + // thread is still using, so the unwind is caught and the handles + // are abandoned instead. + let served = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| { + while let Some(cmd) = rx.blocking_recv() { + match cmd { + WorkerCommand::Provision { config, reply } => { + worker.begin_provision(config, reply, &worker_tx); + } + WorkerCommand::PullFinished { token, outcome } => { + worker.finish_provision(token, outcome); + } + WorkerCommand::Container(work) => { + worker.dispatch(work, &worker_tx); + } + WorkerCommand::ExecFinished { sandbox_id, report } => { + worker.finish_exec(&sandbox_id, report, &worker_tx); + } + WorkerCommand::Retire { sandbox_id, reply } => { + worker.retire(&sandbox_id); + let _ = reply.send(()); + } + WorkerCommand::ContainerCount { reply } => { + let _ = reply.send(worker.live_container_count()); + } + WorkerCommand::Shutdown { reply } => { + worker.shutdown(&mut rx, EXEC_DRAIN_TIMEOUT); + let _ = reply.send(()); + break; } - } - WorkerCommand::Stop { config, reply } => { - let _ = reply.send(worker.stop(config)); - } - WorkerCommand::Deprovision { config, reply } => { - let _ = reply.send(worker.deprovision(config)); - } - WorkerCommand::Retire { sandbox_id, reply } => { - worker.retire(&sandbox_id); - let _ = reply.send(()); - } - WorkerCommand::ContainerCount { reply } => { - let _ = reply.send(worker.live_container_count()); - } - WorkerCommand::Shutdown { reply } => { - worker.shutdown(); - let _ = reply.send(()); - break; } } + })); + + if served.is_err() { + let in_flight = worker.execs_in_flight.load(Ordering::SeqCst); + worker.abandon_if_execs_running(in_flight); } }) .map_err(|e| anyhow::anyhow!("spawn WSLc worker thread: {e}"))?; @@ -1635,7 +2045,549 @@ mod tests { "expected the callback to be split, got {chunks}" ); } - // + + /// A started container, so [`Worker::validate_exec`] admits an exec on it. + fn started_entry() -> ContainerEntry { + ContainerEntry { + started: true, + ..test_entry(false) + } + } + + /// A container whose handle release is observable, so a test can tell a + /// deliberate leak from a normal drop. + fn flagged_entry( + release: unsafe extern "C" fn(mxc_sdk::wslc_common::wslc_bindings::WslcContainer) -> i32, + ) -> ContainerEntry { + let sentinel = std::ptr::dangling_mut(); + ContainerEntry { + started: false, + quarantined: false, + retired: false, + + // 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) }, + } + } + + /// An exec occupying a container's in-flight slot, with the completion + /// receiver a test uses to observe what the worker answers. + fn in_flight_exec( + exec_id: &str, + ) -> ( + InFlightExec, + oneshot::Receiver>, + ) { + let (done, done_rx) = oneshot::channel(); + let active_execs: ActiveExecs = Arc::new(Mutex::new(HashMap::new())); + let cancellation = Arc::new(AtomicBool::new(false)); + let registration = + Arc::new(register_exec(&active_execs, exec_id, "run", &cancellation).unwrap()); + ( + InFlightExec { + container: std::ptr::dangling_mut(), + done, + _registration: registration, + parked: Vec::new(), + }, + done_rx, + ) + } + + /// An exec command and the handles a test needs to observe it. + struct TestExec { + work: ContainerWork, + admit: oneshot::Receiver>, + done: oneshot::Receiver>, + cancellation: Arc, + } + + fn exec_work(sandbox_id: &str, exec_id: &str) -> TestExec { + let (admit, admit_rx) = oneshot::channel(); + let (done, done_rx) = oneshot::channel(); + let active_execs: ActiveExecs = Arc::new(Mutex::new(HashMap::new())); + let cancellation = Arc::new(AtomicBool::new(false)); + let registration = + Arc::new(register_exec(&active_execs, exec_id, "run", &cancellation).unwrap()); + let sink: OutputSink = Box::new(|_, _| {}); + TestExec { + work: ContainerWork::Exec(ExecRequest { + config: ExecConfig { + exec_id: exec_id.to_string(), + run_token: "run".to_string(), + sandbox_id: sandbox_id.to_string(), + script_code: "echo hi".to_string(), + working_directory: String::new(), + env: Vec::new(), + env_scope: EnvScope::Merge, + timeout_ms: 0, + }, + sink, + cancellation: Arc::clone(&cancellation), + registration, + admit, + done, + }), + admit: admit_rx, + done: done_rx, + cancellation, + } + } + + /// A worker holding one started container with an exec already in flight. + fn worker_with_exec_in_flight() -> (Worker, oneshot::Receiver>) + { + let mut worker = Worker::new(); + worker + .containers + .insert("wslc:busy".to_string(), started_entry()); + let (in_flight, done) = in_flight_exec("running"); + worker + .exec_in_flight + .insert("wslc:busy".to_string(), in_flight); + (worker, done) + } + + fn stop_work(sandbox_id: &str) -> (ContainerWork, oneshot::Receiver>) { + let (reply, reply_rx) = oneshot::channel(); + ( + ContainerWork::Stop { + config: StopConfig { + sandbox_id: sandbox_id.to_string(), + }, + reply, + }, + reply_rx, + ) + } + + fn deprovision_work( + sandbox_id: &str, + ) -> (ContainerWork, oneshot::Receiver>) { + let (reply, reply_rx) = oneshot::channel(); + ( + ContainerWork::Deprovision { + config: DeprovisionConfig { + sandbox_id: sandbox_id.to_string(), + }, + reply, + }, + reply_rx, + ) + } + + /// Deleting the container would free the handle the run thread is holding. + #[test] + fn deprovision_parks_behind_an_exec_using_the_same_container() { + let (tx, _rx) = mpsc::unbounded_channel(); + let (mut worker, _done) = worker_with_exec_in_flight(); + let (work, mut reply) = deprovision_work("wslc:busy"); + + worker.dispatch(work, &tx); + + assert!( + reply.try_recv().is_err(), + "a deprovision must not run while an exec holds the container" + ); + assert!( + worker.containers.contains_key("wslc:busy"), + "the container handle must survive until the run reports back" + ); + assert_eq!(worker.exec_in_flight["wslc:busy"].parked.len(), 1); + } + + #[test] + fn an_exec_does_not_park_work_for_another_container() { + let (tx, _rx) = mpsc::unbounded_channel(); + let (mut worker, _done) = worker_with_exec_in_flight(); + worker + .containers + .insert("wslc:idle".to_string(), test_entry(false)); + let (work, mut reply) = deprovision_work("wslc:idle"); + + worker.dispatch(work, &tx); + + assert!( + reply.try_recv().expect("idle container answered").is_ok(), + "a command for an idle container must not wait on another container's exec" + ); + assert!(!worker.containers.contains_key("wslc:idle")); + } + + #[test] + fn finishing_an_exec_releases_the_work_parked_behind_it() { + let (tx, _rx) = mpsc::unbounded_channel(); + let (mut worker, mut done) = worker_with_exec_in_flight(); + let (work, mut reply) = deprovision_work("wslc:busy"); + worker.dispatch(work, &tx); + + worker.finish_exec( + "wslc:busy", + ExecReport::Finished(Ok(ExecTerminal::Exited(0))), + &tx, + ); + + assert_eq!(done.try_recv().unwrap().unwrap(), ExecTerminal::Exited(0)); + assert!(reply + .try_recv() + .expect("the parked deprovision must be answered") + .is_ok()); + assert!(!worker.containers.contains_key("wslc:busy")); + } + + #[test] + fn a_second_exec_on_the_same_container_waits_for_the_first() { + 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); + + assert!( + second.admit.try_recv().is_err(), + "a second exec must not be admitted while the first holds the container" + ); + } + + /// 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() { + let (tx, _rx) = mpsc::unbounded_channel(); + let (mut worker, _done) = worker_with_exec_in_flight(); + let mut queued = exec_work("wslc:busy", "queued"); + worker.dispatch(queued.work, &tx); + + queued.cancellation.store(true, Ordering::Release); + worker.finish_exec( + "wslc:busy", + ExecReport::Finished(Ok(ExecTerminal::Exited(0))), + &tx, + ); + + assert!(queued.admit.try_recv().unwrap().is_ok()); + assert_eq!( + queued.done.try_recv().unwrap().unwrap(), + ExecTerminal::Cancelled + ); + } + + #[test] + fn an_unconfirmed_exec_is_quarantined_on_the_worker() { + let (tx, _rx) = mpsc::unbounded_channel(); + let (mut worker, mut done) = worker_with_exec_in_flight(); + + worker.finish_exec( + "wslc:busy", + ExecReport::Unconfirmed("the exit callback never fired".to_string()), + &tx, + ); + + let err = done.try_recv().unwrap().unwrap_err(); + assert!(err.to_string().contains("quarantined"), "got {err}"); + assert!(worker.containers["wslc:busy"].quarantined); + } + + /// A report can outlive its slot: shutdown drains the run and clears the + /// map before the posted command is serviced. + #[test] + fn a_report_for_a_container_with_no_exec_in_flight_is_a_no_op() { + let (tx, _rx) = mpsc::unbounded_channel(); + let mut worker = Worker::new(); + + worker.finish_exec( + "wslc:never-ran", + ExecReport::Finished(Ok(ExecTerminal::Exited(0))), + &tx, + ); + + assert_eq!(worker.live_container_count(), 0); + } + + /// A run that has not reported back is still holding its container handle, + /// so a teardown that gives up waiting must abandon the handles rather than + /// release them. + #[test] + fn a_timed_out_shutdown_answers_clients_without_releasing_live_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 (tx, mut rx) = mpsc::unbounded_channel(); + let mut worker = Worker::new(); + worker + .containers + .insert("wslc:busy".to_string(), flagged_entry(release)); + let (in_flight, mut done) = in_flight_exec("running"); + worker + .exec_in_flight + .insert("wslc:busy".to_string(), in_flight); + let (work, mut parked_reply) = stop_work("wslc:busy"); + worker.dispatch(work, &tx); + + // The run never reports, so the drain can only expire. + let _counted = ExecCount::enter(&worker.execs_in_flight); + + worker.shutdown(&mut rx, Duration::from_millis(100)); + + assert!( + !RELEASED.load(Ordering::SeqCst), + "a container handle the run thread still holds must not be released" + ); + assert!( + done.try_recv() + .expect("the running exec must be answered") + .is_err(), + "the client gets a typed error rather than a dropped channel" + ); + assert!(parked_reply + .try_recv() + .expect("parked work must be answered") + .is_err()); + } + + /// A command that reaches the queue ahead of the exec report must not cost + /// the exec its result, nor be swallowed on its way past. + #[test] + fn shutdown_answers_a_command_queued_ahead_of_the_exec_report() { + let (tx, mut rx) = mpsc::unbounded_channel(); + let (mut worker, mut done) = worker_with_exec_in_flight(); + + let (queued, mut queued_reply) = stop_work("wslc:other"); + tx.send(WorkerCommand::Container(queued)).unwrap(); + tx.send(WorkerCommand::ExecFinished { + sandbox_id: "wslc:busy".to_string(), + report: ExecReport::Finished(Ok(ExecTerminal::Exited(0))), + }) + .unwrap(); + + worker.shutdown(&mut rx, EXEC_DRAIN_TIMEOUT); + + assert_eq!( + done.try_recv() + .expect("the finished exec must be answered") + .unwrap(), + ExecTerminal::Exited(0), + "a command ahead of the report must not cost the exec its exit code" + ); + assert!( + queued_reply + .try_recv() + .expect("the queued command must be answered, not dropped") + .is_err(), + "a command consumed during teardown gets a typed error" + ); + } + + /// A run that finished before teardown keeps its own exit code: the drain + /// found its report already queued. + #[test] + fn shutdown_delivers_a_finished_execs_real_outcome() { + let (tx, mut rx) = mpsc::unbounded_channel(); + let (mut worker, mut done) = worker_with_exec_in_flight(); + let (work, mut reply) = stop_work("wslc:busy"); + worker.dispatch(work, &tx); + tx.send(WorkerCommand::ExecFinished { + sandbox_id: "wslc:busy".to_string(), + report: ExecReport::Finished(Ok(ExecTerminal::Exited(0))), + }) + .unwrap(); + + worker.shutdown(&mut rx, EXEC_DRAIN_TIMEOUT); + + assert_eq!( + done.try_recv() + .expect("the finished exec must be answered") + .unwrap(), + ExecTerminal::Exited(0), + "a run that completed before teardown must keep its exit code" + ); + assert!( + reply + .try_recv() + .expect("the parked stop must be answered rather than dropped") + .is_err(), + "parked work is refused once the daemon is shutting down" + ); + } + + /// The synthetic error is for a run whose report never arrived. + #[test] + fn shutdown_answers_an_exec_whose_report_never_arrived() { + let (tx, mut rx) = mpsc::unbounded_channel(); + let (mut worker, mut done) = worker_with_exec_in_flight(); + let (work, mut reply) = stop_work("wslc:busy"); + worker.dispatch(work, &tx); + + worker.shutdown(&mut rx, EXEC_DRAIN_TIMEOUT); + + assert!(done + .try_recv() + .expect("the running exec must be answered") + .is_err()); + assert!(reply + .try_recv() + .expect("the parked stop must be answered rather than dropped") + .is_err()); + } + + /// A thread that unwinds never reaches its own release statement, so the + /// count has to come off in `Drop` or teardown waits out the full budget. + #[test] + fn a_panicking_run_still_releases_its_in_flight_count() { + let count = Arc::new(AtomicUsize::new(0)); + + let panicked = std::panic::catch_unwind({ + let count = Arc::clone(&count); + move || { + let _counted = ExecCount::enter(&count); + assert_eq!(count.load(Ordering::SeqCst), 1); + panic!("the run exploded"); + } + }); + + assert!(panicked.is_err()); + assert_eq!( + count.load(Ordering::SeqCst), + 0, + "a panicking run must not leave itself counted forever" + ); + } + + /// A client awaiting a reply learns the worker is gone from its channel + /// closing, so an unwind must not take the reply channels with it. + #[test] + fn an_unwind_closes_the_reply_channels_it_abandons() { + let (mut worker, mut done) = worker_with_exec_in_flight(); + let (work, mut parked_reply) = stop_work("wslc:busy"); + worker + .exec_in_flight + .get_mut("wslc:busy") + .unwrap() + .parked + .push(work); + let (reply, mut pending_reply) = oneshot::channel(); + worker.pending.insert( + 1, + PendingProvision { + config: ProvisionConfig { + image: "alpine:latest".to_string(), + image_tar_path: None, + volumes: Vec::new(), + network: Default::default(), + port_mappings: Vec::new(), + }, + reply, + _retire_deadline: std::sync::mpsc::channel().0, + }, + ); + + worker.abandon_if_execs_running(1); + + assert!( + matches!(done.try_recv(), Err(oneshot::error::TryRecvError::Closed)), + "the running exec's client must see its channel close" + ); + assert!( + matches!( + parked_reply.try_recv(), + Err(oneshot::error::TryRecvError::Closed) + ), + "parked work's client must see its channel close" + ); + assert!( + matches!( + pending_reply.try_recv(), + Err(oneshot::error::TryRecvError::Closed) + ), + "a parked provision's client must see its channel close" + ); + } + + /// Releasing a container handle is an SDK call; a run thread still holding + /// one must outlive the worker that owned it. + #[test] + fn an_unwind_with_a_run_still_counted_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:busy".to_string(), flagged_entry(release)); + + worker.abandon_if_execs_running(1); + + assert!( + !RELEASED.load(Ordering::SeqCst), + "a handle the run thread is still using must not be released" + ); + } + + #[test] + fn an_unwind_with_no_run_counted_releases_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:idle".to_string(), flagged_entry(release)); + + worker.abandon_if_execs_running(0); + + assert!( + RELEASED.load(Ordering::SeqCst), + "with nothing in flight the handles must be released normally" + ); + } + + /// Teardown asks whether *its own* handles are free, so one worker's run + /// must not hold another's drain open. + #[test] + fn one_workers_run_does_not_count_against_anothers_drain() { + let busy = Worker::new(); + let idle = Worker::new(); + + let _counted = ExecCount::enter(&busy.execs_in_flight); + + assert_eq!(busy.execs_in_flight.load(Ordering::SeqCst), 1); + assert_eq!( + idle.execs_in_flight.load(Ordering::SeqCst), + 0, + "a run on one worker must not be counted against another's teardown" + ); + } + + /// Teardown must not free a container handle a run thread is still using. + #[test] + fn a_counted_exec_keeps_the_drain_from_reporting_clear() { + let in_flight = AtomicUsize::new(1); + + assert!(!wait_for_execs_in_flight( + &in_flight, + Duration::from_millis(100) + )); + + in_flight.store(0, Ordering::SeqCst); + assert!(wait_for_execs_in_flight( + &in_flight, + Duration::from_millis(100) + )); + } + // 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 @@ -1817,4 +2769,69 @@ mod tests { .unwrap(); handle.shutdown().await.unwrap(); } + + /// A deprovision that reaches the worker mid-run would delete the container + /// and free the handle the run thread is using. + #[tokio::test] + #[ignore = "requires a WSL2 host with alpine:latest already in the daemon session cache"] + async fn deprovision_during_an_exec_waits_for_the_run() { + 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 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, + }) + .await + .unwrap(); + + // Admission has returned, so the run thread already holds the handle. + let deprovision = tokio::spawn({ + let handle = handle.clone(); + let sandbox_id = id.clone(); + async move { handle.deprovision(DeprovisionConfig { sandbox_id }).await } + }); + + tokio::time::sleep(Duration::from_millis(500)).await; + assert!( + !deprovision.is_finished(), + "deprovision must wait for the run rather than delete the container under it" + ); + + let mut stdout = Vec::new(); + while let Some((kind, data)) = exec.output.recv().await { + if kind == OutStream::Stdout { + stdout.extend_from_slice(&data); + } + } + assert_eq!(exec.done.await.unwrap().unwrap(), ExecTerminal::Exited(0)); + assert_eq!(String::from_utf8_lossy(&stdout).trim(), "survived"); + + deprovision.await.unwrap().unwrap(); + assert_eq!(count(&handle).await, 0); + + handle.shutdown().await.unwrap(); + } }