gpu-handle-types 0.2.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
// SPDX-License-Identifier: MIT OR Apache-2.0
//
// Waiter-thread-fallback architecture.
//
// Coverage matrix:
//
//  1. `WaiterThread::enqueue` resolves a future when the slice closure
//     returns [`SliceOutcome::Signaled`].
//  2. Deadline expiry surfaces as `Err(Error::Timeout)`.
//  3. Slice-closure `Failed` outcome propagates verbatim.
//  4. Dropping the returned future flips the cancellation flag the
//     thread observes at the next slice boundary; the thread is then
//     available for a follow-up request (proves cancellation didn't
//     break the loop).
//  5. Dropping the `WaiterThread` joins cleanly with no in-flight
//     request hanging — the thread-wide shutdown flag exits the slice
//     loop even when the request would otherwise run forever — and the
//     in-flight future is *resolved* rather than left pending forever.
//  6. A request still queued behind a busy slice when the thread is
//     torn down is resolved too.
//  7. Requests are serviced on the same OS thread across calls
//     (`ThreadId` stable — proves the thread is permanent, not
//     respawned per request).
//  8. Two independent `WaiterThread` instances run on distinct
//     `ThreadId`s.
//
// Tests run under `pollster::block_on` — the minimum-footprint
// executor. The default `wait_async` impl on
// `SyncWaiter` covers the same hybrid path under `pollster`, `tokio`,
// `smol`, and `futures` by construction (it uses only
// `core::future::poll_fn` and this crate's own
// [`gpu_handle_types::yield_once`] — no per-runtime types touch the
// future).
//
// ## Why the lifecycle tests contain no sleeps
//
// Cases 4-6 are about *where in the slice loop* a flag is observed:
// "cancellation is observed at the next slice boundary", "shutdown
// breaks the loop instead of running the request to completion".
// Expressing them as wall-clock budgets around
// `std::thread::sleep(WAITER_SLICE)` (a window of observed slices, a
// join budget) would measure the host's scheduler, not the loop. On a
// saturated machine `sleep(10ms)` overshoots by hundreds of
// milliseconds and the thread may not even reach its first `recv()`
// inside a tight window, which makes "0 slices ran" a legitimate
// outcome that any lower bound on the slice count fails on.
//
// Instead they are driven by a two-channel handshake with the slice
// closure: the closure announces that it entered a slice and then parks
// until the test releases it. That pins the loop in a known state, so
// the cancellation assertion is *exact* ("no second slice is ever
// issued") rather than ranged, and the shutdown assertions are
// liveness properties ("`drop` returns at all against a request that
// never terminates"). No assertion reads a clock. `HANG_BUDGET` below
// is only a deadlock detector — it never distinguishes pass from fail,
// and every wait it guards resolves immediately on a correct run.

use std::future::Future;
use std::sync::Arc;
use std::sync::atomic::{AtomicUsize, Ordering};
use std::sync::mpsc;
use std::task::{Context, Poll, Waker};
use std::thread::ThreadId;
use std::time::{Duration, Instant};

use gpu_handle_types::{BackendKind, Error, SliceFn, SliceOutcome, WaiterThread};

/// Deadline for the blocking handshake `recv`s. It is a **deadlock
/// detector**, not a latency budget: on every correct path the value
/// being waited for is already available or the channel is already
/// disconnected, so the wait returns immediately regardless of machine
/// load. It exists purely so a regression that deadlocks the slice loop
/// fails the test instead of hanging the suite forever.
const HANG_BUDGET: Duration = Duration::from_secs(30);

/// Poll a future exactly once and report whether it was ready.
///
/// Used by the shutdown tests, where the waiter thread has already been
/// joined before the poll: the completion cell is either populated or
/// nothing in the process will ever populate it, so `Pending` is a
/// terminal *failure* rather than a state worth waiting on.
/// `pollster::block_on` would park forever on the regression these
/// tests exist to catch — a test that hangs the suite reports nothing.
fn poll_once<F: Future>(fut: F) -> Option<F::Output> {
    let fut = std::pin::pin!(fut);
    match fut.poll(&mut Context::from_waker(Waker::noop())) {
        Poll::Ready(v) => Some(v),
        Poll::Pending => None,
    }
}

/// Slice closure that parks inside **every** slice until the test
/// releases it, so the test can pin the waiter loop at a known point
/// and count slices exactly.
///
/// Returns `(slice_fn, entered_rx, release_tx)`:
/// * `entered_rx` yields one message per slice entered — the slice
///   *count*, observed exactly rather than inferred from elapsed time.
///   It disconnects when the waiter thread drops the request, which is
///   itself an observable event.
/// * `release_tx` lets the current slice return `TimedOut`.
///
/// Only usable when the test thread itself publishes the flag the
/// waiter is expected to observe (the cancellation case): the send that
/// unparks the slice is then the happens-before edge that makes the
/// observation certain. When the flag is set by a *third* thread, use
/// [`park_once_then_free_run`] instead — see its docs.
fn parking_slice() -> (SliceFn, mpsc::Receiver<()>, mpsc::Sender<()>) {
    let (entered_tx, entered_rx) = mpsc::channel::<()>();
    let (release_tx, release_rx) = mpsc::channel::<()>();
    let slice_fn: SliceFn = Box::new(move |_slice| {
        // A disconnected receiver means the test finished early; the
        // waiter thread must not panic on it.
        let _ = entered_tx.send(());
        let _ = release_rx.recv();
        SliceOutcome::TimedOut
    });
    (slice_fn, entered_rx, release_tx)
}

/// Slice closure that parks inside its **first** slice and free-runs
/// afterwards — every later slice yields and returns `TimedOut`
/// immediately, so the request never terminates on its own.
///
/// The park gives the test a deterministic "the waiter thread is busy
/// and cannot reach `recv()`" point. The free-run afterwards is what
/// makes the *shutdown* tests race-free: `WaiterThread::drop` stores the
/// flag from a third thread and there is no way for the test to observe
/// that store, so a closure that parks on every slice could park again
/// before the store lands and deadlock the join. Free-running instead
/// means the loop keeps reaching its flag checks until the store
/// arrives, whenever that is.
///
/// What it costs: the exact "at most one further slice" count is not
/// observable, so the shutdown tests assert the liveness property
/// instead — `drop` returns *at all* against a request that never
/// terminates, which is only possible if the loop tests `shutdown` at
/// the slice boundary rather than only at `recv()`. That is the whole
/// regression this guards; the ≤ 1-slice latency then follows from the
/// loop's structure, and the exact boundary-check placement is pinned
/// separately (and deterministically) by
/// [`cancellation_observed_within_one_slice`].
fn park_once_then_free_run() -> (SliceFn, mpsc::Receiver<()>, mpsc::Sender<()>) {
    let (entered_tx, entered_rx) = mpsc::channel::<()>();
    let (release_tx, release_rx) = mpsc::channel::<()>();
    let mut park = Some((entered_tx, release_rx));
    let slice_fn: SliceFn = Box::new(move |_slice| {
        match park.take() {
            Some((entered_tx, release_rx)) => {
                let _ = entered_tx.send(());
                let _ = release_rx.recv();
            }
            // Stay runnable but cooperative: the loop re-checks
            // `shutdown` / `cancelled` on every pass.
            None => std::thread::yield_now(),
        }
        SliceOutcome::TimedOut
    });
    (slice_fn, entered_rx, release_tx)
}

#[test]
fn signaled_outcome_resolves_future() {
    let thread = WaiterThread::new("test-signal");
    let probes = Arc::new(AtomicUsize::new(0));
    let probes_c = probes.clone();
    let slice_fn: SliceFn = Box::new(move |_slice| {
        let n = probes_c.fetch_add(1, Ordering::SeqCst);
        if n >= 2 { SliceOutcome::Signaled } else { SliceOutcome::TimedOut }
    });
    let fut = thread.enqueue(slice_fn, None);
    let result = pollster::block_on(fut);
    assert!(result.is_ok(), "expected Ok, got {result:?}");
    assert!(probes.load(Ordering::SeqCst) >= 3);
}

#[test]
fn deadline_expiry_surfaces_timeout() {
    let thread = WaiterThread::new("test-timeout");
    let slice_fn: SliceFn = Box::new(|_| SliceOutcome::TimedOut);
    let deadline = Instant::now().checked_add(Duration::from_millis(50));
    let fut = thread.enqueue(slice_fn, deadline);
    let result = pollster::block_on(fut);
    assert!(matches!(result, Err(Error::Timeout)), "expected Err(Timeout), got {result:?}");
}

#[test]
fn failed_outcome_surfaces_verbatim() {
    let thread = WaiterThread::new("test-fail");
    let slice_fn: SliceFn = Box::new(|_| SliceOutcome::Failed(Error::DeviceLost { backend: BackendKind::Vulkan }));
    let fut = thread.enqueue(slice_fn, None);
    let result = pollster::block_on(fut);
    match result {
        Err(Error::DeviceLost { backend: BackendKind::Vulkan }) => {}
        other => panic!("expected Err(DeviceLost {{ Vulkan }}), got {other:?}"),
    }
}

/// Dropping the future must stop the request at the **next slice
/// boundary** — not mid-slice, and not one slice later.
///
/// The handshake makes that exact. The cancel flag is stored (with
/// `Release`) while slice 1 is parked, and the `release_tx.send` that
/// unparks it is the happens-before edge that publishes the flag to the
/// waiter thread, so the `Acquire` load at the top of the next
/// iteration is *guaranteed* to observe it. Slice 2 must therefore
/// never be entered — a property, not a stopwatch reading.
#[test]
fn cancellation_observed_within_one_slice() {
    let thread = WaiterThread::new("test-cancel");
    let (slice_fn, entered_rx, release_tx) = parking_slice();
    let fut = thread.enqueue(slice_fn, None);

    // Deterministic "the waiter thread is now inside slice 1".
    entered_rx.recv_timeout(HANG_BUDGET).expect("waiter thread never entered its first slice");

    drop(fut); // flips the cancellation flag
    release_tx.send(()).expect("waiter thread must still be parked inside slice 1");

    // No slice 2. On the correct path the waiter abandons the request
    // and drops it — taking `entered_tx` with it — so this resolves as
    // `Disconnected` immediately; it never actually waits out the
    // budget. `Ok(())` would mean a slice ran after the cancel flag was
    // published, i.e. the flag is checked in the wrong place (or not at
    // all).
    assert!(
        entered_rx.recv_timeout(HANG_BUDGET).is_err(),
        "a slice was issued after the future was dropped; cancellation is not observed at the slice boundary",
    );

    // The thread must be idle on `recv()` again — cancellation abandons
    // the request without breaking the loop.
    let done = Arc::new(AtomicUsize::new(0));
    let done_c = done.clone();
    let slice_fn: SliceFn = Box::new(move |_| {
        done_c.fetch_add(1, Ordering::SeqCst);
        SliceOutcome::Signaled
    });
    let result = pollster::block_on(thread.enqueue(slice_fn, None));
    assert!(result.is_ok(), "follow-up enqueue must succeed; cancellation broke the loop?");
    assert_eq!(done.load(Ordering::SeqCst), 1);
}

#[test]
fn drop_joins_cleanly_when_idle() {
    // Spawn + drop with nothing in flight. The channel-close path
    // alone suffices here (the slice loop never started).
    let thread = WaiterThread::new("test-shutdown-idle");
    drop(thread); // joins on this line
}

/// The thread is mid-slice on a request that would never finish. The
/// thread-wide shutdown flag set by `WaiterThread::drop` must break the
/// slice loop at the next boundary; otherwise `join()` hangs forever.
///
/// Three properties, none of them a stopwatch reading:
///
/// * `drop` really **joins** (it does not detach) — proved by observing
///   that it cannot return while a slice is parked.
/// * It returns **at all** against a request that never terminates.
///   The outer `recv()` is never reached on this path, so the only way
///   out is the `shutdown` check at the slice boundary; a regression
///   that dropped it turns this into a permanent hang, which is the
///   whole point of the flag.
/// * The in-flight future is **resolved**, not abandoned.
///   [`gpu_handle_types::BackendWaitFuture`] borrows nothing from the
///   thread, so holding one across the drop is legal and an unresolved
///   completion would pend forever with no waker left alive.
#[test]
fn drop_joins_cleanly_with_in_flight_request() {
    let thread = WaiterThread::new("test-shutdown-inflight");
    let (slice_fn, entered_rx, release_tx) = park_once_then_free_run();
    let fut = thread.enqueue(slice_fn, None);

    entered_rx.recv_timeout(HANG_BUDGET).expect("waiter thread never entered its first slice");

    // `drop` blocks in `join`, so it has to run off the test thread.
    let (joined_tx, joined_rx) = mpsc::channel::<()>();
    let dropper = std::thread::spawn(move || {
        drop(thread); // shutdown flag + sender close + join
        let _ = joined_tx.send(());
    });

    // Negative assertion that cannot fail spuriously: the waiter thread
    // is provably parked inside its slice, so a correct `join` cannot
    // have returned no matter how the two threads are scheduled.
    // `Ok(())` here would mean `drop` gave up on the thread instead of
    // joining it.
    assert!(
        joined_rx.recv_timeout(Duration::from_millis(200)).is_err(),
        "WaiterThread::drop returned while the waiter thread was still inside a slice — it must join, not detach",
    );

    release_tx.send(()).expect("waiter thread must still be parked inside slice 1");
    joined_rx
        .recv_timeout(HANG_BUDGET)
        .expect("WaiterThread::drop deadlocked — the slice loop never observed the shutdown flag");
    dropper.join().expect("dropper thread panicked");

    match poll_once(fut) {
        Some(Err(Error::Cancelled)) => {}
        other => panic!("an in-flight future must resolve to Err(Cancelled) on shutdown, got {other:?}"),
    }
}

/// A request that never reached a slice — it was still sitting in the
/// channel behind a busy one when the thread was torn down — must be
/// resolved as well. Same reasoning as the in-flight case: its future
/// is alive, holds no reference to the thread, and nothing else will
/// ever complete it.
///
/// Parking the first request's opening slice is what makes "still
/// queued" deterministic: the waiter cannot reach `recv()` while it is
/// parked, so the second `enqueue` provably lands in the channel and
/// stays there.
#[test]
fn drop_resolves_a_queued_request_instead_of_hanging() {
    let thread = WaiterThread::new("test-shutdown-queued");
    let (slice_fn, entered_rx, release_tx) = park_once_then_free_run();
    let busy = thread.enqueue(slice_fn, None);

    entered_rx.recv_timeout(HANG_BUDGET).expect("waiter thread never entered its first slice");
    let queued = thread.enqueue(Box::new(|_| SliceOutcome::Signaled), None);

    let (joined_tx, joined_rx) = mpsc::channel::<()>();
    let dropper = std::thread::spawn(move || {
        drop(thread);
        let _ = joined_tx.send(());
    });
    release_tx.send(()).expect("waiter thread must still be parked inside slice 1");
    joined_rx.recv_timeout(HANG_BUDGET).expect("WaiterThread::drop deadlocked");
    dropper.join().expect("dropper thread panicked");

    match poll_once(queued) {
        Some(Err(Error::Cancelled)) => {}
        other => panic!("a queued request's future must resolve to Err(Cancelled) on shutdown, got {other:?}"),
    }
    match poll_once(busy) {
        Some(Err(Error::Cancelled)) => {}
        other => panic!("the in-flight request's future must resolve to Err(Cancelled) on shutdown, got {other:?}"),
    }
}

#[test]
fn requests_share_a_single_thread() {
    let thread = WaiterThread::new("test-reuse");
    let observed: Arc<std::sync::Mutex<Vec<ThreadId>>> = Arc::new(std::sync::Mutex::new(Vec::new()));
    for _ in 0..5 {
        let observed_c = observed.clone();
        let slice_fn: SliceFn = Box::new(move |_| {
            observed_c.lock().unwrap().push(std::thread::current().id());
            SliceOutcome::Signaled
        });
        pollster::block_on(thread.enqueue(slice_fn, None)).expect("Signaled outcome should resolve to Ok");
    }
    let ids = observed.lock().unwrap();
    assert_eq!(ids.len(), 5);
    let first = ids[0];
    for (i, id) in ids.iter().enumerate().skip(1) {
        assert_eq!(*id, first, "slice {i} ran on a different ThreadId; thread respawned?");
    }
}

#[test]
fn distinct_instances_run_on_distinct_threads() {
    let a = WaiterThread::new("test-distinct-a");
    let b = WaiterThread::new("test-distinct-b");
    let (tid_a, tid_b) = {
        let cap_a: Arc<std::sync::Mutex<Option<ThreadId>>> = Arc::new(std::sync::Mutex::new(None));
        let cap_b = cap_a.clone();
        let cap_a_inner = cap_a.clone();
        let slice_fn_a: SliceFn = Box::new(move |_| {
            *cap_a_inner.lock().unwrap() = Some(std::thread::current().id());
            SliceOutcome::Signaled
        });
        pollster::block_on(a.enqueue(slice_fn_a, None)).unwrap();
        let tid_a = cap_a.lock().unwrap().take().unwrap();

        let cap_b_inner = cap_b.clone();
        let slice_fn_b: SliceFn = Box::new(move |_| {
            *cap_b_inner.lock().unwrap() = Some(std::thread::current().id());
            SliceOutcome::Signaled
        });
        pollster::block_on(b.enqueue(slice_fn_b, None)).unwrap();
        let tid_b = cap_b.lock().unwrap().take().unwrap();
        (tid_a, tid_b)
    };
    assert_ne!(tid_a, tid_b, "two WaiterThread instances must spawn distinct OS threads");
}