otsdaq-mu2e-stm  5.02.01
cpu_utils.hh
1 #ifndef CPU_UTILS_HH
2 #define CPU_UTILS_HH
3 
4 #include <pthread.h>
5 #include <sched.h>
6 
7 #include <atomic>
8 #include <iostream>
9 #include <memory>
10 #include <mutex>
11 #include <string>
12 #include <unordered_set>
13 #include <vector>
14 
15 // Config code
16 #include "Mu2e-STMDAQ/config/config.hh"
17 
18 // Forward declaration
19 class AsyncLogger;
20 
21 // NUMA socket identifier
22 enum class Socket {
23  Socket0,
24  Socket1
25 };
26 
27 // Application role used to select the appropriate CPU affinity configuration from the shared XML
28 enum class CpuRole {
29  Standalone,
30  ArtDAQ
31 };
32 
33 // Singleton utility class for managing CPU thread pinning
34 class cpu_utils {
35 
36 private:
37 
38  // Private constructor; only accessible by getInstance()
39  explicit cpu_utils(const Config& cfg_, CpuRole role_);
40 
41  // Low-level function to pin thread to a specified core
42  void pin_thread_to_core(size_t core_id, const std::string& name);
43 
44  // Store reference to the Config instance
45  const Config& cfg;
46 
47  // Application type (Standalone or ArtDAQ) Used to determine which CPU affinity settings should be loaded from the XML configuration
48  CpuRole role;
49 
50  // First logical core to assign from within the configured socket
51  size_t starting_core_id;
52 
53  // Total number of available hardware cores reported by the operating system
54  size_t max_cores;
55 
56  // Atomic counter for round-robin assignment of logical cores within the selected socket
57  std::atomic<size_t> next_core_id;
58 
59  // Mutex to protect access to used_cores set
60  std::mutex tracking_mutex;
61 
62  // Set of all physical CPU cores that have been assigned so far
63  std::unordered_set<size_t> used_cores;
64 
65  // NUMA socket assigned to this application
66  //Socket socket_id;
67  int socket_id;
68 
69  // Physical CPU IDs available to the configured NUMA socket
70  std::vector<int> allowed_cores;
71 
72  // Get logger with mutex to protect setting/appending
73  std::mutex log_mutex;
74  std::shared_ptr<AsyncLogger> logger;
75  std::vector<std::pair<std::string,int>> pending_logs;
76 
77  // Singleton instance pointer
78  static std::shared_ptr<cpu_utils> instance;
79 
80  // Flag to ensure singleton is initialized only once
81  static std::once_flag init_flag;
82 
83  // Verified physical CPU mappings for NUMA socket 0
84  static const std::vector<int> socket0_cores;
85 
86  // Verified physical CPU mappings for NUMA socket 1
87  static const std::vector<int> socket1_cores;
88 
89  void log(const std::string& msg, int level);
90 
91 public:
92 
93  // Delete copy and move constructors/operators to enforce singleton pattern
94  cpu_utils(const cpu_utils&) = delete;
95  cpu_utils& operator=(const cpu_utils&) = delete;
96  cpu_utils(cpu_utils&&) = delete;
97  cpu_utils& operator=(cpu_utils&&) = delete;
98 
99  // Static method to access the singleton instance and initialise the CPU affinity manager for the specified application role
100  static std::shared_ptr<cpu_utils> getInstance(const Config& cfg, CpuRole role);
101 
102  // Assigns the calling thread to the next available CPU core within the configured NUMA socket and returns the physical CPU core ID
103  size_t get_next_core(const std::string& name);
104 
105  void set_logger(std::shared_ptr<AsyncLogger> l);
106 
107  // Clear logger for graceful run transitions
108  void clear_logger() {
109  std::lock_guard<std::mutex> lock(log_mutex);
110  logger = nullptr;
111  pending_logs.clear();
112  }
113 
114  // Get the core checkpoint
115  size_t checkpoint() {
116  return next_core_id.load();
117  }
118 
119  // After we close operations, remove those cores from list
120  void reset_to(size_t checkpoint) {
121  std::lock_guard<std::mutex> lock(tracking_mutex);
122  size_t current = next_core_id.load();
123  for (size_t i = checkpoint; i < current; i++) {
124  size_t logical_core = (starting_core_id + i) % allowed_cores.size();
125  size_t core_id = allowed_cores[logical_core];
126  used_cores.erase(core_id);
127  }
128  next_core_id.store(checkpoint);
129  }
130 
131  // Destructor
132  ~cpu_utils() {
133  log("CPU utils shutting down. Cores assigned: " +
134  std::to_string(next_core_id.load()) +
135  ", unique cores used: " +
136  std::to_string(used_cores.size()),1);
137 
138  //std::cout << "cpu_utils destructor called.\n";
139  }
140 
141 };
142 
143 #endif // CPU_UTILS_HH
Definition: config.hh:19