Build Packages / build:windows:nocuda (push) Successful in 16m8s
Build Packages / build:windows:cuda (push) Successful in 18m58s
Build Packages / build:viewer-tgz:cpu (push) Successful in 20m35s
Build Packages / build:viewer-tgz:cuda (push) Successful in 22m31s
Build Packages / build:rpm (rocky9_nocuda) (push) Successful in 25m9s
Build Packages / build:rpm (ubuntu2404_nocuda) (push) Successful in 25m6s
Build Packages / build:rpm (rocky8_nocuda) (push) Successful in 28m57s
Build Packages / build:rpm (ubuntu2204_nocuda) (push) Successful in 28m58s
Build Packages / build:rpm (rocky8_sls9) (push) Successful in 28m58s
Build Packages / XDS test (durin plugin) (push) Successful in 12m3s
Build Packages / build:rpm (rocky9_sls9) (push) Successful in 22m24s
Build Packages / build:rpm (rocky9) (push) Successful in 21m45s
Build Packages / Generate python client (push) Successful in 53s
Build Packages / build:rpm (rocky8) (push) Successful in 26m9s
Build Packages / Create release (push) Skipped
Build Packages / Build documentation (push) Successful in 1m37s
Build Packages / build:rpm (ubuntu2204) (push) Successful in 25m34s
Build Packages / build:rpm (ubuntu2404) (push) Successful in 22m0s
Build Packages / XDS test (JFJoch plugin) (push) Successful in 10m53s
Build Packages / XDS test (neggia plugin) (push) Successful in 9m29s
Build Packages / DIALS test (push) Successful in 23m40s
Build Packages / Unit tests (push) Successful in 1h20m1s
Three costs before and around the image loop. Every image allocated a fresh buffer for its compressed chunk and resized it, which value-initialises, and the read then overwrote every byte. At a few megabytes a chunk the allocation is large enough to be mapped rather than reused, so the zeroing was page-fault bound and cost more than the read it preceded - twenty gigabytes of it over a long sweep. The buffer now uses an allocator that does not construct, and the two HDF5 read paths are templated on the allocator so every existing caller compiles unchanged. The rebind is deliberate: without it the vector base rebinds to the default allocator and the zeroing quietly returns. The bitshuffle decoder allocated a whole uncompressed frame in its constructor - seventy megabytes a worker, five hundred and fifty across the loop - for the route that decodes the shuffled image separately. That route is taken only when a bitshuffle block is too large for the fused kernel, which neither writer this pipeline reads produces, so on a real frame the buffer is allocated, never touched, and freed. It is now allocated where it is used. The comment two lines below already warned against sizing a buffer from the uncompressed size; the line above it had not been given the same treatment. The first call into cuFFT pays the library's one-time initialisation, and it landed in the middle of the first pass with nothing to overlap it. It is now forced on a background thread at startup, alongside the file open and the mapping build, in the manner the shadow finder already uses. Finally, the detector mask was copied into the start message whether or not a file would carry it, which a merging run does not. It is filled where a writer is constructed - both places one is constructed, the second being the fallback that writes a process file when nothing indexed. Faster on eleven of thirty-eight crystals and slower on none; the whole rotation test set falls from four minutes thirty to four minutes seventeen, with each binary repeating itself to within half a per cent. Space groups thirty-five of thirty-eight and no failures throughout, and every column of the comparison table is identical. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01EGpGdgmJ8MyY9pCGWjktyi
300 lines
12 KiB
C++
300 lines
12 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) {
|
|
tmp_result = nullptr;
|
|
spdlog::error("Indexer thread {} failed: {}", threadid, 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 = nullptr;
|
|
}
|
|
{
|
|
std::unique_lock<std::mutex> lock(m);
|
|
worker_busy[task] = 0;
|
|
worker_free_count++;
|
|
}
|
|
c.notify_one();
|
|
}
|
|
if (result)
|
|
return *result;
|
|
return IndexerResult{.lattice = {}, .indexing_time_s = 0};
|
|
}
|