From de70147c4584df1a568b9b750774df8cdf4752ec Mon Sep 17 00:00:00 2001 From: Melih Elibol Date: Tue, 22 Sep 2026 10:51:30 -0700 Subject: [PATCH] cuda-async: release abandoned in-flight futures through the reactor instead of blocking Dropping a DeviceFuture whose work is still in flight used to synchronize the stream on the dropping thread, which under select!/timeout is an executor thread. The result is a handle by the DeviceOp::execute contract (every device-visible resource is retained by the submission), so it is released immediately; the execution context is parked with the completion reactor behind a flag write enqueued after the abandoned work and released on a dedicated reaper thread when the flag lands. The reactor only hands off, never releases. If the context cannot be parked it is dropped inline, which is the previous blocking wait; faulted and capturing streams still leak the owners with a report (quiet when the fault was already delivered). reaper::parked() and reaper::reaped_total() expose the counts. Signed-off-by: Melih Elibol --- cuda-async/src/device_future.rs | 218 ++++++++--------------------- cuda-async/src/device_operation.rs | 6 + cuda-async/src/lib.rs | 2 + cuda-async/src/reactor.rs | 88 ++++++++++-- cuda-async/src/reaper.rs | 153 ++++++++++++++++++++ cuda-async/src/submission.rs | 25 +++- cuda-async/tests/drop_in_flight.rs | 136 ++++++++++++++---- cutile-rs/CHANGELOG.md | 20 +++ 8 files changed, 445 insertions(+), 203 deletions(-) create mode 100644 cuda-async/src/reaper.rs diff --git a/cuda-async/src/device_future.rs b/cuda-async/src/device_future.rs index 28b408c0c4..814ae70d39 100644 --- a/cuda-async/src/device_future.rs +++ b/cuda-async/src/device_future.rs @@ -18,23 +18,27 @@ //! //! Dropping the future never cancels submitted GPU work; kernels run to //! completion regardless. What the drop decides is *when the host releases -//! the resources that work still uses*. A future dropped while its work is -//! in flight therefore waits for the stream to drain before dropping its -//! undelivered result (the owned output — buffers, DMA targets — plus the -//! execution context's stream and pool handles). If the wait cannot be -//! performed (a faulted context, a stream mid-capture) the result is leaked: -//! releasing memory the device may still write to is the worse failure. The -//! leak is reported on stderr unless the future already resolved with the -//! stream's fault — then the caller has the error and the leak is its -//! documented consequence. See [`DeviceFuture`]'s type docs for why the -//! wait is synchronous. +//! the resources that work still uses*. Those resources live in the +//! execution context's submission (storage leases and other owners the +//! operation retained, see [`ExecutionContext::retain`]); the result itself +//! is a handle by the [`DeviceOp::execute`] contract and is released on the +//! spot. The context is parked with the completion reactor +//! ([`crate::reaper`]) and released on a dedicated thread once a flag write +//! enqueued behind the work lands, so the drop returns immediately and never +//! stalls an executor thread. If the context cannot be parked it is released +//! inline, which waits for the stream; if the stream cannot be waited on (a +//! faulted context, a stream mid-capture) the owners are leaked: releasing +//! memory the device may still write to is the worse failure. The leak is +//! reported on stderr unless the future already resolved with the stream's +//! fault — then the caller has the error and the leak is its documented +//! consequence. use crate::device_operation::{DeviceOp, ExecutionContext}; use crate::error::DeviceError; use cuda_core::{DriverError, Stream}; use futures::task::AtomicWaker; use std::future::Future; -use std::mem::{self, MaybeUninit}; +use std::mem::MaybeUninit; use std::pin::Pin; use std::sync::atomic::{AtomicBool, Ordering}; use std::sync::Arc; @@ -138,31 +142,23 @@ pub(crate) fn probe_stream(stream: &Stream) -> StreamHealth { /// /// Dropping a `DeviceFuture` after its first poll — after `execute` has /// enqueued GPU work — but before it resolved leaves an *undelivered -/// result*: the operation's output, which owns the buffers the GPU is still -/// writing (a tensor, a `Vec` DMA target), or borrows the caller's. The -/// drop waits for the stream to drain and only then drops that result. On a -/// wait failure the result is leaked — loudly, unless the future already -/// delivered the stream's fault to its caller. +/// submission*. The operation's output is released immediately: by the +/// [`DeviceOp::execute`] contract it holds no resource the device still +/// uses, those being retained by the execution context's submission. The +/// context is handed to the reaper (see [`crate::reaper`]), which drops it +/// once the stream has passed the abandoned work, releasing the storage +/// leases and other owners. Dropping is therefore non-blocking; only a +/// context the reactor cannot take falls back to waiting inline. /// -/// cuTile registers owned storage and device access leases with the execution -/// context before enqueueing work. These survive argument recovery, output -/// projection, partial submission failures, and `mem::forget`. Forgetting the -/// future leaks both storage and leases; conflicting cross-stream use remains -/// rejected even after the original Rust borrow ends. Successful completion -/// releases the registered owners before returning the result. -/// -/// The wait is synchronous by necessity, not preference. The alternative — -/// parking the result behind a CUDA event and dropping it later, once -/// `cuEventQuery` passes (the design used by `simt::device_future`) — needs -/// the result to be `'static`, and `DeviceOp::Output` is not: cutile -/// launchers return borrowed inputs (`&'a Tensor`, -/// `Partition<&'a mut Tensor>`) as part of their output. For those the -/// result itself cannot be handed to a background thread. Registered storage -/// owners are independent of these borrows, but arbitrary custom outputs may -/// still require waiting before destruction. Rust cannot specialize on `'static`, so -/// the same rule applies to every output type until `Output: 'static` is a -/// trait-level requirement; at that point the event-gated limbo becomes the -/// default and blocking the last-resort fallback. +/// Outputs may borrow the caller's buffers (`&'a Tensor`, +/// `Partition<&'a mut Tensor>`). That is sound because the borrow is not +/// what keeps the device memory alive: the submission's lease on the +/// storage is, and it is released only after the stream has drained. What +/// the early end of the borrow does allow is a *host* access to the buffer +/// through an unchecked path (a raw device pointer handed to another +/// library) while the device may still be writing; stream-ordered accesses +/// through this crate remain ordered or are rejected by the cross-stream +/// conflict check. #[derive(Debug)] pub struct DeviceFuture> { pub(crate) device_operation: Option, @@ -311,72 +307,31 @@ impl> DeviceFuture { ) && self.result.is_some() } - /// Waits for the submitted work, then drops the undelivered result. + /// Releases an undelivered submission without waiting for the device. /// - /// A stream that is idle costs one query; a busy one is synchronized. A - /// capturing stream cannot be waited on (querying it would invalidate - /// the capture) and a faulted one cannot prove completion; both fall - /// through to the loud leak in - /// [`release_in_flight_result_with`](Self::release_in_flight_result_with). + /// The result is a handle (see the type docs) and drops here. The + /// execution context, whose submission owns everything the device may + /// still touch, is parked with the reaper and released once the stream + /// has passed the abandoned work; when it cannot be parked it is dropped + /// inline, which waits for the stream, or leaks the owners if the stream + /// cannot be waited on. fn release_in_flight_result(&mut self) { if !self.has_undelivered_submission() { return; } - let stream = self - .execution_context - .as_ref() - .map(|ctx| Arc::clone(ctx.get_cuda_stream())); - self.release_in_flight_result_with(move || { - let stream = stream.ok_or_else(|| { - DeviceError::Internal( - "Cannot release an in-flight future without an execution context.".to_string(), - ) - })?; - // The drop may run on an executor thread that never touched - // CUDA; the query and synchronize below need a current context. - stream.device().bind_to_thread()?; - match probe_stream(&stream) { - StreamHealth::Idle => Ok(()), - // SAFETY: the context was bound above; the stream is valid. - StreamHealth::Busy => unsafe { stream.synchronize() }.map_err(DeviceError::Driver), - StreamHealth::Capturing => Err(DeviceError::Internal( - "the future's stream is recording a graph; it cannot be synchronized".into(), - )), - StreamHealth::Faulted(e) => Err(DeviceError::Driver(e)), - } - }); - } - - /// Runs `wait` and drops the stored result on success; on failure the - /// result is leaked loudly, because dropping resources the device may - /// still use is worse than leaking them. - fn release_in_flight_result_with(&mut self, wait: F) - where - F: FnOnce() -> Result<(), DeviceError>, - { - if !self.has_undelivered_submission() { - return; - } - let Some(result) = self.result.take() else { + drop(self.result.take()); + let Some(ctx) = self.execution_context.take() else { return; }; - if let Err(error) = wait() { - if self.fault_delivered { - // The caller already received this stream's fault from - // `poll`; the leak is its documented consequence. A test - // (or program) that handled the error correctly should not - // see an error report for it. - mem::forget(result); - return; - } - crate::leak::report_leak(format_args!( - "cuda-async: leaking the result of a dropped in-flight future; the driver \ - could not prove its GPU work finished: {error}" - )); - mem::forget(result); - return; + if self.fault_delivered { + // The caller already received this stream's fault from `poll`; + // the leak on release is its documented consequence, not news. + ctx.mark_fault_delivered(); } - drop(result); + #[cfg(not(loom))] + crate::reaper::park_or_wait(ctx); + #[cfg(loom)] + drop(ctx); } } @@ -615,70 +570,21 @@ mod release_tests { } } - /// An `Executing` future waits before its result drops. - #[test] - fn release_waits_before_dropping_the_result() { - let events = Arc::new(Mutex::new(Vec::new())); - let tracker = DropTracker { - events: Arc::clone(&events), - }; - let mut future = future_in_state(DeviceFutureState::Executing, Some(tracker)); - - future.release_in_flight_result_with(|| { - events.lock().unwrap().push("wait"); - Ok(()) - }); - - assert_eq!(events.lock().unwrap().as_slice(), ["wait", "drop"]); - assert!(future.result.is_none()); - assert!(!future.has_undelivered_submission()); - } - - /// If the wait fails the result is leaked, never dropped early. + /// Releasing an `Executing` future drops the result handle on the spot: + /// nothing the device uses lives in the result. #[test] - fn release_leaks_when_the_wait_fails() { + fn release_drops_the_result_handle_immediately() { let drops = Arc::new(AtomicUsize::new(0)); let mut future = future_in_state( DeviceFutureState::Executing, Some(CountDrop(Arc::clone(&drops))), ); - let mut capture = crate::leak::capture::start(); - future.release_in_flight_result_with(|| Err(DeviceError::Internal("boom".to_string()))); - let reports = capture.take(); - - assert_eq!(drops.load(Ordering::Relaxed), 0); - assert!(future.result.is_none()); - assert_eq!(reports.len(), 1, "the leak must be reported: {reports:?}"); - assert!(reports[0].contains("dropped in-flight future")); - } - - /// A future that already delivered the stream's fault to its caller - /// leaks its stored result *quietly*: the caller has the error, the - /// leak is its documented consequence. - #[test] - fn release_leaks_quietly_after_the_fault_was_delivered() { - let drops = Arc::new(AtomicUsize::new(0)); - let mut future = future_in_state( - DeviceFutureState::Complete, - Some(CountDrop(Arc::clone(&drops))), - ); - future.fault_delivered = true; - - let mut capture = crate::leak::capture::start(); - future.release_in_flight_result_with(|| Err(DeviceError::Internal("boom".to_string()))); - let reports = capture.take(); + future.release_in_flight_result(); - assert_eq!( - drops.load(Ordering::Relaxed), - 0, - "must still leak, not drop" - ); + assert_eq!(drops.load(Ordering::Relaxed), 1); assert!(future.result.is_none()); - assert!( - reports.is_empty(), - "a delivered fault must not be re-reported: {reports:?}" - ); + assert!(!future.has_undelivered_submission()); } /// The registration-failure shape: `execute` succeeded (work submitted, @@ -694,26 +600,21 @@ mod release_tests { let mut future = future_in_state(DeviceFutureState::Complete, Some(tracker)); assert!(future.has_undelivered_submission()); - future.release_in_flight_result_with(|| { - events.lock().unwrap().push("wait"); - Ok(()) - }); + future.release_in_flight_result(); - assert_eq!(events.lock().unwrap().as_slice(), ["wait", "drop"]); + assert_eq!(events.lock().unwrap().as_slice(), ["drop"]); assert!(!future.has_undelivered_submission()); } - /// A delivered result leaves nothing to release: no wait happens. + /// A delivered result leaves nothing to release. #[test] fn release_is_noop_after_result_delivery() { let mut future: DeviceFuture> = future_in_state(DeviceFutureState::Complete, None); assert!(!future.has_undelivered_submission()); - future.release_in_flight_result_with(|| { - panic!("delivered futures must not wait during release") - }); future.release_in_flight_result(); + assert!(!future.has_undelivered_submission()); } /// An `Idle` future never submitted work: no wait, and dropping it drops @@ -724,8 +625,7 @@ mod release_tests { let mut future = future_in_state(DeviceFutureState::Idle, Some(CountDrop(Arc::clone(&drops)))); - future - .release_in_flight_result_with(|| panic!("idle futures must not wait during release")); + future.release_in_flight_result(); assert_eq!(drops.load(Ordering::Relaxed), 0); assert!(future.result.is_some()); diff --git a/cuda-async/src/device_operation.rs b/cuda-async/src/device_operation.rs index 9c0c51d173..5f28de3043 100644 --- a/cuda-async/src/device_operation.rs +++ b/cuda-async/src/device_operation.rs @@ -186,6 +186,12 @@ impl ExecutionContext { pub(crate) unsafe fn complete(&self) { self.submission.complete(); } + + /// The owning future delivered this stream's fault to its caller; a + /// leak on release is then not reported again. + pub(crate) fn mark_fault_delivered(&self) { + self.submission.mark_fault_delivered(); + } pub fn device(&self) -> &Arc { &self.device } diff --git a/cuda-async/src/lib.rs b/cuda-async/src/lib.rs index 8f61f510a6..5cc953d81d 100644 --- a/cuda-async/src/lib.rs +++ b/cuda-async/src/lib.rs @@ -22,6 +22,8 @@ pub mod prelude; // model-checked through `slot_table`'s mock backend instead. #[cfg(not(loom))] mod reactor; +#[cfg(not(loom))] +pub mod reaper; pub mod scheduling_policies; /// SIMT-model async surface, copied from cuda-oxide for the shared /// host-crate migration. Not re-exported at the root; see the module docs. diff --git a/cuda-async/src/reactor.rs b/cuda-async/src/reactor.rs index 3a29ab208b..7fd36de9e8 100644 --- a/cuda-async/src/reactor.rs +++ b/cuda-async/src/reactor.rs @@ -72,13 +72,42 @@ impl FlagArray for CudaFlags { } } -/// What a slot carries: the waker to fire, and the stream whose flag write +/// What a landed (or retired) slot triggers. +enum Payload { + /// Wake the future registered for this completion. + Wake(Arc), + /// Release the resources of a dropped in-flight future (see + /// [`crate::reaper`]). Never released on the scanner thread. + Reap(crate::reaper::Parked), +} + +/// What a slot carries: the reaction to fire, and the stream whose flag write /// completes the slot (kept alive, and probed if the slot goes stale). struct Registration { - waker_state: Arc, + payload: Payload, stream: Arc, } +impl Registration { + /// Fires the slot's reaction. `faulted` slots were retired by the stale + /// probe instead of their flag: a future is woken without being marked + /// complete so its poll observes the driver error; a parked context is + /// released the same way as a landed one, since its submission probes + /// the stream itself and leaks rather than frees on a fault. + fn fire(self, faulted: bool) { + match self.payload { + Payload::Wake(waker_state) => { + if faulted { + waker_state.wake(); + } else { + waker_state.signal(); + } + } + Payload::Reap(parked) => crate::reaper::release(parked), + } + } +} + struct Reactor { table: SlotTable, /// Device-side alias of the flag slab (CU_MEMHOSTALLOC_DEVICEMAP). @@ -143,10 +172,11 @@ fn scan_loop() { loop { reactor.table.scan_once(&mut woken); if !woken.is_empty() { - // Wakers fire outside any lock the scan held, so a registration - // is never blocked behind a waking phase. + // Reactions fire outside any lock the scan held, so a registration + // is never blocked behind a waking phase. A parked context is + // handed to the reaper thread here, never released in this loop. for reg in woken.drain(..) { - reg.waker_state.signal(); + reg.fire(false); } idle_passes = 0; continue; @@ -179,11 +209,10 @@ fn scan_loop() { last_probe = Instant::now(); probe_stale_slots(&reactor.table, &mut woken, &mut faulted); for reg in woken.drain(..) { - reg.waker_state.signal(); + reg.fire(false); } for reg in faulted.drain(..) { - // Wake without completing: the poll observes the fault. - reg.waker_state.wake(); + reg.fire(true); } } thread::yield_now(); @@ -233,24 +262,53 @@ pub(crate) unsafe fn register( stream: &Arc, waker_state: Arc, ) -> Result<(), DeviceError> { - let reactor = reactor()?; - let slot = reactor - .table - .claim() - .ok_or_else(|| internal("reactor slot pool exhausted".into()))?; + arm(stream, Payload::Wake(waker_state)).map_err(|(error, _)| error) +} + +/// Parks a dropped in-flight future's context until the work submitted +/// before this call on `stream` has completed, then releases it on the +/// reaper thread. On failure the payload is handed back so the caller can +/// fall back to releasing it inline. +/// +/// # Safety +/// As for [`register`]. +pub(crate) unsafe fn park( + stream: &Arc, + parked: crate::reaper::Parked, +) -> Result<(), crate::reaper::Parked> { + arm(stream, Payload::Reap(parked)).map_err(|(_, payload)| match payload { + Payload::Reap(parked) => parked, + Payload::Wake(_) => unreachable!("park arms a Reap payload"), + }) +} + +/// Claims a slot, enqueues its flag write behind the work already on +/// `stream`, and publishes `payload` to the scanner. Returns the payload +/// with the error when the slot cannot be armed. +unsafe fn arm(stream: &Arc, payload: Payload) -> Result<(), (DeviceError, Payload)> { + let reactor = match reactor() { + Ok(reactor) => reactor, + Err(error) => return Err((error, payload)), + }; + let Some(slot) = reactor.table.claim() else { + return Err((internal("reactor slot pool exhausted".into()), payload)); + }; reactor.table.reset_flag(slot); let addr = reactor.dptr + (slot * std::mem::size_of::()) as u64; let code = cuda_bindings::cuStreamWriteValue32_v2(stream.cu_stream(), addr, 1, 0); if code != cuda_bindings::cudaError_enum_CUDA_SUCCESS { reactor.table.release(slot); - return Err(internal(format!("cuStreamWriteValue32 failed: {code}"))); + return Err(( + internal(format!("cuStreamWriteValue32 failed: {code}")), + payload, + )); } // empty→wake: unpark only when this registration transitioned the reactor // from idle to active (the scanner may be parked). At higher registration // rates the scanner is already awake and the skipped unparks avoid // cross-core `Parker` contention (+38% throughput in the A/B). let registration = Registration { - waker_state, + payload, stream: Arc::clone(stream), }; if reactor.table.publish(slot, registration) { diff --git a/cuda-async/src/reaper.rs b/cuda-async/src/reaper.rs new file mode 100644 index 0000000000..3566bf0e7a --- /dev/null +++ b/cuda-async/src/reaper.rs @@ -0,0 +1,153 @@ +/* + * SPDX-FileCopyrightText: Copyright (c) 2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved. + * SPDX-License-Identifier: Apache-2.0 + */ + +//! Non-blocking release of what an abandoned in-flight future still owns. +//! +//! Dropping a [`DeviceFuture`] cannot cancel GPU work that has already been +//! submitted. What the drop decides is *when the host releases the resources +//! that work still uses*: the tensor storage and other owners retained by the +//! submission (see [`ExecutionContext::retain`]). Releasing them early would +//! free memory the device may still write; waiting for the stream inside +//! `Drop` stalls whichever thread runs the drop, which under `select!` or a +//! timeout is an executor thread. +//! +//! The reaper removes the wait without weakening the guarantee. On drop, the +//! future's [`ExecutionContext`] is handed to the completion reactor together +//! with a flag write enqueued *behind* the abandoned work on the same stream. +//! When the flag lands, the reactor passes the context to a dedicated release +//! thread, whose drop of the context releases the submission's owners. The +//! dropping thread returns immediately. +//! +//! ```text +//! drop(future) reactor (flag lands) reaper thread +//! release result handle hand off context -> drop context: +//! arm flag write on the stream owners released +//! park context, return +//! ``` +//! +//! Only the crate's own drops run on the reactor and reaper threads: the +//! parked payload is an [`ExecutionContext`], whose owners are storage +//! leases, stream/pool handles, and whatever an operation registered with +//! [`ExecutionContext::retain`]. The future's *result* is not parked. By the +//! [`DeviceOp::execute`] contract every device-visible resource is retained +//! in the submission, so a result is a handle that may be released at any +//! time, on the dropping thread, with the caller's own `Drop` semantics. +//! +//! # Fallbacks +//! +//! Parking is best effort and never trades safety for latency. If the stream +//! is already idle the context is simply dropped. If the reactor cannot take +//! the payload (slot pool exhausted, stream mem-ops unavailable, reactor +//! failed to start), the context is dropped inline, which is the previous +//! blocking wait. A stream mid-capture must not receive a flag write, and a +//! faulted stream can never land one; both drop inline too, where the +//! submission's own release reports and leaks the owners instead of freeing +//! them. +//! +//! [`DeviceFuture`]: crate::device_future::DeviceFuture +//! [`DeviceOp::execute`]: crate::device_operation::DeviceOp::execute +//! [`ExecutionContext::retain`]: crate::device_operation::ExecutionContext::retain + +use crate::device_future::{probe_stream, StreamHealth}; +use crate::device_operation::ExecutionContext; +use std::sync::atomic::{AtomicU64, AtomicUsize, Ordering}; +use std::sync::mpsc::{self, Sender}; +use std::sync::{Arc, OnceLock}; +use std::thread; + +/// An abandoned submission: dropping it releases the submission's owners. +pub(crate) struct Parked { + /// Held only for its `Drop`: releasing the submission's owners. + _ctx: ExecutionContext, +} + +impl Drop for Parked { + fn drop(&mut self) { + // `_ctx` drops after this body: the submission waits (or, on a stream + // it cannot wait on, leaks) and then releases its owners. + PARKED.fetch_sub(1, Ordering::AcqRel); + REAPED.fetch_add(1, Ordering::Relaxed); + } +} + +/// Contexts currently parked: handed to the reactor, flag not yet landed +/// or release not yet run. +static PARKED: AtomicUsize = AtomicUsize::new(0); +/// Contexts released by the reaper since process start. +static REAPED: AtomicU64 = AtomicU64::new(0); + +/// Number of abandoned submissions whose resources are still held by the +/// reaper. Memory an abandoned future owned now outlives its dropping scope +/// by however long its GPU work takes; this is how much is outstanding. +pub fn parked() -> usize { + PARKED.load(Ordering::Acquire) +} + +/// Number of abandoned submissions the reaper has released so far. +pub fn reaped_total() -> u64 { + REAPED.load(Ordering::Relaxed) +} + +/// Releases `ctx` without blocking the caller when possible. +/// +/// Called from [`DeviceFuture`]'s drop after the result handle has been +/// released. See the module docs for the exact policy. +/// +/// [`DeviceFuture`]: crate::device_future::DeviceFuture +pub(crate) fn park_or_wait(ctx: ExecutionContext) { + let stream = Arc::clone(ctx.get_cuda_stream()); + // Probing and arming need a current context; a thread that cannot bind + // the device falls back to the inline release, which reports the error. + if stream.device().bind_to_thread().is_err() { + drop(ctx); + return; + } + match probe_stream(&stream) { + // Nothing in flight: the inline release costs one more query. + StreamHealth::Idle => drop(ctx), + StreamHealth::Busy => { + PARKED.fetch_add(1, Ordering::AcqRel); + let parked = Parked { _ctx: ctx }; + // SAFETY: the stream is valid and its context is current on this + // thread (bound above); the flag write is ordered behind the + // abandoned work because it is enqueued on the same stream. + if let Err(parked) = unsafe { crate::reactor::park(&stream, parked) } { + // The reactor could not take it: release inline, which is the + // blocking wait. `Parked::drop` keeps the counters honest. + drop(parked); + } + } + // A capturing stream must not receive a flag write, and a faulted one + // can never land it. The submission's release handles both: it + // cannot prove completion, so it reports and leaks the owners. + StreamHealth::Capturing | StreamHealth::Faulted(_) => drop(ctx), + } +} + +/// Passes a landed payload to the release thread, so the reactor's scan loop +/// never runs a release itself. If the thread is gone the payload is +/// released on the caller's thread: still sound, just not isolated. +pub(crate) fn release(parked: Parked) { + if let Err(mpsc::SendError(parked)) = sender().send(parked) { + drop(parked); + } +} + +fn sender() -> &'static Sender { + static SENDER: OnceLock> = OnceLock::new(); + SENDER.get_or_init(|| { + let (tx, rx) = mpsc::channel::(); + // If the thread cannot be spawned the receiver drops here, every + // later `send` fails, and `release` falls back to the inline drop. + let _ = thread::Builder::new() + .name("cuda-async-reaper".into()) + .spawn(move || { + for parked in rx { + drop(parked); + } + }); + tx + }) +} diff --git a/cuda-async/src/submission.rs b/cuda-async/src/submission.rs index 8a2e3c890a..6ae97af795 100644 --- a/cuda-async/src/submission.rs +++ b/cuda-async/src/submission.rs @@ -9,6 +9,7 @@ use crate::device_future::{probe_stream, StreamHealth}; use crate::device_operation::ReplayResource; use crate::error::DeviceError; use cuda_core::Stream; +use std::sync::atomic::{AtomicBool, Ordering}; use std::sync::{Arc, Mutex}; /// Storage whose per-access bookkeeping is released when the submission @@ -90,6 +91,10 @@ pub(crate) struct Submission { stream: Arc, owners: Mutex, recorded: Mutex>>, + /// The future that owned this submission already resolved with the + /// stream's fault: a leak on release is that fault's documented + /// consequence, not news, so it is not reported again. + fault_delivered: AtomicBool, } impl std::fmt::Debug for Submission { @@ -104,9 +109,14 @@ impl Submission { stream, owners: Mutex::new(Owners(Some(OwnerSet::default()))), recorded: Mutex::new(Vec::new()), + fault_delivered: AtomicBool::new(false), } } + pub(crate) fn mark_fault_delivered(&self) { + self.fault_delivered.store(true, Ordering::Relaxed); + } + pub(crate) fn retain(&self, owner: impl Send + 'static) -> Result<(), DeviceError> { self.owners .lock() @@ -147,7 +157,14 @@ impl Submission { } impl Drop for Submission { + /// Releases the owners once the stream has drained. This is the blocking + /// last resort: a dropped in-flight future normally parks its context + /// with [`crate::reaper`], which drops it only after the flag behind the + /// work has landed, so the probe below finds the stream idle. A stream + /// that cannot be waited on (faulted, or mid-capture) leaks the owners, + /// loudly unless the owning future already delivered the fault. fn drop(&mut self) { + let quiet = self.fault_delivered.load(Ordering::Relaxed); self.owners .get_mut() .unwrap_or_else(|e| e.into_inner()) @@ -165,7 +182,13 @@ impl Drop for Submission { )), } })(); - if result.is_err() { + if let Err(error) = &result { + if !quiet { + crate::leak::report_leak(format_args!( + "cuda-async: leaking the resources of an abandoned submission; the \ + driver could not prove its GPU work finished: {error}" + )); + } // Keep the stream identity valid for leaked access leases too. std::mem::forget(self.stream.clone()); } diff --git a/cuda-async/tests/drop_in_flight.rs b/cuda-async/tests/drop_in_flight.rs index 55464399f8..1c85b4b829 100644 --- a/cuda-async/tests/drop_in_flight.rs +++ b/cuda-async/tests/drop_in_flight.rs @@ -148,12 +148,13 @@ impl Drop for Gate { } } -/// What the output's `Drop` observed: whether the device had passed the -/// event recorded after the op's work when the output was released. +/// What the retained operand's `Drop` observed: whether the device had +/// passed the event recorded after the op's work when it was released. type DropLog = Arc>>; -/// The op's output: owns the completion event and reports, on drop, whether -/// the work had completed by then. +/// An operand the op retains in its submission (the thing whose release +/// must wait for the device): owns the completion event and reports, on +/// drop, whether the work had completed by then. struct Tracked { event: Event, log: DropLog, @@ -177,8 +178,8 @@ struct SlowOp { } impl DeviceOp for SlowOp { - type Output = Tracked; - unsafe fn execute(self, context: &ExecutionContext) -> Result { + type Output = (); + unsafe fn execute(self, context: &ExecutionContext) -> Result<(), DeviceError> { let stream = context.get_cuda_stream(); if let Some(gate) = &self.gate { gate.arm(stream.cu_stream())?; @@ -194,16 +195,17 @@ impl DeviceOp for SlowOp { } let event = stream.device().new_event()?; event.record(stream)?; - Ok(Tracked { + context.retain(Tracked { event, log: self.log, - }) + })?; + Ok(()) } } impl IntoFuture for SlowOp { - type Output = Result; - type IntoFuture = cuda_async::device_future::DeviceFuture; + type Output = Result<(), DeviceError>; + type IntoFuture = cuda_async::device_future::DeviceFuture<(), SlowOp>; fn into_future(self) -> Self::IntoFuture { let policy = global_policy(0).expect("global policy"); match self.schedule(&policy) { @@ -240,6 +242,19 @@ fn gated_op(dptr: u64, log: &DropLog, gate: &Arc) -> SlowOp { } } +/// Waits until the reaper holds nothing, so releases can be asserted on. +fn wait_for_reaper(deadline: Duration) { + let start = Instant::now(); + while cuda_async::reaper::parked() != 0 { + assert!( + start.elapsed() < deadline, + "the reaper still holds {} parked submission(s) after {deadline:?}", + cuda_async::reaper::parked() + ); + std::thread::sleep(Duration::from_millis(1)); + } +} + fn block_on_with_deadline(mut future: F, deadline: Duration) -> F::Output { let start = Instant::now(); let waker = noop_waker(); @@ -259,11 +274,13 @@ fn block_on_with_deadline(mut future: F, deadline: Duration) } /// The core regression: poll once with the work provably in flight (held -/// behind a closed gate), drop the future, and check that the output was -/// released only after the device passed the event recorded behind the work. +/// behind a closed gate), drop the future, and check two things: the drop +/// returned before the gate opened (it did not block on the device), and +/// the retained operand was released only after the device passed the event +/// recorded behind the work (the reaper waited for it). /// -/// The gate opens from a helper thread after `GATE_DELAY`, so a drop that -/// waits cannot return before then, and one that does not is caught twice: +/// The gate opens from a helper thread after `GATE_DELAY`, so a blocking +/// drop cannot return before then, and a premature release is caught twice: /// it returns early, and its output's event reports the work still in flight. #[test] fn dropping_in_flight_future_releases_output_after_the_device_finished() { @@ -284,6 +301,7 @@ fn dropping_in_flight_future_releases_output_after_the_device_finished() { for _ in 0..4 { let gate = Gate::closed(); + let released_before = log.lock().unwrap().len(); let mut future = gated_op(dptr, &log, &gate).into_future(); let waker = noop_waker(); let mut cx = Context::from_waker(&waker); @@ -311,18 +329,28 @@ fn dropping_in_flight_future_releases_output_after_the_device_finished() { drop(future); let elapsed = started.elapsed(); assert!( - elapsed >= GATE_DELAY, - "the drop returned after {elapsed:?}, before the gate opened at {GATE_DELAY:?}: \ - it did not wait for the in-flight work" + elapsed < GATE_DELAY, + "the drop returned only after {elapsed:?}, past the gate's {GATE_DELAY:?}: it \ + blocked on the in-flight work instead of parking it" + ); + assert_eq!( + log.lock().unwrap().len(), + released_before, + "an operand was released while the gate still held its work back" ); opener.join().expect("gate opener thread panicked"); + wait_for_reaper(Duration::from_secs(30)); } let log = log.lock().unwrap(); - assert_eq!(log.len(), 4, "every dropped future must release its output"); + assert_eq!( + log.len(), + 4, + "every dropped future must release its operand" + ); assert!( log.iter().all(|&done| done), - "an output was released while its GPU work was still in flight: {log:?}" + "an operand was released while its GPU work was still in flight: {log:?}" ); }); } @@ -343,20 +371,65 @@ fn dropping_unpolled_future_submits_nothing() { }); } -/// After the result has been delivered, dropping the future is a no-op; the -/// delivered output is the caller's and reports completion when dropped. +/// A delivered result completes its submission: the retained operand is +/// released right there, after the work, without involving the reaper. #[test] -fn delivered_result_is_owned_by_the_caller() { +fn delivered_result_releases_its_operands_on_delivery() { on_fresh_thread(|| { init_device_contexts(0, 1).expect("init failed (requires GPU)"); let dptr = alloc_device(BUF); let log: DropLog = Arc::new(Mutex::new(Vec::new())); - let tracked = - block_on_with_deadline(slow_op(dptr, &log).into_future(), Duration::from_secs(30)) - .expect("op failed"); - assert!(log.lock().unwrap().is_empty(), "nothing dropped yet"); - drop(tracked); + block_on_with_deadline(slow_op(dptr, &log).into_future(), Duration::from_secs(30)) + .expect("op failed"); + assert_eq!( + log.lock().unwrap().as_slice(), + [true], + "the operand is released at delivery, after the work" + ); + }); +} + +/// Racing a device future against something that wins first (the `select!` +/// / timeout shape) must not stall the executor: the losing device future +/// is dropped mid-flight on the executor thread, which returns promptly, +/// and its operand is released only after the device finishes. +#[test] +fn losing_a_select_does_not_block_the_executor() { + on_fresh_thread(|| { + init_device_contexts(0, 1).expect("init failed (requires GPU)"); + let dptr = alloc_device(BUF); + let log: DropLog = Arc::new(Mutex::new(Vec::new())); + // Warm up the completion path on an ungated op. + let scratch: DropLog = Default::default(); + block_on_with_deadline( + slow_op(dptr, &scratch).into_future(), + Duration::from_secs(30), + ) + .expect("warm-up op failed"); + + let gate = Gate::closed(); + let opener = gate.open_after(GATE_DELAY); + let started = Instant::now(); + futures::executor::block_on(async { + let device = gated_op(dptr, &log, &gate).into_future(); + let winner = futures::future::ready(()); + match futures::future::select(device, winner).await { + futures::future::Either::Left(_) => panic!("the gated op resolved first"), + futures::future::Either::Right(((), device)) => drop(device), + } + }); + let elapsed = started.elapsed(); + assert!( + elapsed < GATE_DELAY, + "the executor was blocked for {elapsed:?} by the losing future's drop" + ); + assert!( + log.lock().unwrap().is_empty(), + "operand released before the gate opened" + ); + opener.join().expect("gate opener thread panicked"); + wait_for_reaper(Duration::from_secs(30)); assert_eq!(log.lock().unwrap().as_slice(), [true]); }); } @@ -380,6 +453,13 @@ fn later_pipelines_complete_after_cancellations() { block_on_with_deadline(slow_op(dptr, &log).into_future(), Duration::from_secs(30)) .expect("op after cancellations failed"); } - assert!(log.lock().unwrap().iter().all(|&done| done)); + wait_for_reaper(Duration::from_secs(30)); + let log = log.lock().unwrap(); + assert_eq!( + log.len(), + 12, + "8 cancelled + 4 completed operands released: {log:?}" + ); + assert!(log.iter().all(|&done| done), "{log:?}"); }); } diff --git a/cutile-rs/CHANGELOG.md b/cutile-rs/CHANGELOG.md index 4b7ec6eff9..ab04e33e97 100644 --- a/cutile-rs/CHANGELOG.md +++ b/cutile-rs/CHANGELOG.md @@ -100,6 +100,26 @@ paths have lower overhead. ### Changed +- Dropping an in-flight `DeviceFuture` no longer blocks the dropping thread. + The result handle is released immediately; the execution context, whose + submission owns every resource the device may still use, is parked with + the completion reactor and released on a dedicated reaper thread once a + flag write enqueued behind the abandoned work lands. Memory safety is + unchanged (nothing the device may write is freed before the stream + drains); the blocking wait remains the fallback when the reactor cannot + take the context, and the leak-with-report path remains for faulted or + capturing streams. `cuda_async::reaper::parked()` and `reaped_total()` + expose the outstanding and released counts. Code that used a future's + drop as a barrier before touching a buffer through an unchecked path (a + raw device pointer handed to another library) must now synchronize. + +- `Tensor::store` returns its completion `Token`; `Tensor::token` reads it + and unsafe `Tensor::set_token` installs an external dependency. Explicit + installation rejects conditional/loop regions, block-local receivers, + and subsequent shadowing until token updates use control-flow carries. + The raw `set_tensor_token` and `make_partition_view` helpers are unsafe; + safe partition helpers continue to preserve the tensor's own token. + - Tile IR compiler, encoder and example/test prerequisites share a capability registry with checked documentation tables. Standalone bytecode validation uses the JIT's assembler-version negotiation, including the CUDA 13.2