gpu-handle-types 0.1.0

Typed, owned native GPU resource handles (Vulkan, D3D11/12, Metal, OpenGL, CUDA, OpenCL, DMA-BUF, IOSurface, AHardwareBuffer, WebGPU, ...), cross-API sync points and video pixel formats, for passing GPU resources between libraries.
Documentation
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
// SPDX-License-Identifier: MIT OR Apache-2.0

//! Waiter-thread infrastructure for the waiter-thread fallback.
//!
//! The fast-path yield-spin inside [`SyncWaiter::wait_async`] catches the
//! sub-millisecond signal-already-here case with 64 cooperative yields.
//! Anything past that bounded spin must hand off to a **permanently-
//! attached** thread that issues the **native blocking-with-timeout**
//! primitive on 10 ms slices — without `tokio::spawn_blocking`, without
//! per-runtime timers, without per-operation thread spawn.
//!
//! ## Components
//!
//! - [`WaiterThread`] — owns the request channel + the OS thread.
//!   Typically one instance per backend kind (lazy-initialised via
//!   [`OnceLock`] inside each backend's waiter module); per-instance
//!   ownership is supported when callers want a deterministic
//!   shutdown via `Drop`.
//! - [`SliceFn`] — boxed closure each per-backend `wait_async`
//!   override builds at handoff time. The waiter thread calls it
//!   repeatedly with up to 10 ms of slice budget; it returns
//!   [`SliceOutcome::Signaled`] / `TimedOut` / `Failed`.
//! - [`BackendWaitFuture`] — the future returned to the executor.
//!   Its [`Drop`] flips a cancellation flag the waiter thread observes
//!   at slice boundaries; cancellation latency is therefore ≤ one
//!   slice (10 ms) — sized to fit inside a 60 fps frame budget.
//! - [`run_hybrid_wait`] — composes the fast-path yield-spin + the
//!   waiter-thread fallback (thread handoff). Backends with a native blocking primitive call
//!   this from their `wait_async` override and pass the slice
//!   closure.
//!
//! ## Why a per-process (per-backend-kind) static thread
//!
//! "One thread per backend wait-registry" is the natural shape. In
//! practice a process holds only one long-lived wait-registry per
//! backend kind on the hot path; more importantly the waiter thread is
//! **stateless** — every
//! per-instance datum (device handle, semaphore, fence, …) lives
//! inside the [`SliceFn`] closure. A second per-instance thread
//! would spend its life parked on `recv()` doing no work, and a
//! per-`SyncWaiter`-impl thread would spawn thousands per second on
//! steady-state pipelines.
//!
//! The trade-off: with a static singleton the thread is not joined on
//! `Drop` of any individual backend; it lives until process exit and
//! is reaped by the OS. The `Drop`-time join contract
//! still holds when callers construct a [`WaiterThread`] explicitly
//! and own it from a parent registry struct (the type supports that
//! shape — `new` + `Drop` join cleanly), which is how this crate's
//! tests exercise the lifecycle.

use std::future::Future;
use std::pin::Pin;
use std::sync::atomic::{AtomicBool, Ordering};
// mpsc channel + join handle back the native waiter thread; wasm has none.
#[cfg(not(target_family = "wasm"))]
use std::sync::mpsc;
use std::sync::{Arc, Mutex};
use std::task::{Context, Poll, Waker};
#[cfg(not(target_family = "wasm"))]
use std::thread::JoinHandle;
use std::time::{Duration, Instant};

use crate::Error;

/// Per-slice native blocking-with-timeout budget. Sized to:
///
/// - Stay below the 60 fps frame budget (16.6 ms) so a future
///   dropped at the frame boundary releases its waiter slot inside
///   the same frame.
/// - Be coarse enough that the syscall overhead per slice is
///   amortised even on long waits — at 5 s of accumulated wait, the
///   thread re-enters the kernel ~500 times, not 5_000.
/// - Bound cancellation + deadline-exceeded latency to a single
///   slice.
///
/// Tunable as a `Duration` rather than a feature so the value can be
/// re-derived in tests without rebuilding the crate.
pub const WAITER_SLICE: Duration = Duration::from_millis(10);

/// Result of one slice-bounded native wait. Returned by the
/// [`SliceFn`] closure each per-backend `wait_async` override
/// supplies.
#[derive(Debug)]
pub enum SliceOutcome {
    /// The primitive reached its target value within the slice.
    /// Waiter thread resolves the future with `Ok(())`.
    Signaled,
    /// The slice's native timeout elapsed without signalling. Waiter
    /// thread re-issues a fresh slice (or transitions to
    /// `Err(Timeout)` if the per-request `deadline` has now passed).
    TimedOut,
    /// Driver-level failure (`VK_ERROR_DEVICE_LOST`,
    /// `DXGI_ERROR_DEVICE_REMOVED`, `CL_INVALID_EVENT`, …). Waiter
    /// thread resolves the future with this error directly.
    Failed(Error),
}

/// Boxed closure each per-backend `wait_async` override builds at
/// handoff time. The waiter thread calls it on its own thread with
/// the current slice budget; the closure must:
///
/// - Issue the native blocking-with-timeout primitive
///   (`vkWaitSemaphores`, `WaitForSingleObject`,
///   `MTLSharedEvent::waitUntilSignaledValue:timeoutMS:`,
///   `clWaitForEvents`, `device.poll(Wait { … })`, …) with the slice
///   budget.
///
/// - Return one of the three [`SliceOutcome`] variants.
///
/// `Send` is required because the closure runs on the waiter thread,
/// not on the executor's poll thread. `'static` because the slice
/// closure is consumed asynchronously and outlives the `await`
/// position that built it.
// `Send` off wasm (the slice runs on the waiter thread). On wasm there is
// no waiter thread — `enqueue` reports `NotSupported` — and the
// slice closure may close over thread-affine `wgpu` handles, so the bound
// is dropped. An explicit `+ Send` on a `dyn FnMut` cannot be spelled with
// a non-auto marker trait, so the alias is cfg-split directly.
#[cfg(not(target_family = "wasm"))]
pub type SliceFn = Box<dyn FnMut(Duration) -> SliceOutcome + Send + 'static>;
#[cfg(target_family = "wasm")]
pub type SliceFn = Box<dyn FnMut(Duration) -> SliceOutcome + 'static>;

/// Shared state between [`BackendWaitFuture`] (on the executor) and
/// the waiter thread.
struct WaitCompletion {
    result: Option<Result<(), Error>>,
    waker: Option<Waker>,
}

// Native-only: the request that crosses the channel to the waiter thread.
// wasm has no waiter thread, so this is never constructed there.
#[cfg(not(target_family = "wasm"))]
struct WaitRequest {
    slice_fn: SliceFn,
    deadline: Option<Instant>,
    cancelled: Arc<AtomicBool>,
    completion: Arc<Mutex<WaitCompletion>>,
}

/// Permanently-attached waiter thread + request channel.
///
/// Cheap to construct (one `mpsc::channel` + one `thread::spawn`).
/// Joined cleanly on `Drop`: the channel sender is dropped, the
/// thread's `recv()` returns `Err`, the loop exits, and the join
/// handle is consumed.
///
/// Every request that is in flight or still queued when `Drop` runs is
/// resolved with [`Error::Cancelled`] before the thread exits.
/// [`BackendWaitFuture`] borrows nothing from this type, so a caller may
/// hold one across the drop; without that resolution its `await` would
/// pend forever.
///
/// Per-instance ownership is supported, but the usual shape is a
/// per-backend-kind static singleton in each backend's waiter module.
/// See module-level docs.
pub struct WaiterThread {
    // Native-only: the channel + join handle to the waiter thread. On wasm
    // these are absent so `WaiterThread` stays `Send + Sync` (a
    // `Sender<WaitRequest>` carrying a `!Send` `SliceFn` would otherwise be
    // `!Sync`, breaking the process-`static` singletons that hold it).
    #[cfg(not(target_family = "wasm"))]
    sender: Option<mpsc::Sender<WaitRequest>>,
    #[cfg(not(target_family = "wasm"))]
    join: Option<JoinHandle<()>>,
    /// Thread-wide shutdown flag observed at slice boundaries.
    /// Flipped in `Drop` so the thread exits even mid-request — closing
    /// the channel alone wouldn't suffice when a long-running request is
    /// in flight (the slice loop never returns to `recv()`).
    shutdown: Arc<AtomicBool>,
}

impl std::fmt::Debug for WaiterThread {
    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
        #[cfg(not(target_family = "wasm"))]
        let alive = self.sender.is_some();
        // wasm has no waiter thread — never "alive".
        #[cfg(target_family = "wasm")]
        let alive = false;
        f.debug_struct("WaiterThread").field("alive", &alive).finish()
    }
}

impl WaiterThread {
    /// Spawn a new waiter thread. `name` is appended to
    /// `"wgpu-interop-wait-"` for the thread name visible to debuggers
    /// and `ps -L` / `Process Explorer`. Pick a stable
    /// backend-identifying string (`"vulkan"`, `"wgpu"`, `"d3d12"`,
    /// `"metal"`, `"opencl"`, …).
    pub fn new(name: &str) -> Self {
        // wasm has no OS threads. Construct an inert `WaiterThread`
        // (`sender = None`): `enqueue` then takes its existing
        // "already shut down" branch and reports `NotSupported`. Blocking
        // waits are unavailable on wasm — callers use async
        // readback; the `run_hybrid_wait` fast-path spin still runs.
        #[cfg(target_family = "wasm")]
        {
            let _ = name;
            return Self { shutdown: Arc::new(AtomicBool::new(false)) };
        }
        #[cfg(not(target_family = "wasm"))]
        {
            let (tx, rx) = mpsc::channel::<WaitRequest>();
            let shutdown = Arc::new(AtomicBool::new(false));
            let shutdown_thread = shutdown.clone();
            let join = std::thread::Builder::new()
                .name(format!("wgpu-interop-wait-{name}"))
                .spawn(move || waiter_loop(&rx, &shutdown_thread))
                .expect("gpu-handle-types: failed to spawn waiter thread");
            Self { sender: Some(tx), join: Some(join), shutdown }
        }
    }

    /// Enqueue a wait request. Returns a [`BackendWaitFuture`] the
    /// executor can `.await`; `Drop` of the future signals
    /// cancellation to the waiter thread.
    ///
    /// `deadline = None` means "wait forever" — the slice loop never
    /// transitions to `Err(Timeout)` and only exits on
    /// `Signaled` / `Failed` / `cancelled`.
    #[cfg_attr(target_family = "wasm", allow(unused_variables))]
    pub fn enqueue(&self, slice_fn: SliceFn, deadline: Option<Instant>) -> BackendWaitFuture {
        let cancelled = Arc::new(AtomicBool::new(false));
        let completion = Arc::new(Mutex::new(WaitCompletion { result: None, waker: None }));
        // wasm: no waiter thread — blocking waits are unavailable.
        // Resolve the future immediately to `NotSupported`;
        // callers use async readback. The `run_hybrid_wait` fast-path spin
        // already handles the already-signaled case before reaching here.
        #[cfg(target_family = "wasm")]
        finish_completion(
            &completion,
            Err(Error::NotSupported(
                "blocking waits are unavailable on wasm (no waiter thread); drive the async wait \
                 path instead"
                    .into(),
            )),
        );
        #[cfg(not(target_family = "wasm"))]
        {
            let req = WaitRequest { slice_fn, deadline, cancelled: cancelled.clone(), completion: completion.clone() };
            match self.sender.as_ref() {
                Some(s) => {
                    if s.send(req).is_err() {
                        finish_completion(
                            &completion,
                            Err(Error::NotSupported("waiter thread died before request was accepted".into())),
                        );
                    }
                }
                None => {
                    finish_completion(&completion, Err(Error::NotSupported("waiter thread already shut down".into())));
                }
            }
        }
        BackendWaitFuture { completion, cancelled }
    }
}

impl Drop for WaiterThread {
    fn drop(&mut self) {
        // Two-phase shutdown: (1) flip the thread-wide flag so the
        // slice loop exits at the next slice boundary even if the
        // request is still nominally in flight; (2) drop the sender
        // so the outer `recv()` returns `Err` after the current
        // request (if any) finishes. (wasm holds neither the sender nor
        // the join handle — nothing to tear down beyond the flag.)
        self.shutdown.store(true, Ordering::Release);
        #[cfg(not(target_family = "wasm"))]
        {
            self.sender.take();
            if let Some(h) = self.join.take() {
                let _ = h.join();
            }
        }
    }
}

#[cfg(not(target_family = "wasm"))]
fn waiter_loop(rx: &mpsc::Receiver<WaitRequest>, shutdown: &AtomicBool) {
    while let Ok(mut req) = rx.recv() {
        if shutdown.load(Ordering::Acquire) {
            abandon_on_shutdown(&req);
            break;
        }
        // Set when the per-request loop exited because the thread is
        // shutting down, so the outer loop stops pulling new work
        // instead of blocking in `recv()` behind a sender that
        // `WaiterThread::drop` may not have released yet.
        let mut shutting_down = false;
        loop {
            // Thread-wide shutdown — `WaiterThread::drop` was called.
            // Exit the per-request loop and the outer recv loop.
            if shutdown.load(Ordering::Acquire) {
                abandon_on_shutdown(&req);
                shutting_down = true;
                break;
            }
            // Cancellation — future dropped. Abandon the request;
            // the native slice we just issued (if any) has already
            // returned, and the closure may keep state across
            // returns that we'd corrupt by dropping mid-slice.
            // Slice granularity (10 ms) is the cancellation latency
            // bound by design.
            //
            // No `finish_completion` here, deliberately: the only way
            // this flag is set is [`BackendWaitFuture::drop`], so there
            // is no future left to resolve. Contrast the shutdown arm
            // above, where the future is typically still alive.
            if req.cancelled.load(Ordering::Acquire) {
                break;
            }
            // Deadline check.
            let remaining = match req.deadline {
                Some(d) => match d.checked_duration_since(Instant::now()) {
                    Some(r) => r,
                    None => {
                        finish_completion(&req.completion, Err(Error::Timeout));
                        break;
                    }
                },
                None => Duration::MAX,
            };
            let slice = std::cmp::min(WAITER_SLICE, remaining);
            match (req.slice_fn)(slice) {
                SliceOutcome::Signaled => {
                    finish_completion(&req.completion, Ok(()));
                    break;
                }
                SliceOutcome::TimedOut => continue,
                SliceOutcome::Failed(e) => {
                    finish_completion(&req.completion, Err(e));
                    break;
                }
            }
        }
        if shutting_down {
            break;
        }
    }
    // Requests still queued when the thread stops are resolved too, for
    // the same reason: their futures are alive and nothing else will
    // ever complete them. The sender has been dropped (or is about to
    // be) by `WaiterThread::drop`, so this drains and terminates.
    while let Ok(req) = rx.try_recv() {
        abandon_on_shutdown(&req);
    }
}

/// Resolve a request the waiter thread is giving up on because the
/// thread itself is being torn down.
///
/// This is **not** optional bookkeeping. [`BackendWaitFuture`] borrows
/// nothing from [`WaiterThread`] — it holds only the shared completion
/// cell — so `let f = t.enqueue(..); drop(t); f.await` is well-typed and
/// a caller may legitimately outlive the thread it enqueued on. Leaving
/// the completion unresolved makes that `await` pend forever with no
/// waker left in existence to rescue it.
///
/// [`Error::Cancelled`] rather than a failure kind: tearing the waiter
/// thread down is a caller-initiated stop, exactly the "not a failure,
/// the operation was asked to stop" case that
/// [`crate::ErrorKind::Cancelled`] exists to classify. Callers void it
/// instead of surfacing a fault.
#[cfg(not(target_family = "wasm"))]
fn abandon_on_shutdown(req: &WaitRequest) {
    finish_completion(&req.completion, Err(Error::Cancelled));
}

fn finish_completion(completion: &Arc<Mutex<WaitCompletion>>, result: Result<(), Error>) {
    let mut c = completion.lock().expect("WaitCompletion mutex poisoned");
    // Don't clobber an existing terminal outcome (the cancellation
    // path lets the future drop without resolving, but a follow-up
    // poll would still observe the result — keep the first one).
    if c.result.is_none() {
        c.result = Some(result);
    }
    let waker = c.waker.take();
    drop(c);
    if let Some(w) = waker {
        w.wake();
    }
}

/// Future returned by [`WaiterThread::enqueue`]. Resolves when the
/// waiter thread reports `Signaled` / `Failed` / `Err(Timeout)` on
/// the deadline.
///
/// Drop semantics — `Drop` flips the cancellation flag the waiter
/// thread observes at the next slice boundary. The native wait
/// already in flight completes its current slice (up to
/// [`WAITER_SLICE`]) before the slot is released; cancellation
/// latency is bounded.
pub struct BackendWaitFuture {
    completion: Arc<Mutex<WaitCompletion>>,
    cancelled: Arc<AtomicBool>,
}

impl std::fmt::Debug for BackendWaitFuture {
    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
        f.debug_struct("BackendWaitFuture")
            .field("cancelled", &self.cancelled.load(Ordering::Relaxed))
            .field("resolved", &self.completion.lock().map(|c| c.result.is_some()).unwrap_or(false))
            .finish()
    }
}

impl Future for BackendWaitFuture {
    type Output = Result<(), Error>;
    fn poll(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Self::Output> {
        let mut c = self.completion.lock().expect("WaitCompletion mutex poisoned");
        if let Some(r) = c.result.take() {
            return Poll::Ready(r);
        }
        // Refresh the waker on every poll — executors that pass new
        // wakers per re-poll (futures::executor) need this.
        c.waker = Some(cx.waker().clone());
        Poll::Pending
    }
}

impl Drop for BackendWaitFuture {
    fn drop(&mut self) {
        self.cancelled.store(true, Ordering::Release);
    }
}

/// Compose the fast-path yield-spin (bounded) + the waiter-thread
/// fallback (thread handoff) into the canonical hybrid `wait_async` body.
///
/// Per-backend `SyncWaiter::wait_async` overrides call this with:
/// - `is_signaled`: closure issuing the backend's non-blocking probe
///   (`vkGetSemaphoreCounterValue`, `GetCompletedValue`,
///   `signaledValue() >= value`, `clGetEventInfo`, `device.poll`
///   with `Duration::ZERO`, …). Runs on the executor thread during
///   the fast-path yield-spin.
/// - `waiter_thread`: per-backend-kind [`WaiterThread`] reference.
/// - `timeout`: absolute timeout from the caller.
/// - `make_slice_fn`: builder for the [`SliceFn`] consumed on the
///   waiter-thread fallback. Lazy so backends that resolve in the
///   fast-path yield-spin skip closure construction entirely.
///
/// Returns:
/// - `Ok(())` on signal,
/// - `Err(Error::Timeout)` on deadline,
/// - `Err(Error::DeviceLost { .. })` / `Err(Error::NotSupported(_.into()))`
///   on driver-level failure surfaced by the slice closure.
pub async fn run_hybrid_wait<F, M>(
    is_signaled: F,
    waiter_thread: &WaiterThread,
    timeout: Duration,
    make_slice_fn: M,
) -> Result<(), Error>
where
    // `MaybeSend` = `Send` off wasm. On wasm the closures
    // may capture thread-affine `wgpu` handles; the fast-path spin loop
    // still runs on the executor, and the waiter-thread fallback reports
    // `NotSupported`.
    F: Fn() -> Result<bool, Error> + crate::MaybeSend,
    M: FnOnce() -> SliceFn + crate::MaybeSend,
{
    // Fast-path yield-spin — bounded yield-poll. Iteration count, not wall-clock
    // (pollster's re-poll-on-wake collapses wall-clock caps to a
    // CPU burn here).
    const SPIN_ITERATIONS: usize = 64;
    // `None` iff `timeout == Duration::MAX` (the wait-forever sentinel, which
    // `enqueue` also reads as "no deadline"); every finite timeout gets a bounded,
    // always-representable deadline, so an `Instant + Duration` overflow cannot
    // masquerade as "wait forever".
    let deadline = crate::wait_deadline(timeout);
    for _ in 0..SPIN_ITERATIONS {
        match is_signaled() {
            Ok(true) => return Ok(()),
            Ok(false) => {}
            Err(e) => return Err(e),
        }
        if let Some(d) = deadline
            && Instant::now() >= d
        {
            return Err(Error::Timeout);
        }
        crate::yield_once().await;
    }
    // Waiter-thread fallback — hand off to the waiter thread.
    let slice_fn = make_slice_fn();
    waiter_thread.enqueue(slice_fn, deadline).await
}