Merge branch 'par-merge' into rc175

This commit is contained in:
2026-10-08 15:58:29 +02:00
18 changed files with 135 additions and 36 deletions
+1 -1
View File
@@ -53,7 +53,7 @@ AzimuthalIntegrationMapping::AzimuthalIntegrationMapping(const DiffractionExperi
polarization_factor = experiment.GetPolarizationFactor();
if (in_nthreads == 0)
nthreads = std::thread::hardware_concurrency();
nthreads = AvailableCpus();
else
nthreads = in_nthreads;
+2 -1
View File
@@ -4,6 +4,7 @@
#include "ImageBuffer.h"
#include "JFJochException.h"
#include "ThreadAffinity.h"
#include "ZeroCopyReturnValue.h"
#include <algorithm>
@@ -19,7 +20,7 @@ namespace {
// faulting in a 150-200 GB allocation. Replaces numa_alloc_interleaved + a single-threaded
// memset, so this file no longer needs libnuma.
void parallel_first_touch(uint8_t *buffer, size_t buffer_size) {
const unsigned n = std::max(1u, std::thread::hardware_concurrency());
const unsigned n = AvailableCpus();
const size_t chunk = (buffer_size + n - 1) / n;
std::vector<std::thread> threads;
+1 -1
View File
@@ -58,7 +58,7 @@ namespace parallel_detail {
private:
WorkerPool() {
const unsigned hw = std::max(1u, std::thread::hardware_concurrency());
const unsigned hw = AvailableCpus();
workers.reserve(hw - 1);
try {
for (unsigned i = 0; i + 1 < hw; i++) // the submitting thread takes a share too
+55
View File
@@ -5,6 +5,8 @@
#ifdef __linux__
#include <algorithm>
#include <cmath>
#include <cstdio>
#include <filesystem>
#include <fstream>
@@ -26,6 +28,45 @@ namespace {
}
const cpu_set_t start_mask = StartMask();
// The CPUs the cgroup's quota allows, rounded up; 0 where none is set. The process's own cgroup
// first, then the root of the hierarchy as mounted, which is the container's own cgroup where the
// container does not see the host's paths. cgroup v2 (cpu.max: "max 100000" or "200000 100000")
// and v1 (cpu.cfs_quota_us = -1 or the quota, cpu.cfs_period_us).
unsigned CgroupCpuQuota() {
std::ifstream f("/proc/self/cgroup");
std::string line;
while (std::getline(f, line)) {
const size_t c1 = line.find(':'), c2 = line.find(':', c1 + 1);
if (c1 == std::string::npos || c2 == std::string::npos)
continue;
const std::string controllers = "," + line.substr(c1 + 1, c2 - c1 - 1) + ",";
const std::string path = line.substr(c2 + 1);
const bool v2 = line.compare(0, c2 + 1, "0::") == 0;
if (!v2 && controllers.find(",cpu,") == std::string::npos)
continue;
for (const std::string &dir : {"/sys/fs/cgroup" + std::string(v2 ? "" : "/cpu") + path,
"/sys/fs/cgroup" + std::string(v2 ? "" : "/cpu")}) {
double quota = -1, period = 0;
if (v2) {
std::ifstream q(dir + "/cpu.max");
std::string max;
if (!(q >> max >> period))
continue;
if (max != "max")
quota = std::atof(max.c_str());
} else {
std::ifstream q(dir + "/cpu.cfs_quota_us"), p(dir + "/cpu.cfs_period_us");
if (!(q >> quota) || !(p >> period))
continue;
}
if (quota > 0 && period > 0)
return static_cast<unsigned>(std::ceil(quota / period));
break; // found, and no quota set
}
}
return 0;
}
int NumaNodeCount() {
int n = 0;
std::error_code ec;
@@ -76,6 +117,16 @@ int NumaNodeOfPciDevice(const std::string &pci_bus_id) {
return node;
}
unsigned AvailableCpus() {
static const unsigned n = [] {
unsigned cpus = static_cast<unsigned>(std::max(1, CPU_COUNT(&start_mask)));
if (const unsigned quota = CgroupCpuQuota(); quota > 0)
cpus = std::min(cpus, quota);
return cpus;
}();
return n;
}
void PinThreadToNumaNode(int node) {
if (node < 0)
return;
@@ -97,6 +148,10 @@ void RestoreThreadAffinity() {
#else
#include <algorithm>
#include <thread>
unsigned AvailableCpus() { return std::max(1u, std::thread::hardware_concurrency()); }
int NumaNodeOfPciDevice(const std::string &) { return -1; }
void PinThreadToNumaNode(int) {}
void RestoreThreadAffinity() {}
+7
View File
@@ -14,6 +14,13 @@
// not known or the machine has a single node.
int NumaNodeOfPciDevice(const std::string &pci_bus_id);
// How many CPUs this process may use: the CPUs it was started on (its affinity mask - what taskset,
// a cpuset or a Slurm job gives it), and no more than its cgroup's CPU quota where one is set (a
// container started with --cpus). std::thread::hardware_concurrency() counts every CPU of the
// machine, so a job given 16 of a node's 192 would otherwise start 192 threads. Read once, at the
// first call. Outside Linux it is hardware_concurrency(). At least 1.
unsigned AvailableCpus();
// Restrict the calling thread to the CPUs of `node` that the process may run on. Nothing happens for
// node < 0, or where that would leave the thread no CPU.
void PinThreadToNumaNode(int node);
+3 -3
View File
@@ -442,7 +442,7 @@ ShadowFinder::Projection ShadowFinder::Reduce() const {
// A pixel is touched by one worker only and the sums and counts are integers, so the result is
// the same as folding on one thread.
out.frames += host.frames;
ParallelChunks(static_cast<int>(out.max_value.size()), std::thread::hardware_concurrency(), [&](int lo, int hi) {
ParallelChunks(static_cast<int>(out.max_value.size()), AvailableCpus(), [&](int lo, int hi) {
for (int i = lo; i < hi; i++) {
if (host.valid_count[i] == 0)
continue;
@@ -590,7 +590,7 @@ std::vector<float> ShadowFinder::GetMeanProjection() const {
const auto &valid_count = p.valid_count;
std::vector<float> mean(static_cast<size_t>(width) * height);
ParallelChunks(static_cast<int>(mean.size()), std::thread::hardware_concurrency(), [&](int lo, int hi) {
ParallelChunks(static_cast<int>(mean.size()), AvailableCpus(), [&](int lo, int hi) {
for (int i = lo; i < hi; i++)
mean[i] = valid_count[i] > 0 && pixel_mask[i] == 0
? static_cast<float>(static_cast<double>(sum_value[i]) / valid_count[i]) : NAN;
@@ -601,7 +601,7 @@ std::vector<float> ShadowFinder::GetMeanProjection() const {
std::vector<uint32_t> ShadowFinder::GetMask(size_t nthreads) const {
std::unique_lock ul(m);
if (nthreads == 0)
nthreads = std::max(1u, std::thread::hardware_concurrency());
nthreads = AvailableCpus();
#ifdef JFJOCH_USE_CUDA
// Where every frame went to the device the mask is made there too, from the projection as it lies.
if (Gpu() && gpu->GetFrameCount() > 0 && host.frames == 0) {
@@ -91,7 +91,7 @@ FindBeamCenterFromBackground(const DiffractionExperiment &experiment, const Pixe
const std::vector<float> &mean, size_t nthreads,
std::optional<std::pair<float, float>> start, bool allow_device) {
if (nthreads == 0)
nthreads = std::max(1u, std::thread::hardware_concurrency());
nthreads = AvailableCpus();
const auto W = static_cast<int>(experiment.GetXPixelsNumConv());
const auto H = static_cast<int>(experiment.GetYPixelsNumConv());
const size_t n_pixels = static_cast<size_t>(W) * H;
@@ -299,7 +299,7 @@ RotationScaleMerge::RotationScaleMerge(const DiffractionExperiment &experiment,
size_t nthreads, Logger &logger,
std::string observation_dump_path)
: x(experiment), partials_out(partial_outcomes), reference_cell(std::move(reference_cell)),
nthreads(nthreads == 0 ? std::thread::hardware_concurrency() : nthreads), logger(&logger),
nthreads(nthreads == 0 ? AvailableCpus() : nthreads), logger(&logger),
observation_dump_path(std::move(observation_dump_path)) {
const auto s = x.GetScalingSettings();
min_partiality = s.GetMinPartiality();
@@ -389,7 +389,7 @@ void RotationScaleMerge::IngestHost() {
#ifdef JFJOCH_USE_CUDA
// Probed before the observations are built: whether the build lays down the full Obs array or
// only the ingest arrays depends on whether the device pipeline will be active.
gpu_ = std::make_unique<RotationScaleMergeGPU>();
gpu_ = std::make_unique<RotationScaleMergeGPU>(gpu_device);
gpu_active_ = gpu_->Available();
resident_ingest = gpu_active_ && observation_dump_path.empty();
#endif
@@ -154,6 +154,10 @@ public:
// Merge only the first n frames of the outcomes (by default all of them). Set before Ingest().
void SetFrameLimit(int n) { frame_limit = n; }
// The GPU the engine runs on (device 0 by default). Set before Ingest(). The device only says where
// the merge runs: on cards of one model the result is the same on any of them.
void SetGpuDevice(int device) { gpu_device = device; }
// Whether Run() also hands back the merge's fulls unmerged (Result::scaled_fulls). Off by default:
// it is a copy as long as the fulls, wanted only by a caller that writes them.
void SetExportScaledFulls(bool on) { export_scaled_fulls = on; }
@@ -516,6 +520,7 @@ private:
#endif
bool write_back_per_frame_scale = true; // see SetWriteBackPerFrameScale
int frame_limit = -1; // see SetFrameLimit
int gpu_device = 0; // see SetGpuDevice
bool export_scaled_fulls = false; // see SetExportScaledFulls
bool french_wilson = true; // see SetFrenchWilson
@@ -1287,12 +1287,11 @@ struct DeviceGuard {
};
} // namespace
RotationScaleMergeGPU::RotationScaleMergeGPU() : impl_(std::make_unique<Impl>()) {
RotationScaleMergeGPU::RotationScaleMergeGPU(int device) : impl_(std::make_unique<Impl>()) {
if (get_gpu_count() > 0) {
// One instance, one device. It stays on device 0 for now - the merge is a single object and
// nothing else runs beside it - but it is recorded rather than assumed, so the guard below
// can put the caller's device back instead of leaving the thread moved.
impl_->device = 0;
// One instance, one device, recorded so the guard below can put the caller's device back
// instead of leaving the thread moved.
impl_->device = device;
DeviceGuard guard(impl_->device, true);
impl_->stream = std::make_unique<CudaStream>();
impl_->available = true;
@@ -20,7 +20,8 @@
// the pimpl is null without CUDA, and RotationScaleMerge falls back to the CPU loops).
class RotationScaleMergeGPU {
public:
RotationScaleMergeGPU();
// On `device`, which must be one of the visible GPUs.
explicit RotationScaleMergeGPU(int device = 0);
~RotationScaleMergeGPU();
RotationScaleMergeGPU(const RotationScaleMergeGPU &) = delete;
RotationScaleMergeGPU &operator=(const RotationScaleMergeGPU &) = delete;
+2 -1
View File
@@ -4,6 +4,7 @@
#include "ScaleOnTheFly.h"
#include "../../common/Logger.h"
#include "../../common/ThreadAffinity.h"
#include <algorithm>
#include <atomic>
@@ -195,7 +196,7 @@ void ScaleOnTheFly::RejectCollapsedScales(std::vector<IntegrationOutcome> &integ
void ScaleOnTheFly::Scale(std::vector<IntegrationOutcome> &integration, size_t nthreads) const {
if (nthreads == 0)
nthreads = std::thread::hardware_concurrency();
nthreads = AvailableCpus();
if (nthreads <= 1) {
for (auto & i : integration)
@@ -13,6 +13,7 @@
#include <ceres/rotation.h>
#include "Merge.h"
#include "../../common/ThreadAffinity.h"
namespace {
constexpr size_t MIN_FIT_REFLECTIONS = 20;
@@ -350,7 +351,7 @@ double StillsPartialityRefine::RefineOne(IntegrationOutcome &outcome,
double StillsPartialityRefine::Run(std::vector<IntegrationOutcome> &outcomes, size_t nthreads) const {
if (nthreads == 0)
nthreads = std::thread::hardware_concurrency();
nthreads = AvailableCpus();
nthreads = std::max<size_t>(1, nthreads);
double last_mean_tilt = 0.0;
+2 -1
View File
@@ -13,6 +13,7 @@
#include "JFJochTIFF.h"
#include "../common/JFJochException.h"
#include "../common/JFJochMath.h"
#include "../common/ThreadAffinity.h"
#include "../common/DiffractionGeometry.h"
#include "../frame_serialize/CBORStream2Deserializer.h"
#include "../compression/JFJochDecompress.h"
@@ -365,7 +366,7 @@ void PreviewImage::Configure(const DiffractionExperiment &in_experiment, const P
mask.resize(experiment.GetPixelsNum(), 0);
if (nthreads == 0)
nthreads = std::thread::hardware_concurrency();
nthreads = AvailableCpus();
nthreads = std::clamp<size_t>(nthreads, 1, 8);
+2 -1
View File
@@ -14,6 +14,7 @@
#include "../common/JFJochException.h"
#include "../common/Logger.h"
#include "../common/ThreadAffinity.h"
namespace {
@@ -285,7 +286,7 @@ Layout Place(const std::vector<Frame> &frames, const std::string &logger_name) {
}
void ForEachInOrder(size_t n, const std::function<void(size_t)> &fn) {
const size_t nthreads = std::min<size_t>(std::max(1u, std::thread::hardware_concurrency()), 8);
const size_t nthreads = std::min<size_t>(AvailableCpus(), 8);
std::atomic<size_t> next{0};
std::vector<std::future<void>> futures;
for (size_t t = 0; t < nthreads; t++)
+35 -6
View File
@@ -726,6 +726,14 @@ bool Rugnux::ScaleMergeAndSymmetry(PipelineLocals &p) {
// post-refinement, so it runs beside it; its log lines are held and replayed where it is used.
const bool merge_first_part = is_rotation && quality_guard_pass1_ && prepass_images < images_to_process
&& !experiment_.GetRefineRotationWedgeInScaling() && !rot_ss.GetRotationWedgeForScaling().has_value();
// With two GPUs or more the guard's merge has an engine of its own on the second card, and the
// main engine does not wait for it: its ingest and first merge run beside the guard's merge. The
// guard reads the outcomes only while it ingests - the first merge writes per-frame fields back
// into them (mosaicity among them, which the ingest reads) - so that merge waits for the guard's
// ingest and nothing else. Neither merge reads anything the other writes, and on cards of one
// model a merge gives the same bits on either, so this changes only when the guard runs.
const bool guard_on_second_gpu = get_gpu_count() >= 2;
std::promise<void> guard_ingested;
Logger guard_log = Logger::Buffered();
const auto merge_guard = [&] {
TimingMark mark(guard_log, "merge quality guard (pre-pass frames)");
@@ -734,7 +742,13 @@ bool Rugnux::ScaleMergeAndSymmetry(PipelineLocals &p) {
guard_log, "");
first_part_rsm.SetFrameLimit(prepass_images);
first_part_rsm.SetWriteBackPerFrameScale(false);
first_part_rsm.Ingest();
if (guard_on_second_gpu)
first_part_rsm.SetGpuDevice(1);
{
// Set however the ingest ends: a guard that failed rethrows from its future, later.
struct Ingested { std::promise<void> &p; ~Ingested() { p.set_value(); } } ingested{guard_ingested};
first_part_rsm.Ingest();
}
first_part_rsm.SetExportScaledFulls(false);
first_part_rsm.SetFrenchWilson(false);
auto r = first_part_rsm.Run(!experiment_.GetGemmiSpaceGroup().has_value(),
@@ -743,9 +757,14 @@ bool Rugnux::ScaleMergeAndSymmetry(PipelineLocals &p) {
};
std::future<ScaleMergeResult> guard_ahead;
if (merge_first_part && !(postrefine_probe_ && !geometry_prepass && postrefine_probe_only_))
guard_ahead = std::async(std::launch::async, merge_guard);
guard_ahead = std::async(std::launch::async, [&] {
if (guard_on_second_gpu)
set_gpu(1); // for any device work of the merge that is not the engine's own
return merge_guard();
});
// The scaling engine's ingest reads the outcomes and nothing else either, so its host half runs
// beside the post-refinement too; the device half waits for the guard to hand the GPU back.
// beside the post-refinement too; with one GPU the device half waits for the guard to hand the
// card back.
Logger engine_log = Logger::Buffered();
std::future<void> engine_ahead;
if (is_rotation && postrefine_probe_ && !geometry_prepass && !postrefine_probe_only_
@@ -820,8 +839,10 @@ bool Rugnux::ScaleMergeAndSymmetry(PipelineLocals &p) {
if (merge_first_part) {
logger.Info("Two-pass: the quality guard reads this pass's merge of the first {} images, "
"the ones the pre-pass integrated", prepass_images);
guard_merge = guard_ahead.valid() ? guard_ahead.get() : merge_guard();
guard_log.ReplayInto(logger);
if (!guard_ahead.valid() || !guard_on_second_gpu) {
guard_merge = guard_ahead.valid() ? guard_ahead.get() : merge_guard();
guard_log.ReplayInto(logger);
}
}
// A reference MTZ is allowed for rotation: it fixes the space group / cell (on the CLI) and
// resolves the indexing ambiguity (below), but is NOT used to scale - the rotation merge stays
@@ -978,10 +999,18 @@ bool Rugnux::ScaleMergeAndSymmetry(PipelineLocals &p) {
return out;
};
// First pass: P1 when searching, or directly the user-fixed space group.
// First pass: P1 when searching, or directly the user-fixed space group. A guard merging on the
// second GPU (above) is still running; this merge only waits for it to have read the outcomes.
const bool guard_beside = merge_first_part && guard_on_second_gpu && guard_ahead.valid();
if (guard_beside)
guard_ingested.get_future().wait();
const auto initial_sg = experiment_.GetGemmiSpaceGroup();
auto sm = scale_and_merge(initial_sg ? initial_sg->short_name() : "P1", !initial_sg.has_value(),
/*measure_cc=*/true);
if (guard_beside) {
guard_merge = guard_ahead.get();
guard_log.ReplayInto(logger);
}
// The two-pass quality guard judges the second pass against the first on THIS merge, which both
// passes make in the same terms - P1 (or the group the user fixed), full range, no correction
+7 -8
View File
@@ -34,6 +34,7 @@
#include <spdlog/sinks/stdout_color_sinks.h>
#include "../common/Logger.h"
#include "../common/ThreadAffinity.h"
#include "../common/GitInfo.h"
#include "../common/Definitions.h"
#include "../common/DiffractionExperiment.h"
@@ -98,7 +99,7 @@ void print_usage() {
std::cout << " -o, --output-prefix <txt> Output file prefix (default: output). mx and scale runs always write "
"<prefix>_report.txt, the results report; every run writes <prefix>_log.txt, the full timestamped log "
"(debug level); an empty prefix writes nothing at all" << std::endl;
std::cout << " -N, --threads <num> Number of threads (default: all hardware threads, of which the "
std::cout << " -N, --threads <num> Number of threads (default: every CPU the run may use - its affinity and cgroup quota - of which the "
"per-image loop uses at most 16 per GPU; an explicit value is used as given)" << std::endl;
std::cout << " -s, --start-image <num> Start image number (default: 0)" << std::endl;
std::cout << " -e, --end-image <num> End image number (default: all)" << std::endl;
@@ -846,7 +847,7 @@ static int RunRugnux(int argc, char **argv) {
std::string input_file;
std::string output_prefix = "output";
int nthreads = 0; // 0 = auto: resolved to all hardware threads after parsing (see below)
int nthreads = 0; // 0 = auto: resolved to the CPUs the run may use after parsing (see below)
int start_image = 0;
int end_image = -1; // -1 indicates process until end
int image_stride = 1;
@@ -1512,15 +1513,13 @@ static int RunRugnux(int argc, char **argv) {
input_file = argv[optind];
// -N defaults to 0 = "use all hardware threads"; resolve it to a concrete count here so every mode
// -N defaults to 0 = "use every CPU the run may use" (AvailableCpus); resolve it to a concrete count here so every mode
// behaves the same. The scale/merge engines expand 0 on their own, but the per-image processing
// loop (Rugnux) spawns exactly nthreads workers, so passing 0 there would spawn none and process
// nothing - hence resolving it centrally rather than relying on each consumer.
const bool nthreads_auto = nthreads <= 0;
if (nthreads <= 0) {
unsigned int hw = std::thread::hardware_concurrency();
nthreads = hw > 0 ? static_cast<int>(hw) : 1;
}
if (nthreads <= 0)
nthreads = static_cast<int>(AvailableCpus());
// The log file, next to the outputs, and everything logged so far into it.
std::string &log_file = g_log_file;
@@ -1537,7 +1536,7 @@ static int RunRugnux(int argc, char **argv) {
}
}
log_sinks->remove_sink(early_log);
logger.Info("Threads: {}{}", nthreads, nthreads_auto ? " (all hardware threads)" : " (-N)");
logger.Info("Threads: {}{}", nthreads, nthreads_auto ? " (the CPUs this run may use)" : " (-N)");
console->Field("Input", input_file);
console->Field("Output", log_file.empty() ? output_prefix + "_*"
@@ -7,6 +7,7 @@
#include "../../rugnux/RugnuxDefaults.h"
#include "../widgets/ToolbarIcons.h"
#include "../../rugnux/RugnuxCommandLine.h"
#include "../../common/ThreadAffinity.h"
#include <QApplication>
#include <QCheckBox>
@@ -37,15 +38,13 @@
#include <QToolButton>
#include <QVBoxLayout>
#include <thread>
namespace {
// Status (the progress bar) is column 2 so it stays visible in the narrow dock.
enum Column { COL_NAME = 0, COL_STATUS, COL_STARTED, COL_MODE, COL_IMAGES, COL_INDEX, COL_CELL, COL_ACTIONS, COL_COUNT };
int default_threads() {
const unsigned hc = std::thread::hardware_concurrency();
return hc == 0 ? 4 : static_cast<int>(hc);
return static_cast<int>(AvailableCpus());
}
QTableWidgetItem *fixedItem(const QString &text) { // a non-editable cell