Files
Jungfraujoch/common/ThreadAffinity.cpp
leonarski_fandClaude Opus 5.5 50006eff88 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) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01SVmAWnzCmRKAXVUCdc4iNi
2026-10-08 11:15:29 +02:00

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