diff --git a/.gitlab-ci.yml b/.gitlab-ci.yml index 17a6fcb..eb5f82f 100644 --- a/.gitlab-ci.yml +++ b/.gitlab-ci.yml @@ -1,12 +1,13 @@ variables: - # Use the specialized MPCDF HPC image CUDA_IMAGE: "gitlab-registry.mpcdf.mpg.de/mpcdf/ci-module-image/nvhpcsdk_26:2026" PIP_DISABLE_PIP_VERSION_CHECK: "1" + ARRAY_BACKEND: "cupy" + CUNUMPY_REQUIRE_CUDA: "1" stages: - test -gpu_tests: +.gpu_job: stage: test image: ${CUDA_IMAGE} tags: @@ -16,32 +17,34 @@ gpu_tests: - module load python-waterboa/2025.06 - module load nvhpcsdk/26 - module load fftw-serial/3.3.10 - script: - - echo "--- CUDA Sanity Check ---" - nvidia-smi - - - echo "--- Detect Compute Capability ---" - - nvidia-smi --query-gpu=compute_cap --format=csv,noheader - - - echo "--- Tool Versions ---" - - cmake --version - - python3 --version - - git --version - - - echo "--- Pytest Execution ---" - # The MPCDF image likely has a specific python environment. - # We install our dependencies into the user directory or a virtualenv. - # One CUDA version only: nvhpcsdk/26 provides CUDA 13.2 (its headers are used - # when CuPy compiles kernels with NVRTC), so install CuPy for CUDA 13 with the - # CUDA 13.2 libraries and NVRTC. Mixing CUDA 12 (cupy-cuda12x) with the CUDA 13.2 - # headers fails to compile CuPy's own kernels (e.g. CUB reductions). + # Keep NVRTC, CUDA libraries and headers on the same toolkit generation. - python3 -m pip install --user "cupy-cuda13x[ctk]" "cuda-toolkit==13.2.*" - - python3 -m pip install --user -e . - - # Add the user bin to PATH for pytest + - python3 -m pip install --user -e ".[test]" - export PATH="$HOME/.local/bin:$PATH" - - - export ARRAY_BACKEND=cupy - python3 -c "import cupy; cupy.show_config()" - - - pytest -xvs . + +gpu_tests: + extends: .gpu_job + script: + - python3 -m pytest -q tests + +gpu_sanitizers: + extends: .gpu_job + timeout: 45 minutes + parallel: + matrix: + - SANITIZER: [memcheck, racecheck, synccheck] + script: + - compute-sanitizer --version + # These tests intentionally contain no illegal-access demonstration kernels. + - >- + compute-sanitizer --tool "$SANITIZER" --error-exitcode 1 + --target-processes all python3 -m pytest -q + tests/unit/test_cuda_collectives.py + tests/unit/test_cuda_launch_limits.py + tests/unit/test_reusable_streams.py + tests/unit/test_staging.py + tests/unit/test_mirror.py + tests/unit/test_mpi_staging.py + tests/unit/test_segment_plans.py diff --git a/CHANGELOG.md b/CHANGELOG.md index 96e35c7..a3a2539 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -8,6 +8,15 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 ## [Unreleased] ### Added +- Per-device eager CUDA compilation, compiler log streams and explicit `recompile()`. +- Launch validation against actual device/kernel dimensions, thread limits and + static plus dynamic shared memory, with cached per-device opt-in state. +- Transfer payload byte counts, mirror-refresh and kernel-output accounting, and + separate device-only conversion observations. +- Producer `stream=`/`event=` dependencies for output staging and mirror downloads; + device-bound retained storage and CPU-to-pinned staging upgrades. +- GPU CI requires real hardware and runs focused Compute Sanitizer memory, + shared-memory race and synchronization checks. - Reusable `cuda.create_stream`, `create_event`, `record_event`, and `wait_event`, with synchronous host equivalents; `cuda.stream(existing)` selects an existing stream. - `mpi.MPIStaging` reuses host staging storage and rejects overlapping device uses. @@ -35,8 +44,12 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 ### Development - Formatting is checked with `ruff format` and import order with ruff's isort rules (`ruff check`), on `src/` and `tests/`, instead of black and isort; the `dev` extra installs ruff. +- GPU matrix/FFT correctness tests check numerical parity without timing assertions; + performance comparisons belong in scope-profiler. ### Fixed +- Block min/max reductions synchronize warp shared-memory reads before lane 0 + overwrites the result slot, fixing Compute Sanitizer racecheck hazards. - `CudaKernel`'s header hash (`-DCUNUMPY_INCLUDE_HASH`) now covers the headers shipped with cunumpy (`cunumpy/atomic.cuh`, `reduce.cuh`, ...), also when included in angle brackets. Before, an upgrade of cunumpy that changed one of them left CuPy's kernel cache serving the kernel compiled with the old header. `resolve_includes(..., angle_dirs=...)` tracks angle-bracket includes found in the given directories. - `xp.testing.assert_kernels_agree` reads `CudaStructArguments` objects and struct values through their struct fields, so their arrays get the same names as the attributes of the host argument object (before, arrays behind properties were named after the private attribute holding the owner, and the comparison failed with "do not have the same array arguments"). diff --git a/README.md b/README.md index b2dfe0a..2653299 100644 --- a/README.md +++ b/README.md @@ -149,17 +149,20 @@ with xp.use_backend("cupy"): To verify that a block, such as a time step, makes no transfer at all, count them: `count_transfers()` records every `to_numpy()`, `to_cupy()` and -`to_cunumpy()` call that actually copies, every `PyccelKernel` call that -converts device arrays, and every `Kernel` fallback to the host kernel, with -the call site of each. `assert_no_transfers()` raises with that report if -anything was counted. Only transfers made through CuNumpy are seen; raw +`to_cunumpy()` call that actually copies, mirror/staging refreshes, argument +conversion and kernel output copy-back, with call sites and payload byte counts. +Host kernel conversions and fallbacks have separate explanatory markers. +`assert_no_transfers()` rejects host/device movement and permits device-only +conversions. Only CuNumpy execution/conversion helpers are counted; forwarded +backend calls such as `xp.asarray()` and raw `cupy.ndarray.get()` or `cupy.asarray()` calls need a profiler such as `nsys`. ```python with xp.profiling.count_transfers() as counter: propagator(dt) -assert counter.total == 0, counter.report() +assert counter.to_host == counter.to_device == 0, counter.report() +print(counter.bytes_to_host, counter.bytes_to_device) ``` ## Random numbers and dtypes @@ -368,7 +371,8 @@ When building such argument objects, `xp.as_device_array(value, dtype, ndim=None)` applies the "reference or copy once" rule: a CuPy array that already has the dtype and is C-contiguous is returned as it is, anything else (a tuple such as `degree = (3, 3, 3)`, a host array, another dtype, a -non-contiguous view) becomes one device copy. Call it once when the object is +non-contiguous view) is converted; dtype and layout changes can require separate +device copies. Call it once when the object is built, not per kernel call; on the NumPy backend it raises, so host data is never copied to the device implicitly. When the host kernel takes such a group as one object too (e.g. a Pyccel class diff --git a/docs/source/api.md b/docs/source/api.md index 46a777f..861842c 100644 --- a/docs/source/api.md +++ b/docs/source/api.md @@ -283,30 +283,34 @@ with xp.profiling.count_transfers() as counter: assert counter.total == 0, counter.report() ``` -Four kinds of events are recorded: +Five kinds of events are recorded: -* `to_host`: `to_numpy()` (or `to_cunumpy()`) called with a CuPy array; -* `to_device`: `to_cupy()` (or `to_cunumpy()`) called with anything that is not - a CuPy array already; +* `to_host`: an actual device-to-host copy through conversion, mirror, staging, + serial MPI, or host kernel helpers; +* `to_device`: an actual host-to-device copy through those helpers, including + `as_device_array()` and `kernel_output()` copy-back; * `kernel_conversion`: a `PyccelKernel` call that copied device arrays to the host (and back), one event per call, naming the kernel and the number of arrays converted; * `fallback`: a `Kernel` without CUDA kernel calling its host kernel on the CuPy backend (`missing_cuda="fallback"`), one event per call, naming the - kernel. The host copies it makes are counted as one `kernel_conversion` - event in addition. + kernel. Physical copies and a `kernel_conversion` marker are recorded separately; +* `device_copy`: a device-only dtype/layout conversion through CuNumpy helpers. Only real transfers count: `to_numpy()` of a NumPy array or `to_cupy()` of a CuPy array records nothing. The counter has the attributes `to_host`, `to_device`, `kernel_conversions`, `fallbacks` (counts per kind), `total`, -`events` (a list of `TransferEvent(kind, description, where)`, where `where` -is the `file:line` of the caller outside CuNumpy) and -`kernel_conversion_calls` (the `kernel_conversion` events). `report()` returns +`device_copies`, `bytes_to_host`, `bytes_to_device`, and `bytes(kind)`. +`total` includes physical copies and conversion/fallback markers. Byte totals +include only physical copies, so markers do not double-count bytes. +`events` is a list of `TransferEvent(kind, description, where, nbytes=None)`; +`where` is the caller's `file:line` and `nbytes` is None for markers/unknown sizes. +`kernel_conversion_calls` selects the conversion markers. `report()` returns a multi-line string with the events grouped by kind and call site, with counts: ```text -4 transfer(s) through cunumpy (3 to_host, 1 to_device, 0 kernel_conversion, 0 fallback) +4 transfer(s) through cunumpy (3 to_host, 1 to_device, 0 kernel_conversion, 0 fallback, 0 device_copy) to_host (3): /home/me/sim/diagnostics.py:42: to_numpy(shape=(100000,), dtype=float64) (x3) to_device (1): @@ -320,13 +324,15 @@ not thread-safe. **Limitation:** only transfers made through CuNumpy are seen. Raw `cupy.ndarray.get()`, `cupy.asarray(numpy_array)`, `numpy.asarray(cupy_array)`, -`float(device_array)`, and implicit conversions inside other libraries are +`float(device_array)`, forwarded backend calls such as `xp.asarray()`, and +implicit conversions inside other libraries are not counted. Use `nsys` (or CuPy's profiling hooks) to find those. ### `profiling.assert_no_transfers()` Context manager that raises `AssertionError` with the counter's `report()` if -the block makes a transfer through CuNumpy. It yields the `TransferCounter` +the block makes a host/device transfer or host fallback through CuNumpy. +Device-only dtype/layout conversions are allowed. It yields the `TransferCounter` too. An exception raised inside the block propagates as it is: ```python @@ -344,11 +350,14 @@ argument object is built, never per kernel call: * a CuPy array that already has `dtype` (any dtype if `dtype` is `None`) and is C-contiguous is returned unchanged, the same object without a copy, so kernels write into the caller's array; -* anything else becomes one C-contiguous device copy, +* anything else is converted to a C-contiguous device array using `cupy.ascontiguousarray(cupy.asarray(value, dtype))`: tuples and lists (`degree = (3, 3, 3)`), host NumPy arrays (one explicit transfer at build time), device arrays of another dtype, and non-contiguous views. +Dtype and layout conversion can require separate device copies; each copy +through this helper appears in transfer accounting. + The result passes the pointer checks of `CudaKernel` and `CudaStruct`. On the NumPy backend it raises `RuntimeError`: device arguments are only built when running on CuPy, and host data is never copied to the device implicitly. If @@ -368,7 +377,7 @@ class DeviceParticles(xp.cuda.CudaArguments): For the arguments of a `Kernel` with `dispatch="arrays"`, whose choice follows the arrays. `as_kernel_array` returns `value` on the side of `like` (a CuPy array if `like` is one, a NumPy array otherwise), C-contiguous and with `dtype` -(any if `None`): `value` itself if it already is such an array, else one copy, +(any if `None`): `value` itself if it already is such an array, else a conversion, moved across if needed (counted by `count_transfers()`). `kernel_output` is a context manager yielding the buffer for an array the kernel writes: `out` itself if `as_kernel_array` takes it unchanged, else a converted copy whose @@ -871,7 +880,12 @@ Wraps the `__global__` function `name` in the CUDA C `source` (declared through CuPy on the first call (or by `compile()`), and cached, also on disk by CuPy. CuPy is imported only then, so kernels can be created and their signatures parsed without CuPy; `compile()` raises `RuntimeError` without a -GPU. +GPU. `compile(log_stream=None)` compiles eagerly and returns the CuPy raw kernel; +repeated calls reuse successful state on the current device. `is_compiled` +reports that device's state. `recompile(log_stream=None)` refreshes headers and +options on the current device; finish in-flight launches before rebuilding. +Compiler failures remain retryable and a writable `log_stream` receives compiler +output. Catalog compilation therefore reports errors during setup. `from_file` reads the source from a file; the kernel name defaults to the file name without `suffix` (`axpy_cuda.cu` -> `axpy`), and the directory of the file @@ -972,10 +986,10 @@ shape is given either by `n_threads` or by `grid`: * `grid`: number of blocks per dimension, instead of `n_threads`. * `block`: block shape for this call, instead of `block_size`. * `shared_mem`: dynamic shared memory per block in bytes, for - `extern __shared__` arrays. Above 48 KiB (the limit every device has) the - compiled kernel's `max_dynamic_shared_size_bytes` is raised to `shared_mem` - once, up to the device's opt-in limit (`xp.cuda.max_shared_memory_per_block( - opt_in=True)`); a larger request raises `ValueError` before the launch. + `extern __shared__` arrays. Static plus dynamic storage is checked against + the device's opt-in limit. Dynamic storage above the default allowance after + subtracting static storage opts in through `max_dynamic_shared_size_bytes`. + Block/grid dimensions and device/kernel thread limits are also checked. Without `n_threads` and `grid`, a launch uses `n_threads_from(args)` if the kernel has one (constructor argument and settable property): a function of the @@ -1175,7 +1189,7 @@ too); `xp.cuda.get_cuda_debug()` returns the current setting. A kernel created w enabling it also affects kernels created earlier; `debug=True` or `debug=False` fix the mode for one kernel. Only the compile options are fixed at compile time: a kernel compiled before debug mode was enabled keeps its -options, so call `compile()` after enabling, or create the kernels after +options, so call `recompile()` after enabling, or create the kernels after enabling. `kernel.debug_active()` tells whether debug mode applies to a kernel now, and `kernel.compile_options()` returns the options a compilation now would use. @@ -1972,6 +1986,16 @@ buffer is reused `buffers` copies later; a stale result raises the NumPy backend copy at once. Device copies are counted by `count_transfers()`. The arrays must have the staging shape and dtype. +`copy(array, *, stream=None, event=None)` accepts an explicit producer stream or +event (at most one). An event queues a wait before snapshotting on the current +stream. Later writes on another stream must wait for the snapshot/copy; +`copy.result()` is a conservative completion point. Keep the source's device +current when submitting copies. Storage binds to that device and rejects another +device. `ready()`/`result()` temporarily select the owning device and restore the +caller's device. The first GPU use of CPU-initialized storage allocates fresh +pinned slots; previous completed CPU handles keep their snapshots. `shape` and +`dtype` are read-only. + ## `memory.DeviceMirror` ```python @@ -1989,7 +2013,7 @@ raises `TypeError`. transfer). On the NumPy backend it is the host array itself, so the same code runs without any copy on the CPU. * `to_device()`: copies the host array into the existing device array; - `to_host()`: copies the device array into the host array, in place, so the + `to_host(stream=None, event=None)`: copies the device array into the host array, in place, so the host array keeps its identity and the owning library sees the new values. Both are no-ops on the NumPy backend. * `zero()`: zeroes the device array (allocating it empty if needed), or the @@ -2001,6 +2025,14 @@ raises `TypeError`. * `host`, `shape`, `dtype` properties. `to_device()`, `to_host()`, `zero()` and `rebind()` return the mirror, for chaining. +Make the mirror's device current before using its device storage. A mismatched +device raises before copying or zeroing. Pass an explicit producer `stream=` or +`event=` to `to_host()` when production happened outside the current stream; +the host copy is complete on return. Initial uploads and refreshes in both +directions appear in transfer accounting, with payload byte counts. Switching +to NumPy uses the host array; explicitly refresh with `to_device()` after CPU +changes before resuming GPU work. + The transfers are explicit so that one per accumulation is visible and bounded: diff --git a/docs/source/guides/data-movement.md b/docs/source/guides/data-movement.md index ad598e0..07f3f78 100644 --- a/docs/source/guides/data-movement.md +++ b/docs/source/guides/data-movement.md @@ -59,8 +59,8 @@ Move the conversions out of the loop and convert inside the `if` only. ## Count the transfers The classic performance bug of a GPU port is a transfer that sneaks into the -time loop. `count_transfers()` records every copy made through CuNumpy in a -block, with the file and line that caused it: +time loop. `count_transfers()` records copies made through CuNumpy's conversion +and execution helpers, with the file and line that caused them: ```python with xp.profiling.count_transfers() as counter: @@ -72,23 +72,29 @@ print(counter.report()) ``` ```text -20 transfer(s) through cunumpy (10 to_host, 10 to_device, 0 kernel_conversion, 0 fallback) +20 transfer(s) through cunumpy (10 to_host, 10 to_device, 0 kernel_conversion, 0 fallback, 0 device_copy) to_host (10): /home/me/sim/diagnostics.py:42: to_numpy(shape=(100000,), dtype=float64) (x10) to_device (10): /home/me/sim/step.py:17: to_cupy(shape=(100000,), dtype=float64) (x10) ``` -Four kinds of events are recorded: `to_host`, `to_device`, -`kernel_conversion` (a [`PyccelKernel`](../kernels/pyccel-kernel.md) that copied -device arrays to the host and back) and `fallback` (a -[`Kernel`](../kernels/dispatch.md) without CUDA version running its host kernel -on the GPU backend). Calls that do not copy, such as `to_numpy()` of a NumPy -array, are not counted, so a `count_transfers()` block on the NumPy backend -reports zero. +Physical copies are recorded as `to_host`, `to_device`, or `device_copy` +(device-only dtype/layout conversions). Mirror refreshes, staging, and kernel +output copies are included. Each copy has an `event.nbytes` payload size; +`counter.bytes_to_host`, `counter.bytes_to_device`, and `counter.bytes(kind)` +sum the known sizes. -In tests, `assert_no_transfers()` turns this into a check that fails with the -report: +`kernel_conversion` marks a [`PyccelKernel`](../kernels/pyccel-kernel.md) call +that copied device arrays to the host and back, and `fallback` marks a +[`Kernel`](../kernels/dispatch.md) without a CUDA version running its host kernel +on the GPU backend. Physical copies are recorded separately from these markers, +which add no bytes. `counter.total` includes markers; `counter.to_host + +counter.to_device` counts physical host/device copies. Reference-only calls, +such as `to_numpy()` of a NumPy array, record nothing. + +In tests, `assert_no_transfers()` rejects host/device copies and host execution +on the GPU backend, with the report below. It allows device-only conversions: ```python def test_step_stays_on_device(): @@ -97,7 +103,8 @@ def test_step_stays_on_device(): step(state, dt) ``` -The counter only sees transfers made through CuNumpy. Raw `cupy.asarray(host)`, +The counter only sees transfers made through CuNumpy helpers. Forwarded backend +operations such as `xp.asarray(host)`, raw `cupy.asarray(host)`, `device_array.get()`, `float(device_scalar)` and conversions inside other libraries are invisible to it; use `nsys` to find those (see [Timing and profiling](profiling.md)). @@ -136,6 +143,13 @@ place. On the NumPy backend `mirror.device` *is* the host array and the copies are no-ops. See [Accumulation kernels](../kernels/accumulation.md) for the full pattern. +Both directions appear in transfer accounting. A retained device buffer is +bound to its CUDA device; make that device current before using it. Pass +`mirror.to_host(stream=producer)` or `mirror.to_host(event=completion)` when +the kernel ran on another stream. The refresh blocks until the host data is +ready. After CPU work changes the host buffer, call `mirror.to_device()` before +resuming GPU work. + ## Pinned memory Host-to-device copies from page-locked ("pinned") host memory are faster and diff --git a/docs/source/guides/execution-helpers.md b/docs/source/guides/execution-helpers.md index 819ca22..fc67400 100644 --- a/docs/source/guides/execution-helpers.md +++ b/docs/source/guides/execution-helpers.md @@ -92,6 +92,60 @@ Without either, CuNumpy synchronizes all work on the buffer's device. `synchronize_for_mpi(*arrays, stream=..., event=...)` follows the same rule; pass at most one dependency, covering all reads/writes of the supplied buffers. +## Order output snapshots and mirror refreshes + +```python +staging = xp.memory.HostStaging(field.shape, field.dtype, buffers=2) +copy = staging.copy(field, stream=producer) +# Alternatively: copy = staging.copy(field, event=produced) +host_field = copy.result() # complete host data for an output writer +``` + +The snapshot runs on the supplied producer stream, or on the current stream +after waiting for the supplied event. Without either, the current stream must +already cover production. Subsequent writes on that stream are ordered after +the snapshot; another stream must wait before overwriting the source. Waiting +for `copy.result()` is a conservative way to establish completion. + +Make the source device current when submitting copies. The staging instance +binds to that device and rejects another. Completion checks temporarily select +the owning device and restore the caller's device. CPU-initialized storage gets +fresh pinned slots at first GPU use; earlier CPU snapshots stay valid. Reusing +a slot still invalidates its earlier handle, so copy a returned host buffer if +the writer must retain it. + +`DeviceMirror.to_host(stream=producer)` and `.to_host(event=produced)` provide the +same explicit producer choices for synchronous refreshes into another library's +host array. Make the mirror's device current. After CPU work changes its host +array, explicitly call `.to_device()` before resuming GPU work. + +## Check copies and prepare kernels + +```python +with xp.profiling.count_transfers() as copies: + mirror.to_device() + mirror.to_host() +print(copies.bytes_to_device, copies.bytes_to_host) +``` + +Transfers through execution/conversion helpers carry payload byte counts, +including mirror refreshes, MPI/output staging, argument conversion, and kernel +output copy-back. Conversion/fallback markers describe calls and add no bytes. +`total` includes these markers; use `to_host + to_device` for physical boundary +copy counts. Device-only conversions are recorded separately as `device_copy` +and are permitted by `assert_no_transfers()`. Forwarded backend operations such +as `xp.asarray()` and external library copies are outside this accounting. + +`kernel.compile(log_stream=log)` compiles eagerly and caches successful state +per device; catalog setup catches CUDA compiler failures before timesteps. +`kernel.recompile()` explicitly refreshes headers/debug options. Each launch +checks actual block/grid, kernel thread, and static-plus-dynamic shared-memory +limits. These execution checks are independent of scope-profiler. + +GPU CI requires a real CUDA device (`CUNUMPY_REQUIRE_CUDA=1`) and runs focused +`memcheck`, `racecheck`, and `synccheck` jobs with nonzero sanitizer error exits. +Numerical parity tests are separate from performance comparisons. + ## Prepare cell ranges and segment reductions ```python diff --git a/docs/source/kernels/cuda-kernel.md b/docs/source/kernels/cuda-kernel.md index fd93217..df2abe5 100644 --- a/docs/source/kernels/cuda-kernel.md +++ b/docs/source/kernels/cuda-kernel.md @@ -98,6 +98,11 @@ kernel(*args, n_threads=None, grid=None, block=None, shared_mem=0, stream=None) * **Stream**: `stream=s` queues the launch on a stream, see [Devices, memory and streams](../guides/gpu-devices.md). +Launches validate block/grid dimensions, threads per block against both device +and compiled-kernel limits, and static plus dynamic shared memory. Dynamic +storage above the default allowance opts in where supported. Queries and opt-in +state are cached per device. Invalid configurations raise before launching. + A zero-sized launch (`n_threads=0`) launches nothing. `launch_shape()` returns the `(grid, block)` a call would use, which is how to size per-block outputs: @@ -248,8 +253,16 @@ The first call of each kernel compiles it, which takes from a fraction of a second to several seconds. Call `kernel.compile()` (or `catalog.compile_all(jobs=8)` for a whole catalog, in parallel threads) during setup so the first time step is not slower than the rest and compilation errors -appear before the simulation starts. `kernel.is_compiled` tells whether it has -happened. Compiling requires a GPU (`RuntimeError` otherwise). +appear before the simulation starts. `kernel.is_compiled` reports successful +compilation on the current device. Failed compilation can be retried. Compiling +requires a GPU (`RuntimeError` otherwise). Parallel catalog compilation preserves +the device selected by its caller. + +`kernel.compile(log_stream=log)` accepts a writable compiler log stream. +`kernel.recompile(log_stream=log)` refreshes headers and debug options on the +current device; other devices keep their compiled versions. Finish in-flight +launches before rebuilding. Scope-profiler can inspect the returned CuPy kernel's +existing `attributes`; CuNumpy does not add a separate resource-reporting API. ## Tips for writing kernels diff --git a/docs/source/kernels/debugging.md b/docs/source/kernels/debugging.md index 73753ba..4333eb1 100644 --- a/docs/source/kernels/debugging.md +++ b/docs/source/kernels/debugging.md @@ -46,8 +46,9 @@ with xp.cuda.cuda_debug(): Compile options are fixed when a kernel is compiled. A kernel that was already compiled before debug mode was enabled keeps its options (the synchronization -still applies). Enable debug mode before the kernels are first called, or set -`CUNUMPY_CUDA_DEBUG=1` in the environment. `kernel.debug_active()` and +still applies). Enable debug mode before the kernels are first called, call +`kernel.recompile()` after enabling it, or set `CUNUMPY_CUDA_DEBUG=1` in the +environment. Finish in-flight launches before rebuilding. `kernel.debug_active()` and `kernel.compile_options()` show what applies to a kernel. Use `-DCUNUMPY_BOUNDS_CHECK` in your own code too: @@ -108,7 +109,7 @@ on for the kernel under suspicion, not in production. | correct on small inputs, wrong on large ones | `int` overflow in index computations; use `long long`. | | results differ slightly between runs | floating-point atomics in a different order; compare with a tolerance. | | results occasionally garbage after MPI | missing `synchronize_for_mpi()` before the MPI call. | -| a header change has no effect | the header is included with angle brackets (not hashed) or the kernel object was compiled before the change; see "Headers and the compile cache" in [Writing CUDA kernels](cuda-kernel.md). | +| a header change has no effect | the kernel object was compiled before the change; call `kernel.recompile()`. Headers must resolve through the configured include directories to be hashed; see "Headers and the compile cache" in [Writing CUDA kernels](cuda-kernel.md). | After an illegal memory access the CUDA context of the process is unusable: all later CUDA calls fail. Restart the process (or the pytest run) after fixing the diff --git a/src/cunumpy/LLM_GUIDE.md b/src/cunumpy/LLM_GUIDE.md index 50589f7..05b1fa2 100644 --- a/src/cunumpy/LLM_GUIDE.md +++ b/src/cunumpy/LLM_GUIDE.md @@ -131,15 +131,19 @@ xp.to_numpy(a) # -> numpy.ndarray; copies only if a is CuPy xp.to_cupy(a) # -> cupy.ndarray; copies only if a is not CuPy; ImportError w/o CuPy xp.to_cunumpy(a) # -> array of the active backend xp.as_device_array(value, dtype=None, ndim=None, *, name=None) - # contiguous CuPy array of dtype: returned as is; else one copy; + # contiguous CuPy array of dtype: returned as is; else convert; # RuntimeError on the NumPy backend with xp.profiling.count_transfers() as c: ... # c.total, c.to_host, c.to_device, # c.kernel_conversions, c.fallbacks, c.events, c.report() -with xp.profiling.assert_no_transfers(): ... # AssertionError with report if anything copied +with xp.profiling.assert_no_transfers(): ... # rejects host/device copies and host fallback ``` -Only transfers through cunumpy are counted (not raw `cupy.asarray`, `.get()`, -`float(device_scalar)`, or `DeviceMirror.to_host()/to_device()`). +Mirror refreshes, MPI/output staging, argument conversion and kernel output +copy-back are counted with `event.nbytes`, `c.bytes_to_host`, and +`c.bytes_to_device`. Conversion/fallback markers add no bytes; `c.total` +includes markers. Device-only conversions are separate `device_copy` events +and are allowed by `assert_no_transfers()`. Forwarded `xp.asarray`, raw CuPy +copies, and external library/scalar conversions are not counted. MPI, accumulation and versions: @@ -264,7 +268,9 @@ k = xp.cuda.CudaKernel(source, name, *, block_size=128, options=(), include_dirs k = xp.cuda.CudaKernel.from_file("push/push_cuda.cu") # name "push" ks = xp.cuda.CudaKernel.all_from_file("ops.cu") # dict name -> kernel k(*args, n_threads=None, grid=None, block=None, shared_mem=0, stream=None) -k.compile(); k.is_compiled; k.launch_shape(n_threads) -> (grid, block) +k.compile(log_stream=None); k.recompile(log_stream=None) +k.is_compiled # successful compilation on the current CUDA device +k.launch_shape(n_threads) -> (grid, block) k.included_headers; k.compile_options(); k.debug_active() xp.cuda.CudaKernelVariants(factory).get(*key); .compile_all(keys, jobs=1) xp.cuda.ctype_of(np.float64) == "double" @@ -284,6 +290,10 @@ xp.cuda.cuda_include_dir() `ArrayND` params: CuPy arrays of dtype T and ndim N, any strides. * Compiled lazily with NVRTC on first call; cached on disk by CuPy; quoted `#include "..."` headers are hashed into the options so edits recompile. + `compile()` compiles eagerly and reuses successful per-device state; + `recompile()` refreshes that state after header/debug changes. Each launch + validates actual device/kernel dimensions, threads and static+dynamic shared + memory. Use the returned CuPy raw kernel's attributes for resource inspection. * Creating a `CudaKernel` does not import CuPy; compiling needs a GPU. Shipped CUDA headers (always on the include path): @@ -375,10 +385,20 @@ m = xp.memory.DeviceMirror(host_numpy_array) # TypeError if not numpy.ndarray m.device # CuPy copy (lazy) on CuPy; the host array itself on NumPy m.zero() m.to_device() -m.to_host() # to_host copies in place; no-ops on NumPy +m.to_host(stream=None, event=None) # in place; pass at most one dependency m.rebind(new_host_array) # after the owner reallocates ``` +Make a retained mirror's CUDA device current before using it. `to_host()` +blocks until the copy completes; supply the producer stream or event when work +ran elsewhere. Refresh with `to_device()` after CPU work modifies the host. +`HostStaging.copy(a, stream=None, event=None)` snapshots on the producer stream +(or current stream after an event wait). Later writes must be ordered after +the snapshot. Staging binds to the first source device and upgrades ordinary +CPU slots to pinned slots on first GPU use. `ready()`/`result()` temporarily +select the owning device and restore the caller's device. Copy the returned +host buffer to retain it beyond slot reuse. + Debugging: ```python @@ -508,8 +528,8 @@ def test_scale(): * Comparing CPU and GPU results with exact equality where the GPU uses atomics or a different reduction order; use a tolerance. * Assuming `xp.random.seed(s)` gives the same numbers on both backends. -* Expecting `DeviceMirror` copies or raw CuPy conversions to show up in - `count_transfers()`. +* Expecting raw CuPy or forwarded backend operations to show up in + `count_transfers()`; mirror copies and execution helpers are counted. * Writing code that requires CuPy at import time. `import cunumpy` and creating `CudaKernel` objects work without CuPy; import `cupy` lazily, only on the GPU path, if you need it at all. diff --git a/src/cunumpy/_cuda_kernel.py b/src/cunumpy/_cuda_kernel.py index e022bed..1a37f52 100644 --- a/src/cunumpy/_cuda_kernel.py +++ b/src/cunumpy/_cuda_kernel.py @@ -50,6 +50,7 @@ class is the one definition of the arguments. import hashlib import inspect import math +import operator import os import re import sys @@ -57,6 +58,7 @@ class is the one definition of the arguments. from collections.abc import Callable, Hashable, Iterable, Iterator, Mapping, Sequence from concurrent.futures import ThreadPoolExecutor from contextlib import nullcontext +from dataclasses import dataclass from pathlib import Path from typing import Any, NamedTuple @@ -84,6 +86,32 @@ class is the one definition of the arguments. # CUDA limit on the number of threads per block _MAX_THREADS_PER_BLOCK = 1024 +# Runtime queries are setup work, never repeated for a prepared launch. +_DEVICE_LIMITS: dict[int, dict[str, Any]] = {} + + +def _current_device() -> int: + import cupy as cp + + return cp.cuda.runtime.getDevice() + + +def _device_limits(device: int) -> dict[str, Any]: + if device not in _DEVICE_LIMITS: + import cupy as cp + + _DEVICE_LIMITS[device] = cp.cuda.runtime.getDeviceProperties(device) + return _DEVICE_LIMITS[device] + + +@dataclass +class _CompiledKernel: + raw: Any + limits: dict[str, Any] + attributes: dict[str, int] + dynamic_shared: int + + # Headers shipped with cunumpy: #include etc. _CUDA_INCLUDE_DIR = Path(__file__).resolve().parent / "cuda" / "include" _ARRAY_VIEW_INCLUDE = '#include "cunumpy/array_view.cuh"' @@ -433,8 +461,22 @@ def _compile_in_threads( compile() return list(compilers) + from cunumpy.xp import cupy_available + + device = _current_device() if cupy_available() else None + + def on_device(compiler: Callable[[], Any]) -> Any: + if device is None: + return compiler() + import cupy as cp + + with cp.cuda.Device(device): + return compiler() + with ThreadPoolExecutor(max_workers=min(jobs, len(compilers))) as pool: - futures = {name: pool.submit(compile) for name, compile in compilers.items()} + futures = { + name: pool.submit(on_device, compile) for name, compile in compilers.items() + } compiled, error = [], None for name, future in futures.items(): exc = future.exception() @@ -1362,7 +1404,7 @@ def verify_layout( RuntimeError If CuPy is not available. """ - from cunumpy.xp import cupy_available + from cunumpy.xp import cupy_available, to_numpy if not cupy_available(): raise RuntimeError("verify_layout() compiles a CUDA kernel and needs CuPy") @@ -1378,7 +1420,7 @@ def verify_layout( ) out = cp.zeros(len(self._fields) + 2, dtype=cp.uint64) kernel(out, n_threads=1) - measured = [int(v) for v in out.get()] + measured = [int(v) for v in to_numpy(out)] layout = {"sizeof": measured[0], "alignof": measured[1]} layout.update({f.name: n for f, n in zip(self._fields, measured[2:])}) @@ -1931,8 +1973,6 @@ def __init__( ) -> None: self._block = self._check_block(_as_shape(block_size, "block_size")) self.n_threads_from = n_threads_from - # dynamic shared memory the compiled kernel is set up for (see __call__) - self._shared_mem_opt_in = 48 * 1024 self._check_finite = bool(check_finite) self._debug = None if debug is None else bool(debug) self._source = source @@ -1958,7 +1998,7 @@ def __init__( if self._signature is None else [_checker(p, i) for i, p in enumerate(self._signature)] ) - self._raw_kernel = None + self._compiled: dict[int, _CompiledKernel] = {} @classmethod def from_file( @@ -2213,51 +2253,80 @@ def signature(self) -> tuple[CudaParameter, ...] | None: @property def is_compiled(self) -> bool: - """Whether the kernel has been compiled in this process.""" - return self._raw_kernel is not None + """Whether compilation succeeded on the current CUDA device.""" + return bool(self._compiled) and _current_device() in self._compiled - def compile(self) -> Any: + def compile(self, *, log_stream: Any = None) -> Any: """Compile the kernel now (it is otherwise compiled on the first call). The options are :meth:`compile_options`, evaluated now: a kernel - compiled before debug mode was enabled keeps its options. + compiled before debug mode was enabled keeps its options. Call + :meth:`recompile` to refresh headers/options explicitly. Compilation + and launch resources are cached separately on each CUDA device; + failed compilation remains retryable. ``log_stream`` receives compiler + output (a writable file object, or None). Returns ------- cupy.RawKernel The compiled kernel; compiled once and cached (also on disk by CuPy, - keyed on the source and :meth:`compile_options`). + keyed on the source and :meth:`compile_options`). No launch is needed + to surface compiler errors. Raises ------ RuntimeError If CuPy or a GPU is not available. """ - if self._raw_kernel is None: - from cunumpy.xp import cupy_available + from cunumpy.xp import cupy_available - if not cupy_available(): - raise RuntimeError( - f"cannot compile CUDA kernel {self.expression!r}: " - "CuPy is not installed or no GPU is available", - ) - import cupy as cp + if not cupy_available(): + raise RuntimeError( + f"cannot compile CUDA kernel {self.expression!r}: " + "CuPy is not installed or no GPU is available", + ) + import cupy as cp + device = _current_device() + if device not in self._compiled: options = self.compile_options() if self._template_args is None: - self._raw_kernel = cp.RawKernel( + raw = cp.RawKernel( self._source, self._name, options=options, ) + raw.compile(log_stream=log_stream) else: module = cp.RawModule( code=self._source, options=options, name_expressions=[self.expression], ) - self._raw_kernel = module.get_function(self.expression) - return self._raw_kernel + module.compile(log_stream=log_stream) + raw = module.get_function(self.expression) + limits = _device_limits(device) + attributes = dict(raw.attributes) + static = max(0, attributes.get("shared_size_bytes", 0)) + dynamic = attributes.get("max_dynamic_shared_size_bytes", -1) + if dynamic < 0: + dynamic = max(0, int(limits["sharedMemPerBlock"]) - static) + dynamic = min(dynamic, max(0, int(limits["sharedMemPerBlock"]) - static)) + # Do not publish failed compilation or incomplete setup as compiled. + self._compiled[device] = _CompiledKernel(raw, limits, attributes, dynamic) + return self._compiled[device].raw + + def recompile(self, *, log_stream: Any = None) -> Any: + """Recompile on the current device using current headers/debug options. + + Other devices keep their compiled kernels. A failed rebuild remains + uncompiled and can be retried; in-flight launches must finish first. + """ + from cunumpy.xp import cupy_available + + if cupy_available(): + self._compiled.pop(_current_device(), None) + return self.compile(log_stream=log_stream) def prepare_args(self, *args: Any) -> tuple[Any, ...]: """The arguments as passed to ``cupy.RawKernel``: flattened and checked. @@ -2331,7 +2400,7 @@ def launch_shape( "numbers of dimensions", ) block_shape = block_shape + (1,) * (len(threads) - 1) - grid_shape = tuple(math.ceil(n / b) for n, b in zip(threads, block_shape)) + grid_shape = tuple((n + b - 1) // b for n, b in zip(threads, block_shape)) return grid_shape, block_shape def __call__( @@ -2374,10 +2443,14 @@ def __call__( illegal memory access; the CuPy error is chained. Without debug mode, such an error surfaces at a later synchronization (a ``.get()``, an MPI call, ...), not necessarily in this kernel. + ValueError + Dimensions, threads or static plus dynamic shared memory exceed + this device/kernel's limits, or arrays/stream belong to another device. """ if n_threads is None and grid is None and self._n_threads_from is not None: n_threads = self._n_threads_from(args) grid_shape, block_shape = self.launch_shape(n_threads, grid=grid, block=block) + shared_mem = operator.index(shared_mem) if shared_mem < 0: raise ValueError(f"shared_mem must be non-negative, got {shared_mem}") values = self.prepare_args(*args) @@ -2385,8 +2458,23 @@ def __call__( return kernel = self.compile() - if shared_mem > self._shared_mem_opt_in: - self._opt_in_shared_memory(kernel, shared_mem) + device = _current_device() + state = self._compiled[device] + self._validate_launch(state, device, grid_shape, block_shape, shared_mem) + if stream is not None and getattr(stream, "device_id", device) not in ( + device, + -1, + ): + raise ValueError( + f"kernel {self.expression!r}: stream belongs to another device" + ) + for label, array in _device_arrays_in(args): + array_device = getattr(getattr(array, "device", None), "id", None) + if array_device is not None and array_device != device: + raise ValueError( + f"kernel {self.expression!r} on CUDA device {device}: " + f"{label} belongs to CUDA device {array_device}", + ) debug = self.debug_active() with stream if stream is not None else nullcontext(): kernel(grid_shape, block_shape, values, shared_mem=shared_mem) @@ -2395,29 +2483,50 @@ def __call__( if self._check_finite: self._check_finite_arrays(args) - def _opt_in_shared_memory(self, kernel: Any, shared_mem: int) -> None: - """Allow `shared_mem` bytes of dynamic shared memory above the default limit. - - Up to :data:`~cunumpy.cuda.DEFAULT_SHARED_MEMORY_PER_BLOCK` (48 KiB) every - device launches without setup. Above it, newer GPUs need the kernel - attribute ``max_dynamic_shared_size_bytes``; it is set once (and again - for a larger request) up to the device's opt-in limit. - """ - from cunumpy._device import ( - DEFAULT_SHARED_MEMORY_PER_BLOCK, - max_shared_memory_per_block, + def _validate_launch( + self, + state: _CompiledKernel, + device: int, + grid: tuple[int, ...], + block: tuple[int, ...], + shared_mem: int, + ) -> None: + limits, attributes = state.limits, state.attributes + name = limits.get("name", "unknown") + if isinstance(name, bytes): + name = name.decode(errors="replace") + context = f"kernel {self.expression!r} on CUDA device {device} ({name})" + for label, shape, maximum in ( + ("block", block, limits["maxThreadsDim"]), + ("grid", grid, limits["maxGridSize"]), + ): + for axis, (requested, limit) in enumerate(zip(shape, maximum)): + if requested > limit: + raise ValueError( + f"{context}: {label}[{axis}]={requested} exceeds the {limit} limit", + ) + threads_limit = int(limits["maxThreadsPerBlock"]) + kernel_limit = attributes.get("max_threads_per_block", -1) + if kernel_limit > 0: + threads_limit = min(threads_limit, kernel_limit) + if math.prod(block) > threads_limit: + raise ValueError( + f"{context}: block {block} has {math.prod(block)} threads, " + f"exceeding the {threads_limit} device/kernel limit", + ) + static = max(0, attributes.get("shared_size_bytes", 0)) + maximum_shared = max( + int(limits["sharedMemPerBlock"]), + int(limits.get("sharedMemPerBlockOptin", 0)), ) - - if shared_mem <= DEFAULT_SHARED_MEMORY_PER_BLOCK: - return - limit = max_shared_memory_per_block(opt_in=True) - if shared_mem > limit: + if shared_mem + static > maximum_shared: raise ValueError( - f"kernel {self.expression!r}: shared_mem={shared_mem} bytes exceeds " - f"the {limit} bytes a block may use on this device", + f"{context}: static shared memory {static} + shared_mem={shared_mem} " + f"exceeds the {maximum_shared} bytes a block may use on this device", ) - kernel.max_dynamic_shared_size_bytes = shared_mem - self._shared_mem_opt_in = shared_mem + if shared_mem > state.dynamic_shared: + state.raw.max_dynamic_shared_size_bytes = shared_mem + state.dynamic_shared = shared_mem def _check_finite_arrays(self, args: tuple[Any, ...]) -> None: """Raise if a floating-point array among `args` holds NaN or inf.""" diff --git a/src/cunumpy/_fake_cupy_impl.py b/src/cunumpy/_fake_cupy_impl.py index 5a043d0..0a4a13a 100644 --- a/src/cunumpy/_fake_cupy_impl.py +++ b/src/cunumpy/_fake_cupy_impl.py @@ -401,6 +401,9 @@ def __init__(self, *args, **kwargs): def __call__(self, *args, **kwargs): raise NotImplementedError("fake CuPy cannot run CUDA kernels") + def compile(self, log_stream=None): + raise NotImplementedError("fake CuPy cannot compile CUDA kernels") + def get_function(self, name): return _NoGPU() diff --git a/src/cunumpy/_kernel.py b/src/cunumpy/_kernel.py index 1c74c1a..d626273 100644 --- a/src/cunumpy/_kernel.py +++ b/src/cunumpy/_kernel.py @@ -35,7 +35,7 @@ import numpy as np from cunumpy._transfers import _ACTIVE as _COUNTERS -from cunumpy._transfers import _record +from cunumpy._transfers import _describe, _is_device_copy, _nbytes, _record from cunumpy.xp import _cupy_backend, _to_cupy, _to_numpy, is_gpu, to_cupy, to_numpy __all__ = ["CompiledHostKernel", "KernelArguments", "PyccelKernel", "resolve_host_args"] @@ -278,6 +278,12 @@ def _convert_to_numpy( if _is_device_array(value): value_np = _device_to_host(value) + if _COUNTERS: + _record( + "to_host", + f"PyccelKernel {self.name!r} input ({_describe(value)})", + nbytes=_nbytes(value_np), + ) memo[key] = value_np converted.append((value, value_np)) return value_np @@ -323,13 +329,23 @@ def _convert_to_numpy( def _convert_from_numpy(self, value: Any) -> Any: """Move host arrays returned by the kernel back to the device.""" if self._is_array(value): - return _host_to_device(value) + return self._copy_to_device(value) if isinstance(value, tuple): return tuple(self._convert_from_numpy(item) for item in value) if isinstance(value, list): return [self._convert_from_numpy(item) for item in value] return value + def _copy_to_device(self, value: Any) -> Any: + result = _host_to_device(value) + if _COUNTERS: + _record( + "to_device", + f"PyccelKernel {self.name!r} output ({_describe(value)})", + nbytes=_nbytes(value), + ) + return result + def _collect_host_arrays(self, value: Any, found: set[int], seen: set[int]) -> None: """Record the id of every host array reachable from `value`. @@ -479,7 +495,7 @@ def __call__(self, *args: Any, **kwargs: Any) -> Any: # Copy in-place kernel updates back to the device arrays. for device_array, host_array in converted: if writeable is None or id(host_array) in writeable: - device_array[...] = _host_to_device(host_array) + device_array[...] = self._copy_to_device(host_array) return self._convert_from_numpy(result) @@ -819,7 +835,14 @@ def as_kernel_array(value: Any, like: Any, dtype: Any = None) -> Any: if not is_gpu(value): value = to_cupy(value) - return cupy.ascontiguousarray(value, dtype=dtype) + result = cupy.ascontiguousarray(value, dtype=dtype) + if _COUNTERS and _is_device_copy(value, result): + _record( + "device_copy", + f"as_kernel_array({_describe(value)})", + nbytes=_nbytes(result), + ) + return result return np.ascontiguousarray(to_numpy(value), dtype=dtype) @@ -838,4 +861,9 @@ def kernel_output(out: Any, like: Any, dtype: Any = None) -> Generator[Any]: buffer = as_kernel_array(out, like, dtype) yield buffer if buffer is not out: - out[...] = buffer if is_gpu(out) or not is_gpu(buffer) else to_numpy(buffer) + if is_gpu(out) and not is_gpu(buffer): + out[...] = to_cupy(buffer) + elif not is_gpu(out) and is_gpu(buffer): + out[...] = to_numpy(buffer) + else: + out[...] = buffer diff --git a/src/cunumpy/_mirror.py b/src/cunumpy/_mirror.py index 147735c..fbcd650 100644 --- a/src/cunumpy/_mirror.py +++ b/src/cunumpy/_mirror.py @@ -18,6 +18,8 @@ import numpy as np +from cunumpy._transfers import _ACTIVE as _COUNTERS +from cunumpy._transfers import _describe, _nbytes, _record from cunumpy.xp import _cupy_backend, to_cupy __all__ = ["DeviceMirror"] @@ -104,13 +106,24 @@ def device(self) -> Any: """ if not _cupy_backend(): return self._host + self._check_bound() + self._check_device() if self._device is None: - self._check_bound() # The only host -> device allocation; it goes through to_cupy so # that it is visible wherever transfers are accounted for. self._device = to_cupy(self._host) return self._device + def _check_device(self) -> None: + if self._device is not None: + import cupy as cp + + if cp.cuda.runtime.getDevice() != self._device.device.id: + raise ValueError( + f"DeviceMirror is bound to CUDA device {self._device.device.id}; " + "make that device current before using the mirror", + ) + def _check_bound(self) -> None: """Raise if the host array no longer has the shape/dtype it was bound with.""" if self._host.shape != self._shape or self._host.dtype != self._dtype: @@ -156,37 +169,78 @@ def to_device(self) -> DeviceMirror: self._check_bound() if not _cupy_backend(): return self + self._check_device() if self._device is None: _ = self.device # allocates as a copy of the host, through to_cupy else: # host -> device copy into the existing device buffer self._device.set(self._host) + if _COUNTERS: + _record( + "to_device", + f"DeviceMirror.to_device({_describe(self._host)})", + nbytes=_nbytes(self._host), + ) return self - def to_host(self) -> DeviceMirror: + def to_host(self, *, stream: Any = None, event: Any = None) -> DeviceMirror: """Copy the device array into the host array, in place. The host array keeps its identity, so a library that holds the buffer sees the new values. No-op on the NumPy backend, and if no device array has been created yet. + Make the mirror's device current. Pass ``stream=`` to copy on a known + producer stream, or ``event=`` to wait for production on the current + stream. With neither, only the current stream's work is ordered. + The method blocks until the host copy is complete. + Raises ------ ValueError If the host array was reallocated with another shape or dtype. """ + if stream is not None and event is not None: + raise ValueError("pass only one of stream and event") self._check_bound() if not _cupy_backend() or self._device is None: return self + self._check_device() + from cunumpy._streams import HostEvent, HostStream + + if isinstance(event, HostEvent) or isinstance(stream, HostStream): + raise TypeError("device copies require a CUDA producer stream or event") + import cupy as cp + + if stream is not None and getattr( + stream, "device_id", self._device.device.id + ) not in ( + self._device.device.id, + -1, + ): + raise ValueError("DeviceMirror producer stream belongs to another device") + if event is not None: + cp.cuda.get_current_stream().wait_event(event) # device -> host copy into the existing host buffer - self._device.get(out=self._host) + if stream is None: + self._device.get(out=self._host) + else: + self._device.get(out=self._host, stream=stream) + if _COUNTERS: + _record( + "to_host", + f"DeviceMirror.to_host({_describe(self._host)})", + nbytes=_nbytes(self._host), + ) return self def zero(self) -> DeviceMirror: """Zero the array kernels write into (the device array, or the host on NumPy).""" + self._check_bound() if not _cupy_backend(): self._host.fill(0) return self + self._check_device() if self._device is None: self._check_bound() import cupy as cp diff --git a/src/cunumpy/_mpi.py b/src/cunumpy/_mpi.py index 9f9541f..1ed6387 100644 --- a/src/cunumpy/_mpi.py +++ b/src/cunumpy/_mpi.py @@ -15,7 +15,7 @@ _LOCAL_RANK_VARIABLES, # noqa: F401 - re-exported ) from cunumpy._transfers import _ACTIVE as _COUNTERS -from cunumpy._transfers import _describe, _record +from cunumpy._transfers import _describe, _nbytes, _record from cunumpy.xp import array_backend, cupy_available, to_numpy _logger = logging.getLogger(__name__) @@ -243,17 +243,25 @@ def mpi_buffer( synchronize_for_mpi(array, stream=stream, event=event) transfer_array = staging._transfer_array(array) if send: - if _COUNTERS: - _record("to_host", f"mpi_buffer({_describe(array)}) staging for send") if transfer_array is not array: transfer_array[...] = array transfer_array.get(out=host) + if _COUNTERS: + _record( + "to_host", + f"mpi_buffer({_describe(array)}) staging for send", + nbytes=_nbytes(array), + ) yield host if recv: + transfer_array.set(host) if _COUNTERS: - _record("to_device", f"mpi_buffer({_describe(array)}) staging for recv") + _record( + "to_device", + f"mpi_buffer({_describe(array)}) staging for recv", + nbytes=_nbytes(array), + ) # Ensure MPI's host buffer can be reused immediately on context exit. - transfer_array.set(host) if transfer_array is not array: array[...] = transfer_array cp.cuda.get_current_stream().synchronize() diff --git a/src/cunumpy/_mpi_serial.py b/src/cunumpy/_mpi_serial.py index 5cf29a2..855f21f 100644 --- a/src/cunumpy/_mpi_serial.py +++ b/src/cunumpy/_mpi_serial.py @@ -236,7 +236,11 @@ def _copy(source: Any, target: Any, offset: int = 0) -> None: if source is _IN_PLACE or source is None or target is None: return if hasattr(source, "get") and not hasattr(target, "get"): + from cunumpy._transfers import _ACTIVE, _nbytes, _record + source = source.get() + if _ACTIVE: + _record("to_host", "SerialComm receive from device", nbytes=_nbytes(source)) if not target.flags.c_contiguous: raise ValueError("the receive buffer must be C-contiguous") flat = target.reshape(-1) @@ -247,6 +251,11 @@ def _copy(source: Any, target: Any, offset: int = 0) -> None: f"at offset {offset}", ) flat[offset : offset + source.size] = source + if not hasattr(source, "get") and hasattr(target, "get"): + from cunumpy._transfers import _ACTIVE, _nbytes, _record + + if _ACTIVE: + _record("to_device", "SerialComm receive from host", nbytes=_nbytes(source)) def _check_rank(rank: int, what: str) -> None: diff --git a/src/cunumpy/_staging.py b/src/cunumpy/_staging.py index b465c9b..1fd8b7d 100644 --- a/src/cunumpy/_staging.py +++ b/src/cunumpy/_staging.py @@ -29,13 +29,15 @@ from __future__ import annotations +import operator +from contextlib import nullcontext from typing import Any import array_api_compat import numpy as np from cunumpy._transfers import _ACTIVE as _COUNTERS -from cunumpy._transfers import _record +from cunumpy._transfers import _nbytes, _record __all__ = ["HostStaging", "StagedCopy"] @@ -56,11 +58,19 @@ def _empty_pinned(shape: tuple[int, ...], dtype: Any) -> np.ndarray: class _Slot: - def __init__(self, host: np.ndarray) -> None: + def __init__(self, host: np.ndarray, device_id: int | None = None) -> None: self.host = host self.device: Any = None # the snapshot, allocated on the first device copy self.event: Any = None # recorded after the copy to the host self.generation = 0 + self.device_id = device_id + + def context(self) -> Any: + return ( + nullcontext() + if self.device_id is None + else _cupy().cuda.Device(self.device_id) + ) class StagedCopy: @@ -81,7 +91,10 @@ def _check(self) -> None: def ready(self) -> bool: """Whether the copy has finished (never waits).""" self._check() - return self._slot.event is None or bool(self._slot.event.done) + if self._slot.event is None: + return True + with self._slot.context(): + return bool(self._slot.event.done) def result(self) -> np.ndarray: """Wait for the copy and return the host array. @@ -95,7 +108,8 @@ def result(self) -> np.ndarray: """ self._check() if self._slot.event is not None: - self._slot.event.synchronize() + with self._slot.context(): + self._slot.event.synchronize() return self._slot.host @@ -112,6 +126,12 @@ class HostStaging: Number of buffers (and device snapshots) used in turn: how many copies may be in flight at once. 2 (double buffering) lets one copy run while the previous result is written out. + + Notes + ----- + GPU storage is bound to the first source device; make it current when + submitting copies. A CPU-initialized instance allocates new pinned slots + on its first GPU copy; earlier CPU results retain their completed snapshots. """ def __init__( @@ -120,14 +140,24 @@ def __init__( dtype: Any, buffers: int = 2, ) -> None: + buffers = operator.index(buffers) if buffers < 1: raise ValueError(f"buffers must be at least 1, got {buffers}") - self.shape = (shape,) if isinstance(shape, int) else tuple(shape) - self.dtype = np.dtype(dtype) + self._shape = ( + (operator.index(shape),) + if isinstance(shape, (int, np.integer)) + else tuple(operator.index(n) for n in shape) + ) + if any(n < 0 for n in self.shape): + raise ValueError("staging shape must be non-negative") + self._dtype = np.dtype(dtype) + if self.dtype.hasobject: + raise TypeError("HostStaging does not support object dtype") self._slots: list[_Slot] = [] self._n_buffers = buffers self._next = 0 self._stream: Any = None + self._device_id: int | None = None def __repr__(self) -> str: return ( @@ -135,29 +165,50 @@ def __repr__(self) -> str: f"buffers={self._n_buffers})" ) + @property + def shape(self) -> tuple[int, ...]: + """The fixed shape of each snapshot.""" + return self._shape + + @property + def dtype(self) -> np.dtype: + """The fixed dtype of each snapshot.""" + return self._dtype + @property def buffers(self) -> int: """Number of buffers used in turn.""" return self._n_buffers def _slot(self, device: bool) -> _Slot: - if not self._slots: + if not self._slots or (device and self._slots[0].device_id is None): allocate = _empty_pinned if device else np.empty + # On the first GPU copy, replace ordinary host storage with pinned + # slots. Existing CPU handles keep their own completed snapshots. self._slots = [ - _Slot(allocate(self.shape, self.dtype)) for _ in range(self._n_buffers) + _Slot(allocate(self.shape, self.dtype), self._device_id) + for _ in range(self._n_buffers) ] + self._next = 0 slot = self._slots[self._next] self._next = (self._next + 1) % self._n_buffers return slot - def copy(self, array: Any) -> StagedCopy: + def copy(self, array: Any, *, stream: Any = None, event: Any = None) -> StagedCopy: """Start copying `array` to the host and return at once. Parameters ---------- array : cupy.ndarray | numpy.ndarray - An array of the staging shape and dtype. It may be overwritten - right after this call: a device array is snapshotted first. + An array of the staging shape and dtype. Later writes on the + snapshot's stream are ordered after it; another stream must wait + for the snapshot/copy before overwriting the source. + stream : cupy.cuda.Stream | None + Producer stream on which to snapshot, current stream by default. + event : cupy.cuda.Event | None + Producer completion to wait for on the current stream before the + snapshot. Pass at most one of stream and event. Host copies ignore + these dependencies because host input is assumed ready. Returns ------- @@ -165,45 +216,77 @@ def copy(self, array: Any) -> StagedCopy: ``ready()`` tells whether the copy has finished, ``result()`` waits for it and returns the host array. """ + if stream is not None and event is not None: + raise ValueError("pass only one of stream and event") if tuple(array.shape) != self.shape or np.dtype(array.dtype) != self.dtype: raise ValueError( f"HostStaging for {self.shape} {self.dtype} got an array of shape " f"{tuple(array.shape)} and dtype {array.dtype}", ) device = _is_device_array(array) + if device: + from cunumpy._streams import HostEvent, HostStream + + if isinstance(event, HostEvent) or isinstance(stream, HostStream): + raise TypeError("device copies require a CUDA producer stream or event") + cp = _cupy() + device_id = array.device.id + if cp.cuda.runtime.getDevice() != device_id: + raise ValueError("make the source array's CUDA device current") + if self._device_id is not None and self._device_id != device_id: + raise ValueError("HostStaging is bound to another CUDA device") + if stream is not None and getattr(stream, "device_id", device_id) not in ( + device_id, + -1, + ): + raise ValueError( + "HostStaging producer stream belongs to another device" + ) + self._device_id = device_id slot = self._slot(device) if slot.event is not None: - slot.event.synchronize() # the buffer's previous copy is done + with slot.context(): + slot.event.synchronize() # the buffer's previous copy is done slot.generation += 1 if not device: np.copyto(slot.host, np.asarray(array)) slot.event = None return StagedCopy(slot, slot.generation) - cp = _cupy() if self._stream is None: self._stream = cp.cuda.Stream(non_blocking=True) - if slot.device is None: - slot.device = cp.empty(self.shape, dtype=self.dtype) - if _COUNTERS: - _record("to_host", f"HostStaging.copy({self.shape} {self.dtype})") - # snapshot on the producer's stream, after the kernels that wrote `array` - slot.device[...] = array - snapshot_done = cp.cuda.get_current_stream().record() + # Snapshot on the producer stream, or wait for a supplied event on the + # current stream. Later writes must be ordered after this snapshot. + producer = cp.cuda.get_current_stream() if stream is None else stream + with producer: + if event is not None: + producer.wait_event(event) + if slot.device is None: + slot.device = cp.empty(self.shape, dtype=self.dtype) + slot.device[...] = array + snapshot_done = producer.record() # copy the snapshot to the pinned buffer on the staging stream self._stream.wait_event(snapshot_done) try: - # blocking=False (CuPy >= 13): return once the copy is enqueued - slot.device.get(stream=self._stream, out=slot.host, blocking=False) - except ( - TypeError - ): # older CuPy: a copy on a stream into pinned memory is asynchronous - slot.device.get(stream=self._stream, out=slot.host) - slot.event = self._stream.record() + try: + # blocking=False (CuPy >= 13): return once the copy is enqueued + slot.device.get(stream=self._stream, out=slot.host, blocking=False) + except TypeError: # older CuPy + slot.device.get(stream=self._stream, out=slot.host) + if _COUNTERS: + _record( + "to_host", + f"HostStaging.copy({self.shape} {self.dtype})", + nbytes=_nbytes(array), + ) + finally: + # Retain completion even if a copy failed after enqueuing work. + slot.event = self._stream.record() return StagedCopy(slot, slot.generation) def synchronize(self) -> None: """Wait for all copies in flight.""" for slot in self._slots: if slot.event is not None: - slot.event.synchronize() + with slot.context(): + slot.event.synchronize() diff --git a/src/cunumpy/_transfers.py b/src/cunumpy/_transfers.py index 509b63e..f2d642e 100644 --- a/src/cunumpy/_transfers.py +++ b/src/cunumpy/_transfers.py @@ -11,7 +11,7 @@ propagator(dt) assert counter.total == 0, counter.report() -or, equivalently:: +or, to reject host/device copies and host execution on the GPU backend:: with xp.profiling.assert_no_transfers(): propagator(dt) @@ -25,13 +25,21 @@ * ``kernel_conversion``: a :class:`~cunumpy.kernels.PyccelKernel` call that copied device arrays to the host (and back), one event per call; * ``fallback``: a :class:`~cunumpy.kernels.Kernel` without CUDA kernel calling its host - kernel on the CuPy backend (``missing_cuda="fallback"``), one event per call. + kernel on the CuPy backend (``missing_cuda="fallback"``), one event per call; +* ``device_copy``: device-only dtype/layout conversions in CuNumpy helpers. + +Mirror refreshes, argument conversions, staging, serial MPI and kernel output +copy-back are also counted. Each physical host/device copy is recorded with its +payload size; conversion/fallback markers have no byte count. ``total`` counts +all observations, including markers. ``assert_no_transfers`` allows device-only +copies and rejects host/device copies and host fallback/conversion markers. Limitations ----------- Only transfers made *through cunumpy* are seen. Raw ``cupy.ndarray.get()``, ``cupy.asarray(numpy_array)``, ``numpy.asarray(cupy_array)``, ``float(device_array)``, -and implicit conversions inside other libraries are not counted; use ``nsys`` +forwarded backend operations such as ``xp.asarray`` and implicit conversions +inside other libraries are not counted; use ``nsys`` (or CuPy's own profiling hooks) to find those. Like the backend selection, the set of active counters is process-wide state @@ -41,6 +49,7 @@ from __future__ import annotations +import math import os import sys from collections.abc import Generator @@ -56,7 +65,7 @@ ] #: Event kinds, in the order they are reported. -KINDS = ("to_host", "to_device", "kernel_conversion", "fallback") +KINDS = ("to_host", "to_device", "kernel_conversion", "fallback", "device_copy") # The currently active counters, innermost last. Instrumented code checks # ``if _ACTIVE:`` before doing any work, so the overhead of an inactive counter @@ -75,17 +84,20 @@ class TransferEvent: ---------- kind : str One of ``"to_host"``, ``"to_device"``, ``"kernel_conversion"`` or - ``"fallback"``. + ``"fallback"``, or ``"device_copy"``. description : str What was transferred, e.g. ``"to_numpy(shape=(1000,), dtype=float64)"`` or ``"PyccelKernel 'push': 3 array(s) copied to the host"``. where : str The call site outside cunumpy, as ``"file:line"``. + nbytes : int | None + Payload bytes of a physical copy. None for markers or unknown sizes. """ kind: str description: str where: str + nbytes: int | None = None def __str__(self) -> str: return f"{self.where}: {self.kind}: {self.description}" @@ -141,9 +153,28 @@ def fallbacks(self) -> int: @property def total(self) -> int: - """Number of recorded transfers of all kinds.""" + """Number of observations, including conversion/fallback markers.""" return len(self.events) + def bytes(self, kind: str) -> int: + """Known bytes copied for `kind`; markers and unknown sizes add zero.""" + return sum(e.nbytes or 0 for e in self.events if e.kind == kind) + + @property + def bytes_to_host(self) -> int: + """Known payload bytes copied from device to host.""" + return self.bytes("to_host") + + @property + def bytes_to_device(self) -> int: + """Known payload bytes copied from host to device.""" + return self.bytes("to_device") + + @property + def device_copies(self) -> int: + """Device-only conversions recorded by CuNumpy helpers.""" + return self.count("device_copy") + def report(self) -> str: """A readable summary: events grouped by kind and call site, with counts.""" counts = ", ".join(f"{self.count(kind)} {kind}" for kind in KINDS) @@ -174,9 +205,9 @@ def _caller() -> str: return "" -def _record(kind: str, description: str) -> None: +def _record(kind: str, description: str, *, nbytes: int | None = None) -> None: """Record a transfer in every active counter (call only ``if _ACTIVE:``).""" - event = TransferEvent(kind, description, _caller()) + event = TransferEvent(kind, description, _caller(), nbytes) for counter in _ACTIVE: counter._add(event) @@ -190,6 +221,31 @@ def _describe(array: Any) -> str: return f"shape={tuple(shape)}, dtype={dtype}" +def _nbytes(array: Any) -> int | None: + """Array payload size, without copying data or reading a device scalar.""" + shape, dtype = getattr(array, "shape", None), getattr(array, "dtype", None) + if shape is None or dtype is None: + return None + return math.prod(shape) * dtype.itemsize + + +def _device_pointer(array: Any) -> int | None: + pointer = getattr(getattr(array, "data", None), "ptr", None) + if pointer is not None: + return pointer + interface = getattr(array, "__cuda_array_interface__", {}) + data = interface.get("data") + return None if data is None else data[0] + + +def _is_device_copy(source: Any, result: Any) -> bool: + """Distinguish a newly allocated conversion from a view of the same storage.""" + if result is source: + return False + pointer = _device_pointer(source) + return pointer is None or pointer != _device_pointer(result) + + @contextmanager def count_transfers() -> Generator[TransferCounter, None, None]: """Count the host/device transfers made through cunumpy in the block. @@ -225,7 +281,8 @@ def assert_no_transfers() -> Generator[TransferCounter, None, None]: """Raise ``AssertionError`` if the block makes a transfer through cunumpy. A :func:`count_transfers` block that, on exit, raises with the counter's - :meth:`~TransferCounter.report` if anything was counted. Only checked if + :meth:`~TransferCounter.report` if a host/device copy or host conversion/ + fallback was counted. Device-only conversions are allowed. Only checked if the block exits normally; an exception raised inside propagates as it is. Examples @@ -235,7 +292,7 @@ def assert_no_transfers() -> Generator[TransferCounter, None, None]: """ with count_transfers() as counter: yield counter - if counter.total: + if any(event.kind != "device_copy" for event in counter.events): raise AssertionError( "host/device transfers inside a block that must not transfer:\n" + counter.report(), diff --git a/src/cunumpy/cuda/include/cunumpy/reduce.cuh b/src/cunumpy/cuda/include/cunumpy/reduce.cuh index 06cf58b..ddcc422 100644 --- a/src/cunumpy/cuda/include/cunumpy/reduce.cuh +++ b/src/cunumpy/cuda/include/cunumpy/reduce.cuh @@ -136,6 +136,9 @@ __device__ T cunumpy_block_reduce(T v, Op op) __syncthreads(); if (warp == 0) { v = lane < n_warps ? partial[lane] : Op::fill(partial); + // Min/max padding lanes read partial[0], which lane 0 overwrites below. + // Shuffles do not order shared memory: finish every warp read first. + __syncwarp(mask); v = cunumpy_warp_reduce(v, op, mask); if (lane == 0) partial[0] = v; } diff --git a/src/cunumpy/xp.py b/src/cunumpy/xp.py index 3da3621..272c068 100644 --- a/src/cunumpy/xp.py +++ b/src/cunumpy/xp.py @@ -14,7 +14,13 @@ import array_api_compat.numpy as np from cunumpy._transfers import _ACTIVE as _COUNTERS -from cunumpy._transfers import _describe, _record +from cunumpy._transfers import ( + _describe, + _device_pointer, + _is_device_copy, + _nbytes, + _record, +) if os.environ.get("CUNUMPY_FAKE_CUPY", "").strip().lower() in ("1", "true", "yes"): # tests without a GPU: a strict host stand-in for CuPy, see cunumpy._fake_cupy @@ -281,20 +287,31 @@ def to_numpy(array: Any) -> np.ndarray: A CuPy array is copied to the host, which `count_transfers()` counts as a ``to_host`` transfer; anything else is passed through `numpy.asarray`. """ + result = _to_numpy(array) if _COUNTERS and get_array_backend(array) == "cupy": - _record("to_host", f"to_numpy({_describe(array)})") - return _to_numpy(array) + _record("to_host", f"to_numpy({_describe(array)})", nbytes=_nbytes(result)) + return result def to_cupy(array: Any) -> Any: """Convert an array to a CuPy array. - Anything that is not a CuPy array already is copied to the device, which - `count_transfers()` counts as a ``to_device`` transfer. + Host input is copied to the device and counted as a ``to_device`` transfer. + CUDA-array-interface inputs may be referenced without a copy; device-only + copies are recorded separately as ``device_copy``. """ - if _COUNTERS and get_array_backend(array) != "cupy": - _record("to_device", f"to_cupy({_describe(array)})") - return _to_cupy(array) + result = _to_cupy(array) + if _COUNTERS: + device = get_array_backend(array) == "cupy" or hasattr( + array, "__cuda_array_interface__" + ) + if not device: + _record("to_device", f"to_cupy({_describe(array)})", nbytes=_nbytes(result)) + elif _device_pointer(array) is not None and _is_device_copy(array, result): + _record( + "device_copy", f"to_cupy({_describe(array)})", nbytes=_nbytes(result) + ) + return result def as_device_array( @@ -304,18 +321,20 @@ def as_device_array( *, name: str | None = None, ) -> Any: - """Reference `value` on the device, or make one device copy of it. + """Reference `value` on the device, or convert it to a contiguous device array. The "reference or copy once" rule for building CUDA argument objects (`CudaArguments` subclasses, `CudaStruct` values): call it once when the argument object is built, never per kernel call. A CuPy array that already has the requested `dtype` (any dtype if `dtype` is None) and is C-contiguous is returned unchanged, the same object without a copy, so kernels write - into the caller's array. Anything else is converted with one device copy, + into the caller's array. Anything else is converted with ``cupy.ascontiguousarray(cupy.asarray(value, dtype))``: a tuple or list (e.g. ``degree = (3, 3, 3)``), a host NumPy array (one explicit transfer at build time), a device array of another dtype, or a non-contiguous view. The result passes the pointer checks of `CudaKernel` and `CudaStruct`. + Dtype and layout conversion can require separate device copies; each copy + through this helper is visible in transfer accounting. Raises on the NumPy backend: device argument objects are only built when running on CuPy, and host data is never copied to the device implicitly. @@ -361,7 +380,25 @@ def as_device_array( ): result = value else: - result = cp.ascontiguousarray(cp.asarray(value, dtype=dtype)) + converted = cp.asarray(value, dtype=dtype) + if _COUNTERS: + device_only = isinstance(value, cp.ndarray) or hasattr( + value, + "__cuda_array_interface__", + ) + if not device_only or _is_device_copy(value, converted): + _record( + "device_copy" if device_only else "to_device", + f"as_device_array({_describe(value)})", + nbytes=_nbytes(converted), + ) + result = cp.ascontiguousarray(converted) + if _COUNTERS and _is_device_copy(converted, result): + _record( + "device_copy", + f"as_device_array({_describe(converted)}) layout", + nbytes=_nbytes(result), + ) if ndim is not None and result.ndim != ndim: raise ValueError( f"{what} must have {ndim} dimension(s), got {result.ndim} " diff --git a/tests/conftest.py b/tests/conftest.py new file mode 100644 index 0000000..cb4ce45 --- /dev/null +++ b/tests/conftest.py @@ -0,0 +1,23 @@ +"""CUDA CI must fail if its hardware is unavailable rather than skip coverage.""" + +import os + +import pytest + + +def pytest_configure(config): + if os.environ.get("CUNUMPY_REQUIRE_CUDA", "").lower() not in {"1", "true", "yes"}: + return + import cunumpy as xp + + try: + xp.set_backend("cupy", strict=True) + import cupy as cp + + if getattr(cp, "__cunumpy_fake__", False): + raise RuntimeError("fake CuPy does not provide CUDA hardware coverage") + cp.zeros(1).get() # exercise allocation, execution and synchronization + except Exception as error: + raise pytest.UsageError( + f"CUDA CI requires a usable real GPU: {error}" + ) from error diff --git a/tests/unit/test_backend_numerics.py b/tests/unit/test_backend_numerics.py new file mode 100644 index 0000000..59ab52c --- /dev/null +++ b/tests/unit/test_backend_numerics.py @@ -0,0 +1,25 @@ +"""Backend numerical parity; performance comparisons belong in scope-profiler.""" + +import numpy as np +import pytest + +import cunumpy as xp + +pytestmark = pytest.mark.skipif(not xp.cupy_available(), reason="requires CUDA") + + +def test_matmul_matches_numpy(): + rng = np.random.default_rng(123) + a = rng.normal(size=(64, 96)).astype(np.float32) + b = rng.normal(size=(96, 32)).astype(np.float32) + with xp.use_backend("cupy", strict=True): + actual = xp.to_numpy(xp.asarray(a) @ xp.asarray(b)) + np.testing.assert_allclose(actual, a @ b, rtol=2e-5, atol=2e-5) + + +def test_fft_matches_numpy(): + rng = np.random.default_rng(123) + values = (rng.normal(size=1024) + 1j * rng.normal(size=1024)).astype(np.complex64) + with xp.use_backend("cupy", strict=True): + actual = xp.to_numpy(xp.fft.fft(xp.asarray(values))) + np.testing.assert_allclose(actual, np.fft.fft(values), rtol=2e-5, atol=2e-5) diff --git a/tests/unit/test_benchmarks.py b/tests/unit/test_benchmarks.py deleted file mode 100644 index e89a0e9..0000000 --- a/tests/unit/test_benchmarks.py +++ /dev/null @@ -1,81 +0,0 @@ -import time - -import pytest - -import cunumpy as xp - - -@pytest.mark.skipif( - not xp.cupy_available(), - reason="CuPy/GPU not available or not functional", -) -def test_benchmark_matmul(): - """Benchmark matrix multiplication to show CuPy performance gain.""" - size = 2000 - - # --- Benchmark NumPy --- - with xp.use_backend("numpy"): - a_np = xp.random.rand(size, size).astype(xp.float32) - b_np = xp.random.rand(size, size).astype(xp.float32) - - start_np = time.perf_counter() - _ = a_np @ b_np - # No sync needed for NumPy as it is synchronous - end_np = time.perf_counter() - t_np = end_np - start_np - - # --- Benchmark CuPy --- - with xp.use_backend("cupy"): - a_cp = xp.random.rand(size, size).astype(xp.float32) - b_cp = xp.random.rand(size, size).astype(xp.float32) - - # Warm up - _ = a_cp @ b_cp - xp.synchronize() - - start_cp = time.perf_counter() - _ = a_cp @ b_cp - xp.synchronize() # CRITICAL for benchmarking GPU - end_cp = time.perf_counter() - t_cp = end_cp - start_cp - - print(f"\n[Benchmark] Size: {size}x{size}") - print(f"NumPy time: {t_np:.4f}s") - print(f"CuPy time: {t_cp:.4f}s") - print(f"Speedup: {t_np / t_cp:.2f}x") - - # On a real GPU (A100/A30), CuPy should be significantly faster - # We use a conservative threshold of 1.5x for the test to pass on various hardware - assert t_cp < t_np, f"CuPy ({t_cp:.4f}s) was not faster than NumPy ({t_np:.4f}s)" - - -@pytest.mark.skipif( - not xp.cupy_available(), - reason="CuPy/GPU not available or not functional", -) -def test_benchmark_fft(): - """Benchmark FFT performance.""" - size = 2**22 # ~4 million elements - - with xp.use_backend("numpy"): - data_np = xp.random.rand(size).astype(xp.complex64) - start = time.perf_counter() - _ = xp.fft.fft(data_np) - t_np = time.perf_counter() - start - - with xp.use_backend("cupy"): - data_cp = xp.random.rand(size).astype(xp.complex64) - # Warm up - _ = xp.fft.fft(data_cp) - xp.synchronize() - - start = time.perf_counter() - _ = xp.fft.fft(data_cp) - xp.synchronize() - t_cp = time.perf_counter() - start - - print(f"\n[Benchmark] FFT Size: {size}") - print(f"NumPy time: {t_np:.4f}s") - print(f"CuPy time: {t_cp:.4f}s") - print(f"Speedup: {t_np / t_cp:.2f}x") - assert t_cp < t_np diff --git a/tests/unit/test_cuda_collectives.py b/tests/unit/test_cuda_collectives.py index 30b0fc5..3abfdc1 100644 --- a/tests/unit/test_cuda_collectives.py +++ b/tests/unit/test_cuda_collectives.py @@ -83,7 +83,9 @@ def test_partial_warps_multidimensional_blocks_and_repeated_collectives(block): data = cp.asarray(host) inc, exc, again = (cp.empty_like(data) for _ in range(3)) sums, lo, hi = (cp.empty(3) for _ in range(3)) - kernel = CudaKernel(BLOCK_SOURCE, "collectives", block_size=block) + kernel = CudaKernel( + BLOCK_SOURCE, "collectives", block_size=block, options=("-lineinfo",) + ) kernel(data, inc, exc, again, sums, lo, hi, grid=3) rows = host.reshape(3, count) expected_inc = np.cumsum(rows, axis=1) @@ -95,6 +97,54 @@ def test_partial_warps_multidimensional_blocks_and_repeated_collectives(block): np.testing.assert_array_equal(cp.asnumpy(got), expected) +REPEATED_EXTREMA_SOURCE = r""" +#include +template +__global__ void repeated_extrema(const T* x, T* minima, T* maxima, int rounds) { + int t = cunumpy_block_thread(); + int count = cunumpy_block_threads(); + int i = blockIdx.x * count + t; + int stride = gridDim.x * count; + for (int r = 0; r < rounds; ++r) { + T v = ((r & 1) ? -x[i] : x[i]) + T(r); + T lo = cunumpy_block_min(v); + T hi = cunumpy_block_max(v); + minima[r * stride + i] = lo; + maxima[r * stride + i] = hi; + } +} +""" + + +@requires_cuda +@pytest.mark.parametrize("dtype", [np.int32, np.float64]) +@pytest.mark.parametrize("block", [17, 32, 33, 128, 1024, (7, 5)]) +def test_repeated_extrema_are_returned_to_every_thread(dtype, block): + import cupy as cp + + count, blocks, rounds = int(np.prod(block)), 3, 8 + host = np.random.default_rng(123).integers(-100, 100, blocks * count).astype(dtype) + data = cp.asarray(host) + minima, maxima = (cp.empty((rounds, data.size), dtype=dtype) for _ in range(2)) + kernel = CudaKernel( + REPEATED_EXTREMA_SOURCE, + "repeated_extrema", + block_size=block, + template_args=(dtype,), + options=("-lineinfo",), + ) + kernel(data, minima, maxima, rounds, grid=blocks) + rows = host.reshape(blocks, count) + values = np.stack([(rows if r % 2 == 0 else -rows) + r for r in range(rounds)]) + for actual, expected in ( + (minima, values.min(axis=2)), + (maxima, values.max(axis=2)), + ): + np.testing.assert_array_equal( + cp.asnumpy(actual), np.repeat(expected, count, axis=1) + ) + + MASKED_SOURCE = r""" #include extern "C" __global__ void masked(const double* x, unsigned int mask, @@ -223,7 +273,7 @@ def test_scan_header_is_resolved_and_emulation_refuses_warp_intrinsics(): ] if emulation_compiler() is None: pytest.skip("requires a C++ compiler") - with pytest.raises(NotImplementedError, match="warp shuffles"): + with pytest.raises(NotImplementedError, match="warp (shuffles|synchronization)"): emulate_cuda_kernel(kernel, *(np.zeros(32) for _ in range(7)), n_threads=32) @@ -240,6 +290,7 @@ def test_collective_headers_compile_with_all_supported_arithmetic_types(): inline int __ffs(unsigned v) { return v ? __builtin_ctz(v) + 1 : 0; } inline int __clz(unsigned v) { return __builtin_clz(v); } inline void __syncthreads() {} +inline void __syncwarp(unsigned) {} template T __shfl_sync(unsigned, T v, int) { return v; } template T __shfl_up_sync(unsigned, T v, int) { return v; } template T __shfl_xor_sync(unsigned, T v, int) { return v; } @@ -270,6 +321,7 @@ def test_collective_headers_compile_with_all_supported_arithmetic_types(): input=_STUBS + stubs + BLOCK_SOURCE + + REPEATED_EXTREMA_SOURCE + MASKED_SOURCE + ATOMICS_SOURCE + instantiate, diff --git a/tests/unit/test_cuda_kernel.py b/tests/unit/test_cuda_kernel.py index 83eed65..68b0841 100644 --- a/tests/unit/test_cuda_kernel.py +++ b/tests/unit/test_cuda_kernel.py @@ -13,7 +13,7 @@ import pytest import cunumpy as xp -import cunumpy._device as device_module +import cunumpy._cuda_kernel as cuda_module from cunumpy import as_device_array from cunumpy.cuda import ( CudaArguments, @@ -72,6 +72,7 @@ class FakeDeviceArray: def __init__(self, dtype, ptr=0x1000, shape=(1,), strides=None, flags=None): self.dtype = np.dtype(dtype) + self.device = SimpleNamespace(id=0) self.data = SimpleNamespace(ptr=ptr) self.shape = tuple(shape) self.ndim = len(self.shape) @@ -2186,6 +2187,21 @@ def recorded(monkeypatch): """A CudaKernel whose launches are recorded instead of run.""" kernel = CudaKernel(AXPY, "axpy", block_size=128) raw = RecordingRawKernel() + limits = { + "name": b"recording GPU", + "maxThreadsPerBlock": 1024, + "maxThreadsDim": (1024, 1024, 64), + "maxGridSize": (2**31 - 1, 65535, 65535), + "sharedMemPerBlock": 48 * 1024, + "sharedMemPerBlockOptin": 100_000, + } + kernel._compiled[0] = cuda_module._CompiledKernel( + raw, + limits, + {"shared_size_bytes": 0, "max_threads_per_block": 1024}, + 48 * 1024, + ) + monkeypatch.setattr(cuda_module, "_current_device", lambda: 0) monkeypatch.setattr(kernel, "compile", lambda: raw) monkeypatch.setattr(kernel, "debug_active", lambda: False) return kernel, raw @@ -2208,13 +2224,8 @@ def test_n_threads_from_first_array(recorded): as_option.n_threads_from((1.0, 2)) -def test_shared_memory_above_the_default_is_opted_in(recorded, monkeypatch): +def test_shared_memory_above_the_default_is_opted_in(recorded): kernel, raw = recorded - monkeypatch.setattr( - device_module, - "max_shared_memory_per_block", - lambda opt_in=False: 100_000, - ) x, y = FakeDeviceArray(np.float64), FakeDeviceArray(np.float64) kernel(1.0, x, y, 1, n_threads=1, shared_mem=40_000) # below 48 KiB: no setup assert raw.max_dynamic_shared_size_bytes == 48 * 1024 diff --git a/tests/unit/test_cuda_launch_limits.py b/tests/unit/test_cuda_launch_limits.py new file mode 100644 index 0000000..4834bc5 --- /dev/null +++ b/tests/unit/test_cuda_launch_limits.py @@ -0,0 +1,329 @@ +"""Compilation/launch contracts with fake runtime limits and real GPU coverage.""" + +import io +import sys +import threading +import types + +import pytest + +import cunumpy as xp +from cunumpy import _cuda_kernel as implementation +from cunumpy.cuda import CudaKernel +from cunumpy.kernels import Kernel, KernelCatalog + +EMPTY = 'extern "C" __global__ void empty() {}' +requires_gpu = pytest.mark.skipif(not xp.cupy_available(), reason="requires CUDA") + + +@pytest.fixture +def runtime(monkeypatch): + current = threading.local() + current.device = 0 + raw_kernels, queries, compilations = [], [], [] + properties = { + i: { + "name": f"test GPU {i}".encode(), + "maxThreadsDim": (1024, 1024, 64), + "maxGridSize": (2**31 - 1, 65535, 65535), + "maxThreadsPerBlock": 1024, + "sharedMemPerBlock": 49152, + "sharedMemPerBlockOptin": 98304, + } + for i in range(2) + } + + def device_id(): + return getattr(current, "device", 0) + + class Device: + def __init__(self, device): + self.id = device + + def __enter__(self): + self.previous = device_id() + current.device = self.id + return self + + def __exit__(self, *exc): + current.device = self.previous + + class Raw: + def __init__(self, source, name, options=()): + self.source, self.name, self.options = source, name, options + self.device = device_id() + self.launches = [] + self.attributes = { + "max_threads_per_block": 256, + "shared_size_bytes": 4096, + "max_dynamic_shared_size_bytes": 49152, + } + self.max_dynamic_shared_size_bytes = 49152 + raw_kernels.append(self) + + def compile(self, log_stream=None): + compilations.append((self.device, self.name, self.options)) + if "bad_syntax" in self.source: + raise RuntimeError("NVRTC rejected the source") + if log_stream is not None: + log_stream.write("compiled\n") + + def __call__(self, grid, block, args, shared_mem=0): + assert device_id() == self.device + self.launches.append((grid, block, shared_mem)) + + class Module: + def __init__(self, code, options, name_expressions): + self.raw = Raw(code, name_expressions[0], options) + self.compiled = False + + def compile(self, log_stream=None): + self.raw.compile(log_stream) + self.compiled = True + + def get_function(self, name): + assert self.compiled + assert name == self.raw.name + return self.raw + + def get_properties(device): + queries.append(device) + return properties[device] + + cp = types.SimpleNamespace( + RawKernel=Raw, + RawModule=Module, + cuda=types.SimpleNamespace( + Device=Device, + runtime=types.SimpleNamespace( + getDevice=device_id, getDeviceProperties=get_properties + ), + ), + ) + monkeypatch.setitem(sys.modules, "cupy", cp) + monkeypatch.setattr(xp.xp, "cupy_available", lambda: True) + monkeypatch.setattr(implementation, "_DEVICE_LIMITS", {}) + return types.SimpleNamespace( + cp=cp, + current=current, + properties=properties, + raw=raw_kernels, + queries=queries, + compilations=compilations, + ) + + +def test_compile_is_eager_cached_and_device_specific(runtime): + kernel = CudaKernel(EMPTY, "empty") + assert not kernel.is_compiled + log = io.StringIO() + first = kernel.compile(log_stream=log) + assert log.getvalue() == "compiled\n" + assert kernel.is_compiled + assert kernel.compile() is first + with runtime.cp.cuda.Device(1): + assert not kernel.is_compiled + second = kernel.compile() + assert second is not first + assert kernel.is_compiled + assert kernel.compile() is first + assert [c[0] for c in runtime.compilations] == [0, 1] + assert runtime.queries == [0, 1] + + +def test_failed_compile_is_not_cached_and_can_be_retried(runtime, monkeypatch): + kernel = CudaKernel(EMPTY, "empty") + original = runtime.cp.RawKernel.compile + + def fail(self, log_stream=None): + raise RuntimeError("compiler failed") + + monkeypatch.setattr(runtime.cp.RawKernel, "compile", fail) + with pytest.raises(RuntimeError, match="compiler failed"): + kernel.compile() + assert not kernel.is_compiled + monkeypatch.setattr(runtime.cp.RawKernel, "compile", original) + assert kernel.compile() is runtime.raw[-1] + assert kernel.is_compiled + + +def test_recompile_refreshes_headers_and_debug_on_only_current_device( + runtime, tmp_path +): + header = tmp_path / "value.cuh" + header.write_text("#define VALUE 1\n") + kernel = CudaKernel( + '#include "value.cuh"\n' + EMPTY, "empty", include_dirs=[tmp_path] + ) + first = kernel.compile() + with runtime.cp.cuda.Device(1): + second = kernel.compile() + header.write_text("#define VALUE 2\n") + with xp.cuda.cuda_debug(): + rebuilt = kernel.recompile() + assert rebuilt is not first + assert rebuilt.options != first.options + assert "-lineinfo" in rebuilt.options + with runtime.cp.cuda.Device(1): + assert kernel.compile() is second + + +def test_template_is_compiled_before_function_is_returned(runtime): + kernel = CudaKernel( + "template __global__ void empty() {}", + "empty", + template_args=["double"], + ) + assert kernel.compile().name == "empty" + assert len(runtime.compilations) == 1 + + +def test_catalog_fails_during_setup_and_workers_keep_selected_device(runtime): + good = CudaKernel(EMPTY, "empty") + bad = CudaKernel(EMPTY + "\nbad_syntax", "empty") + catalog = KernelCatalog( + {"good": Kernel(lambda: None, good), "bad": Kernel(lambda: None, bad)} + ) + with runtime.cp.cuda.Device(1): + with pytest.raises(RuntimeError, match="NVRTC rejected"): + catalog.compile_all(jobs=2) + assert good.is_compiled + assert not bad.is_compiled + assert all(c[0] == 1 for c in runtime.compilations) + + +@pytest.mark.parametrize( + "launch, expected", + [ + ({"grid": (1, 65536)}, "grid\\[1\\]=65536"), + ({"grid": 1, "block": (1, 1, 65)}, "block\\[2\\]=65"), + ({"grid": 1, "block": 512}, "256 device/kernel limit"), + ({"grid": 1, "shared_mem": 98304 - 4096 + 1}, "static shared memory 4096"), + ], +) +def test_actual_limits_reject_invalid_launch_before_execution( + runtime, launch, expected +): + kernel = CudaKernel(EMPTY, "empty") + with pytest.raises(ValueError, match=expected) as error: + kernel(**launch) + assert "test GPU 0" in str(error.value) + assert runtime.raw[0].launches == [] + + +def test_device_limits_can_be_stricter_than_kernel_limits(runtime): + runtime.properties[0]["maxThreadsPerBlock"] = 64 + with pytest.raises(ValueError, match="64 device/kernel limit"): + CudaKernel(EMPTY, "empty")(grid=1) + + +def test_static_memory_requires_opt_in_even_with_large_initial_dynamic_attribute( + runtime, +): + kernel = CudaKernel(EMPTY, "empty") + raw = kernel.compile() + state = kernel._compiled[0] + assert state.dynamic_shared == 45056 + kernel(grid=1, shared_mem=49152 - 4096 + 1) + assert raw.max_dynamic_shared_size_bytes == 45057 + + +def test_opt_in_accounts_for_static_memory_and_is_per_device(runtime): + kernel = CudaKernel(EMPTY, "empty") + kernel(grid=1, shared_mem=45057) + first = runtime.raw[-1] + assert first.max_dynamic_shared_size_bytes == 45057 + kernel(grid=1, shared_mem=60000) + kernel(grid=1, shared_mem=50000) + assert first.max_dynamic_shared_size_bytes == 60000 + with runtime.cp.cuda.Device(1): + kernel(grid=1, shared_mem=45057) + assert runtime.raw[-1].max_dynamic_shared_size_bytes == 45057 + assert len(runtime.queries) == 2 + assert len(runtime.compilations) == 2 + + +def test_stream_on_another_device_is_rejected(runtime): + with pytest.raises(ValueError, match="stream belongs to another device"): + CudaKernel(EMPTY, "empty")(grid=1, stream=types.SimpleNamespace(device_id=1)) + + +@pytest.mark.parametrize("wrapped", [False, True]) +def test_array_on_another_device_is_rejected_before_launch(runtime, wrapped): + array = types.SimpleNamespace( + __cuda_array_interface__={}, device=types.SimpleNamespace(id=1) + ) + argument = ( + types.SimpleNamespace(__cuda_args__=lambda: (array,)) if wrapped else array + ) + kernel = CudaKernel(EMPTY, "empty", check_signature=False) + with pytest.raises(ValueError, match="belongs to CUDA device 1"): + kernel(argument, grid=1) + assert runtime.raw[0].launches == [] + + +def test_empty_launch_does_not_compile(runtime): + CudaKernel(EMPTY, "empty")(n_threads=0) + assert runtime.compilations == [] + + +@requires_gpu +def test_invalid_cuda_fails_at_compile_not_first_launch(): + kernel = CudaKernel( + 'extern "C" __global__ void broken() { missing_symbol(); }', "broken" + ) + with pytest.raises(Exception, match="missing_symbol"): + kernel.compile() + assert not kernel.is_compiled + + +@requires_gpu +def test_kernel_launch_bounds_are_enforced(): + kernel = CudaKernel( + 'extern "C" __global__ __launch_bounds__(64) void limited() {}', + "limited", + check_signature=False, + block_size=64, + ) + raw = kernel.compile() + assert raw.attributes["max_threads_per_block"] == 64 + with pytest.raises(ValueError, match="64 device/kernel limit"): + kernel(grid=1, block=128) + kernel(grid=1) + xp.synchronize() + + +@requires_gpu +def test_same_kernel_compiles_on_multiple_devices(): + import cupy as cp + + if cp.cuda.runtime.getDeviceCount() < 2: + pytest.skip("requires two devices") + kernel = CudaKernel(EMPTY, "empty") + for device in (0, 1, 0): + with cp.cuda.Device(device): + kernel(grid=1) + assert kernel.is_compiled + cp.cuda.get_current_stream().synchronize() + + +@requires_gpu +def test_static_shared_memory_is_included_in_actual_gpu_limit(): + import cupy as cp + + source = r"""extern "C" __global__ void scratch(const double* x, double* out) { + __shared__ double values[128]; + values[threadIdx.x] = x[threadIdx.x]; + __syncthreads(); + if (threadIdx.x == 0) out[0] = values[127]; + }""" + kernel = CudaKernel(source, "scratch") + raw = kernel.compile() + static = raw.attributes["shared_size_bytes"] + assert static >= 128 * 8 + properties = cp.cuda.runtime.getDeviceProperties(cp.cuda.runtime.getDevice()) + limit = max( + properties["sharedMemPerBlock"], properties.get("sharedMemPerBlockOptin", 0) + ) + with pytest.raises(ValueError, match="static shared memory"): + kernel(cp.ones(128), cp.empty(1), n_threads=128, shared_mem=limit - static + 1) diff --git a/tests/unit/test_emulation.py b/tests/unit/test_emulation.py index ff21718..2bc3640 100644 --- a/tests/unit/test_emulation.py +++ b/tests/unit/test_emulation.py @@ -182,7 +182,7 @@ def test_warp_shuffles_in_an_included_header_are_refused(): " double s = cunumpy_block_sum(1.0); if (threadIdx.x == 0) out[0] = s; }", "k", ) - with pytest.raises(NotImplementedError, match="warp shuffles"): + with pytest.raises(NotImplementedError, match="warp (shuffles|synchronization)"): emulate_cuda_kernel(kernel, np.zeros(1), 1, n_threads=32) diff --git a/tests/unit/test_execution_transfers.py b/tests/unit/test_execution_transfers.py new file mode 100644 index 0000000..da5ba0a --- /dev/null +++ b/tests/unit/test_execution_transfers.py @@ -0,0 +1,241 @@ +"""Physical copies, fallback markers, and retained-buffer ownership without CUDA.""" + +import sys +import types + +import numpy as np +import pytest + +import cunumpy as xp +from cunumpy import _mirror as mirrors +from cunumpy.memory import DeviceMirror + + +@pytest.fixture +def device(monkeypatch): + cp = types.SimpleNamespace(current=0, copies=[], waits=[]) + + class Array: + def __init__(self, data): + self.array = np.asarray(data) + self.data = types.SimpleNamespace(ptr=self.array.ctypes.data) + self.device = types.SimpleNamespace(id=cp.current) + self.shape, self.dtype, self.flags = ( + self.array.shape, + self.array.dtype, + self.array.flags, + ) + self.ndim = self.array.ndim + + def set(self, host): + cp.copies.append("set") + np.copyto(self.array, host) + + def get(self, out=None, stream=None): + cp.copies.append(("get", stream)) + if out is None: + return self.array.copy() + np.copyto(out, self.array) + return out + + def fill(self, value): + self.array.fill(value) + + def __setitem__(self, index, value): + self.array[index] = value.array if isinstance(value, Array) else value + + def asarray(value, dtype=None): + if isinstance(value, Array): + if dtype is None or value.dtype == np.dtype(dtype): + return value + value = value.array + return Array(np.array(value, dtype=dtype, copy=True)) + + def contiguous(value, dtype=None): + value = asarray(value, dtype=dtype) + return ( + value + if value.flags.c_contiguous + else Array(np.ascontiguousarray(value.array)) + ) + + cp.ndarray = Array + cp.asarray = asarray + cp.ascontiguousarray = contiguous + cp.zeros = lambda shape, dtype: Array(np.zeros(shape, dtype=dtype)) + cp.cuda = types.SimpleNamespace( + runtime=types.SimpleNamespace(getDevice=lambda: cp.current), + get_current_stream=lambda: types.SimpleNamespace(wait_event=cp.waits.append), + ) + monkeypatch.setitem(sys.modules, "cupy", cp) + monkeypatch.setattr(xp.xp.array_backend, "_backend", "cupy") + monkeypatch.setattr(xp.xp, "cupy_available", lambda: True) + monkeypatch.setattr(mirrors, "_cupy_backend", lambda: xp.get_backend() == "cupy") + monkeypatch.setattr(xp.xp, "_to_cupy", asarray) + monkeypatch.setattr( + xp.xp, + "get_array_backend", + lambda a: "cupy" if isinstance(a, Array) else "numpy", + ) + return cp + + +def test_mirror_initialization_and_refreshes_count_each_copy_once(device): + host = np.arange(4.0) + mirror = DeviceMirror(host) + with xp.profiling.count_transfers() as counter: + mirror.to_device() + mirror.to_device() + mirror.to_host() + assert counter.to_device == 2 and counter.to_host == 1 + assert counter.bytes_to_device == 2 * host.nbytes + assert counter.bytes_to_host == host.nbytes + assert [e.nbytes for e in counter.events] == [host.nbytes] * 3 + assert mirror.host is host + + +def test_mirror_refresh_breaks_no_transfer_assertion(device): + mirror = DeviceMirror(np.ones(2)) + mirror.to_device() + with ( + pytest.raises(AssertionError, match="DeviceMirror.to_device"), + xp.profiling.assert_no_transfers(), + ): + mirror.to_device() + with ( + pytest.raises(AssertionError, match="DeviceMirror.to_host"), + xp.profiling.assert_no_transfers(), + ): + mirror.to_host() + + +def test_zero_allocates_without_transferring_initial_host_contents(device): + mirror = DeviceMirror(np.ones(3)) + with xp.profiling.assert_no_transfers() as counter: + mirror.zero() + assert mirror.device is not mirror.host + assert counter.total == 0 + np.testing.assert_array_equal(mirror.device.array, np.zeros(3)) + + +@pytest.mark.parametrize("operation", ["device", "to_device", "to_host", "zero"]) +def test_mirror_rejects_another_current_device_before_touching_storage( + device, operation +): + mirror = DeviceMirror(np.ones(3)) + mirror.to_device() + device.current = 1 + with pytest.raises(ValueError, match="bound to CUDA device 0"): + getattr(mirror, operation)() if operation != "device" else mirror.device + assert device.copies == [] + + +def test_mirror_explicit_dependencies_and_wrong_stream(device): + mirror = DeviceMirror(np.ones(2)) + mirror.to_device() + producer = types.SimpleNamespace(device_id=0) + mirror.to_host(stream=producer) + assert device.copies[-1] == ("get", producer) + event = object() + mirror.to_host(event=event) + assert device.waits == [event] + with pytest.raises(ValueError, match="only one"): + mirror.to_host(stream=producer, event=event) + with pytest.raises(ValueError, match="another device"): + mirror.to_host(stream=types.SimpleNamespace(device_id=1)) + with pytest.raises(TypeError, match="CUDA producer"): + mirror.to_host(event=xp.cuda.HostEvent()) + + +def test_cpu_backend_refresh_does_not_touch_retained_device_buffer(device): + mirror = DeviceMirror(np.ones(2)) + mirror.to_device() + with xp.use_backend("numpy"), xp.profiling.assert_no_transfers(): + assert mirror.device is mirror.host + mirror.zero() + mirror.to_device().to_host() + assert device.copies == [] + mirror.to_device() + np.testing.assert_array_equal(mirror.device.array, [0.0, 0.0]) + + +def test_host_argument_conversion_counts_bytes_and_identity_does_not(device): + with xp.profiling.count_transfers() as counter: + array = xp.as_device_array([1, 2, 3], dtype=np.float32) + assert xp.as_device_array(array, dtype=np.float32) is array + assert counter.to_device == 1 and counter.bytes_to_device == 12 + with ( + pytest.raises(AssertionError, match="as_device_array"), + xp.profiling.assert_no_transfers(), + ): + xp.as_device_array(np.ones(2)) + + +def test_device_only_dtype_and_layout_copies_are_distinguished(device): + source = device.ndarray(np.zeros((4, 3))[:, ::2]) + with xp.profiling.assert_no_transfers() as counter: + packed = xp.as_device_array(source) + cast = xp.as_device_array(packed, dtype=np.float32) + assert packed.flags.c_contiguous and cast.dtype == np.float32 + assert counter.to_host == counter.to_device == 0 + assert counter.device_copies == 2 + assert counter.bytes("device_copy") == packed.array.nbytes + cast.array.nbytes + + +def test_failed_transfer_is_not_recorded(device, monkeypatch): + def fail(value): + raise RuntimeError("copy failed") + + monkeypatch.setattr(xp.xp, "_to_cupy", fail) + with ( + xp.profiling.count_transfers() as counter, + pytest.raises(RuntimeError, match="copy failed"), + ): + xp.to_cupy(np.zeros(2)) + assert counter.total == 0 + + +def test_kernel_output_counts_host_back_to_device_copy(device): + out = device.ndarray(np.zeros(4)) + with ( + xp.profiling.count_transfers() as counter, + xp.kernels.kernel_output(out, like=np.zeros(4)) as buffer, + ): + buffer[:] = 3.0 + assert counter.to_host == counter.to_device == 1 + assert counter.bytes_to_host == counter.bytes_to_device == out.array.nbytes + np.testing.assert_array_equal(out.array, np.full(4, 3.0)) + + +def test_interoperable_device_reference_does_not_report_host_transfer(device): + array = device.ndarray(np.ones(3)) + + class Exporter: + def __init__(self): + self.__cuda_array_interface__ = {"data": (array.data.ptr, False)} + + exported = Exporter() + device.asarray = lambda value, dtype=None: array + with xp.profiling.assert_no_transfers() as counter: + assert xp.as_device_array(exported) is array + assert counter.total == 0 + + +@pytest.mark.parametrize("helper", ["as_device_array", "as_kernel_array"]) +def test_reference_only_layout_view_is_not_counted_as_copy(device, helper): + scalar = device.ndarray(np.array(3.0)) + view = device.ndarray(scalar.array.reshape(1)) + device.ascontiguousarray = lambda value, dtype=None: view + if helper == "as_device_array": + # An interoperable exporter takes the conversion path. The resulting + # one-dimensional view still references the original scalar's storage. + exported = types.SimpleNamespace( + __cuda_array_interface__={"data": (scalar.data.ptr, False)} + ) + device.asarray = lambda value, dtype=None: scalar + convert, value, kwargs = xp.as_device_array, exported, {} + else: + convert, value, kwargs = xp.kernels.as_kernel_array, scalar, {"like": scalar} + with xp.profiling.count_transfers() as counter: + assert convert(value, **kwargs) is view + assert counter.total == 0 diff --git a/tests/unit/test_mirror.py b/tests/unit/test_mirror.py index 769fba2..4b67b6e 100644 --- a/tests/unit/test_mirror.py +++ b/tests/unit/test_mirror.py @@ -19,6 +19,30 @@ reason="CuPy/GPU not available or not functional", ) + +@requires_gpu +@pytest.mark.parametrize("dependency", ["stream", "event"]) +def test_mirror_to_host_after_nondefault_producer(dependency): + import cupy as cp + + host = np.zeros(1000) + with xp.use_backend("cupy"): + mirror = DeviceMirror(host) + mirror.zero() + xp.synchronize() + producer = cp.cuda.Stream(non_blocking=True) + event = cp.cuda.Event(disable_timing=True) + with producer: + mirror.device.fill(7.0) + event.record() + with xp.profiling.count_transfers() as counter: + mirror.to_host( + **{dependency: producer if dependency == "stream" else event} + ) + np.testing.assert_array_equal(host, np.full(1000, 7.0)) + assert counter.bytes_to_host == host.nbytes + + BIN_ADD = r""" #include diff --git a/tests/unit/test_porting_helpers.py b/tests/unit/test_porting_helpers.py index dbdfaa8..3e89c33 100644 --- a/tests/unit/test_porting_helpers.py +++ b/tests/unit/test_porting_helpers.py @@ -274,18 +274,29 @@ def __init__(self): # --------------------------------------------------------------------------- -def _launching(kernel): +def _launching(kernel, monkeypatch): """Replace compilation by a recorder of (grid, block) launches.""" launches = [] - kernel.compile = lambda: ( - lambda grid, block, values, shared_mem: launches.append(grid) + raw = lambda grid, block, values, shared_mem: launches.append(grid) + kernel.compile = lambda: raw + kernel._compiled[0] = cuda_kernel_module._CompiledKernel( + raw, + { + "maxThreadsDim": (1024, 1024, 64), + "maxGridSize": (2**31 - 1, 65535, 65535), + "maxThreadsPerBlock": 1024, + "sharedMemPerBlock": 49152, + }, + {"max_threads_per_block": 1024, "shared_size_bytes": 0}, + 49152, ) + monkeypatch.setattr(cuda_kernel_module, "_current_device", lambda: 0) return launches -def test_n_threads_from(): +def test_n_threads_from(monkeypatch): kernel = CudaKernel(SCALE, "scale", n_threads_from=lambda args: args[2]) - launches = _launching(kernel) + launches = _launching(kernel, monkeypatch) kernel(FakeDeviceArray(np.float64), 2.0, 300) assert launches == [(3,)] kernel(FakeDeviceArray(np.float64), 2.0, 300, n_threads=10) # explicit wins @@ -300,7 +311,7 @@ def test_n_threads_from(): def test_check_finite(monkeypatch): kernel = CudaKernel(SCALE, "scale", check_finite=True) assert kernel.check_finite - _launching(kernel) + _launching(kernel, monkeypatch) kernel._synchronize_after_launch = lambda *a: None fake_cupy = ModuleType("cupy") fake_cupy.isfinite = np.isfinite diff --git a/tests/unit/test_staging.py b/tests/unit/test_staging.py index 679e02a..78f4ff3 100644 --- a/tests/unit/test_staging.py +++ b/tests/unit/test_staging.py @@ -1,6 +1,7 @@ """Tests for `xp.memory.HostStaging`: background copies of device arrays to the host.""" import types +from contextlib import nullcontext import numpy as np import pytest @@ -49,6 +50,10 @@ class DeviceArray(np.ndarray): log = None + @property + def device(self): + return types.SimpleNamespace(id=getattr(self, "_device_id", 0)) + def get(self, stream=None, out=None, blocking=True): assert blocking is False # an asynchronous copy DeviceArray.log.append(("get", stream.name)) @@ -70,6 +75,13 @@ class FakeStream: def __init__(self, name): self.name = name self.count = 0 + self.device_id = 0 + + def __enter__(self): + return self + + def __exit__(self, *exc): + pass def record(self): self.count += 1 @@ -88,6 +100,8 @@ def fake_device(monkeypatch): cuda = types.SimpleNamespace( Stream=lambda non_blocking=False: FakeStream("staging"), get_current_stream=lambda: current, + Device=lambda device: nullcontext(), + runtime=types.SimpleNamespace(getDevice=lambda: 0), ) def empty(shape, dtype): @@ -145,6 +159,101 @@ def test_copies_are_counted_as_transfers(fake_device): with xp.profiling.count_transfers() as counter: staging.copy(np.zeros(2).view(DeviceArray)) assert counter.total == 1 + assert counter.bytes_to_host == 16 + + +def test_explicit_producer_stream_and_event_order_snapshot(fake_device): + log = fake_device + source = np.ones(3).view(DeviceArray) + producer = FakeStream("producer") + staging = HostStaging(3, source.dtype) + staging.copy(source, stream=producer) + assert ("record", "producer#1") in log + log.clear() + event = FakeEvent("produced") + staging.copy(source, event=event) + assert log.index(("wait", "compute", "produced")) < log.index( + ("record", "compute#1") + ) + assert ("synchronize", "produced") not in log + with pytest.raises(ValueError, match="only one"): + staging.copy(source, stream=producer, event=event) + + +def test_cpu_initialized_staging_upgrades_to_pinned_and_preserves_old_copy( + fake_device, monkeypatch +): + allocations = [] + + def pinned(shape, dtype): + allocations.append(shape) + return np.empty(shape, dtype) + + monkeypatch.setattr(staging_module, "_empty_pinned", pinned) + staging = HostStaging(2, np.float64) + old = staging.copy(np.array([1.0, 2.0])) + device_copy = staging.copy(np.array([3.0, 4.0]).view(DeviceArray)) + assert len(allocations) == staging.buffers + np.testing.assert_array_equal(old.result(), [1.0, 2.0]) + np.testing.assert_array_equal(device_copy.result(), [3.0, 4.0]) + assert old.result() is not device_copy.result() + + +def test_device_binding_is_checked_before_reusing_any_buffer(fake_device, monkeypatch): + staging = HostStaging(2, np.float64) + source = np.ones(2).view(DeviceArray) + copy = staging.copy(source) + source._device_id = 1 + with pytest.raises(ValueError, match="device current"): + staging.copy(source) + monkeypatch.setattr(staging_module._cupy().cuda.runtime, "getDevice", lambda: 1) + with pytest.raises(ValueError, match="bound to another"): + staging.copy(source) + assert copy.result().tolist() == [1.0, 1.0] + + +def test_wrong_stream_and_host_dependencies_are_rejected(fake_device): + staging = HostStaging(2, np.float64) + source = np.ones(2).view(DeviceArray) + stream = FakeStream("wrong") + stream.device_id = 1 + with pytest.raises(ValueError, match="stream belongs"): + staging.copy(source, stream=stream) + with pytest.raises(TypeError, match="CUDA producer"): + staging.copy(source, event=xp.cuda.HostEvent()) + + +def test_failed_copy_keeps_completion_before_slot_reuse(fake_device, monkeypatch): + staging = HostStaging(1, np.float64, buffers=1) + source = np.ones(1).view(DeviceArray) + original = DeviceArray.get + + def fail(self, **kwargs): + raise RuntimeError("copy failed after enqueue") + + monkeypatch.setattr(DeviceArray, "get", fail) + with pytest.raises(RuntimeError, match="copy failed"): + staging.copy(source) + monkeypatch.setattr(DeviceArray, "get", original) + result = staging.copy(source) + assert ("synchronize", "staging#1") in fake_device + assert result.result().tolist() == [1.0] + + +def test_host_copy_waits_for_previous_device_slot(fake_device): + staging = HostStaging(1, np.float64, buffers=1) + old = staging.copy(np.ones(1).view(DeviceArray)) + result = staging.copy(np.array([2.0])) + assert ("synchronize", "staging#1") in fake_device + with pytest.raises(RuntimeError, match="reused"): + old.result() + assert result.result().tolist() == [2.0] + + +@pytest.mark.parametrize("shape,dtype", [((-1,), float), ((2,), object)]) +def test_invalid_staging_storage(shape, dtype): + with pytest.raises((ValueError, TypeError)): + HostStaging(shape, dtype) def test_on_gpu(): @@ -157,3 +266,57 @@ def test_on_gpu(): copy = staging.copy(data) data *= 0.0 # overwritten right away on the compute stream np.testing.assert_array_equal(copy.result(), np.arange(1000.0)) + + +@pytest.mark.parametrize("dependency", ["stream", "event"]) +def test_nondefault_producer_on_gpu(dependency): + if not xp.cupy_available(): + pytest.skip("requires CUDA") + import cupy as cp + + source = cp.zeros(1000) + cp.cuda.get_current_stream().synchronize() + producer = cp.cuda.Stream(non_blocking=True) + produced = cp.cuda.Event(disable_timing=True) + with producer: + source.fill(7.0) + produced.record() + staging = HostStaging(source.shape, source.dtype, buffers=1) + host_copy = staging.copy(np.full(1000, 2.0)) + copy = staging.copy( + source, **{dependency: producer if dependency == "stream" else produced} + ) + np.testing.assert_array_equal(host_copy.result(), np.full(1000, 2.0)) + np.testing.assert_array_equal(copy.result(), np.full(1000, 7.0)) + # Copy completion makes the source safe to overwrite on another stream. + source.fill(9.0) + cp.cuda.get_current_stream().synchronize() + again = staging.copy(source) + with pytest.raises(RuntimeError, match="reused"): + copy.result() + np.testing.assert_array_equal(again.result(), np.full(1000, 9.0)) + + +def test_staging_storage_description_cannot_change(): + staging = HostStaging(2, np.float64) + with pytest.raises(AttributeError): + staging.shape = (4,) + with pytest.raises(AttributeError): + staging.dtype = np.dtype(np.float32) + + +def test_staged_result_restores_callers_device_on_gpu(): + if not xp.cupy_available(): + pytest.skip("requires CUDA") + import cupy as cp + + if cp.cuda.runtime.getDeviceCount() < 2: + pytest.skip("requires two devices") + with cp.cuda.Device(0): + staging = HostStaging(3, np.float64) + copy = staging.copy(cp.arange(3.0)) + with cp.cuda.Device(1): + np.testing.assert_array_equal(copy.result(), np.arange(3.0)) + assert cp.cuda.runtime.getDevice() == 1 + with pytest.raises(ValueError, match="bound to another"): + staging.copy(cp.zeros(3)) diff --git a/tests/unit/test_transfers.py b/tests/unit/test_transfers.py index a431c69..773cc11 100644 --- a/tests/unit/test_transfers.py +++ b/tests/unit/test_transfers.py @@ -84,7 +84,7 @@ def test_empty_counter(): assert counter.events == [] and counter.kernel_conversion_calls == [] assert counter.report().startswith("0 transfer(s) through cunumpy") assert repr(counter) == ( - "TransferCounter(to_host=0, to_device=0, kernel_conversion=0, fallback=0)" + "TransferCounter(to_host=0, to_device=0, kernel_conversion=0, fallback=0, device_copy=0)" ) @@ -189,7 +189,7 @@ def test_report_groups_events_by_kind_and_call_site(fake_device): lines = report.splitlines() assert lines[0] == ( "4 transfer(s) through cunumpy " - "(3 to_host, 1 to_device, 0 kernel_conversion, 0 fallback)" + "(3 to_host, 1 to_device, 0 kernel_conversion, 0 fallback, 0 device_copy)" ) assert " to_host (3):" in lines assert " to_device (1):" in lines @@ -278,7 +278,9 @@ def scale(x, y, factor): assert np.array_equal(x.data, np.full(3, 18.0)) # 1 * 2 * 3 * 3 (aliased) assert np.array_equal(y.data, np.full(3, 2.0)) assert counter.kernel_conversions == 2 - assert counter.total == 2, "the copies inside the call are not counted twice" + assert counter.total == 8 # six copies and two conversion markers + assert counter.to_host == counter.to_device == 3 + assert counter.bytes_to_host == counter.bytes_to_device == 3 * 3 * 8 first, second = counter.kernel_conversion_calls assert first.kind == "kernel_conversion" assert first.description == ( @@ -377,8 +379,9 @@ def scale(x, factor): PyccelKernel(scale)(x, 2.0) assert cp.all(x == 2.0) - assert counter.kernel_conversions == 1 and counter.total == 1 - assert counter.events[0].description == ( + assert counter.kernel_conversions == 1 and counter.total == 3 + assert counter.bytes_to_host == counter.bytes_to_device == x.nbytes + assert counter.kernel_conversion_calls[0].description == ( "PyccelKernel 'scale': 1 device array(s) copied to the host" ) @@ -403,7 +406,7 @@ def scale(x, factor, n): assert cp.all(x == 2.0) assert counter.fallbacks == 1 assert counter.kernel_conversions == 1, "the fallback converts on the host" - assert counter.total == 2 + assert counter.total == 4 # two copies, conversion and fallback markers @requires_cupy