Files
Jungfraujoch/tests/ParallelForTest.cpp
T
leonarski_fandClaude Opus 5.5 c20cd016dc Fixed-partition reductions: sums that do not depend on the thread count
ParallelChunks cuts a range into one chunk per worker, so a floating-point
sum folded per chunk and then added up rounds differently at another -N.
ParallelBlocks cuts it by n alone (n / min_per_block blocks, at most 256);
each block folds into its own slot and the slots are added in block order,
so the sum has the same bits at any thread count.

Converted:
- RotationScaleMerge::ApplyCellSurface: the per-cell cross/ref2 sums of the
  surface fit (modulation, absorption), previously per-thread partials cut by
  ThreadsForWork(idx_all) threads.
- PostRefine: the scale-scan cost grid (per-chunk slots cut by -N).
- PostRefine joint solve: Ceres at a fixed 16 threads (as the rotation
  indexer's chain) instead of -N. Ceres sums cost and gradient in
  4 * num_threads pieces, so the split no longer follows -N. The pieces are
  handed to its threads in scheduling order, which only one thread makes
  exact; one thread was measured at +1.0-1.4 s of post-refinement on myob and
  lyso (0.5 -> 1.9 s, 1.1 -> 2.1 s), so that channel is left.
- IndexAndRefine supercell probe: per-frame probes kept by image number and
  summed in frame order, instead of added under a mutex in completion order;
  the primitive is the highest probed frame's instead of the last finisher's.

Integer reductions and per-item passes are unchanged; the three ParallelSort
callers already break ties on the index.

Validation (myob, cytc, lyso, sparse; GPU and CPU builds; -N 8/16/32):
p.hkl, p.mtz, p_P1.mtz and p_unmerged.mtz are md5-identical across -N, and
identical to the base 351be7de0 at every -N. The base was already -N
invariant on these four sets (GPU -N 1..32, CPU lyso -N 2..32), so no output
moved: the converted sums differ from the old ones only in last bits that
never reached a written float.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01D1G8gJVAy6gp1K5Dz3NE5C
2026-09-28 18:57:02 +02:00

189 lines
7.7 KiB
C++

// SPDX-FileCopyrightText: 2026 Filip Leonarski, Paul Scherrer Institute <filip.leonarski@psi.ch>
// SPDX-License-Identifier: GPL-3.0-only
#include <catch2/catch_all.hpp>
#include <algorithm>
#include <cmath>
#include <mutex>
#include <numeric>
#include <stdexcept>
#include <vector>
#include "../common/ParallelFor.h"
namespace {
const std::vector<size_t> THREAD_COUNTS = {0, 1, 2, 3, 8, 64};
const std::vector<int> SIZES = {0, 1, 2, 7, 64, 1000};
}
// The split is fixed and contiguous, which is what lets a pass whose per-element work is independent
// give the serial answer bit for bit. Every element has to land in exactly one slice, the slices have
// to tile [0, n) in order, and there must never be an empty one.
TEST_CASE("ParallelChunks_SlicesTileTheRange", "[ParallelFor]") {
for (int n: SIZES) {
for (size_t nthreads: THREAD_COUNTS) {
CAPTURE(n, nthreads);
std::mutex m;
std::vector<std::pair<int, int>> slices;
std::vector<int> visits(static_cast<size_t>(std::max(n, 1)), 0);
ParallelChunks(n, nthreads, [&](int lo, int hi) {
{
std::scoped_lock lock(m);
slices.emplace_back(lo, hi);
}
for (int i = lo; i < hi; i++)
visits[static_cast<size_t>(i)]++; // slices are disjoint, so no lock needed
});
std::sort(slices.begin(), slices.end());
int expected_lo = 0;
for (const auto &[lo, hi]: slices) {
CHECK(lo == expected_lo);
CHECK(hi > lo); // never an empty slice
expected_lo = hi;
}
CHECK(expected_lo == n); // and they reach the end
CHECK(slices.size() <= static_cast<size_t>(std::max(n, 0)));
for (int i = 0; i < n; i++)
CHECK(visits[static_cast<size_t>(i)] == 1);
}
}
}
// Work stealing, so a worker takes whatever is next rather than a fixed slice - but every item still
// has to be done exactly once, whatever order they come out in.
TEST_CASE("ParallelFor_VisitsEveryItemOnce", "[ParallelFor]") {
for (int n: SIZES) {
for (size_t nthreads: THREAD_COUNTS) {
CAPTURE(n, nthreads);
std::vector<std::atomic<int>> visits(static_cast<size_t>(std::max(n, 1)));
for (auto &v: visits)
v.store(0);
ParallelFor(n, nthreads, [&](int i) { visits[static_cast<size_t>(i)].fetch_add(1); });
for (int i = 0; i < n; i++)
CHECK(visits[static_cast<size_t>(i)].load() == 1);
}
}
}
// One thread means the caller's loop, in order - the property the serial fallbacks rely on.
TEST_CASE("ParallelFor_IsSerialAndInOrderForOneThread", "[ParallelFor]") {
for (size_t nthreads: {size_t{0}, size_t{1}}) {
CAPTURE(nthreads);
std::vector<int> order;
ParallelFor(16, nthreads, [&](int i) { order.push_back(i); });
std::vector<int> expected(16);
std::iota(expected.begin(), expected.end(), 0);
CHECK(order == expected);
}
}
// Splitting the work must not change the answer. Each element is written by exactly one worker, so
// the result has to match the serial loop element for element, at every thread count.
TEST_CASE("ParallelFor_SplitDoesNotChangeTheResult", "[ParallelFor]") {
constexpr int N = 1000;
std::vector<double> serial(N);
for (int i = 0; i < N; i++)
serial[static_cast<size_t>(i)] = std::sin(i * 0.001) * 1e6 + i;
for (size_t nthreads: THREAD_COUNTS) {
CAPTURE(nthreads);
std::vector<double> chunked(N, 0.0), stolen(N, 0.0);
ParallelChunks(N, nthreads, [&](int lo, int hi) {
for (int i = lo; i < hi; i++)
chunked[static_cast<size_t>(i)] = std::sin(i * 0.001) * 1e6 + i;
});
ParallelFor(N, nthreads, [&](int i) {
stolen[static_cast<size_t>(i)] = std::sin(i * 0.001) * 1e6 + i;
});
CHECK(chunked == serial); // bit for bit, not approximately
CHECK(stolen == serial);
}
}
// The blocks depend on n alone: the same tiling of [0, n), block b always the same range, whatever the
// thread count.
TEST_CASE("ParallelBlocks_PartitionIgnoresTheThreadCount", "[ParallelFor]") {
for (int n: {0, 1, 7, 4095, 4096, 100000, 5000000}) {
const int nb = ReductionBlocks(n);
CHECK(nb <= MAX_REDUCTION_BLOCKS);
std::vector<std::pair<int, int>> reference;
for (size_t nthreads: THREAD_COUNTS) {
CAPTURE(n, nthreads);
std::vector<std::pair<int, int>> range(static_cast<size_t>(nb), {-1, -1});
ParallelBlocks(n, nthreads, [&](int b, int lo, int hi) { range[static_cast<size_t>(b)] = {lo, hi}; });
int expected_lo = 0;
for (const auto &[lo, hi]: range) {
CHECK(lo == expected_lo);
CHECK(hi > lo);
expected_lo = hi;
}
CHECK(expected_lo == n);
if (reference.empty()) reference = range;
CHECK(range == reference);
}
}
}
// What the blocks are for: a floating-point sum folded per block and added up in block order has the
// same bits at any thread count. The terms span many decades, so any change of grouping would show.
TEST_CASE("ParallelBlocks_SumIsTheSameAtAnyThreadCount", "[ParallelFor]") {
constexpr int N = 1000003;
constexpr int NCELL = 37;
std::vector<double> term(N);
for (int i = 0; i < N; i++)
term[static_cast<size_t>(i)] = std::exp(std::sin(i * 0.37) * 20.0) * (i % 3 == 0 ? -1.0 : 1.0);
auto blocked_sums = [&](size_t nthreads) {
const int nb = ReductionBlocks(N);
std::vector<double> part(static_cast<size_t>(nb), 0.0);
std::vector<std::vector<double>> cell_part(static_cast<size_t>(nb), std::vector<double>(NCELL, 0.0));
ParallelBlocks(N, nthreads, [&](int b, int lo, int hi) {
for (int i = lo; i < hi; i++) {
part[static_cast<size_t>(b)] += term[static_cast<size_t>(i)];
cell_part[static_cast<size_t>(b)][static_cast<size_t>(i % NCELL)] += term[static_cast<size_t>(i)];
}
});
std::vector<double> out(NCELL + 1, 0.0);
for (int b = 0; b < nb; b++) {
out[NCELL] += part[static_cast<size_t>(b)];
for (int c = 0; c < NCELL; c++)
out[static_cast<size_t>(c)] += cell_part[static_cast<size_t>(b)][static_cast<size_t>(c)];
}
return out;
};
const std::vector<double> one = blocked_sums(1);
for (size_t nthreads: {size_t{3}, size_t{7}, size_t{32}}) {
CAPTURE(nthreads);
CHECK(blocked_sums(nthreads) == one); // bit for bit, not approximately
}
}
// A negative or zero count is a no-op rather than an error - callers pass a computed size.
TEST_CASE("ParallelFor_DoesNothingForAnEmptyRange", "[ParallelFor]") {
int calls = 0;
for (int n: {0, -1, -1000}) {
ParallelChunks(n, 8, [&](int, int) { calls++; });
ParallelFor(n, 8, [&](int) { calls++; });
ParallelBlocks(n, 8, [&](int, int, int) { calls++; });
}
CHECK(calls == 0);
}
// An exception thrown in a worker reaches the caller rather than terminating: the futures are waited
// on, so the other workers finish first and only then does it propagate.
TEST_CASE("ParallelFor_PropagatesAnException", "[ParallelFor]") {
CHECK_THROWS_AS(ParallelChunks(64, 4, [](int lo, int) {
if (lo == 0) throw std::runtime_error("from a chunk");
}), std::runtime_error);
CHECK_THROWS_AS(ParallelFor(64, 4, [](int i) {
if (i == 0) throw std::runtime_error("from an item");
}), std::runtime_error);
}