rugnux: GPU merge waits for the probe pass running beside it instead of failing

The speculative geometry probe (StartSpeculativeGeometryProbe) runs an
indexing-only pass on a copy of the run beside the pass's scaling merge.
It holds ~4.5 GB of engines (25 image-sized spot-finding buffers, the
per-engine tables, FFT indexers) plus what its streams' pool retains, for a
few seconds. When a merge allocation landed in that window and did not fit,
RotationScaleMergeGPU::Impl::Alloc reported "needs more GPU memory than this
card has ... too large for GPU scaling", although the per-observation arrays
were 1.3 GB on a 16.6 GB card and the set runs fine alone. Timing-dependent:
one full-battery failure, not reproduced in ~15 plain reruns; reproduced
deterministically by letting the probe hold extra device memory.

- GPUWorkBeside (common/CUDAWrapper): a process-wide count of GPU work
  running beside the main line, with a condition variable signalled when the
  last one ends. The speculative probe holds one for its whole pass.
- Alloc: on failure, wait for that work to end (bounded, 10 min, then a
  "GPU busy" error), then ask for the same buffer once more. Nothing else in
  flight means no wait, so a genuine shortage still fails at once. The
  computation is the same whenever the allocation succeeds, so results do
  not depend on the wait.
- "Too large for this card" is now said only by the early check of the
  per-observation arrays against total device memory; a later shortage with
  nothing running beside gets its own message (buffer size, free/total, and
  that another program or the set's size is the cause).

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01D1G8gJVAy6gp1K5Dz3NE5C
This commit is contained in:
2026-09-29 07:44:01 +02:00
co-authored by Claude Opus 5.5
parent 528af5e155
commit c644c4e38a
5 changed files with 112 additions and 2 deletions
+27
View File
@@ -3,6 +3,9 @@
#include "CUDAWrapper.h"
#include <condition_variable>
#include <mutex>
// Build-independent: the CUDA build gets get_gpu_names() from CUDAWrapper.cu, the CPU-only build from
// the stub below, and this collapses whichever list came back. Four identical cards read better as
// "4x <name>" than as the same name four times, and a mixed machine keeps one group per model.
@@ -24,6 +27,30 @@ std::string get_gpu_description() {
return out;
}
namespace {
std::mutex gpu_work_beside_mutex;
std::condition_variable gpu_work_beside_done;
int gpu_work_beside = 0;
}
GPUWorkBeside::GPUWorkBeside() {
std::lock_guard lock(gpu_work_beside_mutex);
++gpu_work_beside;
}
GPUWorkBeside::~GPUWorkBeside() {
{
std::lock_guard lock(gpu_work_beside_mutex);
--gpu_work_beside;
}
gpu_work_beside_done.notify_all();
}
bool wait_for_gpu_work_beside(std::chrono::seconds timeout) {
std::unique_lock lock(gpu_work_beside_mutex);
return gpu_work_beside_done.wait_for(lock, timeout, [] { return gpu_work_beside == 0; });
}
#ifndef JFJOCH_USE_CUDA
int32_t get_gpu_count() {
+15
View File
@@ -3,6 +3,7 @@
#pragma once
#include <chrono>
#include <cstdint>
#include <string>
#include <vector>
@@ -50,3 +51,17 @@ void cuda_clear_error();
// fallback from that, and falling back only moves the failure somewhere less clear. A handled,
// non-sticky failure passes. No-op without CUDA.
void cuda_throw_if_context_lost();
// GPU work that runs beside the main line of a process and gives all of its device memory back when it
// ends: a probe pass run on a copy of the run while the run itself scales and merges. Hold one for as
// long as that work runs. Build-independent - it only counts.
class GPUWorkBeside {
public:
GPUWorkBeside();
~GPUWorkBeside();
GPUWorkBeside(const GPUWorkBeside &) = delete;
GPUWorkBeside &operator=(const GPUWorkBeside &) = delete;
};
// Wait until no GPUWorkBeside is held, for at most `timeout`. True once none is - at once if none was.
bool wait_for_gpu_work_beside(std::chrono::seconds timeout);
@@ -4,6 +4,7 @@
#include "RotationScaleMergeGPU.h"
#include <algorithm>
#include <chrono>
#include <cmath>
#include <cstdio>
#include <string>
@@ -666,6 +667,39 @@ namespace {
double(total_bytes) / 1e9);
throw JFJochException(JFJochExceptionCategory::GPUCUDAError, msg);
}
// A buffer that did not fit although the per-observation arrays fit the card, with nothing else in
// this process running on the GPU any more: what holds the rest is this run or another program.
[[noreturn]] void ThrowOutOfGPUMemory(int64_t n_obs, size_t bytes) {
size_t free_bytes = 0, total_bytes = 0;
cudaMemGetInfo(&free_bytes, &total_bytes);
cudaGetLastError();
char msg[768];
snprintf(msg, sizeof(msg),
"Scaling %lld partial observations ran out of GPU memory: a buffer of %.1f MB did not fit, "
"with %.1f GB of the %.1f GB card free (the per-observation arrays take %.1f GB). Either "
"another program is using this GPU, or this data set is too large for GPU scaling on this "
"card - run it on the CPU (CUDA_VISIBLE_DEVICES= rugnux ..., or a CPU-only build) on a "
"machine with plenty of host memory",
static_cast<long long>(n_obs), double(bytes) / 1e6, double(free_bytes) / 1e9,
double(total_bytes) / 1e9, double(n_obs) * OBS_BYTES / 1e9);
throw JFJochException(JFJochExceptionCategory::GPUCUDAError, msg);
}
// How long a merge waits for GPU work running beside it (GPUWorkBeside) to give its memory back.
// That work is a probe pass of seconds, so this long means it is stuck, not slow.
constexpr auto GPU_BUSY_TIMEOUT = std::chrono::minutes(10);
[[noreturn]] void ThrowGPUBusy(int64_t n_obs) {
char msg[512];
snprintf(msg, sizeof(msg),
"GPU busy: scaling %lld partial observations waited %lld minutes for GPU work running "
"beside it in this process to give back its memory, and it did not. The data set fits "
"this card; run it again, or on the CPU (CUDA_VISIBLE_DEVICES= rugnux ...)",
static_cast<long long>(n_obs),
static_cast<long long>(std::chrono::duration_cast<std::chrono::minutes>(GPU_BUSY_TIMEOUT).count()));
throw JFJochException(JFJochExceptionCategory::GPUCUDAError, msg);
}
}
struct RotationScaleMergeGPU::Impl {
@@ -673,12 +707,24 @@ struct RotationScaleMergeGPU::Impl {
bool available = false;
int n_obs = 0, n_frames = 0, n_groups = 0;
// The arrays fit the card (SetPartialsLayout checks that first), but not necessarily beside other
// GPU work of this process: a probe pass started on a copy of the run holds several GB of engines
// for a few seconds while this merge runs. So a buffer that does not fit waits for that work to
// end and is asked for again - the same buffer, so what the merge computes does not depend on the
// wait. It is asked for again also when nothing is running beside it any more: work that ended
// between the failure and the check has just given its memory back.
template <typename T>
CudaDevicePtr<T> Alloc(size_t n) const {
try {
return CudaDevicePtr<T>(n, ALLOC);
} catch (const JFJochException &) {
ThrowTooLargeForGPU(n_obs);
}
if (!wait_for_gpu_work_beside(GPU_BUSY_TIMEOUT))
ThrowGPUBusy(n_obs);
try {
return CudaDevicePtr<T>(n, ALLOC);
} catch (const JFJochException &) {
ThrowOutOfGPUMemory(n_obs, n * sizeof(T));
}
}
+3 -1
View File
@@ -3253,8 +3253,10 @@ void Rugnux::StartSpeculativeGeometryProbe(const std::array<float, 5> &geometry)
auto sp = std::make_shared<SpeculativeProbe>();
sp->run = run;
// A probe that fails leaves no memo, and the walk's own probe then runs it again and meets the
// failure itself.
// failure itself. It runs beside the merge of the pass that started it, on the same card: a merge
// that does not fit beside it waits for it to end (GPUWorkBeside).
sp->done = std::async(std::launch::async, [run] {
GPUWorkBeside beside;
try {
run->RunPipeline(nullptr, /*write_output=*/false, /*geometry_prepass=*/false);
} catch (const std::exception &) {
+20
View File
@@ -4,6 +4,9 @@
#include <catch2/catch_all.hpp>
#include "../common/CUDAWrapper.h"
#include <optional>
#include <thread>
#ifdef JFJOCH_USE_CUDA
#include <future>
@@ -101,3 +104,20 @@ TEST_CASE("CudaDevicePtr_HandledPoolFailureLeavesNoError", "[CUDAMemHelpers]") {
}
#endif
// A merge that finds the card full waits for the GPU work running beside it: until the last one
// ends, and not forever.
TEST_CASE("GPUWorkBeside_WaitEndsWithTheWork", "[CUDAMemHelpers]") {
CHECK(wait_for_gpu_work_beside(std::chrono::seconds(0)));
std::optional<GPUWorkBeside> work;
work.emplace();
CHECK_FALSE(wait_for_gpu_work_beside(std::chrono::seconds(0)));
std::thread ends([&work] {
std::this_thread::sleep_for(std::chrono::milliseconds(100));
work.reset();
});
CHECK(wait_for_gpu_work_beside(std::chrono::seconds(60)));
ends.join();
}