otsdaq-mu2e-stm  5.02.01
dqm_manager.cc
1 #include <atomic>
2 #include <chrono>
3 #include <iostream>
4 #include <sstream>
5 #include <thread>
6 
7 // Include dqm manager code
8 #include "Mu2e-STMDAQ/processing/dqm_manager.hh"
9 // Include dqm structs code
10 #include "Mu2e-STMDAQ/dqm/dqm_structs.hh"
11 // Operations manager header
12 #include "Mu2e-STMDAQ/processing/operation_manager.hh"
13 // Include UDP code
14 #include "Mu2e-STMDAQ/processing/udp.hh"
15 
16 // Constructor
17 DQM::DQM(const Config& cfg_,
18  const std::shared_ptr<AsyncLogger>& logger_,
19  const std::shared_ptr<STMdata>& stm_,
20  const std::shared_ptr<SignalHandler>& signal_,
21  OperationManager* op_man_)
22  : cfg(cfg_)
23  , logger(logger_)
24  , stm(stm_)
25  , signal(signal_)
26  , op_man(op_man_)
27  , baseline_mean_prev(stm->baseline_config.prev_mean)
28  , baseline_sigma_prev(stm->baseline_config.prev_sigma)
29  , baseline_nbins(stm->baseline_config.hist_bin_num)
30  , noise_len(stm->buffer_config.baseline_len)
31  , raw_len(stm->dqm_config.raw_len)
32  , peak_nbins(stm->dqm_config.peak_nbins)
33  , min_ph(stm->dqm_config.min_ph)
34  , inv_wid(stm->dqm_config.inv_bin_wid)
35  , peak_hist_window(peak_nbins, 0)
36  , peak_hist_all(peak_nbins, 0)
37  , window_size_buffers(stm->baseline_config.hist_window_buffers)
38  , per_buffer_hists(window_size_buffers, std::vector<uint64_t>(peak_nbins, 0))
39  , num_dropped_packets(0)
40  , op_num(0)
41  , running_(false)
42 {
43  size_t baseline_bytes =
44  sizeof(dqm_data_baseline) +
45  ((baseline_nbins - 1) * sizeof(uint64_t)) + // -1 for array initialised with 1
46  (noise_len * sizeof(int16_t)) + 7 + // Alignment padding
47  sizeof(uint64_t); // End counter
48  size_t raw_bytes =
49  sizeof(dqm_data_raw) +
50  ((raw_len - 1) * sizeof(int16_t)) + // -1 for array initialised with size 1
51  7 + // Alignment padding
52  sizeof(uint64_t); // End counter
53  size_t peak_bytes =
54  sizeof(dqm_data_peak) +
55  (((2 * peak_nbins) - 1) *
56  sizeof(uint64_t)) + // Two histograms -1 for array initialised with size 1
57  sizeof(uint64_t); // End counter
58 
59  // Register shared memory blocks
60  registerBlock<dqm_data_raw>(DQMPageType::RAW, "/dqm_raw_data", raw_bytes);
61  registerBlock<dqm_data_baseline>(
62  DQMPageType::BASELINE, "/dqm_adc_baseline", baseline_bytes);
63  registerBlock<dqm_data_peak>(DQMPageType::PEAKS, "/dqm_peak_data", peak_bytes);
64 
65  // Register operations for OperationManager
66  register_operation("update_dqm", [this](auto& b) { update_dqm(b); });
67 }
68 
69 void DQM::init_shm()
70 {
71  // Op num only known after op_manager has run
72  this->op_num = op_man->getUseOps().size() - 1;
73 
74  size_t daq_bytes =
75  sizeof(dqm_data_daq) +
76  ((op_num - 1) * sizeof(op_entry)) + // -1 for array initialised with size 1
77  7 + // Alignment padding
78  sizeof(uint64_t); // End counter
79 
80  registerBlock<dqm_data_daq>(DQMPageType::SPEEDS, "/dqm_daq_data", daq_bytes);
81 }
82 
83 auto t0 = std::chrono::steady_clock::now();
84 auto prev_baseline = t0;
85 auto prev_raw = t0;
86 auto prev_peak = t0;
87 auto prev_cpu = t0;
88 auto prev_alarm = t0;
89 
90 constexpr auto baseline_period = std::chrono::milliseconds{200};
91 constexpr auto raw_period = std::chrono::milliseconds{200};
92 constexpr auto peak_period = std::chrono::milliseconds{200};
93 constexpr auto cpu_period = std::chrono::milliseconds{1000};
94 constexpr auto alarm_period = std::chrono::milliseconds{200};
95 
96 // Main update loop: runs once per second
97 void DQM::update_dqm(std::shared_ptr<DataStruct>& buffer)
98 {
99  // Do the buffer by buffer operations
100  alarm_info(buffer);
101  histogram_peaks(buffer);
102 
103  // Get the time now
104  auto now = std::chrono::steady_clock::now();
105 
106  // Update baseline DQM (if baseline period has passed)
107  if(now - prev_baseline >= baseline_period)
108  {
109  // Update baseline DQM
110  update_dqm_baseline(buffer);
111  // Store last baseline update time
112  prev_baseline = now;
113  }
114 
115  // Update raw data DQM (if raw data period has passed)
116  if(now - prev_raw >= raw_period)
117  {
118  // Update raw data DQM
119  update_dqm_raw(buffer);
120  // Store last raw data update time
121  prev_raw = now;
122  }
123 
124  // Update peak data DQM (if peak data period has passed)
125  if(now - prev_peak >= peak_period)
126  {
127  // Update peak data DQM
128  update_dqm_peak(buffer);
129  // Store last peak data update time
130  prev_peak = now;
131  }
132 
133  // Update cpu performance DQM (if cpu period has passed)
134  if(now - prev_cpu >= cpu_period)
135  {
136  // Update cpu performance DQM
137  update_cpu_performance(buffer);
138  // Store last raw data update time
139  prev_cpu = now;
140  }
141 }
142 
143 void DQM::alarm_info(std::shared_ptr<DataStruct>& buffer)
144 {
145  // Check for dropped packets (want to check all buffers)
146  if(buffer->dropped_packet_count > 0)
147  {
148  num_dropped_packets += buffer->dropped_packet_count;
149  }
150 }
151 
152 // Fill histogram of peaks in every buffer
153 void DQM::histogram_peaks(std::shared_ptr<DataStruct>& buffer)
154 {
155  // Get updated histogram of peak heights
156  // Remove the oldest buffer from the window histogram
157  for(int b = 0; b < peak_nbins; b++)
158  {
159  peak_hist_window[b] -= per_buffer_hists[window_index][b];
160  total_window -= per_buffer_hists[window_index][b];
161  per_buffer_hists[window_index][b] = 0; // reset this buffer histogram
162  }
163 
164  // Loop through EWTs and peaks
165  for(const auto& ewt : buffer->EWTs)
166  {
167  for(const auto& pulse : ewt.ph)
168  {
169  // Get pulse height
170  int16_t height = pulse.height;
171 
172  if(height < min_ph || height > 0)
173  continue;
174 
175  // Find bin
176  int bin = int((height - min_ph) * inv_wid);
177 
178  // Add to histograms
179  ++per_buffer_hists[window_index][bin];
180  ++peak_hist_window[bin];
181  ++peak_hist_all[bin];
182  }
183  }
184 
185  // Advance circular index
186  ++window_index;
187  if(window_index == window_size_buffers)
188  window_index = 0;
189 }
190 
191 // Update the DQM baseline
192 void DQM::update_dqm_baseline(std::shared_ptr<DataStruct>& buffer)
193 {
194  // Get the DQM baseline shm block
195  auto* baseline = get<dqm_data_baseline>(DQMPageType::BASELINE);
196 
197  // Begin safe write of shm
198  baseline->gen_start.fetch_add(1, std::memory_order_release);
199 
200  // Update baseline values
201  baseline->timestamp_ns = get_current_time_ns();
202  baseline->baseline_mean_prev = baseline_mean_prev;
203  baseline->baseline_sigma_prev = baseline_sigma_prev;
204  baseline->baseline_mean_avg = buffer->baseline_mean_avg;
205  baseline->baseline_sigma_avg = buffer->baseline_sigma_avg;
206  baseline->baseline_mean_current = buffer->baseline_mean_current;
207  baseline->baseline_sigma_current = buffer->baseline_sigma_current;
208  baseline->baseline_fit_mean_all = buffer->baseline_all.mu0;
209  baseline->baseline_fit_sigma_all = buffer->baseline_all.sigma0;
210  baseline->baseline_fit_mean_wind = buffer->baseline_window.mu0;
211  baseline->baseline_fit_sigma_wind = buffer->baseline_window.sigma0;
212  baseline->baseline_bins = baseline_nbins;
213  baseline->noise_len = noise_len;
214 
215  // Copy hist data
216  uint64_t* hist_dest = reinterpret_cast<uint64_t*>(baseline->baseline_data);
217  std::memcpy(
218  hist_dest, buffer->hist_counts_all.data(), baseline_nbins * sizeof(uint64_t));
219 
220  // Copy noise data
221  size_t copy_len = std::min(buffer->noise_data.size(), static_cast<size_t>(noise_len));
222  int16_t* noise_dest = reinterpret_cast<int16_t*>(baseline->baseline_data +
223  (baseline_nbins * sizeof(uint64_t)));
224  std::memcpy(noise_dest, buffer->noise_data.data(), copy_len * sizeof(int16_t));
225  // Fill rest with zeroes
226  if(copy_len < noise_len)
227  std::memset(noise_dest + copy_len, 0, (noise_len - copy_len) * sizeof(int16_t));
228 
229  // Align with 8 byte for gen end write
230  uintptr_t noise_end_addr =
231  reinterpret_cast<uintptr_t>(noise_dest) + (noise_len * sizeof(int16_t));
232  uintptr_t aligned_gen_end_addr = (noise_end_addr + 7) & ~7;
233  auto* gen_end_ptr = reinterpret_cast<std::atomic<uint64_t>*>(aligned_gen_end_addr);
234 
235  // End safe write of shm
236  gen_end_ptr->store(baseline->gen_start.load(std::memory_order_acquire),
237  std::memory_order_release);
238 }
239 
240 // Update raw data display
241 void DQM::update_dqm_raw(std::shared_ptr<DataStruct>& buffer)
242 {
243  // Get raw DQM SHM block
244  auto* raw = get<dqm_data_raw>(DQMPageType::RAW);
245 
246  // Write to shared memory safely
247  raw->gen_start.fetch_add(1, std::memory_order_release);
248 
249  // Timestamp and rate
250  raw->timestamp_ns = get_current_time_ns(); // Implement this helper
251 
252  raw->raw_len = raw_len;
253 
254  // Copy raw data
255  int16_t* raw_dest = reinterpret_cast<int16_t*>(raw->raw_data);
256  size_t copy_len = std::min(buffer->raw_len, static_cast<size_t>(raw_len));
257  std::memcpy(raw_dest, buffer->raw.data(), copy_len * sizeof(int16_t));
258  // Fill rest with zeroes
259  if(copy_len < raw_len)
260  std::memset(raw_dest + copy_len, 0, (raw_len - copy_len) * sizeof(int16_t));
261 
262  // Align with 8 byte for gen end write
263  uintptr_t raw_end_addr =
264  reinterpret_cast<uintptr_t>(raw_dest) + (raw_len * sizeof(int16_t));
265  uintptr_t aligned_gen_end_addr = (raw_end_addr + 7) & ~7;
266  auto* gen_end_ptr = reinterpret_cast<std::atomic<uint64_t>*>(aligned_gen_end_addr);
267 
268  // Complete safe write
269  gen_end_ptr->store(raw->gen_start.load(std::memory_order_acquire),
270  std::memory_order_release);
271 }
272 
273 // Update peak peak display
274 void DQM::update_dqm_peak(std::shared_ptr<DataStruct>& buffer)
275 {
276  // Get peak DQM SHM block
277  auto* peaks = get<dqm_data_peak>(DQMPageType::PEAKS);
278 
279  // Write to shared memory safely
280  peaks->gen_start.fetch_add(1, std::memory_order_release);
281 
282  // Timestamp and rate
283  peaks->timestamp_ns = get_current_time_ns(); // Implement this helper
284 
285  peaks->peak_nbins = peak_nbins;
286 
287  // Copy data
288  std::memcpy(peaks->peak_data, peak_hist_all.data(), peak_nbins * sizeof(uint64_t));
289 
290  std::memcpy(peaks->peak_data + peak_nbins,
291  peak_hist_window.data(),
292  peak_nbins * sizeof(uint64_t));
293 
294  // Find gen end write location
295  auto* gen_end_ptr =
296  reinterpret_cast<std::atomic<uint64_t>*>(peaks->peak_data + 2 * peak_nbins);
297 
298  // Complete safe write
299  gen_end_ptr->store(peaks->gen_start.load(std::memory_order_acquire),
300  std::memory_order_release);
301 
302  //uint64_t* gen_end_ptr = peaks->peak_data + (2 * peak_nbins);
303  //uint64_t gen_end = peaks->gen_start.load(std::memory_order_acquire);
304  //__atomic_store_n(gen_end_ptr, gen_end, __ATOMIC_RELEASE);
305 }
306 
307 // Update CPU performace DQM
308 void DQM::update_cpu_performance(std::shared_ptr<DataStruct>& buffer)
309 {
310  auto daq_perf = get<dqm_data_daq>(DQMPageType::SPEEDS);
311 
312  // Write to shared memory safely
313  daq_perf->gen_start.fetch_add(1, std::memory_order_release);
314 
315  // Timestamp
316  daq_perf->timestamp_ns = get_current_time_ns();
317 
318  // Dropped packets counter
319  daq_perf->num_dropped_packets = num_dropped_packets;
320 
321  // Get ADC temp
322  double temp = 60.0;
323  daq_perf->adc_temp = temp;
324 
325  // Number of operations
326  daq_perf->num_ops = op_num;
327 
328  // Start print line
329  std::ostringstream oss;
330  oss << "Operation speed (Gbit/s): ";
331 
332  bool first_entry = true;
333  for(size_t i = 0; i < buffer->cpu_performance.size(); ++i)
334  {
335  const auto& name = buffer->cpu_performance[i].first;
336  double perf = buffer->cpu_performance[i].second;
337 
338  // Print
339  if(name.find("DQM") != std::string::npos)
340  continue;
341 
342  if(!first_entry)
343  oss << " | ";
344 
345  oss << name << "=" << std::fixed << std::setprecision(2) << perf;
346  first_entry = false;
347 
348  // Copy name up to 31 characters with null end
349  std::memset(daq_perf->ops_data[i].name, 0, 32);
350  std::strncpy(daq_perf->ops_data[i].name, name.c_str(), 31);
351 
352  // Copy speed
353  daq_perf->ops_data[i].speed = perf;
354  }
355 
356  logger->log(oss.str(), 1);
357 
358  uintptr_t data_end_addr =
359  reinterpret_cast<uintptr_t>(daq_perf->ops_data) + (op_num * sizeof(op_entry));
360  uintptr_t aligned_gen_end_addr = (data_end_addr + 7) & ~7;
361  auto* gen_end_ptr = reinterpret_cast<std::atomic<uint64_t>*>(aligned_gen_end_addr);
362 
363  gen_end_ptr->store(daq_perf->gen_start.load(std::memory_order_acquire),
364  std::memory_order_release);
365 }
366 
367 uint64_t DQM::get_current_time_ns()
368 {
369  return std::chrono::duration_cast<std::chrono::nanoseconds>(
370  std::chrono::system_clock::now().time_since_epoch())
371  .count();
372 }
Definition: config.hh:19
Definition: peaks.hh:6
Definition: dqm_structs.hh:53