// SPDX-FileCopyrightText: 2025 Filip Leonarski, Paul Scherrer Institute // SPDX-License-Identifier: GPL-3.0-only #pragma once #include #include #include #include #include #include "../common/JFJochException.h" class CudaStream { cudaStream_t stream_ = nullptr; public: // Non-blocking by default: a stream created with cudaStreamDefault synchronises against the legacy // NULL stream, so any NULL-stream operation anywhere in the process serialises every worker's GPU // work against every other's. With one engine per worker thread that costs most of the parallelism. CudaStream(unsigned int flags = cudaStreamNonBlocking) { if (cudaStreamCreateWithFlags(&stream_, flags) != cudaSuccess) throw JFJochException(JFJochExceptionCategory::GPUCUDAError, "Failed to create CUDA stream"); } ~CudaStream() { if (stream_) cudaStreamDestroy(stream_); } // Move-only type CudaStream(CudaStream&& other) noexcept : stream_(other.stream_) { other.stream_ = nullptr; } CudaStream& operator=(CudaStream&& other) noexcept { if (this != &other) { if (stream_) cudaStreamDestroy(stream_); stream_ = other.stream_; other.stream_ = nullptr; } return *this; } CudaStream(const CudaStream&) = delete; CudaStream& operator=(const CudaStream&) = delete; operator cudaStream_t() const { return stream_; } cudaStream_t get() const { return stream_; } }; // A timing event, so a phase that is queued on a stream can still report how long the device spent // on it. cudaEventDisableTiming is deliberately NOT used - timing is the whole point here. class CudaEvent { cudaEvent_t event_ = nullptr; public: CudaEvent() { if (cudaEventCreate(&event_) != cudaSuccess) throw JFJochException(JFJochExceptionCategory::GPUCUDAError, "Failed to create CUDA event"); } ~CudaEvent() { if (event_) cudaEventDestroy(event_); } CudaEvent(CudaEvent&& other) noexcept : event_(other.event_) { other.event_ = nullptr; } CudaEvent& operator=(CudaEvent&& other) noexcept { if (this != &other) { if (event_) cudaEventDestroy(event_); event_ = other.event_; other.event_ = nullptr; } return *this; } CudaEvent(const CudaEvent&) = delete; CudaEvent& operator=(const CudaEvent&) = delete; operator cudaEvent_t() const { return event_; } cudaEvent_t get() const { return event_; } }; class CudaFFTPlan { cufftHandle plan_ = 0; public: CudaFFTPlan() = default; CudaFFTPlan( int rank, const int* n, const int* inembed, int istride, int idist, const int* onembed, int ostride, int odist, cufftType type, int batch) { if (cufftPlanMany(&plan_, rank, const_cast(n), const_cast(inembed), istride, idist, const_cast(onembed), ostride, odist, type, batch) != CUFFT_SUCCESS) throw JFJochException(JFJochExceptionCategory::GPUCUDAError, "Failed to create CUFFT plan with cufftPlanMany"); } // Convenience overload for vector input CudaFFTPlan( int rank, const std::vector& n, const std::vector& inembed, int istride, int idist, const std::vector& onembed, int ostride, int odist, cufftType type, int batch) : CudaFFTPlan(rank, n.data(), inembed.data(), istride, idist, onembed.data(), ostride, odist, type, batch) {} ~CudaFFTPlan() { if (plan_) cufftDestroy(plan_); } // Move-only type CudaFFTPlan(CudaFFTPlan&& other) noexcept : plan_(other.plan_) { other.plan_ = 0; } CudaFFTPlan& operator=(CudaFFTPlan&& other) noexcept { if (this != &other) { if (plan_) cufftDestroy(plan_); plan_ = other.plan_; other.plan_ = 0; } return *this; } CudaFFTPlan(const CudaFFTPlan&) = delete; CudaFFTPlan& operator=(const CudaFFTPlan&) = delete; operator cufftHandle() const { return plan_; } cufftHandle get() const { return plan_; } }; // How much freed memory the pool keeps rather than returning it to the driver. Returning it puts the // next allocation straight back on the path this exists to avoid, so hold the per-worker engine // buffers; but the card also has to fit the merge, which asks for several gigabytes of its own after // the image loop, so do not hold everything. constexpr uint64_t CUDA_MEM_POOL_RELEASE_THRESHOLD = 1ull << 30; // The stream device allocations are ordered on: one per (thread, device), created on first use and // never destroyed. // // cudaMalloc and cudaFree are on CUDA's implicit-synchronization list - each one synchronises the // device across every stream - so one analysis engine per worker thread, each making a few dozen of // them, serialises every worker against every other AND stalls the workers already processing // images. cudaMallocAsync/cudaFreeAsync take the stream-ordered path instead and do not. // // The stream is deliberately not owned by the engine: cudaFreeAsync must be ordered on a stream that // is still alive, and an engine's own stream can be destroyed before the buffers it allocated. It is // keyed by device because a thread pinned to one GPU must not order work on another's stream. inline cudaStream_t cuda_allocation_stream() { int device = 0; if (cudaGetDevice(&device) != cudaSuccess) return nullptr; thread_local std::map streams; const auto it = streams.find(device); if (it != streams.end()) return it->second; cudaStream_t stream = nullptr; if (cudaStreamCreateWithFlags(&stream, cudaStreamNonBlocking) != cudaSuccess) return nullptr; cudaMemPool_t pool = nullptr; if (cudaDeviceGetDefaultMemPool(&pool, device) == cudaSuccess) { uint64_t threshold = CUDA_MEM_POOL_RELEASE_THRESHOLD; cudaMemPoolSetAttribute(pool, cudaMemPoolAttrReleaseThreshold, &threshold); } streams.emplace(device, stream); return stream; } // Which allocator a device buffer uses. Pooled is what every per-worker buffer wants. Synchronous is // for a buffer that is freed by a thread other than the one whose work read it: cudaFreeAsync orders // the free on the freeing side's allocation stream, which says nothing about kernels queued // elsewhere, while cudaFree synchronises the whole device and therefore cannot be early. enum class CudaAlloc { Pooled, Synchronous }; template class CudaDevicePtr { T* ptr_ = nullptr; // The stream the allocation was ordered on, and the one the free must be ordered on. Null when // the pool was not used, so the two always pair up. cudaStream_t stream_ = nullptr; void release() { if (!ptr_) return; if (stream_) cudaFreeAsync(ptr_, stream_); else cudaFree(ptr_); ptr_ = nullptr; } public: CudaDevicePtr() = default; explicit CudaDevicePtr(size_t count, CudaAlloc alloc = CudaAlloc::Pooled) { if (alloc == CudaAlloc::Pooled) stream_ = cuda_allocation_stream(); if (stream_ && cudaMallocAsync(&ptr_, count * sizeof(T), stream_) == cudaSuccess) { // The allocation is ordered on this stream, so anything that wants to use the memory from // another stream has to be ordered after it. Waiting for it once here is that ordering, // and it leaves the pointer as freely usable as cudaMalloc's would have been. cudaStreamSynchronize(stream_); return; } stream_ = nullptr; if (cudaMalloc(&ptr_, count * sizeof(T)) != cudaSuccess) throw JFJochException(JFJochExceptionCategory::GPUCUDAError, "Failed to allocate device memory"); } ~CudaDevicePtr() { release(); } // Move-only type CudaDevicePtr(CudaDevicePtr&& other) noexcept : ptr_(other.ptr_), stream_(other.stream_) { other.ptr_ = nullptr; other.stream_ = nullptr; } CudaDevicePtr& operator=(CudaDevicePtr&& other) noexcept { if (this != &other) { release(); ptr_ = other.ptr_; stream_ = other.stream_; other.ptr_ = nullptr; other.stream_ = nullptr; } return *this; } CudaDevicePtr(const CudaDevicePtr&) = delete; CudaDevicePtr& operator=(const CudaDevicePtr&) = delete; T* get() const { return ptr_; } operator T*() const { return ptr_; } }; template class CudaHostPtr { T* ptr_ = nullptr; public: CudaHostPtr() = default; explicit CudaHostPtr(size_t count) { if (cudaMallocHost(&ptr_, count * sizeof(T)) != cudaSuccess) throw JFJochException(JFJochExceptionCategory::GPUCUDAError, "Failed to allocate pinned host memory"); } ~CudaHostPtr() { if (ptr_) cudaFreeHost(ptr_); } // Move-only type CudaHostPtr(CudaHostPtr&& other) noexcept : ptr_(other.ptr_) { other.ptr_ = nullptr; } CudaHostPtr& operator=(CudaHostPtr&& other) noexcept { if (this != &other) { if (ptr_) cudaFreeHost(ptr_); ptr_ = other.ptr_; other.ptr_ = nullptr; } return *this; } CudaHostPtr(const CudaHostPtr&) = delete; CudaHostPtr& operator=(const CudaHostPtr&) = delete; T* get() const { return ptr_; } operator T*() const { return ptr_; } }; template class CudaRegisteredVector { std::vector* vec_ = nullptr; bool registered_ = false; static void registerPtr(void* ptr, size_t bytes, unsigned int flags) { cudaError_t err = cudaHostRegister(ptr, bytes, flags); if (err != cudaSuccess) throw JFJochException(JFJochExceptionCategory::GPUCUDAError, "cudaHostRegister failed"); } static void unregisterPtr(void* ptr) { cudaError_t err = cudaHostUnregister(ptr); if (err != cudaSuccess) throw JFJochException(JFJochExceptionCategory::GPUCUDAError, "cudaHostUnregister failed"); } // Unchecked, for the destructor and the move-assignment: both are noexcept, so throwing out of // them aborts the process instead of reporting the failure - and teardown is exactly where // cudaHostUnregister fails (after a device reset, or while another exception is unwinding). // The other destructors in this header ignore their teardown status for the same reason. static void unregisterPtrNoThrow(void* ptr) noexcept { cudaHostUnregister(ptr); } public: // Non-owning wrapper. Does NOT provide accessors to the vector. CudaRegisteredVector() = default; CudaRegisteredVector(std::vector& vec, unsigned int flags = cudaHostRegisterDefault) : vec_(&vec) { if (!vec.empty()) { registerPtr(vec.data(), vec.size() * sizeof(T), flags); registered_ = true; } } ~CudaRegisteredVector() { if (registered_ && vec_ && !vec_->empty()) { unregisterPtrNoThrow(vec_->data()); } } // Move-only CudaRegisteredVector(CudaRegisteredVector&& other) noexcept : vec_(other.vec_), registered_(other.registered_) { other.vec_ = nullptr; other.registered_ = false; } CudaRegisteredVector& operator=(CudaRegisteredVector&& other) noexcept { if (this != &other) { // Clean current registration if (registered_ && vec_ && !vec_->empty()) { unregisterPtrNoThrow(vec_->data()); } vec_ = other.vec_; registered_ = other.registered_; other.vec_ = nullptr; other.registered_ = false; } return *this; } CudaRegisteredVector(const CudaRegisteredVector&) = delete; CudaRegisteredVector& operator=(const CudaRegisteredVector&) = delete; // Re-register after vector capacity/size change. Caller must ensure // the vector is not registered at the moment of mutation. void rebind(std::vector& vec, unsigned int flags = cudaHostRegisterDefault) { // Unregister previous if needed if (registered_ && vec_ && !vec_->empty()) { unregisterPtr(vec_->data()); } vec_ = &vec; if (!vec.empty()) { registerPtr(vec.data(), vec.size() * sizeof(T), flags); registered_ = true; } else { registered_ = false; } } // Explicit unregister (optional). void unregister() { if (registered_ && vec_ && !vec_->empty()) { unregisterPtr(vec_->data()); registered_ = false; } } bool isRegistered() const { return registered_; } };