Files
Jungfraujoch/image_analysis/indexing/IndexerThreadPool.cpp
T
leonarski_fandjungfrau 4dc2534dbf
Build Packages / build:rpm (rocky9_sls9) (push) Successful in 18m57s
Build Packages / Unit tests (push) Skipped
Build Packages / build:windows:nocuda (push) Successful in 16m55s
Build Packages / build:windows:cuda (push) Successful in 18m48s
Build Packages / build:viewer-tgz:cpu (push) Successful in 13m10s
Build Packages / build:viewer-tgz:cuda (push) Successful in 14m45s
Build Packages / build:rpm (rocky8_nocuda) (push) Successful in 22m23s
Build Packages / build:rpm (rocky9_nocuda) (push) Successful in 20m12s
Build Packages / build:rpm (ubuntu2204_nocuda) (push) Successful in 23m7s
Build Packages / build:rpm (ubuntu2404_nocuda) (push) Successful in 20m43s
Build Packages / build:rpm (rocky8_sls9) (push) Successful in 23m9s
Build Packages / XDS test (durin plugin) (push) Successful in 12m26s
Build Packages / build:rpm (rocky9) (push) Successful in 24m58s
Build Packages / Generate python client (push) Successful in 50s
Build Packages / build:rpm (ubuntu2404) (push) Successful in 23m20s
Build Packages / Create release (push) Skipped
Build Packages / XDS test (JFJoch plugin) (push) Successful in 12m37s
Build Packages / build:rpm (rocky8) (push) Successful in 27m58s
Build Packages / build:rpm (ubuntu2204) (push) Successful in 25m38s
Build Packages / Build documentation (push) Successful in 59s
Build Packages / DIALS test (push) Successful in 23m16s
Build Packages / XDS test (neggia plugin) (push) Successful in 6m38s
v1.0.0.rc-162 (#72)
**Files written by Jungfraujoch now import correctly in DIALS, XDS and pyFAI.** A tilted detector, a grid scan, a still recorded at a goniometer position, and saturated or unreadable pixels were each described in a way that a third-party program acted on wrongly. If you process Jungfraujoch data outside Jungfraujoch, prefer this release to any earlier one.

* HDF5: the detector tilt (`rot1`/`rot2`/`rot3`) is exported correctly in the NXmx transformation chain; untilted geometries are unaffected.
* HDF5: a still recorded at a goniometer position is no longer read back as a single image, and a grid scan records a stationary spindle so a program that requires a rotation axis can open it.
* HDF5: the sample transformation chain is written in mounting order, with a Smargon head position told apart from the spindle, one entry per image, `module_offset` as a float unit vector, and `offset_units` on every offset.
* HDF5: saturated, underloaded and unreadable pixels are described so a downstream program masks them - `saturation_value`, `underload_value`, `error_value` and `bit_depth_readout` are written correctly, and a data file missing next to a VDS master reads as the error marker rather than as zero counts.
* HDF5: the rotation axis is read back under whatever name it carries, and `mirror_y` records whether the assembled image is mirrored in Y relative to the detector's raw readout.
* A grid scan and a goniometer axis can both be set; they are no longer alternatives.
* `images_per_file` is chosen from the acquisition when it is not given: a rotation sweep of at most 20000 images goes into a single data file, a grid scan splits on whole fast-axis rows, and stills and serial keep 1000.
* The writer refuses a stream whose start message declares a different pixel format than its images carry, and a DECTRIS detector sending signed images is no longer declared unsigned.
* The image stream can carry the sample transformation chain (`transformations`, in the END message); a producer that does not send it gets the same chain built by the writer.
* rugnux: fixing the space group with `-S` no longer prevents the lattice from being found - a lattice indexed in a different setting is reindexed into that group's own setting, and a run whose crystal does not have that group's lattice stops and names the cell it indexed as, rather than reporting statistics that cannot describe it.
* rugnux: the per-image resolution estimate now predicts the resolution the merged data reach rather than the highest-resolution spot found, and is reported as `SPOT_RESOLUTION_ESTIMATE`.
* rugnux: two runs of the same command on the same images produce the same merged intensities; the azimuthal profile written alongside them is not yet reproducible in the same way.
* rugnux: the offline lattice refinement is bounded by iterations rather than by a wall clock, so a loaded machine can no longer refine to a different lattice; a live acquisition keeps its real-time bound.
* rugnux: the detector-frame modulation correction is fitted on a grid spanning the detector, so whether it is applied no longer depends on how far integration reached.
* rugnux: the geometry pre-pass no longer writes `<prefix>_01.mtz`, `_01.cif`, `_01.hkl` and `_01_image.dat`; the refined second pass writes those files under `<prefix>`, and that is the result to use.
* rugnux: `_process.h5` describes the pixel format of the images it links to, and is written on a thread of its own.
* rugnux: the detector geometry is also logged in XDS's convention (`ORGX`/`ORGY`, detector axis vectors, rotation axis), so it can be compared with an XDS refinement.
* rugnux: an image integrated in pyFAI through the `.poni` file written by `--mode calibration` comes out with the correct azimuth, and the file declares pyFAI's `orientation`, which needs pyFAI 2024.01 or newer. Radial integration is unchanged.
* rugnux: a rotation run is substantially faster throughout - beam-stop detection, first-pass indexing, geometry refinement, integration, scaling and merging - and observations outside the scaling resolution range are dropped as they are ingested. The refined geometry, the space group chosen and the merged statistics are unchanged.
* Faster spot finding and indexing, on the broker as well as in rugnux; the spots found and the lattices indexed are unchanged.
* A run reserves substantially less GPU memory: nothing is allocated for buffers that are never read, and a worker builds only the engines it uses.
* rugnux: with `-N` left at its default the per-image loop of `--mode mx` uses at most 16 workers per GPU, rather than one per hardware thread; an explicit `-N` is obeyed as given.
* CUDA 12 builds now contain device code for Volta, so the RHEL 8 packages and the portable Linux `.tgz` run on a V100; the CUDA 13 artefacts (RHEL 9, Ubuntu, Windows) remain Turing and newer.
* The build resolves a single Eigen for the whole project, and refuses to configure if Ceres picks up a different one; a build that mixed two Eigen versions was undefined behaviour and crashed at -O2.
* Documentation: a security page, and the supported GPU generations and minimum NVIDIA driver version of every released artefact.

**Breaking change to OpenAPI** - regenerate the client (`jfjoch-client` 1.0.0-rc.162, `frontend/src/client`):
* `dataset_settings.images_per_file` is no longer `default: 1000` and no longer accepts `0`; it is optional, and its minimum is 1. A client sending `0` (previously "one file for the whole run") is now rejected - omit the field instead, which for a rotation sweep gives the same single file.
* `file_writer_format` now defaults to `NXmxVDS`, matching the server's own default and the layout recommended for DIALS, XDS and CrystFEL. A generated client that fills in schema defaults and does not set the format explicitly will write VDS masters where it previously wrote legacy ones; set `NXmxLegacy` explicitly to keep them.

---------

Co-authored-by: jungfrau <jungfrau@mx-aare-test.psi.ch>
Reviewed-on: #72
Co-authored-by: Filip Leonarski <filip.leonarski@psi.ch>
2026-08-25 08:21:39 +02:00

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};
}