diff --git a/common/CMakeLists.txt b/common/CMakeLists.txt index 8a4b2ca8d..69eab22f4 100644 --- a/common/CMakeLists.txt +++ b/common/CMakeLists.txt @@ -61,6 +61,7 @@ ADD_LIBRARY(JFJochCommon STATIC DetectorModuleGeometry.cpp DetectorModuleGeometry.h DetectorSetup.h DetectorSetup.cpp ZeroCopyReturnValue.h Histogram.h DiffractionGeometry.h CUDAWrapper.cpp CUDAWrapper.h + ThreadAffinity.cpp ThreadAffinity.h ADUHistogram.cpp ADUHistogram.h RawToConvertedGeometryCore.h Plot.h diff --git a/common/CUDAWrapper.cpp b/common/CUDAWrapper.cpp index 22d12c221..e57813c12 100644 --- a/common/CUDAWrapper.cpp +++ b/common/CUDAWrapper.cpp @@ -38,6 +38,12 @@ void set_gpu(int32_t dev_id) {} void pin_gpu() {} +void pin_gpu(int32_t dev_id) {} + +void enable_gpu_numa_binding() {} + +void set_gpu_blocking_sync() {} + void cuda_clear_error() {} #endif diff --git a/common/CUDAWrapper.cu b/common/CUDAWrapper.cu index 5cc0669db..6be8b0c81 100644 --- a/common/CUDAWrapper.cu +++ b/common/CUDAWrapper.cu @@ -3,8 +3,11 @@ #include #include +#include +#include #include "CUDAWrapper.h" +#include "ThreadAffinity.h" #include "JFJochException.h" inline void cuda_err(cudaError_t val) { @@ -55,11 +58,57 @@ void set_gpu(int32_t dev_id) { } } +namespace { + std::atomic gpu_numa_binding{false}; + + // The NUMA node of each device, read once. + int GpuNumaNode(int32_t dev_id) { + static std::mutex m; + static std::vector node; + std::lock_guard lock(m); + if (node.empty()) { + const int32_t count = get_gpu_count(); + node.assign(count, -1); + for (int32_t i = 0; i < count; i++) { + char bus_id[32] = {}; + if (cudaDeviceGetPCIBusId(bus_id, sizeof(bus_id), i) == cudaSuccess) + node[i] = NumaNodeOfPciDevice(bus_id); + } + } + return dev_id < static_cast(node.size()) ? node[dev_id] : -1; + } +} + +void enable_gpu_numa_binding() { + gpu_numa_binding = true; +} + +void pin_gpu(int32_t dev_id) { + if (get_gpu_count() == 0) + return; + set_gpu(dev_id); + if (gpu_numa_binding) + PinThreadToNumaNode(GpuNumaNode(dev_id)); +} + void pin_gpu() { static std::atomic counter{0}; auto dev_count = get_gpu_count(); if (dev_count > 0) - set_gpu(counter.fetch_add(1) % dev_count); + pin_gpu(static_cast(counter.fetch_add(1) % dev_count)); +} + +void set_gpu_blocking_sync() { + const int32_t count = get_gpu_count(); + for (int32_t i = 0; i < count; i++) { + if (cudaSetDevice(i) != cudaSuccess) + continue; + // cudaErrorSetOnActiveProcess where a context exists already: it keeps its flags. + if (cudaSetDeviceFlags(cudaDeviceScheduleBlockingSync) != cudaSuccess) + cudaGetLastError(); + } + if (count > 0) + cudaSetDevice(0); } void cuda_clear_error() { diff --git a/common/CUDAWrapper.h b/common/CUDAWrapper.h index 2b078ba3f..3feda3068 100644 --- a/common/CUDAWrapper.h +++ b/common/CUDAWrapper.h @@ -23,6 +23,20 @@ void set_gpu(int32_t dev_id); // is visible. Honours CUDA_VISIBLE_DEVICES via get_gpu_count(). void pin_gpu(); +// From here on, pin_gpu() also keeps the calling thread on the CPUs of the NUMA node its GPU is +// attached to (ThreadAffinity.h) - on a machine with one node, or without CUDA, nothing changes. +// Off unless a program asks for it. +void enable_gpu_numa_binding(); + +// pin_gpu() onto a given device: set_gpu(dev_id), and the NUMA binding above where it is enabled. +// For a worker that takes its card by index rather than round-robin. +void pin_gpu(int32_t dev_id); + +// Have every GPU's host threads BLOCK in a synchronisation (cudaDeviceScheduleBlockingSync) instead +// of spinning on a core until the device finishes. Must be called before anything creates a CUDA +// context; a device that already has one keeps the flags it was created with. No-op without CUDA. +void set_gpu_blocking_sync(); + // Drop the error CUDA has recorded for the calling thread. Call it where a CUDA failure has been // HANDLED - a device route that fell back to the host, an indexing attempt whose failure was turned // into a result - because the error otherwise stays as the thread's last error and the next diff --git a/common/ParallelFor.h b/common/ParallelFor.h index a40e118c7..09a39b09d 100644 --- a/common/ParallelFor.h +++ b/common/ParallelFor.h @@ -13,6 +13,8 @@ #include #include +#include "ThreadAffinity.h" + // 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. @@ -58,7 +60,12 @@ namespace parallel_detail { 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(); }); + workers.emplace_back([this] { + // It inherited whatever CPUs the thread that first used the pool was kept to. + RestoreThreadAffinity(); + in_worker = true; + Loop(); + }); } ~WorkerPool() { diff --git a/common/ThreadAffinity.cpp b/common/ThreadAffinity.cpp new file mode 100644 index 000000000..9a85305da --- /dev/null +++ b/common/ThreadAffinity.cpp @@ -0,0 +1,104 @@ +// SPDX-FileCopyrightText: 2026 Filip Leonarski, Paul Scherrer Institute +// SPDX-License-Identifier: GPL-3.0-only + +#include "ThreadAffinity.h" + +#ifdef __linux__ + +#include +#include +#include +#include + +#include +#include + +namespace { + // The process's CPUs, taken before any thread has been pinned: static initialisation runs on the + // main thread before main(). + cpu_set_t StartMask() { + cpu_set_t mask; + CPU_ZERO(&mask); + if (sched_getaffinity(0, sizeof(mask), &mask) != 0) + for (int c = 0; c < CPU_SETSIZE; c++) + CPU_SET(c, &mask); + return mask; + } + const cpu_set_t start_mask = StartMask(); + + int NumaNodeCount() { + int n = 0; + std::error_code ec; + for (const auto &e : std::filesystem::directory_iterator("/sys/devices/system/node", ec)) { + const std::string name = e.path().filename().string(); + if (name.rfind("node", 0) == 0 && name.size() > 4 + && name.find_first_not_of("0123456789", 4) == std::string::npos) + n++; + } + return n; + } + + // "0-23,48-71" -> the CPUs it names. + bool ParseCpuList(const std::string &list, cpu_set_t &out) { + CPU_ZERO(&out); + std::stringstream ss(list); + std::string range; + bool any = false; + while (std::getline(ss, range, ',')) { + int lo, hi; + if (std::sscanf(range.c_str(), "%d-%d", &lo, &hi) == 2) { + } else if (std::sscanf(range.c_str(), "%d", &lo) == 1) { + hi = lo; + } else + continue; + for (int c = lo; c <= hi && c < CPU_SETSIZE; c++) { + CPU_SET(c, &out); + any = true; + } + } + return any; + } +} + +int NumaNodeOfPciDevice(const std::string &pci_bus_id) { + if (NumaNodeCount() < 2) + return -1; + unsigned domain, bus, device, function; + if (std::sscanf(pci_bus_id.c_str(), "%x:%x:%x.%x", &domain, &bus, &device, &function) != 4) + return -1; + char path[128]; + std::snprintf(path, sizeof(path), "/sys/bus/pci/devices/%04x:%02x:%02x.%x/numa_node", + domain, bus, device, function); + std::ifstream f(path); + int node = -1; + if (!(f >> node)) + return -1; + return node; +} + +void PinThreadToNumaNode(int node) { + if (node < 0) + return; + std::ifstream f("/sys/devices/system/node/node" + std::to_string(node) + "/cpulist"); + std::string list; + cpu_set_t node_cpus; + if (!std::getline(f, list) || !ParseCpuList(list, node_cpus)) + return; + cpu_set_t mask; + CPU_AND(&mask, &node_cpus, &start_mask); + if (CPU_COUNT(&mask) == 0) + return; + pthread_setaffinity_np(pthread_self(), sizeof(mask), &mask); +} + +void RestoreThreadAffinity() { + pthread_setaffinity_np(pthread_self(), sizeof(start_mask), &start_mask); +} + +#else + +int NumaNodeOfPciDevice(const std::string &) { return -1; } +void PinThreadToNumaNode(int) {} +void RestoreThreadAffinity() {} + +#endif diff --git a/common/ThreadAffinity.h b/common/ThreadAffinity.h new file mode 100644 index 000000000..07a7e2c1c --- /dev/null +++ b/common/ThreadAffinity.h @@ -0,0 +1,24 @@ +// SPDX-FileCopyrightText: 2026 Filip Leonarski, Paul Scherrer Institute +// SPDX-License-Identifier: GPL-3.0-only + +#pragma once + +#include + +// Keeping a GPU's worker threads on the socket the card hangs off, on a machine with more than one +// NUMA node, so the pinned host buffers they allocate and the copies to the card stay on that +// socket's memory controller and PCIe root. Read from /sys - no libnuma - and Linux only: elsewhere, +// and on a machine with a single node, every call is a no-op. + +// The NUMA node of a PCI device, from its bus id as CUDA gives it ("0000:41:00.0"); -1 where it is +// not known or the machine has a single node. +int NumaNodeOfPciDevice(const std::string &pci_bus_id); + +// 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); + +// Give the calling thread back the CPUs the process started with. A thread inherits the affinity of +// the thread that creates it, so a thread pool first reached from a pinned worker would otherwise +// run every later parallel pass of the process on one socket. +void RestoreThreadAffinity(); diff --git a/rugnux/rugnux_cli.cpp b/rugnux/rugnux_cli.cpp index ffee0a521..c50bcea33 100644 --- a/rugnux/rugnux_cli.cpp +++ b/rugnux/rugnux_cli.cpp @@ -2937,6 +2937,11 @@ static int RunRugnux(int argc, char **argv) { int main(int argc, char **argv) { try { + // Before anything creates a CUDA context, which fixes the flags for the device. A thread + // waiting on the GPU then sleeps instead of spinning on a core - measured on a 16M rotation + // run, a fifth of all CPU time went to that spinning, with no change in wall time either way. + set_gpu_blocking_sync(); + enable_gpu_numa_binding(); return RunRugnux(argc, argv); } catch (const std::exception &e) { Logger("rugnux").Error("{}", e.what());