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
**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>
187 lines
7.6 KiB
C++
187 lines
7.6 KiB
C++
// SPDX-FileCopyrightText: 2026 Filip Leonarski, Paul Scherrer Institute <filip.leonarski@psi.ch>
|
|
// SPDX-License-Identifier: GPL-3.0-only
|
|
|
|
#pragma once
|
|
|
|
#include <algorithm>
|
|
#include <atomic>
|
|
#include <condition_variable>
|
|
#include <deque>
|
|
#include <exception>
|
|
#include <functional>
|
|
#include <mutex>
|
|
#include <thread>
|
|
#include <vector>
|
|
|
|
// Two shapes of "run this over a range on several threads", used by the analysis code. Both take the
|
|
// worker count from the caller rather than asking the hardware, so a run that was told how many
|
|
// threads to use keeps to it.
|
|
|
|
// How many workers a pass over `n` cheap items should use: enough that each gets at least
|
|
// `min_per_thread` of them, and never more than the caller was given. A pass whose data is small can
|
|
// otherwise spend more on splitting the work than on doing it. That is not a large machine's problem:
|
|
// it is what makes the same code right on an 8-core laptop, a 16-core desktop and a two-socket node,
|
|
// none of which should be handed 48 chunks of a few thousand items.
|
|
inline size_t ThreadsForWork(size_t n, size_t nthreads, size_t min_per_thread = 32768) {
|
|
if (nthreads <= 1 || n == 0) return 1;
|
|
return std::clamp<size_t>(n / min_per_thread, 1, nthreads);
|
|
}
|
|
|
|
namespace parallel_detail {
|
|
// The threads both helpers below run on. They are made once and kept, because the analysis code
|
|
// repeats some of its passes thousands of times in a run and a thread costs tens of microseconds
|
|
// to create and join - more, on a short pass, than the pass itself.
|
|
class WorkerPool {
|
|
public:
|
|
static WorkerPool &Instance() {
|
|
static WorkerPool pool;
|
|
return pool;
|
|
}
|
|
|
|
// True on a thread the pool owns. A parallel pass reached from inside one runs inline instead
|
|
// of queueing: the workers are already occupied by the outer pass, so waiting for one of them
|
|
// to pick up the inner work could wait forever.
|
|
static bool InWorker() { return in_worker; }
|
|
|
|
size_t WorkerCount() const { return workers.size(); }
|
|
|
|
void Submit(std::function<void()> job) {
|
|
{
|
|
std::lock_guard lock(m);
|
|
queue.push_back(std::move(job));
|
|
}
|
|
cv.notify_one();
|
|
}
|
|
|
|
private:
|
|
WorkerPool() {
|
|
const unsigned hw = std::max(1u, std::thread::hardware_concurrency());
|
|
workers.reserve(hw - 1);
|
|
for (unsigned i = 0; i + 1 < hw; i++) // the submitting thread takes a share too
|
|
workers.emplace_back([this] { in_worker = true; Loop(); });
|
|
}
|
|
|
|
~WorkerPool() {
|
|
{
|
|
std::lock_guard lock(m);
|
|
stop = true;
|
|
}
|
|
cv.notify_all();
|
|
for (auto &t: workers) t.join();
|
|
}
|
|
|
|
void Loop() {
|
|
for (;;) {
|
|
std::function<void()> job;
|
|
{
|
|
std::unique_lock lock(m);
|
|
cv.wait(lock, [this] { return stop || !queue.empty(); });
|
|
if (stop) return;
|
|
job = std::move(queue.front());
|
|
queue.pop_front();
|
|
}
|
|
job();
|
|
}
|
|
}
|
|
|
|
std::mutex m;
|
|
std::condition_variable cv;
|
|
std::deque<std::function<void()> > queue;
|
|
std::vector<std::thread> workers;
|
|
bool stop = false;
|
|
static thread_local bool in_worker;
|
|
};
|
|
|
|
inline thread_local bool WorkerPool::in_worker = false;
|
|
|
|
// What the tasks of one pass share: the body to call, how many of them are still outstanding, and
|
|
// the first exception any of them threw.
|
|
struct RunState {
|
|
const std::function<void(int)> *body = nullptr;
|
|
std::atomic<int> remaining{0};
|
|
std::mutex done_m;
|
|
std::condition_variable done_cv;
|
|
bool done = false; // guarded by done_m; see RunOneTask
|
|
std::mutex err_m;
|
|
std::exception_ptr error;
|
|
};
|
|
|
|
inline void RunOneTask(RunState &s, int t) {
|
|
try {
|
|
(*s.body)(t);
|
|
} catch (...) {
|
|
std::lock_guard lock(s.err_m);
|
|
if (!s.error) s.error = std::current_exception();
|
|
}
|
|
if (s.remaining.fetch_sub(1, std::memory_order_acq_rel) == 1) {
|
|
// The flag the waiter tests is set UNDER done_m, and the counter is not that flag. If the
|
|
// waiter watched the counter it could see zero the instant the decrement above lands -
|
|
// before this thread has taken the lock - find its predicate already true, never block,
|
|
// and return from RunTasks. RunState is a local of that frame, so the lock and the notify
|
|
// below would then run on a destroyed mutex and condition variable, writing pthread state
|
|
// into a stack frame the submitting thread has already reused. Watching a flag set under
|
|
// the lock means completion cannot be observed until this thread has released it.
|
|
std::lock_guard lock(s.done_m);
|
|
s.done = true;
|
|
s.done_cv.notify_all();
|
|
}
|
|
}
|
|
|
|
// Call body(t) for every t in [0, ntasks) on the pool and on this thread, and return once they have
|
|
// all finished. An exception from any of them is held until then and rethrown here, so the others
|
|
// still run to completion - which is what waiting on futures used to give.
|
|
inline void RunTasks(int ntasks, const std::function<void(int)> &body) {
|
|
if (ntasks <= 0) return;
|
|
WorkerPool &pool = WorkerPool::Instance();
|
|
if (ntasks == 1 || WorkerPool::InWorker() || pool.WorkerCount() == 0) {
|
|
for (int t = 0; t < ntasks; t++) body(t);
|
|
return;
|
|
}
|
|
RunState s;
|
|
s.body = &body;
|
|
s.remaining.store(ntasks, std::memory_order_relaxed);
|
|
RunState *sp = &s;
|
|
for (int t = 1; t < ntasks; t++)
|
|
pool.Submit([sp, t] { RunOneTask(*sp, t); });
|
|
RunOneTask(s, 0);
|
|
{
|
|
std::unique_lock lock(s.done_m);
|
|
s.done_cv.wait(lock, [sp] { return sp->done; });
|
|
}
|
|
if (s.error) std::rethrow_exception(s.error);
|
|
}
|
|
}
|
|
|
|
// Chunked: each worker gets one contiguous [lo, hi) range and there is no per-item synchronisation.
|
|
// Right for millions of cheap uniform items - the CPU stand-in for a flat CUDA grid-stride kernel.
|
|
// The split is fixed and deterministic, so a pass whose per-element work is independent gives the
|
|
// same answer as the serial loop, bit for bit.
|
|
template <typename Fn>
|
|
void ParallelChunks(int n, size_t nthreads, Fn fn) {
|
|
if (n <= 0) return;
|
|
const int nt = static_cast<int>(std::max<size_t>(1, std::min(nthreads, static_cast<size_t>(n))));
|
|
if (nt == 1) { fn(0, n); return; }
|
|
const int chunk = (n + nt - 1) / nt;
|
|
parallel_detail::RunTasks(nt, [&](int t) {
|
|
const int lo = t * chunk, hi = std::min(n, lo + chunk);
|
|
if (lo < hi) fn(lo, hi);
|
|
});
|
|
}
|
|
|
|
// Work-stealing per item, off a shared atomic counter: one atomic per item, so use it only where the
|
|
// per-item work is heavy and uneven (per-frame fits, per-ring selections) and the atomic amortises.
|
|
// For millions of tiny uniform items a per-item atomic is pure contention - use ParallelChunks.
|
|
template <typename Fn>
|
|
void ParallelFor(int n, size_t nthreads, Fn fn) {
|
|
if (n <= 0) return;
|
|
if (nthreads <= 1 || n == 1) {
|
|
for (int i = 0; i < n; i++) fn(i);
|
|
return;
|
|
}
|
|
const size_t local = std::min(nthreads, static_cast<size_t>(n));
|
|
std::atomic<int> next = 0;
|
|
parallel_detail::RunTasks(static_cast<int>(local), [&](int) {
|
|
for (int i = next.fetch_add(1); i < n; i = next.fetch_add(1)) fn(i);
|
|
});
|
|
}
|