Files
Jungfraujoch/image_analysis/indexing/IndexerThreadPool.cpp
T
leonarski_fandClaude Opus 5 6319d5c600 indexing: tell an indexer that failed apart from one that found nothing
IndexerThreadPool collapsed two different outcomes into the same empty reply. A worker that threw
set result = nullptr, and so did a dispatch that never found a free worker; both then became a
default-constructed IndexerResult, indistinguishable downstream from an indexer that ran and found
no lattice. The first says nothing about the frame at all - the indexer never looked, and it will
fail again on the next one - while the second is a real negative result about the crystal.

The visible cost was the advice a failed run gave. With the card full, rugnux printed "Indexer
thread 0 failed: CUDA (GPU) error" and then ended with "Two-pass rotation indexing found no lattice.
Check the beam centre (--beam-x / --beam-y), raise --max-spots ..." - sending the operator to look
at geometry that was never wrong, for a machine that was simply out of memory. None of those
remedies can help when the frames were not examined.

IndexerResult gains an optional error, set by the worker and by the pool's own catch; the remaining
nullptr path keeps its meaning of "not attempted" and deliberately carries no error.
RotationIndexer records it and exposes GetIndexerError(), and rugnux's first pass branches on it at
the throw site, naming the resource failure instead.

The error travels as DATA through image_analysis/ rather than as an exception, because
RotationIndexer::RunIndexing() is on the online path - IndexAndRefine calls it on a schedule from
the broker and the receiver, where dropping a frame is the right failure and killing a live
acquisition is not. Only rugnux, which owns the "this run is over" decision, turns it into one. The
failure result is byte-identical to the default-constructed one it replaces, so any_executed and the
online frame-drop behaviour are unchanged.

The new IndexerError category exists because the category is only the display prefix on what() -
Category() is read nowhere - and the old line read "Processing failed: Input parameter invalid" for
a GPU fault, which is the same defect one layer up. SpotFinderError is the precedent.

msg.indexing_result is deliberately left alone: of its three states, "not attempted" is the honest
one for a frame the indexer never examined, and asserting false would be the same collapse again.

Verified on a squeezed card (246 MB free, -N 4): exit 1, zero dropped frames, the new message, and
no mention of --beam-x. On a free card, against the previous binary at -N 6 on the same input, the
logs differ only in paths and timing - same space group, cell and merge statistics. Targeted Catch2
cases pass, including three that drive the pool through the online path.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01Lc5JG6kJqZoCWaoZ43JGTW
2026-08-29 15:30:28 +02:00

307 lines
13 KiB
C++

// SPDX-FileCopyrightText: 2025 Filip Leonarski, Paul Scherrer Institute <filip.leonarski@psi.ch>
// SPDX-License-Identifier: GPL-3.0-only
#include "IndexerThreadPool.h"
#include "../common/CUDAWrapper.h"
#include "../common/Logger.h"
#ifdef JFJOCH_USE_CUDA
#include "FFBIDXIndexer.h"
#include "FFTIndexerGPU.h"
#endif
#ifdef JFJOCH_USE_FFTW
#include "FFTIndexerCPU.h"
#endif
void WarmUpCuFFT() {
#ifdef JFJOCH_USE_CUDA
if (get_gpu_count() == 0)
return;
cufftHandle plan = 0;
if (cufftPlan1d(&plan, 1024, CUFFT_C2C, 1) == CUFFT_SUCCESS)
cufftDestroy(plan);
#endif
}
// The indexer for one RESOLVED algorithm, or nullptr if this build/host cannot serve it.
static std::unique_ptr<Indexer> MakeIndexer(IndexingAlgorithmEnum algorithm, const IndexingSettings &settings) {
#ifdef JFJOCH_USE_CUDA
if (get_gpu_count() > 0) {
if (algorithm == IndexingAlgorithmEnum::FFT)
return std::make_unique<FFTIndexerGPU>(settings);
if (algorithm == IndexingAlgorithmEnum::FFBIDX)
return std::make_unique<FFBIDXIndexer>();
}
#endif
#ifdef JFJOCH_USE_FFTW
if (algorithm == IndexingAlgorithmEnum::FFTW)
return std::make_unique<FFTIndexerCPU>(settings);
#endif
return nullptr;
}
IndexerThread::IndexerThread(const IndexingSettings &settings, int threadid, IndexerConstruction construction)
: settings_(settings), construction_(construction) {
std::unique_lock<std::mutex> lock(m);
state = TaskState::STARTING;
worker_thread = std::thread(&IndexerThread::Worker, this, threadid);
c_running.wait(lock, [this] { return state != TaskState::STARTING; });
if (state == TaskState::ERROR) {
worker_thread.join();
throw JFJochException(JFJochExceptionCategory::InputParameterInvalid,
"Indexer thread initialization failed");
}
}
void IndexerThread::Worker(int threadid) {
try {
pin_gpu();
} catch (const std::exception &e) {
spdlog::error("Failed to pin to GPU {}", e.what());
} catch (...) {
// GPU pinning errors are not critical and should be ignored for the time being.
}
std::unique_ptr<Indexer> fft_indexer, ffbidx_indexer, fftw_indexer;
// Preconstruct: build every indexer the requested algorithm could resolve to before the pool
// reports ready, so no cuFFT planning happens once frames are flowing, and a failure is fatal
// for the pool instead of being met frame by frame. OnFirstUse skips this and builds in the
// dispatch below.
if (construction_ == IndexerConstruction::Preconstruct) {
try {
const auto requested = settings_.GetAlgorithm();
if (requested == IndexingAlgorithmEnum::Auto || requested == IndexingAlgorithmEnum::FFT)
fft_indexer = MakeIndexer(IndexingAlgorithmEnum::FFT, settings_);
if (requested == IndexingAlgorithmEnum::Auto || requested == IndexingAlgorithmEnum::FFBIDX)
ffbidx_indexer = MakeIndexer(IndexingAlgorithmEnum::FFBIDX, settings_);
if ((requested == IndexingAlgorithmEnum::Auto && get_gpu_count() == 0)
|| requested == IndexingAlgorithmEnum::FFTW)
fftw_indexer = MakeIndexer(IndexingAlgorithmEnum::FFTW, settings_);
} catch (const std::exception &e) {
spdlog::error("Failed to initialize indexer: {}", e.what());
{
std::unique_lock<std::mutex> lock(m);
state = TaskState::ERROR;
}
c_running.notify_all();
return;
} catch (...) {
spdlog::error("Failed to initialize indexer");
{
std::unique_lock<std::mutex> lock(m);
state = TaskState::ERROR;
}
c_running.notify_all();
return;
}
}
{
std::unique_lock<std::mutex> lock(m);
state = TaskState::IDLE;
}
c_running.notify_all();
while (true) {
std::unique_ptr<TaskInput> input;
// Look for task + handle stop
{
std::unique_lock<std::mutex> lock(m);
c_start.wait(lock, [this] { return stop || state == TaskState::READY; });
if (stop && (state != TaskState::READY))
return;
state = TaskState::RUNNING;
input = std::move(task_input);
}
if (input) {
std::unique_ptr<IndexerResult> tmp_result;
try {
auto algorithm = input->experiment.GetIndexingAlgorithm();
std::unique_ptr<Indexer> *slot = nullptr;
switch (algorithm) {
case IndexingAlgorithmEnum::FFT: slot = &fft_indexer; break;
case IndexingAlgorithmEnum::FFBIDX: slot = &ffbidx_indexer; break;
case IndexingAlgorithmEnum::FFTW: slot = &fftw_indexer; break;
default: break;
}
// A preconstructing worker already holds it; an OnFirstUse worker builds it here,
// on the first frame that resolves to this algorithm.
if (slot && !*slot)
*slot = MakeIndexer(algorithm, settings_);
if (!slot || !*slot) {
// Algorithm is already resolved here (never Auto/None - see
// IndexerThreadPool::Run, which also checked this host can serve it). Reaching
// this means the resolved algorithm has no matching indexer in this build -
// fail loudly instead of silently not indexing.
throw JFJochException(JFJochExceptionCategory::InputParameterInvalid,
"Internal error: no indexer available for the resolved "
"indexing algorithm");
}
Indexer &indexer = **slot;
indexer.Setup(input->experiment);
tmp_result = std::make_unique<IndexerResult>(indexer.Run(input->recip));
} catch (std::exception &e) {
// Hand the failure back as a result carrying the reason. A nullptr here was
// indistinguishable from a worker that was never dispatched, and both then read
// downstream as "this frame did not index".
spdlog::error("Indexer thread {} failed: {}", threadid, e.what());
tmp_result = std::make_unique<IndexerResult>(IndexerResult{
.lattice = {}, .indexing_time_s = 0, .executed = false, .error = e.what()});
}
{
std::unique_lock<std::mutex> lock(m);
state = TaskState::COMPLETED;
result = std::move(tmp_result);
}
c_done.notify_all();
}
}
}
void IndexerThread::Finalize() {
{
std::unique_lock<std::mutex> lock(m);
stop = true;
}
c_start.notify_all();
if (worker_thread.joinable())
worker_thread.join();
}
std::unique_ptr<IndexerResult> IndexerThread::Run(const DiffractionExperiment &experiment,
const std::vector<Coord> &recip) {
std::unique_ptr<IndexerResult> tmp_result;
{
std::unique_lock<std::mutex> lock(m);
if (stop)
return nullptr;
if (state != TaskState::IDLE)
return nullptr;
task_input = std::make_unique<TaskInput>(std::cref(experiment), std::cref(recip));
state = TaskState::READY;
}
c_start.notify_one();
{
std::unique_lock<std::mutex> lock(m);
c_done.wait(lock, [this] { return state == TaskState::COMPLETED; });
tmp_result = std::move(result);
state = TaskState::IDLE;
}
return tmp_result;
}
IndexerThread::~IndexerThread() {
Finalize();
}
IndexerThreadPool::IndexerThreadPool(const IndexingSettings &settings, IndexerConstruction construction)
: worker_busy(settings.GetIndexingThreads(), 0),
worker_free_count(settings.GetIndexingThreads()),
viable_cell_min_spots(settings.GetViableCellMinSpots()),
blocking(settings.GetBlockingBehavior()) {
for (size_t i = 0; i < settings.GetIndexingThreads(); ++i)
tasks.emplace_back(std::make_unique<IndexerThread>(std::cref(settings), i, construction));
}
int IndexerThreadPool::GetFreeWorker() {
std::unique_lock<std::mutex> lock(m);
if (tasks.size() == 0)
return -1;
if (blocking)
c.wait(lock, [this] { return worker_free_count > 0; });
for (int i = 0; i < tasks.size(); i++) {
if (worker_busy[i] == 0) {
worker_busy[i] = 1;
worker_free_count--;
return i;
}
}
return -1;
}
IndexerResult IndexerThreadPool::Run(const DiffractionExperiment &experiment, const std::vector<Coord> &recip) {
const auto algorithm = experiment.GetIndexingAlgorithm();
if (algorithm == IndexingAlgorithmEnum::None)
return IndexerResult{.lattice = {}, .indexing_time_s = 0, .executed = false};
// GetIndexingAlgorithm() must already have resolved Auto to a concrete algorithm;
// the pool has no policy to resolve it, so Auto here is an upstream contract bug.
if (algorithm == IndexingAlgorithmEnum::Auto)
throw JFJochException(JFJochExceptionCategory::InputParameterInvalid,
"Internal error: indexing algorithm must be resolved (not Auto) "
"before reaching the indexer pool");
// The workers built their indexers from the raw requested algorithm, but the algorithm actually
// dispatched is the RESOLVED one (rotation, for instance, always resolves to the GPU FFT indexer
// when a GPU is present, ignoring the request). If the resolution lands on an algorithm this host
// did not build an indexer for, fail here with an explanation instead of the opaque "no indexer
// available for the resolved algorithm" from deep inside a worker.
const auto requested = experiment.GetIndexingSettings().GetAlgorithm();
const bool have_gpu = get_gpu_count() > 0;
#ifdef JFJOCH_USE_FFTW
constexpr bool fftw_built = true;
#else
constexpr bool fftw_built = false;
#endif
const bool servable =
(algorithm == IndexingAlgorithmEnum::FFT && have_gpu &&
(requested == IndexingAlgorithmEnum::Auto || requested == IndexingAlgorithmEnum::FFT)) ||
(algorithm == IndexingAlgorithmEnum::FFBIDX && have_gpu &&
(requested == IndexingAlgorithmEnum::Auto || requested == IndexingAlgorithmEnum::FFBIDX)) ||
(algorithm == IndexingAlgorithmEnum::FFTW && fftw_built &&
((requested == IndexingAlgorithmEnum::Auto && !have_gpu) || requested == IndexingAlgorithmEnum::FFTW));
if (!servable) {
std::string msg;
if (requested == IndexingAlgorithmEnum::FFTW && have_gpu)
msg = "FFTW is the CPU indexer and is not available on a node with a GPU. Rotation indexing "
"always uses the GPU FFT indexer here; select FFT or Auto, or run FFTW on a CPU-only node.";
else if (algorithm == IndexingAlgorithmEnum::FFT && !have_gpu)
msg = "FFT is the GPU indexer but no GPU is available. Select FFTW or Auto for CPU indexing.";
else if (algorithm == IndexingAlgorithmEnum::FFBIDX && !have_gpu)
msg = "FFBIDX is a GPU indexer but no GPU is available. Select FFTW or Auto for CPU indexing.";
else if (algorithm == IndexingAlgorithmEnum::FFTW)
msg = "FFTW (CPU) indexing was requested but this build has no FFTW indexer.";
else
msg = "the requested indexing algorithm resolved to one with no indexer available on this host.";
throw JFJochException(JFJochExceptionCategory::InputParameterInvalid,
"Cannot index: " + msg);
}
// Check if there is available worker
const int task = GetFreeWorker();
std::unique_ptr<IndexerResult> result;
if (task >= 0) {
try {
result = tasks[task]->Run(experiment, recip);
} catch (const std::exception &e) {
spdlog::error("Indexer thread failed: {}", e.what());
result = std::make_unique<IndexerResult>(IndexerResult{
.lattice = {}, .indexing_time_s = 0, .executed = false, .error = e.what()});
}
{
std::unique_lock<std::mutex> lock(m);
worker_busy[task] = 0;
worker_free_count++;
}
c.notify_one();
}
if (result)
return *result;
// No free worker, or the pool is stopping: indexing was not attempted. Distinct from both a
// frame that did not index and an indexer that failed, and left without an error for that reason.
return IndexerResult{.lattice = {}, .indexing_time_s = 0};
}