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) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01SVmAWnzCmRKAXVUCdc4iNi
160 lines
5.4 KiB
C++
160 lines
5.4 KiB
C++
// SPDX-FileCopyrightText: 2026 Filip Leonarski, Paul Scherrer Institute <filip.leonarski@psi.ch>
|
|
// SPDX-License-Identifier: GPL-3.0-only
|
|
|
|
#include "ThreadAffinity.h"
|
|
|
|
#ifdef __linux__
|
|
|
|
#include <algorithm>
|
|
#include <cmath>
|
|
#include <cstdio>
|
|
#include <filesystem>
|
|
#include <fstream>
|
|
#include <sstream>
|
|
|
|
#include <pthread.h>
|
|
#include <sched.h>
|
|
|
|
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();
|
|
|
|
// 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;
|
|
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;
|
|
}
|
|
|
|
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;
|
|
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
|
|
|
|
#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() {}
|
|
|
|
#endif
|