Files
Jungfraujoch/image_analysis/indexing/IndexerThreadPool.cpp
T
leonarski_fandClaude Opus 5.5 a6bb16ce8c RotationIndexer: keep what RunIndexing computed, keyed by everything it reads
A rotation run asks the rotation indexer the same question many times: the
canonical pass's first pass repeats the rotation-scale probe at the stored
angles exactly (identical validation evidence on every set checked), and on
a crystal that does not index, pass 2 and every probe repeat pass 1's rescue
ladder rung for rung. Each such RunIndexing is an FFT search plus a serial
Ceres fixed-point chain, ~1-3 s on the GPU build.

RunIndexing is deterministic in its inputs, so its outcome - every member it
sets - is now kept process-wide under a key of all of them: the accumulated
spots (every field) and their angles, both geometries and the axis, the
experiment's indexing settings, cell and space group, and the settings of the
pool it indexes with (IndexerThreadPool::Settings). A RotationIndexer asking
with the same key takes the outcome. RUGNUX_VERIFY_FIRST_PASS_MEMO recomputes
and throws on a difference.

md5-identical output on four sets; GPU wall 28.9 -> 26.1 s, 43.9 -> 41.0 s,
30.3 -> 26.9 s and 143.6 -> 120.5 s (the set that does not index); the
verify mode found no difference on the last.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01D1G8gJVAy6gp1K5Dz3NE5C
2026-09-26 19:52:46 +02:00

313 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, input->severity_only));
} 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()});
// This thread goes on to the next frame, so a CUDA error left behind here would be
// reported by the first cudaGetLastError() of every frame that follows - one failed
// attempt would read as an indexer that never works again.
cuda_clear_error();
}
{
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, bool severity_only) {
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), severity_only);
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()),
settings_(settings) {
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,
bool severity_only) {
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, severity_only);
} 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};
}