Files
Jungfraujoch/image_analysis/indexing/CUDAMemHelpers.h
T
leonarski_fandClaude Fable 5.1 060a89cd7d CUDA: allocation streams are handed back, not leaked; a handled pool failure leaves no error
A long-running broker began cancelling every data collection with

    Device decoding failed (CUDA (GPU) error (out of memory)), falling back to host decompression
    CUDA (GPU) error (out of memory)

while nvidia-smi showed the cards less than a fifth full. Two defects, both in CUDAMemHelpers.h and
both from the pooled allocator (cudaMallocAsync) that came with rc.162.

The leak. cuda_allocation_stream() kept one stream per (thread, device) in a thread_local map of raw
cudaStream_t and never destroyed them. That was written against rugnux, where the worker threads live
as long as the process. The broker starts fresh std::async threads for every data collection - 16 or
64 of them - so every collection left that many streams behind. Measured: 0.56 MB of device memory
per leaked stream, linear to 4928 streams, never returned, with nothing on the host side growing.

The RAII wrapper (CudaStream) was there but not used at this site, and using it as-is - a stream
destroyed when its thread exits - would not have been safe: ShadowFinder builds its GPU accumulator
on a throw-away std::async thread and frees it from another thread long after, and that free is
ordered on the allocating thread's stream. So the streams are still never destroyed, but a thread
now only borrows one: CudaStream objects live in a process-wide per-device idle list, a thread takes
one on first use and hands it back when it exits. Their number is bounded by the threads that were
ever alive at once instead of by the threads ever started. The list itself is deliberately leaked, so
that nothing calls into CUDA during static destruction.

Replaying the broker's pattern against the real header, 60 collections of 64 threads:
    before   3840 streams, 260 -> 2424 MB of device memory
    after      64 streams, 260 ->  358 MB

The stale error. Every helper here throws a named message, yet the log carried the raw CUDA string,
so the failure came through a cuda_err() and not from an allocation. CudaDevicePtr falls back to
cudaMalloc when cudaMallocAsync fails, silently - but the failed call stays behind as the thread's
last error, and the cudaGetLastError() that follows the next kernel launch reports it. The buffers
were all allocated; the frame was lost anyway, once on the device-decode route (caught, hence the
warning) and once more on the host fallback (fatal). The pooled attempt failing, and a stream that
cannot be created, are both handled by falling back, so both now clear the error they leave.

What finite resource the production cards ran out of at under 4 GB used was not established - no
cap on the number of streams was found up to 4928 on the card this was measured on. The leak is the
only thing on this path that grows with uptime.

tests/CUDAMemHelpersTest.cpp: later threads end up on the same stream, concurrent threads on
different ones, a buffer is freed cleanly after its allocating thread has exited (and another has
borrowed its stream), and a pool that cannot serve a request leaves no error behind - the last by
capping a memory pool at 4 MB so that the pooled attempt fails and the fallback succeeds.

Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
2026-09-20 18:45:03 +02:00

401 lines
15 KiB
C++

// SPDX-FileCopyrightText: 2025 Filip Leonarski, Paul Scherrer Institute <filip.leonarski@psi.ch>
// SPDX-License-Identifier: GPL-3.0-only
#pragma once
#include <cuda_runtime.h>
#include <cufft.h>
#include <map>
#include <mutex>
#include <stdexcept>
#include <vector>
#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<int*>(n),
const_cast<int*>(inembed), istride, idist,
const_cast<int*>(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<int>& n,
const std::vector<int>& inembed, int istride, int idist,
const std::vector<int>& 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), borrowed on first use and
// handed back when the thread exits. The streams themselves are 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.
//
// Nor is it owned by the thread, for the same reason: a buffer can outlive the thread that allocated
// it - the shadow finder builds its accumulator on a throw-away thread and frees it from another,
// long after - so the thread only borrows the stream. But it does hand it back. The broker starts
// fresh worker threads for every data collection, and a stream holds about half a megabyte of device
// memory, so a stream per thread ever started is a leak that a long-running broker does not survive.
// Handed back, the streams number as many as the threads that were ever alive at once.
namespace jfjoch_cuda_allocation_streams {
struct Idle {
std::mutex m;
std::map<int, std::vector<CudaStream>> streams; // by device
};
// Never destroyed: a thread can still be handing a stream back while the process exits, and
// destroying streams during static destruction would call into a runtime that may be gone.
inline Idle &idle() {
static Idle *i = new Idle;
return *i;
}
struct Borrowed {
std::map<int, CudaStream> streams; // by device
~Borrowed() {
std::lock_guard lock(idle().m);
for (auto &[device, stream] : streams)
idle().streams[device].push_back(std::move(stream));
}
};
}
inline cudaStream_t cuda_allocation_stream() {
int device = 0;
if (cudaGetDevice(&device) != cudaSuccess)
return nullptr;
thread_local jfjoch_cuda_allocation_streams::Borrowed borrowed;
const auto it = borrowed.streams.find(device);
if (it != borrowed.streams.end())
return it->second;
auto &idle = jfjoch_cuda_allocation_streams::idle();
{
std::lock_guard lock(idle.m);
auto &streams = idle.streams[device];
if (!streams.empty()) {
CudaStream stream = std::move(streams.back());
streams.pop_back();
return borrowed.streams.emplace(device, std::move(stream)).first->second;
}
}
try {
CudaStream stream;
cudaMemPool_t pool = nullptr;
if (cudaDeviceGetDefaultMemPool(&pool, device) == cudaSuccess) {
uint64_t threshold = CUDA_MEM_POOL_RELEASE_THRESHOLD;
cudaMemPoolSetAttribute(pool, cudaMemPoolAttrReleaseThreshold, &threshold);
}
return borrowed.streams.emplace(device, std::move(stream)).first->second;
} catch (const JFJochException &) {
// No stream: the caller allocates synchronously instead, so this is handled - see
// CudaDevicePtr for why it must then not stay behind as the thread's last error.
cudaGetLastError();
return nullptr;
}
}
// 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 <typename T>
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_) {
if (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;
}
// The pool failing is handled - by the synchronous allocation below - so it must not stay
// behind as the thread's last error. Every kernel launch is followed by cudaGetLastError(),
// which would report it there: an "out of memory" thrown over an image whose buffers were
// all allocated and whose kernel launched.
cudaGetLastError();
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 <typename T>
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 <typename T>
class CudaRegisteredVector {
std::vector<T>* 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<T>& 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<T>& 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_; }
};