From 50006eff889021306b3f7d639e8378766dd63cda Mon Sep 17 00:00:00 2001 From: Filip Leonarski Date: Thu, 8 Oct 2026 08:40:09 +0200 Subject: [PATCH 1/2] Thread count from the CPUs the process may use, not the machine's std::thread::hardware_concurrency() counts every CPU of the node, so a job given a CPU set (taskset, a cpuset, a Slurm allocation) or a container CPU quota started one thread per node CPU: the default -N, the shared worker pool and every "0 = all threads" default. AvailableCpus() (ThreadAffinity) counts the start affinity mask, capped by the cgroup CPU quota (v2 cpu.max, v1 cfs_quota/period, the process's own cgroup and then the mounted root), read once. Outside Linux it is hardware_concurrency(), so the viewer tree stays portable. It replaces every hardware_concurrency() call outside the tests and the vendored pocketfft. Only how many threads run changes; the passes split their work by n alone, so the results do not depend on it. Co-Authored-By: Claude Opus 5.5 (1M context) Claude-Session: https://claude.ai/code/session_01SVmAWnzCmRKAXVUCdc4iNi --- common/AzimuthalIntegrationMapping.cpp | 2 +- common/ImageBuffer.cpp | 3 +- common/ParallelFor.h | 2 +- common/ThreadAffinity.cpp | 55 +++++++++++++++++++ common/ThreadAffinity.h | 7 +++ image_analysis/beam_stop/ShadowFinder.cpp | 6 +- .../BeamCenterFromBackground.cpp | 2 +- .../scale_merge/RotationScaleMerge.cpp | 2 +- image_analysis/scale_merge/ScaleOnTheFly.cpp | 3 +- .../scale_merge/StillsPartialityRefine.cpp | 3 +- preview/PreviewImage.cpp | 3 +- reader/SweepLayout.cpp | 3 +- rugnux/rugnux_cli.cpp | 15 +++-- viewer/windows/JFJochProcessingJobsWindow.cpp | 5 +- 14 files changed, 88 insertions(+), 23 deletions(-) diff --git a/common/AzimuthalIntegrationMapping.cpp b/common/AzimuthalIntegrationMapping.cpp index efebb322b..f33ba2ca6 100644 --- a/common/AzimuthalIntegrationMapping.cpp +++ b/common/AzimuthalIntegrationMapping.cpp @@ -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; diff --git a/common/ImageBuffer.cpp b/common/ImageBuffer.cpp index 5b2809b3e..fff84beb2 100644 --- a/common/ImageBuffer.cpp +++ b/common/ImageBuffer.cpp @@ -4,6 +4,7 @@ #include "ImageBuffer.h" #include "JFJochException.h" +#include "ThreadAffinity.h" #include "ZeroCopyReturnValue.h" #include @@ -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 threads; diff --git a/common/ParallelFor.h b/common/ParallelFor.h index 276cbd35a..65b9e262c 100644 --- a/common/ParallelFor.h +++ b/common/ParallelFor.h @@ -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 diff --git a/common/ThreadAffinity.cpp b/common/ThreadAffinity.cpp index 9a85305da..cd4d9d52c 100644 --- a/common/ThreadAffinity.cpp +++ b/common/ThreadAffinity.cpp @@ -5,6 +5,8 @@ #ifdef __linux__ +#include +#include #include #include #include @@ -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(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(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 +#include + +unsigned AvailableCpus() { return std::max(1u, std::thread::hardware_concurrency()); } int NumaNodeOfPciDevice(const std::string &) { return -1; } void PinThreadToNumaNode(int) {} void RestoreThreadAffinity() {} diff --git a/common/ThreadAffinity.h b/common/ThreadAffinity.h index 07a7e2c1c..3af3232a4 100644 --- a/common/ThreadAffinity.h +++ b/common/ThreadAffinity.h @@ -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); diff --git a/image_analysis/beam_stop/ShadowFinder.cpp b/image_analysis/beam_stop/ShadowFinder.cpp index 3d62d6d6d..879d3e6ec 100644 --- a/image_analysis/beam_stop/ShadowFinder.cpp +++ b/image_analysis/beam_stop/ShadowFinder.cpp @@ -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(out.max_value.size()), std::thread::hardware_concurrency(), [&](int lo, int hi) { + ParallelChunks(static_cast(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 ShadowFinder::GetMeanProjection() const { const auto &valid_count = p.valid_count; std::vector mean(static_cast(width) * height); - ParallelChunks(static_cast(mean.size()), std::thread::hardware_concurrency(), [&](int lo, int hi) { + ParallelChunks(static_cast(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(static_cast(sum_value[i]) / valid_count[i]) : NAN; @@ -601,7 +601,7 @@ std::vector ShadowFinder::GetMeanProjection() const { std::vector 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) { diff --git a/image_analysis/geom_refinement/BeamCenterFromBackground.cpp b/image_analysis/geom_refinement/BeamCenterFromBackground.cpp index 0b99a030c..3f7512114 100644 --- a/image_analysis/geom_refinement/BeamCenterFromBackground.cpp +++ b/image_analysis/geom_refinement/BeamCenterFromBackground.cpp @@ -91,7 +91,7 @@ FindBeamCenterFromBackground(const DiffractionExperiment &experiment, const Pixe const std::vector &mean, size_t nthreads, std::optional> start, bool allow_device) { if (nthreads == 0) - nthreads = std::max(1u, std::thread::hardware_concurrency()); + nthreads = AvailableCpus(); const auto W = static_cast(experiment.GetXPixelsNumConv()); const auto H = static_cast(experiment.GetYPixelsNumConv()); const size_t n_pixels = static_cast(W) * H; diff --git a/image_analysis/scale_merge/RotationScaleMerge.cpp b/image_analysis/scale_merge/RotationScaleMerge.cpp index bc164861c..d48061fb9 100644 --- a/image_analysis/scale_merge/RotationScaleMerge.cpp +++ b/image_analysis/scale_merge/RotationScaleMerge.cpp @@ -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(); diff --git a/image_analysis/scale_merge/ScaleOnTheFly.cpp b/image_analysis/scale_merge/ScaleOnTheFly.cpp index 6d4ccf030..823b42706 100644 --- a/image_analysis/scale_merge/ScaleOnTheFly.cpp +++ b/image_analysis/scale_merge/ScaleOnTheFly.cpp @@ -4,6 +4,7 @@ #include "ScaleOnTheFly.h" #include "../../common/Logger.h" +#include "../../common/ThreadAffinity.h" #include #include @@ -195,7 +196,7 @@ void ScaleOnTheFly::RejectCollapsedScales(std::vector &integ void ScaleOnTheFly::Scale(std::vector &integration, size_t nthreads) const { if (nthreads == 0) - nthreads = std::thread::hardware_concurrency(); + nthreads = AvailableCpus(); if (nthreads <= 1) { for (auto & i : integration) diff --git a/image_analysis/scale_merge/StillsPartialityRefine.cpp b/image_analysis/scale_merge/StillsPartialityRefine.cpp index acdf8e81c..44910f6cb 100644 --- a/image_analysis/scale_merge/StillsPartialityRefine.cpp +++ b/image_analysis/scale_merge/StillsPartialityRefine.cpp @@ -13,6 +13,7 @@ #include #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 &outcomes, size_t nthreads) const { if (nthreads == 0) - nthreads = std::thread::hardware_concurrency(); + nthreads = AvailableCpus(); nthreads = std::max(1, nthreads); double last_mean_tilt = 0.0; diff --git a/preview/PreviewImage.cpp b/preview/PreviewImage.cpp index b7f941727..b53b94cef 100644 --- a/preview/PreviewImage.cpp +++ b/preview/PreviewImage.cpp @@ -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(nthreads, 1, 8); diff --git a/reader/SweepLayout.cpp b/reader/SweepLayout.cpp index d43609df9..343fbd171 100644 --- a/reader/SweepLayout.cpp +++ b/reader/SweepLayout.cpp @@ -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 &frames, const std::string &logger_name) { } void ForEachInOrder(size_t n, const std::function &fn) { - const size_t nthreads = std::min(std::max(1u, std::thread::hardware_concurrency()), 8); + const size_t nthreads = std::min(AvailableCpus(), 8); std::atomic next{0}; std::vector> futures; for (size_t t = 0; t < nthreads; t++) diff --git a/rugnux/rugnux_cli.cpp b/rugnux/rugnux_cli.cpp index 962c17e45..4545ae1a4 100644 --- a/rugnux/rugnux_cli.cpp +++ b/rugnux/rugnux_cli.cpp @@ -34,6 +34,7 @@ #include #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 Output file prefix (default: output). mx and scale runs always write " "_report.txt, the results report; every run writes _log.txt, the full timestamped log " "(debug level); an empty prefix writes nothing at all" << std::endl; - std::cout << " -N, --threads Number of threads (default: all hardware threads, of which the " + std::cout << " -N, --threads 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 Start image number (default: 0)" << std::endl; std::cout << " -e, --end-image 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(hw) : 1; - } + if (nthreads <= 0) + nthreads = static_cast(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 + "_*" diff --git a/viewer/windows/JFJochProcessingJobsWindow.cpp b/viewer/windows/JFJochProcessingJobsWindow.cpp index 0c9b849c7..c63f3ff6a 100644 --- a/viewer/windows/JFJochProcessingJobsWindow.cpp +++ b/viewer/windows/JFJochProcessingJobsWindow.cpp @@ -7,6 +7,7 @@ #include "../../rugnux/RugnuxDefaults.h" #include "../widgets/ToolbarIcons.h" #include "../../rugnux/RugnuxCommandLine.h" +#include "../../common/ThreadAffinity.h" #include #include @@ -37,15 +38,13 @@ #include #include -#include 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(hc); + return static_cast(AvailableCpus()); } QTableWidgetItem *fixedItem(const QString &text) { // a non-editable cell From ef90b933da156dfa0c805e478cb4e2b4b887628c Mon Sep 17 00:00:00 2001 From: Filip Leonarski Date: Thu, 8 Oct 2026 08:40:15 +0200 Subject: [PATCH 2/2] rugnux: the quality guard's merge on a second GPU, beside the main engine With two GPUs or more the two-pass quality guard's merge of the pre-pass frames gets a scaling engine of its own on device 1 (RotationScaleMerge:: SetGpuDevice), and the main engine's device ingest and first merge no longer wait for it to hand the card back. The first merge waits only for the guard's ingest, the last time the guard reads the outcomes: that merge writes per-frame fields (mosaicity among them) back into them. Neither merge reads what the other writes, and a merge gives the same bits on any card of one model, so only when the guard runs changes. With one GPU, and in the CPU build, the order is the one it was. Co-Authored-By: Claude Opus 5.5 (1M context) Claude-Session: https://claude.ai/code/session_01SVmAWnzCmRKAXVUCdc4iNi --- .../scale_merge/RotationScaleMerge.cpp | 2 +- .../scale_merge/RotationScaleMerge.h | 5 +++ .../scale_merge/RotationScaleMergeGPU.cu | 9 ++-- .../scale_merge/RotationScaleMergeGPU.h | 3 +- rugnux/RugnuxScaleMerge.cpp | 41 ++++++++++++++++--- 5 files changed, 47 insertions(+), 13 deletions(-) diff --git a/image_analysis/scale_merge/RotationScaleMerge.cpp b/image_analysis/scale_merge/RotationScaleMerge.cpp index d48061fb9..390d8a28c 100644 --- a/image_analysis/scale_merge/RotationScaleMerge.cpp +++ b/image_analysis/scale_merge/RotationScaleMerge.cpp @@ -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(); + gpu_ = std::make_unique(gpu_device); gpu_active_ = gpu_->Available(); resident_ingest = gpu_active_ && observation_dump_path.empty(); #endif diff --git a/image_analysis/scale_merge/RotationScaleMerge.h b/image_analysis/scale_merge/RotationScaleMerge.h index 1e9f75483..6847ca4e1 100644 --- a/image_analysis/scale_merge/RotationScaleMerge.h +++ b/image_analysis/scale_merge/RotationScaleMerge.h @@ -153,6 +153,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; } @@ -487,6 +491,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 diff --git a/image_analysis/scale_merge/RotationScaleMergeGPU.cu b/image_analysis/scale_merge/RotationScaleMergeGPU.cu index af682a1a6..b2292e4b7 100644 --- a/image_analysis/scale_merge/RotationScaleMergeGPU.cu +++ b/image_analysis/scale_merge/RotationScaleMergeGPU.cu @@ -1173,12 +1173,11 @@ struct DeviceGuard { }; } // namespace -RotationScaleMergeGPU::RotationScaleMergeGPU() : impl_(std::make_unique()) { +RotationScaleMergeGPU::RotationScaleMergeGPU(int device) : impl_(std::make_unique()) { 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(); impl_->available = true; diff --git a/image_analysis/scale_merge/RotationScaleMergeGPU.h b/image_analysis/scale_merge/RotationScaleMergeGPU.h index 68c4de7ab..859f4692c 100644 --- a/image_analysis/scale_merge/RotationScaleMergeGPU.h +++ b/image_analysis/scale_merge/RotationScaleMergeGPU.h @@ -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; diff --git a/rugnux/RugnuxScaleMerge.cpp b/rugnux/RugnuxScaleMerge.cpp index fc4450246..8d1867113 100644 --- a/rugnux/RugnuxScaleMerge.cpp +++ b/rugnux/RugnuxScaleMerge.cpp @@ -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 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 &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 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 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