otsdaq-mu2e-stm  5.02.01
cpu_utils.cc
1 #include <thread>
2 
3 // Include CPU utils header
4 #include "Mu2e-STMDAQ/utils/async_logger.hh"
5 #include "Mu2e-STMDAQ/utils/cpu_utils.hh"
6 
7 // Define static members for singleton management
8 std::shared_ptr<cpu_utils> cpu_utils::instance = nullptr;
9 std::once_flag cpu_utils::init_flag;
10 
11 // Allowed cores for NUMA socket 0
12 const std::vector<int> cpu_utils::socket0_cores = {
13  0, 1, 2, 3, 4, 5, 6, 7, 8, 9, 10, 11, 12, 13, 14, 15, 16, 17, 18, 19, 20, 21, 22, 23};
14 
15 // Allowed cores for NUMA socket 1
16 const std::vector<int> cpu_utils::socket1_cores = {24, 25, 26, 27, 28, 29, 30, 31,
17  32, 33, 34, 35, 36, 37, 38, 39,
18  40, 41, 42, 43, 44, 45, 46, 47};
19 
20 // Allow cpu to log once logger exists
21 void cpu_utils::log(const std::string& msg, int level)
22 {
23  std::lock_guard<std::mutex> lock(log_mutex);
24  if(logger)
25  logger->log(msg, level);
26  else
27  pending_logs.emplace_back(msg, level);
28 }
29 
30 void cpu_utils::set_logger(std::shared_ptr<AsyncLogger> l)
31 {
32  std::lock_guard<std::mutex> lock(log_mutex);
33  logger = l;
34  for(const auto& [msg, level] : pending_logs)
35  logger->log(msg, level);
36  pending_logs.clear();
37 }
38 
39 // Returns the singleton instance; initializes on first call
40 std::shared_ptr<cpu_utils> cpu_utils::getInstance(const Config& cfg, CpuRole role)
41 {
42  std::call_once(init_flag, [&]() {
43  instance = std::shared_ptr<cpu_utils>(new cpu_utils(cfg, role));
44  });
45 
46  return instance;
47 }
48 
49 // Constructor: sets up configuration, starting core, and detects hardware concurrency
50 cpu_utils::cpu_utils(const Config& cfg_, CpuRole role_)
51  : cfg(cfg_), role(role_), next_core_id(0)
52 {
53  // Query total number of CPU cores
54  max_cores = std::thread::hardware_concurrency();
55 
56  // Fallback if detection fails
57  if(max_cores == 0)
58  max_cores = 1;
59 
60  // Set to NUMA node socket
61  if(role == CpuRole::Standalone)
62  {
63  socket_id =
64  EnvVars::expand("${HOSTNAME}") == cfg.getValue<std::string>("stm.ch0_host")
65  ? 1
66  : 0;
67  starting_core_id = cfg.getValue<int>("stm.stmdaq_starting_core");
68  }
69  else
70  {
71  socket_id =
72  EnvVars::expand("${HOSTNAME}") == cfg.getValue<std::string>("stm.ch0_host")
73  ? 0
74  : 1;
75  starting_core_id = cfg.getValue<int>("stm.artdaq_starting_core");
76  }
77 
78  allowed_cores = socket_id == 0 ? socket0_cores : socket1_cores;
79  //allowed_cores = get_socket_cores(socket_id);
80 
81  // Log initialization summary
82  std::string role_name = role == CpuRole::Standalone ? "Standalone" : "ArtDAQ";
83  log("CPU utils initialised: Role: " + role_name +
84  ", Socket: " + std::to_string(socket_id) +
85  ", Starting core offset: " + std::to_string(starting_core_id) +
86  ", Managed cores: " + std::to_string(allowed_cores.size()),
87  1);
88 }
89 
90 // Main interface: assigns a core to the calling thread and returns the core ID
91 size_t cpu_utils::get_next_core(const std::string& name)
92 {
93  // Atomically increment counter and compute core ID using modulo wraparound
94  const size_t local_id = next_core_id.fetch_add(1);
95  const size_t logical_core = (starting_core_id + local_id) % allowed_cores.size();
96  const size_t core_id = allowed_cores[logical_core];
97 
98  {
99  // Lock tracking data structure
100  std::lock_guard<std::mutex> lock(tracking_mutex);
101 
102  // Warn if this core has already been used
103  if(used_cores.find(core_id) != used_cores.end())
104  {
105  log("CPU utils: Warning! Core " + std::to_string(core_id) +
106  " has already been assigned.",
107  2);
108  }
109  else
110  {
111  used_cores.insert(core_id);
112  }
113 
114  // Warn if we've assigned more threads than available cores
115  if(used_cores.size() > allowed_cores.size())
116  {
117  log("CPU utils: Warning! Thread count exceeds available cores on socket.", 2);
118  }
119  }
120 
121  // Actually pin the calling thread to the computed core
122  pin_thread_to_core(core_id, name);
123 
124  return core_id;
125 }
126 
127 // Low-level function to pin thread to specific CPU core
128 void cpu_utils::pin_thread_to_core(size_t core_id, const std::string& name)
129 {
130  // Create and configure a CPU set
131  cpu_set_t cpuset;
132  // Clear the set
133  CPU_ZERO(&cpuset);
134  // Add the specified core to the set
135  CPU_SET(core_id, &cpuset);
136 
137  // Attempt to apply CPU affinity
138  const int result = pthread_setaffinity_np(pthread_self(), sizeof(cpu_set_t), &cpuset);
139 
140  // Log result to user
141  if(result != 0)
142  {
143  log("CPU utils: Error! Failed to pin " + name + " to core " +
144  std::to_string(core_id) + " (errno " + std::to_string(result) + ")",
145  0);
146  }
147  else
148  {
149  log("CPU utils: Pinned " + name + " to core " + std::to_string(core_id), 1);
150  }
151 }
Definition: config.hh:19