use std::sync::Arc;
use std::time::{Duration, Instant};
use onnx_runtime_ep_api::{
ExternalMmapRegion, LazyWeight, LazyWeightBoundary, MmapRegionSource, ResidentWeight,
WeightHandleError,
};
use onnx_runtime_ep_cuda::CudaExecutionProvider;
use onnx_runtime_ep_cuda::prefill_double_buffer::{
LayerTicket, PrefillDoubleBuffer, PrefillLayerRequest, PrefillReject, PrefillSlotStatus,
};
use onnx_runtime_ep_cuda::weight_paging::{
PrefillRoute, PrefillWireDecline, prefill_double_buffer_enabled,
};
use onnx_runtime_ir::DataType;
fn require_cuda() -> CudaExecutionProvider {
match std::panic::catch_unwind(CudaExecutionProvider::new_default) {
Ok(Ok(ep)) => ep,
Ok(Err(error)) => panic!(
"CUDA test requires CUDA device/runtime; CPU-only runs must leave this test ignored: {error}"
),
Err(_) => panic!(
"CUDA test requires CUDA runtime libraries; CPU-only runs must leave this test ignored"
),
}
}
struct LayeredMmap {
mapping_id: usize,
bytes: Vec<u8>,
}
impl MmapRegionSource for LayeredMmap {
fn region_bytes(&self, region: &ExternalMmapRegion) -> Result<&[u8], WeightHandleError> {
if region.mapping_id != self.mapping_id {
return Err(WeightHandleError::DeviceBinding(format!(
"unknown mapping {}",
region.mapping_id
)));
}
let end = region
.offset
.checked_add(region.len)
.ok_or_else(|| WeightHandleError::DeviceBinding("region overflow".into()))?;
self.bytes
.get(region.offset..end)
.ok_or_else(|| WeightHandleError::DeviceBinding("region out of bounds".into()))
}
}
fn whole_layer_weights(
region_bytes: usize,
regions_per_layer: usize,
layers: usize,
) -> (Arc<LayeredMmap>, Vec<LazyWeight>) {
whole_layer_weights_salted(region_bytes, regions_per_layer, layers, 0)
}
fn whole_layer_weights_salted(
region_bytes: usize,
regions_per_layer: usize,
layers: usize,
salt: u8,
) -> (Arc<LayeredMmap>, Vec<LazyWeight>) {
let mapping_id = 42;
let prefix = 4096usize;
let mut backing = vec![0xABu8; prefix];
let mut lazies = Vec::with_capacity(layers);
let layer_bytes = region_bytes * regions_per_layer;
for layer in 0..layers {
let mut regions = Vec::with_capacity(regions_per_layer);
for r in 0..regions_per_layer {
let offset = backing.len();
let pattern = ((1 + layer * regions_per_layer + r) as u8 | 0x40) ^ salt;
backing.resize(offset + region_bytes, pattern);
regions.push(ExternalMmapRegion {
mapping_id,
offset,
len: region_bytes,
});
}
let dtype_shape = vec![layer_bytes];
let materialize_shape = dtype_shape.clone();
let lazy =
LazyWeight::block_quantized_moe(DataType::Uint8, dtype_shape, regions, move || {
ResidentWeight::new(
DataType::Uint8,
materialize_shape.clone(),
vec![0u8; layer_bytes],
)
})
.unwrap();
lazies.push(lazy);
}
(
Arc::new(LayeredMmap {
mapping_id,
bytes: backing,
}),
lazies,
)
}
fn canonical_layer(host: &LayeredMmap, lazy: &LazyWeight) -> Vec<u8> {
let mut out = Vec::with_capacity(lazy.region_bytes_len());
for region in &lazy.regions {
out.extend_from_slice(host.region_bytes(region).unwrap());
}
out
}
fn median(mut samples: Vec<Duration>) -> Duration {
samples.sort_unstable();
samples[samples.len() / 2]
}
#[cfg_attr(
not(feature = "gpu-tests"),
ignore = "requires CUDA device; enable the gpu-tests feature on a CUDA runner"
)]
#[test]
fn whole_layer_prefill_is_byte_identical_across_wraparound() {
let ep = require_cuda();
let runtime = ep.runtime().clone();
let region_bytes = 2 * 1024 * 1024;
let regions_per_layer = 2;
let layers = 4; let layer_bytes = (region_bytes * regions_per_layer) as u64;
let (host, lazies) = whole_layer_weights(region_bytes, regions_per_layer, layers);
let residency = ep.weight_residency(layer_bytes * layers as u64);
let mut db = residency
.build_prefill_double_buffer(layer_bytes)
.expect("capacity-sufficient synthetic layer must build the pipeline");
let mut pending: Option<LayerTicket> = Some(
db.prefetch(0, &PrefillLayerRequest::new(&lazies[0], &*host))
.unwrap(),
);
for layer in 0..layers {
let ticket = pending.take().expect("a prefetch is always pending here");
if layer + 1 < layers {
pending = Some(
db.prefetch(
(layer + 1) as u64,
&PrefillLayerRequest::new(&lazies[layer + 1], &*host),
)
.expect("look-ahead prefetch (incl. slot reuse) must succeed"),
);
}
let view = db.wait(&ticket).expect("ready slot must consume");
assert_eq!(view.len, layer_bytes as usize, "layer {layer} byte count");
runtime.sync_copy_stream().unwrap();
let mut readback = vec![0u8; view.len];
unsafe { runtime.dtoh(&mut readback, view.device_ptr).unwrap() };
assert_eq!(
readback,
canonical_layer(&host, &lazies[layer]),
"layer {layer} device bytes must be byte-identical to its concatenated source regions"
);
db.release(ticket).expect("in-use slot must release");
}
assert!(pending.is_none(), "every prefetched layer was consumed");
let metrics = db.metrics();
assert_eq!(metrics.layers_prefetched, layers as u64);
assert_eq!(metrics.layers_consumed, layers as u64);
assert_eq!(metrics.layers_released, layers as u64);
assert_eq!(metrics.poisoned, 0, "clean path poisons nothing");
assert_eq!(metrics.stale_rejected, 0);
assert_eq!(metrics.cancelled, 0);
assert_eq!(
db.transfer().quarantined_len(),
0,
"no slot may be quarantined on the clean path"
);
drop(db); }
#[cfg_attr(
not(feature = "gpu-tests"),
ignore = "requires CUDA device; enable the gpu-tests feature on a CUDA runner"
)]
#[test]
fn single_and_final_layer_teardown_is_clean() {
let ep = require_cuda();
let runtime = ep.runtime().clone();
let layer_bytes = 4 * 1024 * 1024u64;
let (host, lazies) = whole_layer_weights(layer_bytes as usize, 1, 1);
let residency = ep.weight_residency(layer_bytes);
let mut db = residency
.build_prefill_double_buffer(layer_bytes)
.expect("pipeline builds");
let ticket = db
.prefetch(0, &PrefillLayerRequest::new(&lazies[0], &*host))
.unwrap();
let view = db.wait(&ticket).unwrap();
runtime.sync_copy_stream().unwrap();
let mut readback = vec![0u8; view.len];
unsafe { runtime.dtoh(&mut readback, view.device_ptr).unwrap() };
assert_eq!(readback, canonical_layer(&host, &lazies[0]));
let other = if ticket_slot_is_zero(&db) { 1 } else { 0 };
assert_eq!(db.slot_status(other), PrefillSlotStatus::Free);
db.release(ticket).unwrap();
assert_eq!(db.transfer().quarantined_len(), 0);
drop(db);
}
fn ticket_slot_is_zero(
db: &PrefillDoubleBuffer<impl onnx_runtime_ep_cuda::prefill_double_buffer::PrefillTransfer>,
) -> bool {
!matches!(db.slot_status(0), PrefillSlotStatus::Free)
}
#[cfg_attr(
not(feature = "gpu-tests"),
ignore = "requires CUDA device; enable the gpu-tests feature on a CUDA runner"
)]
#[test]
fn cancellation_frees_slot_for_reuse_without_quarantine() {
let ep = require_cuda();
let runtime = ep.runtime().clone();
let layer_bytes = 4 * 1024 * 1024u64;
let (host, lazies) = whole_layer_weights(layer_bytes as usize, 1, 3);
let residency = ep.weight_residency(layer_bytes * 3);
let mut db = residency
.build_prefill_double_buffer(layer_bytes)
.expect("pipeline builds");
let keep = db
.prefetch(0, &PrefillLayerRequest::new(&lazies[0], &*host))
.unwrap();
let doomed = db
.prefetch(1, &PrefillLayerRequest::new(&lazies[1], &*host))
.unwrap();
let doomed_view = db.wait(&doomed).unwrap();
let doomed_slot = doomed_view.device_ptr;
db.cancel(doomed).unwrap();
assert_eq!(db.metrics().cancelled, 1);
assert_eq!(
db.transfer().quarantined_len(),
0,
"cancel does not quarantine"
);
let reuse = db
.prefetch(2, &PrefillLayerRequest::new(&lazies[2], &*host))
.expect("cancelled slot must be reusable");
let reuse_view = db.wait(&reuse).unwrap();
assert_eq!(
reuse_view.device_ptr, doomed_slot,
"the reused layer must land in the same stable device buffer the cancelled one used"
);
runtime.sync_copy_stream().unwrap();
let mut readback = vec![0u8; reuse_view.len];
unsafe { runtime.dtoh(&mut readback, reuse_view.device_ptr).unwrap() };
assert_eq!(
readback,
canonical_layer(&host, &lazies[2]),
"the reused slot's bytes must be the new layer's, not the cancelled one's"
);
db.release(reuse).unwrap();
db.wait(&keep).unwrap();
db.release(keep).unwrap();
assert_eq!(db.transfer().quarantined_len(), 0);
drop(db);
}
#[cfg_attr(
not(feature = "gpu-tests"),
ignore = "requires CUDA device; enable the gpu-tests feature on a CUDA runner"
)]
#[test]
fn stale_ticket_after_reuse_is_refused() {
let ep = require_cuda();
let layer_bytes = 2 * 1024 * 1024u64;
let (host, lazies) = whole_layer_weights(layer_bytes as usize, 1, 3);
let residency = ep.weight_residency(layer_bytes * 3);
let mut db = residency
.build_prefill_double_buffer(layer_bytes)
.expect("pipeline builds");
let ticket_b = db
.prefetch(0, &PrefillLayerRequest::new(&lazies[0], &*host))
.unwrap();
let ticket_a = db
.prefetch(1, &PrefillLayerRequest::new(&lazies[1], &*host))
.unwrap();
let _ = db.wait(&ticket_a).unwrap();
db.release(ticket_a.clone()).unwrap();
let reuse = db
.prefetch(2, &PrefillLayerRequest::new(&lazies[2], &*host))
.unwrap();
match db.wait(&ticket_a) {
Err(PrefillReject::StaleGeneration { layer_id }) => assert_eq!(layer_id, 1),
other => panic!("stale ticket must be refused, got {other:?}"),
}
assert!(db.metrics().stale_rejected >= 1);
db.wait(&reuse).unwrap();
db.release(reuse).unwrap();
db.wait(&ticket_b).unwrap();
db.release(ticket_b).unwrap();
drop(db);
}
#[cfg_attr(
not(feature = "gpu-tests"),
ignore = "requires CUDA device; enable the gpu-tests feature on a CUDA runner"
)]
#[test]
fn two_pipelines_on_one_device_are_isolated() {
let ep = require_cuda();
let runtime = ep.runtime().clone();
let region_bytes = 2 * 1024 * 1024;
let layer_bytes = (region_bytes * 2) as u64;
let layers = 3usize;
let (host_a, lazies_a) = whole_layer_weights_salted(region_bytes, 2, layers, 0x00);
let (host_b, lazies_b) = whole_layer_weights_salted(region_bytes, 2, layers, 0x24);
let residency = ep.weight_residency(layer_bytes * layers as u64 * 2);
let mut a = residency
.build_prefill_double_buffer(layer_bytes)
.expect("pipeline A builds");
let mut b = residency
.build_prefill_double_buffer(layer_bytes)
.expect("pipeline B builds");
let read = |db: &mut PrefillDoubleBuffer<
onnx_runtime_ep_cuda::prefill_double_buffer::CudaPrefillTransfer,
>,
layer: usize,
lazy: &LazyWeight,
host: &LayeredMmap| {
let t = db
.prefetch(layer as u64, &PrefillLayerRequest::new(lazy, host))
.unwrap();
let view = db.wait(&t).unwrap();
runtime.sync_copy_stream().unwrap();
let mut readback = vec![0u8; view.len];
unsafe { runtime.dtoh(&mut readback, view.device_ptr).unwrap() };
assert_eq!(
readback,
canonical_layer(host, lazy),
"a pipeline must read back only its own layer's bytes, never the other pipeline's"
);
db.release(t).unwrap();
};
for layer in 0..layers {
read(&mut a, layer, &lazies_a[layer], &host_a);
read(&mut b, layer, &lazies_b[layer], &host_b);
}
assert_eq!(a.transfer().quarantined_len(), 0);
assert_eq!(b.transfer().quarantined_len(), 0);
drop(a);
drop(b);
}
#[cfg_attr(
not(feature = "gpu-tests"),
ignore = "requires CUDA device; enable the gpu-tests feature on a CUDA runner"
)]
#[test]
fn oversized_layer_declines_pool_capacity_without_reserving() {
let ep = require_cuda();
let layer_bytes = 300 * 1024 * 1024u64;
let (_host, lazies) = whole_layer_weights(1024, 1, 1);
let residency = ep.weight_residency(layer_bytes);
match residency.build_prefill_double_buffer(layer_bytes) {
Err(PrefillReject::PoolCapacity {
layer_bytes: declined,
}) => assert_eq!(declined, layer_bytes),
other => panic!("oversized layer must decline via PoolCapacity, got {other:?}"),
}
let _ = lazies;
}
#[cfg_attr(
not(feature = "gpu-tests"),
ignore = "requires CUDA device; enable the gpu-tests feature on a CUDA runner"
)]
#[test]
fn disabled_by_default_declines_and_stays_byte_identical() {
let ep = require_cuda();
let layer_bytes = 2 * 1024 * 1024u64;
let (_host, lazies) = whole_layer_weights(layer_bytes as usize, 1, 1);
let residency = ep.weight_residency(layer_bytes);
let gated = residency.prefill_double_buffer(layer_bytes);
if prefill_double_buffer_enabled() {
assert!(
gated.is_ok(),
"with the env gate set, the gated entry must construct the pipeline"
);
} else {
match gated {
Err(PrefillReject::Disabled) => {}
other => panic!("default-off gate must return Disabled, got {other:?}"),
}
}
let _ = lazies;
}
fn assert_paged_bytes_identical(
runtime: &onnx_runtime_ep_cuda::runtime::CudaRuntime,
paged: &onnx_runtime_ep_api::PagedWeight,
host: &LayeredMmap,
lazy: &LazyWeight,
) {
runtime.sync_copy_stream().unwrap();
let mut readback = vec![0u8; paged.len()];
unsafe {
runtime
.dtoh(
&mut readback,
onnx_runtime_ep_cuda::runtime::cuptr(paged.device_ptr()),
)
.unwrap()
};
assert_eq!(readback, canonical_layer(host, lazy), "paged bytes");
}
#[cfg_attr(
not(feature = "gpu-tests"),
ignore = "requires CUDA device; enable the gpu-tests feature on a CUDA runner"
)]
#[test]
fn provider_double_buffers_across_wraparound_byte_identical() {
use onnx_runtime_ep_api::ExecutionProvider;
let mut ep = require_cuda();
let runtime = ep.runtime().clone();
let region_bytes = 2 * 1024 * 1024;
let regions_per_layer = 2;
let layers = 4; let layer_bytes = (region_bytes * regions_per_layer) as u64;
let (host, lazies) = whole_layer_weights(region_bytes, regions_per_layer, layers);
let residency = ep
.weight_residency(layer_bytes * layers as u64)
.with_prefill_double_buffer_enabled(true);
ep.install_residency_for_test(residency);
let source: &dyn MmapRegionSource = &*host;
assert!(
ep.prefetch_lazy_weight(0, &lazies[0], source).unwrap(),
"cold prime prefetch must route through the pipeline"
);
for layer in 0..layers {
if layer + 1 < layers {
assert!(
ep.prefetch_lazy_weight((layer + 1) as u64, &lazies[layer + 1], source)
.unwrap(),
"look-ahead prefetch of layer {} must route through the pipeline",
layer + 1
);
}
let paged = ep
.page_lazy_weight(layer as u64, &lazies[layer], source)
.unwrap()
.expect("an in-pipeline layer must page through the double buffer");
assert_eq!(
paged.len(),
layer_bytes as usize,
"layer {layer} byte count"
);
assert_paged_bytes_identical(&runtime, &paged, &host, &lazies[layer]);
drop(paged);
}
let residency = ep.residency().expect("installed above");
let stats = residency
.prefill_pipeline_stats()
.expect("an eligible layer built the pipeline");
assert_eq!(stats.routed, layers as u64, "every layer routed");
assert_eq!(stats.consumed, layers as u64, "every layer consumed");
assert_eq!(stats.released, layers as u64, "every slot released on drop");
assert_eq!(
residency.prefill_pipeline_quarantined(),
Some(0),
"clean path quarantines nothing"
);
}
#[cfg_attr(
not(feature = "gpu-tests"),
ignore = "requires CUDA device; enable the gpu-tests feature on a CUDA runner"
)]
#[test]
fn provider_single_layer_pages_and_releases() {
use onnx_runtime_ep_api::ExecutionProvider;
let mut ep = require_cuda();
let runtime = ep.runtime().clone();
let layer_bytes = 4 * 1024 * 1024u64;
let (host, lazies) = whole_layer_weights(layer_bytes as usize, 1, 1);
let residency = ep
.weight_residency(layer_bytes)
.with_prefill_double_buffer_enabled(true);
ep.install_residency_for_test(residency);
let source: &dyn MmapRegionSource = &*host;
assert!(ep.prefetch_lazy_weight(0, &lazies[0], source).unwrap());
let paged = ep
.page_lazy_weight(0, &lazies[0], source)
.unwrap()
.expect("single layer pages");
assert_eq!(paged.len(), layer_bytes as usize);
assert_paged_bytes_identical(&runtime, &paged, &host, &lazies[0]);
drop(paged);
let residency = ep.residency().unwrap();
let stats = residency.prefill_pipeline_stats().unwrap();
assert_eq!((stats.routed, stats.consumed, stats.released), (1, 1, 1));
assert_eq!(residency.prefill_pipeline_quarantined(), Some(0));
}
#[cfg_attr(
not(feature = "gpu-tests"),
ignore = "requires CUDA device; enable the gpu-tests feature on a CUDA runner"
)]
#[test]
fn disabled_provider_path_builds_no_pipeline() {
let ep = require_cuda();
let layer_bytes = 2 * 1024 * 1024u64;
let (host, lazies) = whole_layer_weights(layer_bytes as usize, 1, 2);
let residency = ep
.weight_residency(layer_bytes * 2)
.with_prefill_double_buffer_enabled(false);
let source: &dyn MmapRegionSource = &*host;
match residency.prefill_pipeline_prefetch(0, &lazies[0], source) {
PrefillRoute::Declined(PrefillWireDecline::Disabled) => {}
other => panic!("disabled residency must decline with Disabled, got {other:?}"),
}
assert!(
!residency.prefill_pipeline_active(),
"disabled path must not build the pipeline (no extra CUDA work)"
);
assert_eq!(
residency.prefill_pipeline_stats(),
None,
"no pipeline means no wiring stats"
);
}
#[cfg_attr(
not(feature = "gpu-tests"),
ignore = "requires CUDA device; enable the gpu-tests feature on a CUDA runner"
)]
#[test]
fn oversize_later_layer_declines_typed_and_pipeline_survives() {
use onnx_runtime_ep_api::ExecutionProvider;
let ep = require_cuda();
let runtime = ep.runtime().clone();
let small = 2 * 1024 * 1024usize;
let large = 8 * 1024 * 1024usize;
let (host_small, small_layer) = whole_layer_weights(small, 1, 1);
let (_host_large, large_layer) = whole_layer_weights(large, 1, 1);
let residency = Arc::new(
ep.weight_residency(large as u64 * 4)
.with_prefill_double_buffer_enabled(true),
);
let source_small: &dyn MmapRegionSource = &*host_small;
match residency.prefill_pipeline_prefetch(0, &small_layer[0], source_small) {
PrefillRoute::Prefetched => {}
other => panic!("first small layer must route, got {other:?}"),
}
match residency.prefill_pipeline_prefetch(1, &large_layer[0], source_small) {
PrefillRoute::Declined(PrefillWireDecline::Oversize {
layer_bytes,
capacity,
}) => {
assert_eq!(layer_bytes, large as u64);
assert_eq!(capacity, small as u64);
}
other => panic!("oversize later layer must decline via Oversize, got {other:?}"),
}
let stats = residency.prefill_pipeline_stats().unwrap();
assert_eq!(stats.routed, 1);
assert_eq!(stats.declined_oversize, 1);
let paged = residency
.prefill_pipeline_page(0, ep.device_id())
.unwrap()
.expect("the fitting layer still pages");
assert_paged_bytes_identical(&runtime, &paged, &host_small, &small_layer[0]);
drop(paged);
assert_eq!(residency.prefill_pipeline_quarantined(), Some(0));
}
#[cfg_attr(
not(feature = "gpu-tests"),
ignore = "requires CUDA device; enable the gpu-tests feature on a CUDA runner"
)]
#[test]
fn ineligible_boundary_declines_typed_without_building() {
let ep = require_cuda();
let region_bytes = 1024usize;
let (host, _bqmoe) = whole_layer_weights(region_bytes, 1, 1);
let region = ExternalMmapRegion {
mapping_id: 42,
offset: 4096,
len: region_bytes,
};
let materialize_len = region_bytes;
let matmul = LazyWeight::new(
LazyWeightBoundary::MatMul,
onnx_runtime_ir::DataType::Uint8,
vec![region_bytes],
vec![region],
move || {
ResidentWeight::new(
onnx_runtime_ir::DataType::Uint8,
vec![materialize_len],
vec![0u8; materialize_len],
)
},
)
.unwrap();
let residency = ep
.weight_residency(1 << 20)
.with_prefill_double_buffer_enabled(true);
let source: &dyn MmapRegionSource = &*host;
match residency.prefill_pipeline_prefetch(7, &matmul, source) {
PrefillRoute::Declined(PrefillWireDecline::Boundary) => {}
other => panic!("non-BQMoE weight must decline via Boundary, got {other:?}"),
}
assert!(
!residency.prefill_pipeline_active(),
"an ineligible-boundary decline must not build the pipeline"
);
}
#[cfg_attr(
not(feature = "gpu-tests"),
ignore = "requires CUDA device; enable the gpu-tests feature on a CUDA runner"
)]
#[test]
fn two_enabled_residencies_are_isolated() {
use onnx_runtime_ep_api::ExecutionProvider;
let ep = require_cuda();
let runtime = ep.runtime().clone();
let region_bytes = 2 * 1024 * 1024;
let layer_bytes = (region_bytes * 2) as u64;
let layers = 3usize;
let (host_a, lazies_a) = whole_layer_weights_salted(region_bytes, 2, layers, 0x00);
let (host_b, lazies_b) = whole_layer_weights_salted(region_bytes, 2, layers, 0x24);
let res_a = Arc::new(
ep.weight_residency(layer_bytes * layers as u64)
.with_prefill_double_buffer_enabled(true),
);
let res_b = Arc::new(
ep.weight_residency(layer_bytes * layers as u64)
.with_prefill_double_buffer_enabled(true),
);
let device = ep.device_id();
let read = |res: &Arc<onnx_runtime_ep_cuda::weight_paging::CudaWeightResidency>,
layer: usize,
lazy: &LazyWeight,
host: &LayeredMmap| {
let src: &dyn MmapRegionSource = host;
match res.prefill_pipeline_prefetch(layer as u64, lazy, src) {
PrefillRoute::Prefetched => {}
other => panic!("layer {layer} must route, got {other:?}"),
}
let paged = res
.prefill_pipeline_page(layer as u64, device)
.unwrap()
.expect("routed layer pages");
assert_paged_bytes_identical(&runtime, &paged, host, lazy);
drop(paged);
};
for layer in 0..layers {
read(&res_a, layer, &lazies_a[layer], &host_a);
read(&res_b, layer, &lazies_b[layer], &host_b);
}
assert_eq!(res_a.prefill_pipeline_quarantined(), Some(0));
assert_eq!(res_b.prefill_pipeline_quarantined(), Some(0));
assert_eq!(
res_a.prefill_pipeline_stats().unwrap().consumed,
layers as u64
);
assert_eq!(
res_b.prefill_pipeline_stats().unwrap().consumed,
layers as u64
);
}
#[cfg_attr(
not(feature = "gpu-tests"),
ignore = "requires CUDA device; enable the gpu-tests feature on a CUDA runner"
)]
#[test]
fn teardown_with_inflight_and_unconsumed_slots_no_leak() {
use onnx_runtime_ep_api::ExecutionProvider;
let ep = require_cuda();
let layer_bytes = 2 * 1024 * 1024u64;
let (host, lazies) = whole_layer_weights(layer_bytes as usize, 1, 2);
let residency = Arc::new(
ep.weight_residency(layer_bytes * 2)
.with_prefill_double_buffer_enabled(true),
);
let source: &dyn MmapRegionSource = &*host;
match residency.prefill_pipeline_prefetch(0, &lazies[0], source) {
PrefillRoute::Prefetched => {}
other => panic!("layer 0 must route, got {other:?}"),
}
match residency.prefill_pipeline_prefetch(1, &lazies[1], source) {
PrefillRoute::Prefetched => {}
other => panic!("layer 1 must route, got {other:?}"),
}
let paged0 = residency
.prefill_pipeline_page(0, ep.device_id())
.unwrap()
.expect("layer 0 pages");
assert_eq!(
residency.prefill_pipeline_quarantined(),
Some(0),
"an in-flight slot is not quarantined"
);
drop(residency);
drop(paged0);
}
struct ComputeProxy {
a: onnx_runtime_ep_api::DeviceBuffer,
b: onnx_runtime_ep_api::DeviceBuffer,
c: onnx_runtime_ep_api::DeviceBuffer,
kernel: Box<dyn onnx_runtime_ep_api::Kernel>,
m: usize,
k: usize,
n: usize,
}
impl ComputeProxy {
fn new(ep: &CudaExecutionProvider) -> Self {
use onnx_runtime_ep_api::ExecutionProvider;
let (m, k, n) = (2048usize, 4096usize, 4096usize);
let a = ep.allocate(m * k * 4, 256).unwrap();
let b = ep.allocate(k * n * 4, 256).unwrap();
let c = ep.allocate(m * n * 4, 256).unwrap();
let node = onnx_runtime_ir::Node::new(onnx_runtime_ir::NodeId(0), "MatMul", vec![], vec![]);
let kernel = ep.get_kernel(&node, &[vec![m, k], vec![k, n]], 17).unwrap();
Self {
a,
b,
c,
kernel,
m,
k,
n,
}
}
fn run(&mut self, ep: &CudaExecutionProvider) {
use onnx_runtime_ep_api::{
DevicePtr, DevicePtrMut, ExecutionProvider, TensorMut, TensorView,
};
use onnx_runtime_ir::compute_contiguous_strides;
let dev = ep.device_id();
let a_shape = [self.m, self.k];
let b_shape = [self.k, self.n];
let out_shape = [self.m, self.n];
let a_strides = compute_contiguous_strides(&a_shape);
let b_strides = compute_contiguous_strides(&b_shape);
let out_strides = compute_contiguous_strides(&out_shape);
let a_view = TensorView::new(
DevicePtr(self.a.as_ptr()),
DataType::Float32,
&a_shape,
&a_strides,
dev,
);
let b_view = TensorView::new(
DevicePtr(self.b.as_ptr()),
DataType::Float32,
&b_shape,
&b_strides,
dev,
);
let out_view = TensorMut::new(
DevicePtrMut(self.c.as_mut_ptr()),
DataType::Float32,
&out_shape,
&out_strides,
dev,
);
self.kernel
.execute(&[a_view, b_view], &mut [out_view])
.unwrap();
}
fn teardown(self, ep: &CudaExecutionProvider) {
use onnx_runtime_ep_api::ExecutionProvider;
ep.deallocate(self.a).unwrap();
ep.deallocate(self.b).unwrap();
ep.deallocate(self.c).unwrap();
}
}
fn run_serial_arm(
db: &mut PrefillDoubleBuffer<onnx_runtime_ep_cuda::prefill_double_buffer::CudaPrefillTransfer>,
lazies: &[LazyWeight],
source: &dyn MmapRegionSource,
proxy: &mut ComputeProxy,
ep: &CudaExecutionProvider,
) -> Duration {
let start = Instant::now();
for (layer, lazy) in lazies.iter().enumerate() {
let ticket = db
.prefetch(layer as u64, &PrefillLayerRequest::new(lazy, source))
.unwrap();
let _ = db.wait(&ticket).unwrap();
proxy.run(ep);
db.release(ticket).unwrap();
}
start.elapsed()
}
fn run_pipelined_arm(
db: &mut PrefillDoubleBuffer<onnx_runtime_ep_cuda::prefill_double_buffer::CudaPrefillTransfer>,
lazies: &[LazyWeight],
source: &dyn MmapRegionSource,
proxy: &mut ComputeProxy,
ep: &CudaExecutionProvider,
) -> Duration {
let start = Instant::now();
let mut pending = Some(
db.prefetch(0, &PrefillLayerRequest::new(&lazies[0], source))
.unwrap(),
);
for layer in 0..lazies.len() {
let ticket = pending.take().unwrap();
if layer + 1 < lazies.len() {
pending = Some(
db.prefetch(
(layer + 1) as u64,
&PrefillLayerRequest::new(&lazies[layer + 1], source),
)
.unwrap(),
);
}
let _ = db.wait(&ticket).unwrap();
proxy.run(ep);
db.release(ticket).unwrap();
}
start.elapsed()
}
#[ignore = "measurement probe (not a gate): run explicitly with --include-ignored on an idle CUDA GPU"]
#[test]
fn one_slot_vs_two_slot_overlap_probe() {
let ep = require_cuda();
let runtime = ep.runtime().clone();
let region_bytes = 32 * 1024 * 1024;
let regions_per_layer = 2;
let layers = 6usize;
let layer_bytes = (region_bytes * regions_per_layer) as u64;
let (host, lazies) = whole_layer_weights(region_bytes, regions_per_layer, layers);
let source = host.clone() as Arc<dyn MmapRegionSource>;
let mut proxy = ComputeProxy::new(&ep);
let ramp_end = Instant::now() + Duration::from_secs(8);
while Instant::now() < ramp_end {
proxy.run(&ep);
}
runtime.synchronize().unwrap();
let runs = 3usize;
let mut serial_host = Vec::with_capacity(runs);
let mut serial_total = Vec::with_capacity(runs);
let mut pipe_host = Vec::with_capacity(runs);
let mut pipe_total = Vec::with_capacity(runs);
let mut pipe_reuse_wait_ns = 0u64;
let mut cold_reserve_us = 0.0f64;
for run_idx in 0..runs {
let residency = ep.weight_residency(layer_bytes * layers as u64);
let reserve_start = Instant::now();
let mut db = residency
.build_prefill_double_buffer(layer_bytes)
.expect("capacity-sufficient synthetic build");
if run_idx == 0 {
cold_reserve_us = reserve_start.elapsed().as_secs_f64() * 1e6;
}
let host_wall = run_serial_arm(&mut db, &lazies, source.as_ref(), &mut proxy, &ep);
let drain_start = Instant::now();
runtime.synchronize().unwrap();
let total_wall = host_wall + drain_start.elapsed();
serial_host.push(host_wall);
serial_total.push(total_wall);
drop(db);
let residency = ep.weight_residency(layer_bytes * layers as u64);
let mut db = residency
.build_prefill_double_buffer(layer_bytes)
.expect("capacity-sufficient synthetic build");
let host_wall = run_pipelined_arm(&mut db, &lazies, source.as_ref(), &mut proxy, &ep);
let drain_start = Instant::now();
runtime.synchronize().unwrap();
let total_wall = host_wall + drain_start.elapsed();
pipe_host.push(host_wall);
pipe_total.push(total_wall);
pipe_reuse_wait_ns = db.metrics().reuse_wait_ns;
assert_eq!(db.transfer().quarantined_len(), 0);
drop(db);
}
let residency = ep.weight_residency(layer_bytes * layers as u64);
let mut db = residency.build_prefill_double_buffer(layer_bytes).unwrap();
let _ = run_pipelined_arm(&mut db, &lazies, source.as_ref(), &mut proxy, &ep);
runtime.synchronize().unwrap();
let recheck_reuse_wait_ns = db.metrics().reuse_wait_ns;
drop(db);
let serial_host_us = median(serial_host).as_secs_f64() * 1e6;
let serial_total_us = median(serial_total).as_secs_f64() * 1e6;
let pipe_host_us = median(pipe_host).as_secs_f64() * 1e6;
let pipe_total_us = median(pipe_total).as_secs_f64() * 1e6;
println!(
"PREFILL_DB_OVERLAP layer_MiB={:.0} regions_per_layer={} layers={} runs={} \
serial_host_us={:.1} serial_total_us={:.1} pipe_host_us={:.1} pipe_total_us={:.1} \
pipe_reuse_wait_us={:.1} recheck_reuse_wait_us={:.1} cold_reserve_us={:.1} \
note=no_tokps_claim;host_enqueue_and_event_wait_reported_separately",
layer_bytes as f64 / (1024.0 * 1024.0),
regions_per_layer,
layers,
runs,
serial_host_us,
serial_total_us,
pipe_host_us,
pipe_total_us,
pipe_reuse_wait_ns as f64 / 1e3,
recheck_reuse_wait_ns as f64 / 1e3,
cold_reserve_us,
);
proxy.teardown(&ep);
}