From ef90b933da156dfa0c805e478cb4e2b4b887628c Mon Sep 17 00:00:00 2001 From: Filip Leonarski Date: Thu, 8 Oct 2026 08:40:15 +0200 Subject: [PATCH] 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