Allocate off the driver lock, and give the loop the sixteen workers a card wants
Build Packages / build:windows:nocuda (push) Successful in 16m31s
Build Packages / build:windows:cuda (push) Successful in 19m15s
Build Packages / build:viewer-tgz:cpu (push) Successful in 19m58s
Build Packages / build:viewer-tgz:cuda (push) Successful in 22m51s
Build Packages / build:rpm (ubuntu2404_nocuda) (push) Successful in 24m5s
Build Packages / build:rpm (rocky9_nocuda) (push) Successful in 25m3s
Build Packages / build:rpm (ubuntu2204_nocuda) (push) Successful in 28m51s
Build Packages / build:rpm (rocky8_nocuda) (push) Successful in 28m55s
Build Packages / build:rpm (rocky8_sls9) (push) Successful in 28m35s
Build Packages / build:rpm (rocky9_sls9) (push) Successful in 21m19s
Build Packages / XDS test (durin plugin) (push) Successful in 11m17s
Build Packages / build:rpm (rocky9) (push) Successful in 23m30s
Build Packages / Generate python client (push) Successful in 33s
Build Packages / build:rpm (rocky8) (push) Successful in 25m35s
Build Packages / Create release (push) Skipped
Build Packages / Build documentation (push) Successful in 1m37s
Build Packages / build:rpm (ubuntu2404) (push) Successful in 21m40s
Build Packages / XDS test (neggia plugin) (push) Successful in 10m17s
Build Packages / build:rpm (ubuntu2204) (push) Successful in 26m44s
Build Packages / XDS test (JFJoch plugin) (push) Successful in 11m10s
Build Packages / DIALS test (push) Successful in 23m27s
Build Packages / Unit tests (push) Successful in 1h19m59s

Two changes to how the per-image loop is set up, neither of which touches what it
computes.

Device buffers were taken with cudaMalloc and returned with cudaFree, both of which
are on CUDA's implicit-synchronisation list: each one synchronises the device across
every stream. One analysis engine per worker, each making a few dozen of them, means
the workers still constructing stall the workers already processing images, and the
cost grows with the worker count. They are now stream-ordered allocations from the
device's memory pool, with the synchronous pair kept as the fallback where no pool is
available.

Two deliberate limits on that. The pool's release threshold is one gibibyte rather
than unbounded: holding the small per-worker buffers is the whole point, but the card
also has to fit the merge afterwards, which asks for several gigabytes of its own.
And the shared geometry tables keep the synchronous allocator, because their deleter
runs on whichever thread drops the last reference, so an asynchronous free there
would be ordered on a stream that says nothing about the engine streams whose kernels
read the table; they are allocated once per card, so the pool bought them nothing.

The loop's worker cap per card goes from eight to sixteen. The comment beside it
already recorded where the measurement put the minimum - the loop's time falls to
sixteen workers and then rises - and a later measurement on a single card agrees: at
eight the loop waits on the queue rather than on the card. An explicit -N is still
obeyed as given.

Reflection files are byte-identical on four crystals with the worker count doubled,
which is the property the frame-ordered mosaicity smoothing and the deterministic
prediction order were built to give. Thirty consecutive runs of one crystal on the
pooled allocator: no failure, every file identical to the first. Peak device memory
over the whole rotation test set is 4.9 of 16 gibibytes.

Four minutes seventeen to four minutes one over thirty-eight crystals, each binary
repeating itself to within half a per cent; nineteen crystals faster, nineteen level,
none slower, and every column of the comparison table identical.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01EGpGdgmJ8MyY9pCGWjktyi
This commit is contained in:
2026-08-25 02:29:07 +02:00
co-authored by Claude Opus 5
parent 1d16dcd2b9
commit c470bed93a
4 changed files with 89 additions and 13 deletions
+74 -4
View File
@@ -5,6 +5,7 @@
#include <cuda_runtime.h>
#include <cufft.h>
#include <map>
#include <stdexcept>
#include <vector>
#include "../common/JFJochException.h"
@@ -119,27 +120,96 @@ public:
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<int, cudaStream_t> 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 <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) {
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() {
if (ptr_) cudaFree(ptr_);
release();
}
// Move-only type
CudaDevicePtr(CudaDevicePtr&& other) noexcept : ptr_(other.ptr_) { other.ptr_ = nullptr; }
CudaDevicePtr(CudaDevicePtr&& other) noexcept : ptr_(other.ptr_), stream_(other.stream_) {
other.ptr_ = nullptr;
other.stream_ = nullptr;
}
CudaDevicePtr& operator=(CudaDevicePtr&& other) noexcept {
if (this != &other) {
if (ptr_) cudaFree(ptr_);
release();
ptr_ = other.ptr_;
stream_ = other.stream_;
other.ptr_ = nullptr;
other.stream_ = nullptr;
}
return *this;
}
+6 -2
View File
@@ -88,8 +88,12 @@ std::shared_ptr<CudaDevicePtr<T>> SharedDeviceTable(const void *key, size_t coun
it = it->second.expired() ? reg.tables.erase(it) : std::next(it);
// Free on the device that allocated it - the last engine to drop the table may well be a worker
// pinned to a different GPU.
std::shared_ptr<CudaDevicePtr<T>> table(new CudaDevicePtr<T>(count), [device](CudaDevicePtr<T> *p) {
// pinned to a different GPU. And synchronously, for the same reason: the thread that drops the
// last reference is not the one whose kernels read the table, so a stream-ordered free would be
// ordered against the wrong work. The table is allocated once per GPU, so it has nothing to gain
// from the pool anyway.
std::shared_ptr<CudaDevicePtr<T>> table(new CudaDevicePtr<T>(count, CudaAlloc::Synchronous),
[device](CudaDevicePtr<T> *p) {
int current = 0;
cudaGetDevice(&current);
cudaSetDevice(device);
+8 -6
View File
@@ -1863,18 +1863,20 @@ ProcessResult Rugnux::RunPipeline(RugnuxObserver *observer, bool write_output, b
// How many workers the loop actually wants. Every one of them submits its own kernels to a card,
// and a card runs out of room to accept them long before it runs out of work: measured on two
// GPUs, the loop's own time falls to sixteen workers and then rises again, so forty-eight is
// slower than eight on a 16 Mpx set. The cap is per card because that is what the queue depth
// belongs to, and it only applies when -N was left at its default - an explicit -N is a
// deliberate instruction and is obeyed, which is what a previous attempt at this got wrong.
// slower than sixteen on a 16 Mpx set. Sixteen is where that measurement put the minimum, and a
// later one on a single card agrees - eight leaves the loop waiting on the queue rather than the
// card. The cap is per card because that is what the queue depth belongs to, and it only applies
// when -N was left at its default - an explicit -N is a deliberate instruction and is obeyed,
// which is what a previous attempt at this got wrong.
// Only the full-analysis worker drives a card; the azimuthal one preprocesses and integrates on
// the CPU and wants every thread it can have, so the cap must not reach it.
const int gpus = get_gpu_count();
const int image_workers = (per_image_analysis && config_.nthreads_auto && gpus > 0)
? std::min(config_.nthreads, std::max(8, 8 * gpus))
? std::min(config_.nthreads, std::max(16, 16 * gpus))
: config_.nthreads;
if (image_workers < config_.nthreads)
logger.Info("Per-image loop: {} of {} threads ({} GPUs) - more workers than a card can take "
"queue work from make it slower, not faster; pass -N to override",
logger.Info("Per-image loop: {} of {} threads ({} GPUs) - past what a card can take queue "
"work from, more workers make it slower, not faster; pass -N to override",
image_workers, config_.nthreads, gpus);
std::vector<std::future<void> > futures;
futures.reserve(image_workers);
+1 -1
View File
@@ -78,7 +78,7 @@ void print_usage() {
std::cout << " -o, --output-prefix <txt> Output file prefix (default: output). mx and scale runs always write "
"<prefix>_report.txt, the results report; an empty prefix writes nothing at all" << std::endl;
std::cout << " -N, --threads <num> Number of threads (default: all hardware threads, of which the "
"per-image loop uses at most 8 per GPU; an explicit value is used as given)" << std::endl;
"per-image loop uses at most 16 per GPU; an explicit value is used as given)" << std::endl;
std::cout << " -s, --start-image <num> Start image number (default: 0)" << std::endl;
std::cout << " -e, --end-image <num> End image number (default: all)" << std::endl;
std::cout << " -t, --stride <num> Image stride (default: 1)" << std::endl;