From 99e0f1589aaf90679813af7ab86409eaf1bf9715 Mon Sep 17 00:00:00 2001 From: Andy Jost Date: Mon, 17 Aug 2026 13:58:46 -0700 Subject: [PATCH 1/2] test(cuda.core): synchronize IPC buffer initialization Ensure the exporting process completes asynchronous allocation and initialization before an importing child accesses the shared buffer. --- cuda_core/tests/memory_ipc/test_peer_access.py | 2 ++ 1 file changed, 2 insertions(+) diff --git a/cuda_core/tests/memory_ipc/test_peer_access.py b/cuda_core/tests/memory_ipc/test_peer_access.py index 4dc04a8bd0b..c54ace82c54 100644 --- a/cuda_core/tests/memory_ipc/test_peer_access.py +++ b/cuda_core/tests/memory_ipc/test_peer_access.py @@ -82,6 +82,8 @@ def test_main(self, ipc_mempool_device_x2, grant_access_in_parent): buffer = mr.allocate(NBYTES, stream=dev1.default_stream) pgen = PatternGen(dev1, NBYTES) pgen.fill_buffer(buffer, seed=False) + # IPC import does not carry the exporting process's stream ordering. + dev1.sync() # Spawn child process process = mp.Process(target=self.child_main, args=(mr, buffer)) From a35152d9537f07fadad3b983fd0da2bea723e5aa Mon Sep 17 00:00:00 2001 From: Andy Jost Date: Fri, 21 Aug 2026 11:53:13 -0700 Subject: [PATCH 2/2] fix(cuda.core): use consistent streams for async allocations Require PatternGen callers and affected examples/tests to preserve stream ordering, adding explicit synchronization only across host and IPC boundaries. --- cuda_core/examples/memory_pool_resources.py | 2 +- .../strided_memory_view_constructors.py | 7 ++- cuda_core/tests/helpers/buffers.py | 11 ++-- cuda_core/tests/memory/test_managed_ops.py | 15 +++--- cuda_core/tests/memory_ipc/test_errors.py | 18 +++++-- .../memory_ipc/test_ipc_duplicate_import.py | 4 +- cuda_core/tests/memory_ipc/test_leaks.py | 4 +- cuda_core/tests/memory_ipc/test_memory_ipc.py | 52 +++++++++++++------ .../tests/memory_ipc/test_peer_access.py | 17 +++--- .../tests/memory_ipc/test_send_buffers.py | 30 +++++++---- cuda_core/tests/memory_ipc/test_serialize.py | 37 +++++++++---- cuda_core/tests/memory_ipc/test_workerpool.py | 41 +++++++++------ cuda_core/tests/test_helpers.py | 5 +- cuda_core/tests/test_memory.py | 12 ++--- cuda_core/tests/test_memory_peer_access.py | 8 +-- cuda_core/tests/test_module.py | 4 +- 16 files changed, 169 insertions(+), 98 deletions(-) diff --git a/cuda_core/examples/memory_pool_resources.py b/cuda_core/examples/memory_pool_resources.py index aa322ecc206..8b97948c54e 100644 --- a/cuda_core/examples/memory_pool_resources.py +++ b/cuda_core/examples/memory_pool_resources.py @@ -89,13 +89,13 @@ def main(): managed_buffer = managed_mr.allocate(nbytes, stream=stream) pinned_buffer = pinned_mr.allocate(nbytes, stream=stream) + stream.sync() managed_array = np.from_dlpack(managed_buffer).view(np.float32) pinned_array = np.from_dlpack(pinned_buffer).view(np.float32) managed_array[:] = np.arange(size, dtype=dtype) managed_original = managed_array.copy() - stream.sync() managed_buffer.copy_to(pinned_buffer, stream=stream) stream.sync() diff --git a/cuda_core/examples/strided_memory_view_constructors.py b/cuda_core/examples/strided_memory_view_constructors.py index 66820c56e2d..c5799ad998b 100644 --- a/cuda_core/examples/strided_memory_view_constructors.py +++ b/cuda_core/examples/strided_memory_view_constructors.py @@ -18,7 +18,7 @@ import cupy as cp import numpy as np -from cuda.core import Device +from cuda.core import Device, Stream from cuda.core.utils import StridedMemoryView @@ -39,7 +39,8 @@ def main(): device = Device() device.set_current() - stream = device.create_stream() + cupy_stream = cp.cuda.get_current_stream() + stream = Stream.from_handle(cupy_stream.ptr) buffer = None try: @@ -63,7 +64,6 @@ def main(): buffer = device.memory_resource.allocate(gpu_array.nbytes, stream=stream) buffer_array = cp.from_dlpack(buffer).view(dtype=cp.float32).reshape(gpu_array.shape) buffer_array[...] = gpu_array - device.sync() buffer_view = StridedMemoryView.from_buffer( buffer, @@ -77,7 +77,6 @@ def main(): finally: if buffer is not None: buffer.close(stream) - stream.close() if __name__ == "__main__": diff --git a/cuda_core/tests/helpers/buffers.py b/cuda_core/tests/helpers/buffers.py index aefb4f04193..ffd465695a2 100644 --- a/cuda_core/tests/helpers/buffers.py +++ b/cuda_core/tests/helpers/buffers.py @@ -146,8 +146,8 @@ class PatternGen: Provides methods to fill a target buffer with known test patterns and verify the expected values. - If a stream is provided, operations are synchronized with respect to that - stream. Otherwise, they are synchronized over the device. + Operations are submitted to the supplied stream. Verification synchronizes + that stream before comparing results on the host. The test pattern is either a fixed value or a cyclic pattern generated from an 8-bit seed. Only one of `value` or `seed` should be supplied. @@ -158,11 +158,10 @@ class PatternGen: buffer and then perform a comparison. """ - def __init__(self, device, size, stream=None): + def __init__(self, device, size, *, stream): self.device = device self.size = size - self.stream = stream if stream is not None else device.create_stream() - self.sync_target = stream if stream is not None else device + self.stream = Stream_accept(stream) self.pattern_buffers = {} def fill_buffer(self, buffer, seed=None, value=None): @@ -179,7 +178,7 @@ def verify_buffer(self, buffer, seed=None, value=None): pattern_buffer = self._get_pattern_buffer(seed, value) ptr_expected = self._ptr(pattern_buffer) scratch_buffer.copy_from(buffer, stream=self.stream) - self.sync_target.sync() + self.stream.sync() assert libc.memcmp(ptr_test, ptr_expected, self.size) == 0 @staticmethod diff --git a/cuda_core/tests/memory/test_managed_ops.py b/cuda_core/tests/memory/test_managed_ops.py index ed7f44a97f4..2c2ceb6e8e1 100644 --- a/cuda_core/tests/memory/test_managed_ops.py +++ b/cuda_core/tests/memory/test_managed_ops.py @@ -104,6 +104,7 @@ def managed_buffer(request, location_ops_device, location_ops_mr): size = _MANAGED_TEST_ALLOCATION_SIZE if request.param == "pool": buf = location_ops_mr.allocate(size, stream=location_ops_device.default_stream) + location_ops_device.default_stream.sync() yield buf buf.close() else: @@ -215,8 +216,8 @@ def test_same_location(self, location_ops_device, location_ops_mr): from cuda.core.utils import prefetch_batch device = location_ops_device - bufs = [location_ops_mr.allocate(_MANAGED_TEST_ALLOCATION_SIZE, stream=device.default_stream) for _ in range(3)] stream = device.create_stream() + bufs = [location_ops_mr.allocate(_MANAGED_TEST_ALLOCATION_SIZE, stream=stream) for _ in range(3)] prefetch_batch(stream, bufs, device) stream.sync() @@ -230,12 +231,12 @@ def test_per_buffer_location(self, location_ops_device, location_ops_mr): from cuda.core.utils import prefetch_batch device = location_ops_device - bufs = [location_ops_mr.allocate(_MANAGED_TEST_ALLOCATION_SIZE, stream=device.default_stream) for _ in range(2)] + stream = device.create_stream() + bufs = [location_ops_mr.allocate(_MANAGED_TEST_ALLOCATION_SIZE, stream=stream) for _ in range(2)] # Per-buffer prefetch locations are only observable when the buffers sit # on distinct physical pages; assert that here so a pool-packing change # fails loudly instead of silently migrating one shared page. assert _page_base(bufs[0]) != _page_base(bufs[1]) - stream = device.create_stream() prefetch_batch(stream, bufs, [Host(), device]) stream.sync() @@ -257,8 +258,8 @@ def test_basic(self, location_ops_device, location_ops_mr): if not hasattr(driver, "cuMemDiscardBatchAsync"): pytest.skip("cuMemDiscardBatchAsync unavailable") device = location_ops_device - bufs = [location_ops_mr.allocate(_MANAGED_TEST_ALLOCATION_SIZE, stream=device.default_stream) for _ in range(3)] stream = device.create_stream() + bufs = [location_ops_mr.allocate(_MANAGED_TEST_ALLOCATION_SIZE, stream=stream) for _ in range(3)] prefetch_batch(stream, bufs, device) stream.sync() discard_batch(stream, bufs) @@ -276,8 +277,8 @@ def test_same_location(self, location_ops_device, location_ops_mr): if not hasattr(driver, "cuMemDiscardAndPrefetchBatchAsync"): pytest.skip("cuMemDiscardAndPrefetchBatchAsync unavailable") device = location_ops_device - bufs = [location_ops_mr.allocate(_MANAGED_TEST_ALLOCATION_SIZE, stream=device.default_stream) for _ in range(2)] stream = device.create_stream() + bufs = [location_ops_mr.allocate(_MANAGED_TEST_ALLOCATION_SIZE, stream=stream) for _ in range(2)] prefetch_batch(stream, bufs, Host()) stream.sync() discard_prefetch_batch(stream, bufs, device) @@ -486,9 +487,9 @@ def test_instance_discard(self, location_ops_device, managed_buffer): def test_instance_discard_prefetch(self, discard_prefetch_device): device = discard_prefetch_device mr = create_managed_memory_resource_or_skip() - buf = mr.allocate(_MANAGED_TEST_ALLOCATION_SIZE, stream=device.default_stream) + stream = device.create_stream() + buf = mr.allocate(_MANAGED_TEST_ALLOCATION_SIZE, stream=stream) try: - stream = device.create_stream() buf.prefetch(Host(), stream=stream) stream.sync() buf.discard_prefetch(device, stream=stream) diff --git a/cuda_core/tests/memory_ipc/test_errors.py b/cuda_core/tests/memory_ipc/test_errors.py index 0aac9f9a297..9cb5286257a 100644 --- a/cuda_core/tests/memory_ipc/test_errors.py +++ b/cuda_core/tests/memory_ipc/test_errors.py @@ -104,9 +104,11 @@ class TestImportOversizedBufferDescriptorSize(ChildErrorHarness): """Reject peer-supplied sizes larger than the mapped allocation extent.""" def PARENT_ACTION(self, queue): - self.buffer = self.mr.allocate(NBYTES, stream=self.device.default_stream) + stream = self.device.default_stream + self.buffer = self.mr.allocate(NBYTES, stream=stream) payload, _ = self.buffer.ipc_descriptor.__reduce__()[1] oversized = IPCBufferDescriptor._init(payload, NBYTES * 100) + stream.sync() queue.put(oversized) def CHILD_ACTION(self, queue): @@ -141,8 +143,10 @@ def PARENT_ACTION(self, queue): options = DeviceMemoryResourceOptions(max_size=POOL_SIZE, ipc_enabled=True) mr2 = DeviceMemoryResource(self.device, options=options) self._extra_mrs.append(mr2) - buffer = mr2.allocate(NBYTES, stream=self.device.default_stream) - queue.put([self.mr, buffer.ipc_descriptor]) # Note: mr does not own this buffer + stream = self.device.default_stream + self.buffer = mr2.allocate(NBYTES, stream=stream) + stream.sync() + queue.put([self.mr, self.buffer.ipc_descriptor]) # Note: mr does not own this buffer def CHILD_ACTION(self, queue): mr, buffer_desc = queue.get(timeout=CHILD_TIMEOUT_SEC) @@ -159,7 +163,9 @@ class TestImportBuffer(ChildErrorHarness): def PARENT_ACTION(self, queue): # Note: if the buffer is not attached to something to prolong its life, # CUDA_ERROR_INVALID_CONTEXT is raised from Buffer.__del__ - self.buffer = self.mr.allocate(NBYTES, stream=self.device.default_stream) + stream = self.device.default_stream + self.buffer = self.mr.allocate(NBYTES, stream=stream) + stream.sync() queue.put(self.buffer) def CHILD_ACTION(self, queue): @@ -181,8 +187,10 @@ def PARENT_ACTION(self, queue): options = DeviceMemoryResourceOptions(max_size=POOL_SIZE, ipc_enabled=True) mr2 = DeviceMemoryResource(self.device, options=options) self._extra_mrs.append(mr2) - self.buffer = mr2.allocate(NBYTES, stream=self.device.default_stream) + stream = self.device.default_stream + self.buffer = mr2.allocate(NBYTES, stream=stream) buffer_s = pickle.dumps(self.buffer) + stream.sync() queue.put(buffer_s) # Note: mr2 not sent def CHILD_ACTION(self, queue): diff --git a/cuda_core/tests/memory_ipc/test_ipc_duplicate_import.py b/cuda_core/tests/memory_ipc/test_ipc_duplicate_import.py index dc3f5e57c33..771f26a3399 100644 --- a/cuda_core/tests/memory_ipc/test_ipc_duplicate_import.py +++ b/cuda_core/tests/memory_ipc/test_ipc_duplicate_import.py @@ -71,7 +71,9 @@ def test_main(self, ipc_device, ipc_memory_resource): mr = ipc_memory_resource log("allocating buffer") - buffer = mr.allocate(NBYTES, stream=ipc_device.default_stream) + stream = ipc_device.default_stream + buffer = mr.allocate(NBYTES, stream=stream) + stream.sync() # Start the child process. log("starting child") diff --git a/cuda_core/tests/memory_ipc/test_leaks.py b/cuda_core/tests/memory_ipc/test_leaks.py index c6e44824137..bca7948737c 100644 --- a/cuda_core/tests/memory_ipc/test_leaks.py +++ b/cuda_core/tests/memory_ipc/test_leaks.py @@ -102,8 +102,10 @@ def __reduce__(self): def test_pass_object(ipc_device, ipc_memory_resource, launcher, getobject): """Check for fd leaks when an object is sent as a subprocess argument.""" mr = ipc_memory_resource + stream = ipc_device.default_stream with CheckFDLeaks(): - obj = getobject(mr, ipc_device.default_stream) + obj = getobject(mr, stream) + stream.sync() try: launcher(obj, number=2) finally: diff --git a/cuda_core/tests/memory_ipc/test_memory_ipc.py b/cuda_core/tests/memory_ipc/test_memory_ipc.py index 43d356789e7..0eace7f154c 100644 --- a/cuda_core/tests/memory_ipc/test_memory_ipc.py +++ b/cuda_core/tests/memory_ipc/test_memory_ipc.py @@ -26,7 +26,8 @@ def test_main(self, ipc_device, ipc_memory_resource): device = ipc_device mr = ipc_memory_resource assert not mr.is_mapped - pgen = PatternGen(device, NBYTES) + stream = device.default_stream + pgen = PatternGen(device, NBYTES, stream=stream) # Start the child process. queue = mp.Queue() @@ -34,9 +35,10 @@ def test_main(self, ipc_device, ipc_memory_resource): process.start() # Allocate and fill memory. - buffer = mr.allocate(NBYTES, stream=device.default_stream) + buffer = mr.allocate(NBYTES, stream=stream) assert not buffer.is_mapped pgen.fill_buffer(buffer, seed=False) + stream.sync() # Export the buffer via IPC. queue.put(buffer) @@ -50,16 +52,19 @@ def test_main(self, ipc_device, ipc_memory_resource): # Verify that the buffer was modified. pgen.verify_buffer(buffer, seed=True) buffer.close() + stream.sync() def child_main(self, device, mr, queue): device.set_current() assert mr.is_mapped buffer = queue.get(timeout=CHILD_TIMEOUT_SEC) assert buffer.is_mapped - pgen = PatternGen(device, NBYTES) + stream = device.default_stream + pgen = PatternGen(device, NBYTES, stream=stream) pgen.verify_buffer(buffer, seed=False) pgen.fill_buffer(buffer, seed=True) buffer.close() + stream.sync() class TestIPCMempoolMultiple: @@ -70,14 +75,16 @@ def test_main(self, ipc_device, ipc_memory_resource): device = ipc_device mr = ipc_memory_resource q1, q2 = (mp.Queue() for _ in range(2)) + stream = device.default_stream # Allocate memory buffers and export them to each child. - buffer1 = mr.allocate(NBYTES, stream=device.default_stream) + buffer1 = mr.allocate(NBYTES, stream=stream) q1.put(buffer1) q2.put(buffer1) - buffer2 = mr.allocate(NBYTES, stream=device.default_stream) + buffer2 = mr.allocate(NBYTES, stream=stream) q1.put(buffer2) q2.put(buffer2) + stream.sync() # Start the child processes. p1 = mp.Process(target=self.child_main, args=(device, mr, 1, q1)) @@ -94,11 +101,12 @@ def test_main(self, ipc_device, ipc_memory_resource): assert p2.exitcode == 0 # Verify that the buffers were modified. - pgen = PatternGen(device, NBYTES) + pgen = PatternGen(device, NBYTES, stream=stream) pgen.verify_buffer(buffer1, seed=1) pgen.verify_buffer(buffer2, seed=2) buffer1.close() buffer2.close() + stream.sync() def child_main(self, device, mr, seed, queue): # Note: passing the mr registers it so that buffers can be passed @@ -106,13 +114,15 @@ def child_main(self, device, mr, seed, queue): device.set_current() buffer1 = queue.get(timeout=CHILD_TIMEOUT_SEC) buffer2 = queue.get(timeout=CHILD_TIMEOUT_SEC) - pgen = PatternGen(device, NBYTES) + stream = device.default_stream + pgen = PatternGen(device, NBYTES, stream=stream) if seed == 1: pgen.fill_buffer(buffer1, seed=1) elif seed == 2: pgen.fill_buffer(buffer2, seed=2) buffer1.close() buffer2.close() + stream.sync() class TestIPCSharedAllocationHandleAndBufferDescriptors: @@ -126,6 +136,7 @@ def test_main(self, ipc_device, ipc_memory_resource): device = ipc_device mr = ipc_memory_resource alloc_handle = mr.allocation_handle + stream = device.default_stream # Start children. q1, q2 = (mp.Queue() for _ in range(2)) @@ -135,8 +146,9 @@ def test_main(self, ipc_device, ipc_memory_resource): p2.start() # Allocate and share memory. - buffer1 = mr.allocate(NBYTES, stream=device.default_stream) - buffer2 = mr.allocate(NBYTES, stream=device.default_stream) + buffer1 = mr.allocate(NBYTES, stream=stream) + buffer2 = mr.allocate(NBYTES, stream=stream) + stream.sync() q1.put(buffer1.ipc_descriptor) q2.put(buffer2.ipc_descriptor) @@ -149,11 +161,12 @@ def test_main(self, ipc_device, ipc_memory_resource): assert p2.exitcode == 0 # Verify results. - pgen = PatternGen(device, NBYTES) + pgen = PatternGen(device, NBYTES, stream=stream) pgen.verify_buffer(buffer1, seed=False) pgen.verify_buffer(buffer2, seed=True) buffer1.close() buffer2.close() + stream.sync() def child_main(self, device, alloc_handle, seed, queue): """Fills a shared memory buffer.""" @@ -162,10 +175,12 @@ def child_main(self, device, alloc_handle, seed, queue): device.set_current() mr = DeviceMemoryResource.from_allocation_handle(device, alloc_handle) buffer_descriptor = queue.get(timeout=CHILD_TIMEOUT_SEC) - buffer = Buffer.from_ipc_descriptor(mr, buffer_descriptor, stream=device.default_stream) - pgen = PatternGen(device, NBYTES) + stream = device.default_stream + buffer = Buffer.from_ipc_descriptor(mr, buffer_descriptor, stream=stream) + pgen = PatternGen(device, NBYTES, stream=stream) pgen.fill_buffer(buffer, seed=seed) buffer.close() + stream.sync() class TestIPCSharedAllocationHandleAndBufferObjects: @@ -178,6 +193,7 @@ def test_main(self, ipc_device, ipc_memory_resource): device = ipc_device mr = ipc_memory_resource alloc_handle = mr.allocation_handle + stream = device.default_stream # Start children. q1, q2 = (mp.Queue() for _ in range(2)) @@ -187,8 +203,9 @@ def test_main(self, ipc_device, ipc_memory_resource): p2.start() # Allocate and share memory. - buffer1 = mr.allocate(NBYTES, stream=device.default_stream) - buffer2 = mr.allocate(NBYTES, stream=device.default_stream) + buffer1 = mr.allocate(NBYTES, stream=stream) + buffer2 = mr.allocate(NBYTES, stream=stream) + stream.sync() q1.put(buffer1) q2.put(buffer2) @@ -201,11 +218,12 @@ def test_main(self, ipc_device, ipc_memory_resource): assert p2.exitcode == 0 # Verify results. - pgen = PatternGen(device, NBYTES) + pgen = PatternGen(device, NBYTES, stream=stream) pgen.verify_buffer(buffer1, seed=False) pgen.verify_buffer(buffer2, seed=True) buffer1.close() buffer2.close() + stream.sync() def child_main(self, device, alloc_handle, seed, queue): """Fills a shared memory buffer.""" @@ -216,6 +234,8 @@ def child_main(self, device, alloc_handle, seed, queue): # Now get buffers. buffer = queue.get(timeout=CHILD_TIMEOUT_SEC) - pgen = PatternGen(device, NBYTES) + stream = device.default_stream + pgen = PatternGen(device, NBYTES, stream=stream) pgen.fill_buffer(buffer, seed=seed) buffer.close() + stream.sync() diff --git a/cuda_core/tests/memory_ipc/test_peer_access.py b/cuda_core/tests/memory_ipc/test_peer_access.py index c54ace82c54..cd6977a5c64 100644 --- a/cuda_core/tests/memory_ipc/test_peer_access.py +++ b/cuda_core/tests/memory_ipc/test_peer_access.py @@ -79,11 +79,12 @@ def test_main(self, ipc_mempool_device_x2, grant_access_in_parent): assert mr.peer_accessible_by == {dev0} else: assert mr.peer_accessible_by == set() - buffer = mr.allocate(NBYTES, stream=dev1.default_stream) - pgen = PatternGen(dev1, NBYTES) + stream = dev1.default_stream + buffer = mr.allocate(NBYTES, stream=stream) + pgen = PatternGen(dev1, NBYTES, stream=stream) pgen.fill_buffer(buffer, seed=False) # IPC import does not carry the exporting process's stream ordering. - dev1.sync() + stream.sync() # Spawn child process process = mp.Process(target=self.child_main, args=(mr, buffer)) @@ -109,20 +110,22 @@ def child_main(self, mr, buffer): # Test 1: Buffer accessible from resident device (dev1) - should always work dev1 = Device(1) dev1.set_current() - PatternGen(dev1, NBYTES).verify_buffer(buffer, seed=False) + stream1 = dev1.default_stream + PatternGen(dev1, NBYTES, stream=stream1).verify_buffer(buffer, seed=False) # Test 2: Buffer NOT accessible from dev0 initially (peer access not preserved) dev0 = Device(0) dev0.set_current() + stream0 = dev0.default_stream with pytest.raises(CUDAError, match="CUDA_ERROR_INVALID_VALUE"): - PatternGen(dev0, NBYTES).verify_buffer(buffer, seed=False) + PatternGen(dev0, NBYTES, stream=stream0).verify_buffer(buffer, seed=False) # Test 3: Set peer access and verify buffer becomes accessible dev1.set_current() mr.peer_accessible_by = [0] assert mr.peer_accessible_by == {dev0} dev0.set_current() - PatternGen(dev0, NBYTES).verify_buffer(buffer, seed=False) + PatternGen(dev0, NBYTES, stream=stream0).verify_buffer(buffer, seed=False) # Test 4: Revoke peer access and verify buffer becomes inaccessible dev1.set_current() @@ -130,7 +133,7 @@ def child_main(self, mr, buffer): assert mr.peer_accessible_by == set() dev0.set_current() with pytest.raises(CUDAError, match="CUDA_ERROR_INVALID_VALUE"): - PatternGen(dev0, NBYTES).verify_buffer(buffer, seed=False) + PatternGen(dev0, NBYTES, stream=stream0).verify_buffer(buffer, seed=False) buffer.close() # TODO(seberg): 2026-06: mr close may be unsafe with incomplete `buf.close()` diff --git a/cuda_core/tests/memory_ipc/test_send_buffers.py b/cuda_core/tests/memory_ipc/test_send_buffers.py index efa4d8b2abc..50666df2f93 100644 --- a/cuda_core/tests/memory_ipc/test_send_buffers.py +++ b/cuda_core/tests/memory_ipc/test_send_buffers.py @@ -30,13 +30,15 @@ def test_main(self, ipc_device, nmrs): options = DeviceMemoryResourceOptions(max_size=POOL_SIZE, ipc_enabled=True) mrs = [DeviceMemoryResource(device, options=options) for _ in range(nmrs)] buffers = [] + stream = device.default_stream try: # Allocate and fill memory. - buffers = [mr.allocate(NBYTES, stream=device.default_stream) for mr, _ in zip(cycle(mrs), range(NTASKS))] - pgen = PatternGen(device, NBYTES) + buffers = [mr.allocate(NBYTES, stream=stream) for mr, _ in zip(cycle(mrs), range(NTASKS))] + pgen = PatternGen(device, NBYTES, stream=stream) for buffer in buffers: pgen.fill_buffer(buffer, seed=False) + stream.sync() # Start the child process. process = mp.Process(target=self.child_main, args=(device, buffers)) @@ -50,25 +52,26 @@ def test_main(self, ipc_device, nmrs): assert process.exitcode == 0 # Verify that the buffers were modified. - pgen = PatternGen(device, NBYTES) + pgen = PatternGen(device, NBYTES, stream=stream) for buffer in buffers: pgen.verify_buffer(buffer, seed=True) buffer.close() finally: for buffer in buffers: buffer.close() - # TODO(seberg): 2026-06: mr close may be unsafe with incomplete `buf.close()` - device.sync() + stream.sync() for mr in mrs: mr.close() def child_main(self, device, buffers): device.set_current() - pgen = PatternGen(device, NBYTES) + stream = device.default_stream + pgen = PatternGen(device, NBYTES, stream=stream) for buffer in buffers: pgen.verify_buffer(buffer, seed=False) pgen.fill_buffer(buffer, seed=True) buffer.close() + stream.sync() class TestIpcReexport: @@ -93,9 +96,11 @@ def test_main(self, ipc_device, ipc_memory_resource): # Allocate, fill a buffer. mr = ipc_memory_resource - pgen = PatternGen(device, NBYTES) - buffer = mr.allocate(NBYTES, stream=device.default_stream) + stream = device.default_stream + pgen = PatternGen(device, NBYTES, stream=stream) + buffer = mr.allocate(NBYTES, stream=stream) pgen.fill_buffer(buffer, seed=0) + stream.sync() # Set up communication. q_bc = mp.Queue() @@ -123,18 +128,21 @@ def test_main(self, ipc_device, ipc_memory_resource): # Verify that C’s operations are visible. pgen.verify_buffer(buffer, seed=1) buffer.close() + stream.sync() def process_b_main(self, buffer, q_bc, event_b): # Process B: receive buffer from A then forward it to C. device = Device() device.set_current() + stream = device.default_stream # Forward the buffer to C. q_bc.put(buffer) - buffer.close() # Wait for C to receive before exiting. event_b.wait(timeout=CHILD_TIMEOUT_SEC) + buffer.close() + stream.sync() def process_c_main(self, q_bc, event_c): # Process C: receive buffer from B then fill it. @@ -143,9 +151,11 @@ def process_c_main(self, q_bc, event_c): # Get the buffer and fill it. buffer = q_bc.get(timeout=CHILD_TIMEOUT_SEC) - pgen = PatternGen(device, NBYTES) + stream = device.default_stream + pgen = PatternGen(device, NBYTES, stream=stream) pgen.fill_buffer(buffer, seed=1) buffer.close() + stream.sync() # Signal A that the work is complete. event_c.set() diff --git a/cuda_core/tests/memory_ipc/test_serialize.py b/cuda_core/tests/memory_ipc/test_serialize.py index 22596582c49..f2805c31e63 100644 --- a/cuda_core/tests/memory_ipc/test_serialize.py +++ b/cuda_core/tests/memory_ipc/test_serialize.py @@ -30,6 +30,7 @@ class TestObjectSerializationDirect: def test_main(self, ipc_device, ipc_memory_resource): device = ipc_device mr = ipc_memory_resource + stream = device.default_stream # Start the child process. parent_conn, child_conn = mp.Pipe() @@ -41,10 +42,11 @@ def test_main(self, ipc_device, ipc_memory_resource): mp.reduction.send_handle(parent_conn, alloc_handle.handle, process.pid) # Send a buffer. - buffer1 = mr.allocate(NBYTES, stream=device.default_stream) + buffer1 = mr.allocate(NBYTES, stream=stream) parent_conn.send(buffer1) # directly - buffer2 = mr.allocate(NBYTES, stream=device.default_stream) + buffer2 = mr.allocate(NBYTES, stream=stream) + stream.sync() parent_conn.send(buffer2.ipc_descriptor) # by descriptor # Wait for the child process. @@ -54,11 +56,12 @@ def test_main(self, ipc_device, ipc_memory_resource): assert process.exitcode == 0 # Confirm buffers were modified. - pgen = PatternGen(device, NBYTES) + pgen = PatternGen(device, NBYTES, stream=stream) pgen.verify_buffer(buffer1, seed=True) pgen.verify_buffer(buffer2, seed=True) buffer1.close() buffer2.close() + stream.sync() def child_main(self, conn): # Set up the device. @@ -73,14 +76,16 @@ def child_main(self, conn): # Receive the buffers. buffer1 = conn.recv() # directly buffer_desc = conn.recv() - buffer2 = Buffer.from_ipc_descriptor(mr, buffer_desc, stream=device.default_stream) # by descriptor + stream = device.default_stream + buffer2 = Buffer.from_ipc_descriptor(mr, buffer_desc, stream=stream) # by descriptor # Modify the buffers. - pgen = PatternGen(device, NBYTES) + pgen = PatternGen(device, NBYTES, stream=stream) pgen.fill_buffer(buffer1, seed=True) pgen.fill_buffer(buffer2, seed=True) buffer1.close() buffer2.close() + stream.sync() class TestObjectSerializationWithMR: @@ -89,6 +94,7 @@ def test_main(self, ipc_device, ipc_memory_resource): """Test sending IPC memory objects to a child through a queue.""" device = ipc_device mr = ipc_memory_resource + stream = device.default_stream # Start the child process. Sending the memory resource registers it so # that buffers can be handled automatically. @@ -103,7 +109,8 @@ def test_main(self, ipc_device, ipc_memory_resource): assert uuid == mr.uuid # Send a buffer. - buffer = mr.allocate(NBYTES, stream=device.default_stream) + buffer = mr.allocate(NBYTES, stream=stream) + stream.sync() pipe[0].put(buffer) # Wait for the child process. @@ -113,9 +120,10 @@ def test_main(self, ipc_device, ipc_memory_resource): assert process.exitcode == 0 # Confirm buffer was modified. - pgen = PatternGen(device, NBYTES) + pgen = PatternGen(device, NBYTES, stream=stream) pgen.verify_buffer(buffer, seed=True) buffer.close() + stream.sync() def child_main(self, pipe, _): device = Device() @@ -128,9 +136,11 @@ def child_main(self, pipe, _): # Buffer. buffer = pipe[0].get(timeout=CHILD_TIMEOUT_SEC) assert buffer.memory_resource.handle == mr.handle - pgen = PatternGen(device, NBYTES) + stream = device.default_stream + pgen = PatternGen(device, NBYTES, stream=stream) pgen.fill_buffer(buffer, seed=True) buffer.close() + stream.sync() class TestObjectPassing: @@ -148,11 +158,13 @@ def test_main(self, ipc_device, ipc_memory_resource): device = ipc_device mr = ipc_memory_resource alloc_handle = mr.allocation_handle - buffer = mr.allocate(NBYTES, stream=device.default_stream) + stream = device.default_stream + buffer = mr.allocate(NBYTES, stream=stream) buffer_desc = buffer.ipc_descriptor - pgen = PatternGen(device, NBYTES) + pgen = PatternGen(device, NBYTES, stream=stream) pgen.fill_buffer(buffer, seed=False) + stream.sync() # Start the child process. process = mp.Process(target=self.child_main, args=(alloc_handle, mr, buffer_desc, buffer)) @@ -164,6 +176,7 @@ def test_main(self, ipc_device, ipc_memory_resource): pgen.verify_buffer(buffer, seed=True) buffer.close() + stream.sync() def child_main(self, alloc_handle, mr1, buffer_desc, buffer): device = Device() @@ -176,7 +189,8 @@ def child_main(self, alloc_handle, mr1, buffer_desc, buffer): with pytest.raises(TypeError): PinnedMemoryResource.from_allocation_handle(alloc_handle) mr2 = DeviceMemoryResource.from_allocation_handle(device, alloc_handle) - pgen = PatternGen(device, NBYTES) + stream = device.default_stream + pgen = PatternGen(device, NBYTES, stream=stream) # Verify initial content pgen.verify_buffer(buffer, seed=False) @@ -189,3 +203,4 @@ def child_main(self, alloc_handle, mr1, buffer_desc, buffer): # Clean up - only ONE free buffer.close() + stream.sync() diff --git a/cuda_core/tests/memory_ipc/test_workerpool.py b/cuda_core/tests/memory_ipc/test_workerpool.py index e358c043b00..f2a5fb1d50c 100644 --- a/cuda_core/tests/memory_ipc/test_workerpool.py +++ b/cuda_core/tests/memory_ipc/test_workerpool.py @@ -36,30 +36,33 @@ def test_main(self, ipc_device, nmrs): options = DeviceMemoryResourceOptions(max_size=POOL_SIZE, ipc_enabled=True) mrs = [DeviceMemoryResource(device, options=options) for _ in range(nmrs)] buffers = [] + stream = device.default_stream try: - buffers = [mr.allocate(NBYTES, stream=device.default_stream) for mr, _ in zip(cycle(mrs), range(NTASKS))] + buffers = [mr.allocate(NBYTES, stream=stream) for mr, _ in zip(cycle(mrs), range(NTASKS))] + stream.sync() with mp.Pool(NWORKERS) as pool: pool.map(self.process_buffer, buffers) - pgen = PatternGen(device, NBYTES) + pgen = PatternGen(device, NBYTES, stream=stream) for buffer in buffers: pgen.verify_buffer(buffer, seed=True) finally: for buffer in buffers: buffer.close() - # TODO(seberg): 2026-06: mr close may be unsafe with incomplete `buf.close()` - device.sync() + stream.sync() for mr in mrs: mr.close() def process_buffer(self, buffer): device = Device(buffer.memory_resource.device_id) device.set_current() - pgen = PatternGen(device, NBYTES) + stream = device.default_stream + pgen = PatternGen(device, NBYTES, stream=stream) pgen.fill_buffer(buffer, seed=True) buffer.close() + stream.sync() class TestIpcWorkerPoolUsingIPCDescriptors: @@ -82,9 +85,11 @@ def test_main(self, ipc_device, nmrs): options = DeviceMemoryResourceOptions(max_size=POOL_SIZE, ipc_enabled=True) mrs = [DeviceMemoryResource(device, options=options) for _ in range(nmrs)] buffers = [] + stream = device.default_stream try: - buffers = [mr.allocate(NBYTES, stream=device.default_stream) for mr, _ in zip(cycle(mrs), range(NTASKS))] + buffers = [mr.allocate(NBYTES, stream=stream) for mr, _ in zip(cycle(mrs), range(NTASKS))] + stream.sync() with mp.Pool(NWORKERS, initializer=self.init_worker, initargs=(mrs,)) as pool: pool.starmap( @@ -92,14 +97,13 @@ def test_main(self, ipc_device, nmrs): [(mrs.index(buffer.memory_resource), buffer.ipc_descriptor) for buffer in buffers], ) - pgen = PatternGen(device, NBYTES) + pgen = PatternGen(device, NBYTES, stream=stream) for buffer in buffers: pgen.verify_buffer(buffer, seed=True) finally: for buffer in buffers: buffer.close() - # TODO(seberg): 2026-06: mr close may be unsafe with incomplete `buf.close()` - device.sync() + stream.sync() for mr in mrs: mr.close() @@ -107,10 +111,12 @@ def process_buffer(self, mr_idx, buffer_desc): mr = self.mrs[mr_idx] device = Device(mr.device_id) device.set_current() - buffer = Buffer.from_ipc_descriptor(mr, buffer_desc, stream=device.default_stream) - pgen = PatternGen(device, NBYTES) + stream = device.default_stream + buffer = Buffer.from_ipc_descriptor(mr, buffer_desc, stream=stream) + pgen = PatternGen(device, NBYTES, stream=stream) pgen.fill_buffer(buffer, seed=True) buffer.close() + stream.sync() class TestIpcWorkerPoolUsingRegistry: @@ -136,27 +142,30 @@ def test_main(self, ipc_device, nmrs): options = DeviceMemoryResourceOptions(max_size=POOL_SIZE, ipc_enabled=True) mrs = [DeviceMemoryResource(device, options=options) for _ in range(nmrs)] buffers = [] + stream = device.default_stream try: - buffers = [mr.allocate(NBYTES, stream=device.default_stream) for mr, _ in zip(cycle(mrs), range(NTASKS))] + buffers = [mr.allocate(NBYTES, stream=stream) for mr, _ in zip(cycle(mrs), range(NTASKS))] + stream.sync() with mp.Pool(NWORKERS, initializer=self.init_worker, initargs=(mrs,)) as pool: pool.starmap(self.process_buffer, [(device, pickle.dumps(buffer)) for buffer in buffers]) - pgen = PatternGen(device, NBYTES) + pgen = PatternGen(device, NBYTES, stream=stream) for buffer in buffers: pgen.verify_buffer(buffer, seed=True) finally: for buffer in buffers: buffer.close() - # TODO(seberg): 2026-06: mr close may be unsafe with incomplete `buf.close()` - device.sync() + stream.sync() for mr in mrs: mr.close() def process_buffer(self, device, buffer_s): device.set_current() buffer = pickle.loads(buffer_s) # noqa: S301 - pgen = PatternGen(device, NBYTES) + stream = device.default_stream + pgen = PatternGen(device, NBYTES, stream=stream) pgen.fill_buffer(buffer, seed=True) buffer.close() + stream.sync() diff --git a/cuda_core/tests/test_helpers.py b/cuda_core/tests/test_helpers.py index 43dbf8887e2..731924017c8 100644 --- a/cuda_core/tests/test_helpers.py +++ b/cuda_core/tests/test_helpers.py @@ -56,11 +56,12 @@ def test_patterngen_seeds(): device = Device() device.set_current() buffer = make_scratch_buffer(device, 0, NBYTES) + stream = device.default_stream # All seeds are pairwise different. # We test a sampling of values because exhaustive testing is too slow, # especially on Windows. See https://github.com/NVIDIA/cuda-python/issues/1455 - pgen = PatternGen(device, NBYTES) + pgen = PatternGen(device, NBYTES, stream=stream) for i in (ii for ii in range(256) if ii < 5 or ii % 17 == 0): pgen.fill_buffer(buffer, seed=i) pgen.verify_buffer(buffer, seed=i) @@ -77,6 +78,6 @@ def test_patterngen_values(): twos = make_scratch_buffer(device, 2, NBYTES) assert compare_equal_buffers(ones, ones) assert not compare_equal_buffers(ones, twos) - pgen = PatternGen(device, NBYTES) + pgen = PatternGen(device, NBYTES, stream=device.default_stream) pgen.verify_buffer(ones, value=1) pgen.verify_buffer(twos, value=2) diff --git a/cuda_core/tests/test_memory.py b/cuda_core/tests/test_memory.py index 3fdbe98885f..241d9226f3c 100644 --- a/cuda_core/tests/test_memory.py +++ b/cuda_core/tests/test_memory.py @@ -1390,9 +1390,9 @@ def test_device_memory_resource_with_options(init_cuda): buffer.close(stream) # Test memory copying between buffers from same pool - src_buffer = mr.allocate(64, stream=device.default_stream) - dst_buffer = mr.allocate(64, stream=device.default_stream) stream = device.create_stream() + src_buffer = mr.allocate(64, stream=stream) + dst_buffer = mr.allocate(64, stream=stream) src_buffer.copy_to(dst_buffer, stream=stream) device.sync() dst_buffer.close() @@ -1438,9 +1438,9 @@ def test_pinned_memory_resource_with_options(init_cuda): buffer.close(stream) # Test memory copying between buffers from same pool - src_buffer = mr.allocate(64, stream=device.default_stream) - dst_buffer = mr.allocate(64, stream=device.default_stream) stream = device.create_stream() + src_buffer = mr.allocate(64, stream=stream) + dst_buffer = mr.allocate(64, stream=stream) src_buffer.copy_to(dst_buffer, stream=stream) device.sync() dst_buffer.close() @@ -1485,9 +1485,9 @@ def test_managed_memory_resource_with_options(init_cuda): buffer.close(stream) # Test memory copying between buffers from same pool - src_buffer = mr.allocate(64, stream=device.default_stream) - dst_buffer = mr.allocate(64, stream=device.default_stream) stream = device.create_stream() + src_buffer = mr.allocate(64, stream=stream) + dst_buffer = mr.allocate(64, stream=stream) src_buffer.copy_to(dst_buffer, stream=stream) device.sync() dst_buffer.close() diff --git a/cuda_core/tests/test_memory_peer_access.py b/cuda_core/tests/test_memory_peer_access.py index 4763b761ad4..7ba39bb3e99 100644 --- a/cuda_core/tests/test_memory_peer_access.py +++ b/cuda_core/tests/test_memory_peer_access.py @@ -24,9 +24,11 @@ def test_peer_access_basic(mempool_device_x2): zero_on_dev0 = make_scratch_buffer(dev0, 0, NBYTES) one_on_dev0 = make_scratch_buffer(dev0, 1, NBYTES) stream_on_dev0 = dev0.create_stream() + allocation_stream = dev1.create_stream() # Use owned pool to ensure clean initial state (no stale peer access). dmr_on_dev1 = DeviceMemoryResource(dev1, DeviceMemoryResourceOptions(max_size=POOL_SIZE)) - buf_on_dev1 = dmr_on_dev1.allocate(NBYTES, stream=dev1.default_stream) + buf_on_dev1 = dmr_on_dev1.allocate(NBYTES, stream=allocation_stream) + allocation_stream.sync() # No access at first. assert 0 not in dmr_on_dev1.peer_accessible_by @@ -73,11 +75,11 @@ def test_peer_access_transitions(mempool_device_x3): # Allocate per-device resources. streams = [dev.create_stream() for dev in devs] - pgens = [PatternGen(devs[i], NBYTES, streams[i]) for i in range(3)] + pgens = [PatternGen(devs[i], NBYTES, stream=streams[i]) for i in range(3)] # Use owned pools (with options) to ensure clean initial state. # Default pools are shared and may have stale peer access from prior tests. dmrs = [DeviceMemoryResource(dev, DeviceMemoryResourceOptions(max_size=POOL_SIZE)) for dev in devs] - bufs = [dmr.allocate(NBYTES, stream=dev.default_stream) for dmr, dev in zip(dmrs, devs)] + bufs = [dmr.allocate(NBYTES, stream=stream) for dmr, stream in zip(dmrs, streams)] def verify_state(state, pattern_seed): """ diff --git a/cuda_core/tests/test_module.py b/cuda_core/tests/test_module.py index 25cf0e24de4..7a0ba965b5d 100644 --- a/cuda_core/tests/test_module.py +++ b/cuda_core/tests/test_module.py @@ -501,7 +501,7 @@ def test_object_code_load_rdc_with_linker(kind, from_fn, init_cuda): host_buf = cuda.core.LegacyPinnedMemoryResource().allocate(4) result = np.from_dlpack(host_buf).view(np.float32) result[:] = 0.0 - dev_buf = init_cuda.memory_resource.allocate(4, stream=init_cuda.default_stream) + dev_buf = init_cuda.memory_resource.allocate(4, stream=stream) cuda.core.launch( stream, @@ -880,7 +880,7 @@ def get_kernel_only(): result = np.from_dlpack(host_buf).view(np.int32) result[:] = 0 - dev_buf = device.memory_resource.allocate(4, stream=device.default_stream) + dev_buf = device.memory_resource.allocate(4, stream=stream) # Launch kernel config = cuda.core.LaunchConfig(grid=1, block=1)