LDMX Software
eventbuilder::EventBuilder Class Reference

Public Member Functions

 EventBuilder (const std::string &name, framework::Process &proc)
 
 enableLogging ("EventBuilder")
 
void configure (framework::config::Parameters &ps) override
 Callback for the EventProcessor to configure itself from the given set of parameters.
 
void produce (framework::Event &event) override
 Process the event and put new data products into it.
 
- Public Member Functions inherited from framework::Producer
 Producer (const std::string &name, Process &process)
 Class constructor.
 
virtual void process (Event &event) final
 Processing an event for a Producer is calling produce.
 
- Public Member Functions inherited from framework::EventProcessor
 DECLARE_FACTORY (EventProcessor, EventProcessor *, const std::string &, Process &)
 declare that we have a factory for this class
 
 EventProcessor (const std::string &name, Process &process)
 Class constructor.
 
virtual ~EventProcessor ()=default
 Class destructor.
 
virtual void beforeNewRun (ldmx::RunHeader &run_header)
 Callback for Producers to add parameters to the run header before conditions are initialized.
 
virtual void onNewRun (const ldmx::RunHeader &run_header)
 Callback for the EventProcessor to take any necessary action when the run being processed changes.
 
virtual void onFileOpen (EventFile &event_file)
 Callback for the EventProcessor to take any necessary action when a new event input ROOT file is opened.
 
virtual void onFileClose (EventFile &event_file)
 Callback for the EventProcessor to take any necessary action when a event input ROOT file is closed.
 
virtual void onProcessStart ()
 Callback for the EventProcessor to take any necessary action when the processing of events starts, such as creating histograms.
 
virtual void onProcessEnd ()
 Callback for the EventProcessor to take any necessary action when the processing of events finishes, such as calculating job-summary quantities.
 
template<class T >
const T & getCondition (const std::string &condition_name)
 Access a conditions object for the current event.
 
TDirectory * getHistoDirectory ()
 Access/create a directory in the histogram file for this event processor to create histograms and analysis tuples.
 
void setStorageHint (framework::StorageControl::Hint hint)
 Mark the current event as having the given storage control hint from this module_.
 
void setStorageHint (framework::StorageControl::Hint hint, const std::string &purposeString)
 Mark the current event as having the given storage control hint from this module and the given purpose string.
 
int getLogFrequency () const
 Get the current logging frequency from the process.
 
int getRunNumber () const
 Get the run number from the process.
 
std::string getName () const
 Get the processor name.
 
void createHistograms (const std::vector< framework::config::Parameters > &histos)
 Internal function which is used to create histograms passed from the python configuration @parma histos vector of Parameters that configure histograms to create.
 

Private Member Functions

PhysicsEventData assemblePayload (const std::vector< DataFragment > &fragments)
 
void writeEventBinary (const PhysicsEventData &ev, const std::string &path="events.bin")
 

Private Attributes

bool m_verbose_parse_
 
unsigned int m_event_id_
 
FragmentBuffer m_event_buffer_
 
std::string m_input_file_
 
packing::utility::Reader m_reader_
 
std::string m_output_name_ {"BuilderOutput"}
 
long long m_coherence_window_ns_
 
int m_min_subsystems_ {2}
 
std::chrono::steady_clock::time_point m_start_time_
 
std::chrono::steady_clock::time_point m_event_start_time_
 
uint64_t m_total_bytes_read_ {0}
 
uint64_t m_total_events_built_ {0}
 
unsigned long m_events_since_last_report_ {0}
 
std::chrono::steady_clock::time_point m_window_start_time_
 
uint64_t m_window_events_count_ {0}
 
uint64_t m_window_bytes_read_ {0}
 

Static Private Attributes

static constexpr unsigned int WINDOW_SIZE
 

Additional Inherited Members

- Protected Member Functions inherited from framework::EventProcessor
void abortEvent ()
 Abort the event immediately.
 
- Protected Attributes inherited from framework::EventProcessor
HistogramPool histograms_
 helper object for making and filling histograms
 
NtupleManager & ntuple_ {NtupleManager::getInstance()}
 Manager for any ntuples.
 
logging::logger the_log_
 The logger for this EventProcessor.
 

Detailed Description

Definition at line 14 of file EventBuilder.h.

Constructor & Destructor Documentation

◆ EventBuilder()

eventbuilder::EventBuilder::EventBuilder ( const std::string & name,
framework::Process & proc )
inline

Definition at line 16 of file EventBuilder.h.

17 : framework::Producer(name, proc),
18 m_verbose_parse_(false),
19 m_event_id_(0) {}
Base class for a module which produces a data product.

Member Function Documentation

◆ assemblePayload()

PhysicsEventData EventBuilder::assemblePayload ( const std::vector< DataFragment > & fragments)
private

Definition at line 521 of file EventBuilder.cxx.

522 {
523 PhysicsEventData event_data;
524 if (fragments.empty()) return event_data;
525 event_data.event_id_ = m_event_id_;
526 event_data.timestamp_ = fragments.front().header_.timestamp_;
527 for (const auto& fragment : fragments) {
529 g.subsystem_id_ = fragment.header_.subsystem_id_;
530 g.timestamp_ns_ = fragment.header_.timestamp_;
531 g.data_ = fragment.payload_;
532 if (fragment.trailer_.checksum_ != 0)
533 g.checksum_ = fragment.trailer_.checksum_;
534 else
535 g.checksum_ = crc32(fragment.payload_);
536 event_data.blocks_.push_back(std::move(g));
537 event_data.systems_readout_.push_back(fragment.header_.subsystem_id_);
538 }
539 return event_data;
540}

◆ configure()

void EventBuilder::configure ( framework::config::Parameters & parameters)
overridevirtual

Callback for the EventProcessor to configure itself from the given set of parameters.

The parameters a processor has access to are the member variables of the python class in the sequence that has class_name equal to the EventProcessor class name.

For an example, look at MyProcessor.

Parameters
parametersParameters for configuration.

Reimplemented from framework::EventProcessor.

Definition at line 57 of file EventBuilder.cxx.

57 {
58 if (ps.exists("verbose_parse")) {
59 m_verbose_parse_ = ps.get<bool>("verbose_parse");
60 }
61 if (ps.exists("dat_file")) {
62 m_input_file_ = ps.get<std::string>("dat_file");
63 } else {
64 const char* env_input = std::getenv("EVENTBUILDER_INPUT");
65 if (env_input) m_input_file_ = env_input;
66 }
67 if (ps.exists("output_name"))
68 m_output_name_ = ps.get<std::string>("output_name");
69 if (ps.exists("coherence_window_ns")) {
70 m_coherence_window_ns_ = ps.get<double>("coherence_window_ns");
71 }
72 if (ps.exists("min_subsystems")) {
73 m_min_subsystems_ = ps.get<int>("min_subsystems");
74 }
75
76 ldmx_log(info) << "configure(): dat_file='" << m_input_file_
77 << "' output_name='" << m_output_name_
78 << "' verbose_parse=" << (m_verbose_parse_ ? "true" : "false")
79 << " coherence_window_ns=" << m_coherence_window_ns_
80 << " min_subsystems=" << m_min_subsystems_;
81
82 if (!m_input_file_.empty()) {
83 ldmx_log(info) << "configure(): opening file '" << m_input_file_ << "'";
84 m_reader_.open(m_input_file_);
85 if (!m_reader_) {
86 ldmx_log(error) << "failed to open input file: '" << m_input_file_ << "'";
87 } else {
88 ldmx_log(info) << "configure(): file opened successfully";
89 }
90 } else {
91 ldmx_log(error) << "no input file specified";
92 }
93
94 // Initialize performance tracking
95 m_start_time_ = std::chrono::steady_clock::now();
96 m_total_bytes_read_ = 0;
97 m_total_events_built_ = 0;
98 m_events_since_last_report_ = 0;
99
100 // Initialize windowed metrics for real-time DAQ monitoring
101 m_window_start_time_ = m_start_time_;
102 m_window_events_count_ = 0;
103 m_window_bytes_read_ = 0;
104}
void open(const std::string &file_name)
Open a file with this reader.
Definition Reader.h:36

References framework::config::Parameters::exists(), framework::config::Parameters::get(), and packing::utility::Reader::open().

◆ produce()

void EventBuilder::produce ( framework::Event & event)
overridevirtual

Process the event and put new data products into it.

Parameters
eventThe Event to process.

Implements framework::Producer.

Definition at line 106 of file EventBuilder.cxx.

106 {
107 static int produce_call_count = 0;
108 produce_call_count++;
109 ldmx_log(debug) << "produce() called, count=" << produce_call_count;
110 ldmx_log(debug) << "verbose_parse=" << (m_verbose_parse_ ? "true" : "false");
111 if (m_verbose_parse_ || produce_call_count <= 3) {
112 ldmx_log(info) << "produce() call #" << produce_call_count;
113 }
114
115 // Start timing for this event
116 m_event_start_time_ = std::chrono::steady_clock::now();
117
118 // Track errors encountered during event assembly
119 uint32_t current_event_errors = 0;
120
121 while (m_reader_ && !m_reader_.eof()) {
122 // Try to read a RogueFrameHeader
123 packing::RogueFrameHeader frame_header;
124 frame_header.read(m_reader_);
125
126 // Store location of end-of-frame for potential recovery
127 const long long frame_end =
128 static_cast<long long>(m_reader_.tell()) + frame_header.size();
129
130 // Check if this is a data frame (channel 0) and not YAML
131 if (frame_header.channel() != 0 || frame_header.probablyYaml()) {
132 // Skip this frame
133 if (m_verbose_parse_)
134 ldmx_log(debug) << "skipping non-data frame (channel="
135 << frame_header.channel() << ")";
136 m_reader_.seek(frame_end);
137 continue;
138 }
139
140 // We have a valid data frame - try to parse it
141 DataFragment fragment{{}, {}, {.checksum_ = 0}};
142 bool parsed_ok = false;
143
144 // Try RoR format first
145 packing::LDMXRoRHeader ror_header;
146 int pos_before_ror = m_reader_.tell();
147
148 // Attempt to read as RoR header
149 try {
150 ror_header.read(m_reader_);
151 // Successfully parsed RoR header
152 fragment.header_.subsystem_id_ =
153 static_cast<uint64_t>(ror_header.subsystem());
154 fragment.header_.contributor_id_ =
155 static_cast<uint64_t>(ror_header.contributor());
156 fragment.header_.timestamp_ = ror_header.timestamp();
157
158 if (m_verbose_parse_) {
159 ldmx_log(debug) << "parsed RoR header subsys="
160 << (int)ror_header.subsystem()
161 << " contrib=" << (int)ror_header.contributor()
162 << " ts=" << ror_header.timestamp();
163 }
164
165 // Read remaining frame payload after RoR header
166 std::vector<uint8_t> payload_data;
167 long long payload_size = frame_end - m_reader_.tell();
168 if (payload_size > 0) {
169 payload_data.resize(payload_size);
170 m_reader_.read(reinterpret_cast<char*>(payload_data.data()),
171 payload_size);
172 }
173 fragment.payload_ = std::move(payload_data);
174 parsed_ok = true;
175 } catch (...) {
176 // RoR parse failed, try packing subsystem format
177 if (m_verbose_parse_)
178 ldmx_log(debug)
179 << "RoR header parse failed, trying packing subsystem format";
180 current_event_errors |= ldmx::EventSummary::ERROR_PARSE_FAILURE;
181 m_reader_.seek(pos_before_ror);
182 }
183
184 if (!parsed_ok) {
185 // Try packing subsystem format
187 try {
188 pkt.read(m_reader_);
189 fragment.header_.subsystem_id_ = static_cast<uint64_t>(pkt.id());
190 fragment.header_.contributor_id_ =
191 0; // Not available in packing subsystem format
192 fragment.header_.timestamp_ = static_cast<uint64_t>(pkt.header()[1]) *
193 1000000000ULL; // event number to ns
194
195 // Convert packet data to payload bytes
196 const auto& data = pkt.data();
197 fragment.payload_.resize(data.size() * 4);
198 std::memcpy(fragment.payload_.data(),
199 reinterpret_cast<const char*>(data.data()),
200 data.size() * 4);
201
202 if (m_verbose_parse_) {
203 ldmx_log(debug) << "parsed packing subsystem pkt subsys=" << pkt.id()
204 << " data_size=" << data.size();
205 }
206 parsed_ok = true;
207 } catch (...) {
208 // Both parse attempts failed, skip this frame
209 if (m_verbose_parse_)
210 ldmx_log(debug) << "failed to parse as either format, skipping frame";
211 current_event_errors |= ldmx::EventSummary::ERROR_PARSE_FAILURE;
212 m_reader_.seek(frame_end);
213 continue;
214 }
215 }
216
217 if (!parsed_ok) {
218 m_reader_.seek(frame_end);
219 continue;
220 }
221
222 if (m_verbose_parse_)
223 ldmx_log(debug) << "adding fragment subsys="
224 << fragment.header_.subsystem_id_
225 << " ts=" << fragment.header_.timestamp_
226 << " bytes=" << fragment.payload_.size();
227
228 // Check if this fragment is outside the coherence window of the current
229 // event batch If so, it means we should finalize the previous event before
230 // adding this one
231 long long fragment_ts = fragment.header_.timestamp_;
232 long long buffer_ref_time = m_event_buffer_.getReferenceTime();
233
234 if (buffer_ref_time > 0 &&
235 (fragment_ts < buffer_ref_time - m_coherence_window_ns_ ||
236 fragment_ts > buffer_ref_time + m_coherence_window_ns_)) {
237 // This fragment is outside the current window - try to build the previous
238 // event
239 std::vector<DataFragment> assembled_event_fragments;
240 if (m_event_buffer_.tryBuildEvent(m_coherence_window_ns_,
241 m_min_subsystems_,
242 assembled_event_fragments) &&
243 !assembled_event_fragments.empty()) {
244 // Output the complete event and return
245 ++m_event_id_;
246 // Add each subsystem's raw data as vector<uint8_t> directly
247 uint64_t event_timestamp = 0;
248 for (const auto& frag : assembled_event_fragments) {
249 std::string subsys_name = packing::LDMXRoRHeader::getSubsystemName(
250 static_cast<uint8_t>(frag.header_.subsystem_id_),
251 static_cast<uint8_t>(frag.header_.contributor_id_));
252 std::vector<uint8_t> payload_bytes(frag.payload_.begin(),
253 frag.payload_.end());
254 event.add(subsys_name, payload_bytes);
255 if (event_timestamp == 0) {
256 event_timestamp = frag.header_.timestamp_;
257 }
258 }
259 event.getEventHeader().setIntParameter(
260 "RoR Timestamp", static_cast<long long>(event_timestamp));
261 // Still create PhysicsEventData for binary output file
262 PhysicsEventData final_event =
263 assemblePayload(assembled_event_fragments);
264 writeEventBinary(final_event, "events.bin");
265 ldmx_log(info) << "assembled event id=" << m_event_id_
266 << " timestamp=" << final_event.timestamp_
267 << " fragments=" << assembled_event_fragments.size()
268 << " systems=" << final_event.systems_readout_.size();
269
270 ldmx::EventSummary summary;
271 summary.setEventNumber(m_event_id_);
272 summary.setTimestampNs(static_cast<uint64_t>(final_event.timestamp_));
273 std::set<uint64_t> unique_sys;
274 uint64_t total_payload = 0;
275 for (const auto& f : assembled_event_fragments) {
276 unique_sys.insert(f.header_.subsystem_id_);
277 total_payload += f.payload_.size();
278 }
279 // Check for duplicate subsystems
280 if (unique_sys.size() != assembled_event_fragments.size()) {
281 current_event_errors |= ldmx::EventSummary::ERROR_DUPLICATE_SUBSYSTEM;
282 }
283 summary.setNSystems(static_cast<uint32_t>(unique_sys.size()));
284 summary.setSystemIds(
285 std::vector<uint64_t>(unique_sys.begin(), unique_sys.end()));
286 summary.setPayloadSize(total_payload);
287 summary.setErrorFlags(current_event_errors);
288 event.add("EventSummary", summary);
289
290 // Update performance metrics
291 m_total_bytes_read_ += total_payload;
292 m_total_events_built_++;
293 m_events_since_last_report_++;
294 std::chrono::steady_clock::time_point now =
295 std::chrono::steady_clock::now();
296 std::chrono::milliseconds event_build_time_ms =
297 std::chrono::duration_cast<std::chrono::milliseconds>(
298 now - m_event_start_time_);
299 std::chrono::seconds total_elapsed_s =
300 std::chrono::duration_cast<std::chrono::seconds>(now -
301 m_start_time_);
302
303 double events_per_sec =
304 (total_elapsed_s.count() > 0)
305 ? static_cast<double>(m_total_events_built_) /
306 total_elapsed_s.count()
307 : 0.0;
308 double mb_per_sec = (total_elapsed_s.count() > 0)
309 ? static_cast<double>(m_total_bytes_read_) /
310 (1024.0 * 1024.0) /
311 total_elapsed_s.count()
312 : 0.0;
313
314 // Update windowed metrics
315 m_window_events_count_++;
316 m_window_bytes_read_ += total_payload;
317 std::chrono::seconds window_elapsed_s =
318 std::chrono::duration_cast<std::chrono::seconds>(
319 now - m_window_start_time_);
320
321 // Reset window if we've processed WINDOW_SIZE events
322 double window_events_per_sec = 0.0;
323 double window_mb_per_sec = 0.0;
324 if (m_window_events_count_ >= WINDOW_SIZE) {
325 double window_time_sec =
326 std::max(1.0, static_cast<double>(window_elapsed_s.count()));
327 window_events_per_sec =
328 static_cast<double>(m_window_events_count_) / window_time_sec;
329 window_mb_per_sec = static_cast<double>(m_window_bytes_read_) /
330 (1024.0 * 1024.0) / window_time_sec;
331
332 // Reset window
333 m_window_start_time_ = now;
334 m_window_events_count_ = 0;
335 m_window_bytes_read_ = 0;
336 }
337
338 // Write performance metrics to CSV
339 long long event_build_time_ms_val = event_build_time_ms.count();
340 writePerformanceMetric(
341 m_event_id_, static_cast<double>(event_build_time_ms_val),
342 events_per_sec, mb_per_sec, window_events_per_sec,
343 window_mb_per_sec, m_total_events_built_, m_total_bytes_read_);
344
345 if (m_verbose_parse_ || m_events_since_last_report_ % 100 == 0) {
346 ldmx_log(info) << "Performance: "
347 << "total_events=" << m_total_events_built_ << ", "
348 << "events_per_sec=" << std::fixed
349 << std::setprecision(2) << events_per_sec << ", "
350 << "mb_per_sec=" << std::fixed << std::setprecision(3)
351 << mb_per_sec;
352 }
353
354 current_event_errors = 0; // Reset for next event
355 // The fragment that tripped the window-close belongs to the NEXT
356 // event, not the one we just emitted. Buffer it here (the buffer is
357 // empty post-build, so this also resets the reference time) instead
358 // of dropping it -- otherwise every event after the first loses its
359 // triggering subsystem.
360 m_event_buffer_.addFragment(std::move(fragment));
361 return; // Return the event to framework
362 }
363 }
364
365 // Add fragment to buffer (may start a new event batch if buffer was empty)
366 m_event_buffer_.addFragment(std::move(fragment));
367
368 if (m_verbose_parse_)
369 ldmx_log(debug) << "frame added to buffer, searching for more frames...";
370 }
371
372 // Reached EOF - flush any remaining events in the buffer
373 if (m_verbose_parse_)
374 ldmx_log(debug) << "reached EOF, flushing remaining events";
375
376 std::vector<DataFragment> assembled_event_fragments;
377 while (m_event_buffer_.tryBuildEvent(
378 m_coherence_window_ns_, m_min_subsystems_, assembled_event_fragments)) {
379 if (assembled_event_fragments.empty()) break;
380
381 // Mark truncated events (those flushed at EOF)
382 current_event_errors |= ldmx::EventSummary::ERROR_TRUNCATED_EVENT;
383
384 ++m_event_id_;
385 // Add each subsystem's raw data as vector<uint8_t> directly
386 uint64_t event_timestamp = 0;
387 for (const auto& frag : assembled_event_fragments) {
388 std::string subsys_name = packing::LDMXRoRHeader::getSubsystemName(
389 static_cast<uint8_t>(frag.header_.subsystem_id_),
390 static_cast<uint8_t>(frag.header_.contributor_id_));
391 std::vector<uint8_t> payload_bytes(frag.payload_.begin(),
392 frag.payload_.end());
393 event.add(subsys_name, payload_bytes);
394 if (event_timestamp == 0) {
395 event_timestamp = frag.header_.timestamp_;
396 }
397 }
398 event.getEventHeader().setIntParameter(
399 "RoR Timestamp", static_cast<long long>(event_timestamp));
400 // Still create PhysicsEventData for binary output file
401 PhysicsEventData final_event = assemblePayload(assembled_event_fragments);
402 writeEventBinary(final_event, "events.bin");
403 ldmx_log(info) << "assembled event id=" << m_event_id_
404 << " timestamp=" << final_event.timestamp_
405 << " fragments=" << assembled_event_fragments.size()
406 << " systems=" << final_event.systems_readout_.size();
407
408 ldmx::EventSummary summary;
409 summary.setEventNumber(m_event_id_);
410 summary.setTimestampNs(static_cast<uint64_t>(final_event.timestamp_));
411 std::set<uint64_t> unique_sys;
412 uint64_t total_payload = 0;
413 for (const auto& f : assembled_event_fragments) {
414 unique_sys.insert(f.header_.subsystem_id_);
415 total_payload += f.payload_.size();
416 }
417 // Check for duplicate subsystems
418 if (unique_sys.size() != assembled_event_fragments.size()) {
419 current_event_errors |= ldmx::EventSummary::ERROR_DUPLICATE_SUBSYSTEM;
420 }
421 summary.setNSystems(static_cast<uint32_t>(unique_sys.size()));
422 summary.setSystemIds(
423 std::vector<uint64_t>(unique_sys.begin(), unique_sys.end()));
424 summary.setPayloadSize(total_payload);
425 summary.setErrorFlags(current_event_errors);
426 event.add("EventSummary", summary);
427
428 // Update performance metrics
429 m_total_bytes_read_ += total_payload;
430 m_total_events_built_++;
431 m_events_since_last_report_++;
432 std::chrono::steady_clock::time_point now =
433 std::chrono::steady_clock::now();
434 std::chrono::milliseconds event_build_time_ms =
435 std::chrono::duration_cast<std::chrono::milliseconds>(
436 now - m_event_start_time_);
437 std::chrono::seconds total_elapsed_s =
438 std::chrono::duration_cast<std::chrono::seconds>(now - m_start_time_);
439
440 double events_per_sec = (total_elapsed_s.count() > 0)
441 ? static_cast<double>(m_total_events_built_) /
442 total_elapsed_s.count()
443 : 0.0;
444 double mb_per_sec = (total_elapsed_s.count() > 0)
445 ? static_cast<double>(m_total_bytes_read_) /
446 (1024.0 * 1024.0) / total_elapsed_s.count()
447 : 0.0;
448
449 // Update windowed metrics
450 m_window_events_count_++;
451 m_window_bytes_read_ += total_payload;
452 std::chrono::seconds window_elapsed_s =
453 std::chrono::duration_cast<std::chrono::seconds>(now -
454 m_window_start_time_);
455
456 // Calculate window metrics (for final event, may not have full window)
457 double window_events_per_sec = 0.0;
458 double window_mb_per_sec = 0.0;
459 if (m_window_events_count_ >= WINDOW_SIZE || window_elapsed_s.count() > 0) {
460 double window_time_sec =
461 std::max(1.0, static_cast<double>(window_elapsed_s.count()));
462 window_events_per_sec =
463 static_cast<double>(m_window_events_count_) / window_time_sec;
464 window_mb_per_sec = static_cast<double>(m_window_bytes_read_) /
465 (1024.0 * 1024.0) / window_time_sec;
466 }
467
468 // Write performance metrics to CSV
469 long long event_build_time_ms_val = event_build_time_ms.count();
470 writePerformanceMetric(
471 m_event_id_, static_cast<double>(event_build_time_ms_val),
472 events_per_sec, mb_per_sec, window_events_per_sec, window_mb_per_sec,
473 m_total_events_built_, m_total_bytes_read_);
474
475 if (m_verbose_parse_ || m_events_since_last_report_ % 100 == 0) {
476 ldmx_log(info) << "Performance: "
477 << "total_events=" << m_total_events_built_ << ", "
478 << "events_per_sec=" << std::fixed << std::setprecision(2)
479 << events_per_sec << ", "
480 << "mb_per_sec=" << std::fixed << std::setprecision(3)
481 << mb_per_sec;
482 }
483
484 current_event_errors = 0; // Reset for next event
485
486 assembled_event_fragments.clear();
487 return;
488 }
489
490 // No more events - print final summary statistics
491 std::chrono::steady_clock::time_point final_time =
492 std::chrono::steady_clock::now();
493 std::chrono::seconds total_time_s =
494 std::chrono::duration_cast<std::chrono::seconds>(final_time -
495 m_start_time_);
496 double final_events_per_sec =
497 (total_time_s.count() > 0)
498 ? static_cast<double>(m_total_events_built_) / total_time_s.count()
499 : 0.0;
500 double final_mb_per_sec = (total_time_s.count() > 0)
501 ? static_cast<double>(m_total_bytes_read_) /
502 (1024.0 * 1024.0) / total_time_s.count()
503 : 0.0;
504
505 ldmx_log(info) << "\n===== FINAL STATISTICS =====";
506 ldmx_log(info) << "Total events built: " << m_total_events_built_;
507 ldmx_log(info) << "Total bytes read: "
508 << m_total_bytes_read_ / (1024.0 * 1024.0) << " MB";
509 ldmx_log(info) << "Total time: " << total_time_s.count() << " seconds";
510 ldmx_log(info) << "Average throughput: "
511 << "events_per_sec=" << std::fixed << std::setprecision(2)
512 << final_events_per_sec << ", "
513 << "mb_per_sec=" << std::fixed << std::setprecision(3)
514 << final_mb_per_sec;
515 ldmx_log(info) << "=============================\n";
516 ldmx_log(info) << "Event building complete";
517
518 abortEvent();
519}
void abortEvent()
Abort the event immediately.
Provides summary information about an assembled event, including error flags for data quality assessm...
void setSystemIds(const std::vector< uint64_t > &ids)
Set the system IDs.
void setNSystems(uint32_t ns)
Set the number of systems in this event.
void setPayloadSize(uint64_t size)
Set the total payload size.
void setTimestampNs(uint64_t ts)
Set the timestamp in nanoseconds.
void setEventNumber(uint64_t num)
Set the event number.
void setErrorFlags(uint32_t flags)
Set the error flags.
@ ERROR_PARSE_FAILURE
Frame parsing error.
@ ERROR_DUPLICATE_SUBSYSTEM
Duplicate subsystem in single event.
@ ERROR_TRUNCATED_EVENT
Event was truncated (EOF reached)
the header that the LDMX DAQ Firmware block includes in the output data stream at the beginning of ea...
static std::tuple< int, int > subsystem(const std::string &name)
Get the (subsystem, contributor) pair for the input subsystem name.
uint8_t contributor() const
ID number for contributor within subsystem (configured into firmware)
static std::string getSubsystemName(uint8_t subsystem_id, uint8_t contributor_id)
Get the subsystem name from subsystem and contributor IDs.
uint64_t timestamp() const
get timestamp of this RoR
utility::Reader & read(utility::Reader &r)
read the next LDMX RoR header into memory
the header that the Rogue StreamWriter puts includes at the beginning of each frame.
unsigned int size() const
get the size of the frame not including this header
bool probablyYaml() const
check if this frame is probably a yaml dump
int channel() const
get the channel this data was written to
utility::Reader & read(utility::Reader &r)
read the next rogue frame header into memory
std::vector< uint32_t > & data()
Get data.
utility::Reader & read(utility::Reader &r)
read the subsystem packet from the input reader
std::vector< uint32_t > header() const
get the header words
std::streampos tell()
Tell us where the reader is.
Definition Reader.h:88
bool eof()
check if file is done
Definition Reader.h:217
void seek(std::streampos off, std::ios_base::seekdir dir=std::ios::beg)
Go ("seek") a specific position in the file.
Definition Reader.h:62
Reader & read(WordType *w, std::size_t count)
Read the next 'count' words into the input handle.
Definition Reader.h:114

References framework::EventProcessor::abortEvent(), packing::RogueFrameHeader::channel(), packing::LDMXRoRHeader::contributor(), packing::rawdatafile::SubsystemPacket::data(), packing::utility::Reader::eof(), ldmx::EventSummary::ERROR_DUPLICATE_SUBSYSTEM, ldmx::EventSummary::ERROR_PARSE_FAILURE, ldmx::EventSummary::ERROR_TRUNCATED_EVENT, packing::LDMXRoRHeader::getSubsystemName(), packing::rawdatafile::SubsystemPacket::header(), packing::RogueFrameHeader::probablyYaml(), packing::LDMXRoRHeader::read(), packing::rawdatafile::SubsystemPacket::read(), packing::RogueFrameHeader::read(), packing::utility::Reader::read(), packing::utility::Reader::seek(), ldmx::EventSummary::setErrorFlags(), ldmx::EventSummary::setEventNumber(), ldmx::EventSummary::setNSystems(), ldmx::EventSummary::setPayloadSize(), ldmx::EventSummary::setSystemIds(), ldmx::EventSummary::setTimestampNs(), packing::RogueFrameHeader::size(), packing::LDMXRoRHeader::subsystem(), packing::utility::Reader::tell(), and packing::LDMXRoRHeader::timestamp().

◆ writeEventBinary()

void EventBuilder::writeEventBinary ( const PhysicsEventData & ev,
const std::string & path = "events.bin" )
private

Definition at line 542 of file EventBuilder.cxx.

543 {
544 static std::mutex g_out_mutex;
545 std::lock_guard<std::mutex> lg(g_out_mutex);
546 std::ofstream ofs(path, std::ios::binary | std::ios::app);
547 if (!ofs) return;
548 uint64_t event_id_u = static_cast<uint64_t>(ev.event_id_);
549 uint64_t ts = static_cast<uint64_t>(ev.timestamp_);
550 uint32_t nblocks = static_cast<uint32_t>(ev.blocks_.size());
551 ofs.write(reinterpret_cast<const char*>(&event_id_u), sizeof(event_id_u));
552 ofs.write(reinterpret_cast<const char*>(&ts), sizeof(ts));
553 ofs.write(reinterpret_cast<const char*>(&nblocks), sizeof(nblocks));
554 for (const auto& b : ev.blocks_) {
555 uint64_t sid = b.subsystem_id_;
556 uint64_t bts = b.timestamp_ns_;
557 uint32_t psz = static_cast<uint32_t>(b.data_.size());
558 uint32_t csum = b.checksum_;
559 ofs.write(reinterpret_cast<const char*>(&sid), sizeof(sid));
560 ofs.write(reinterpret_cast<const char*>(&bts), sizeof(bts));
561 ofs.write(reinterpret_cast<const char*>(&psz), sizeof(psz));
562 ofs.write(reinterpret_cast<const char*>(&csum), sizeof(csum));
563 if (psz)
564 ofs.write(reinterpret_cast<const char*>(b.data_.data()),
565 static_cast<std::streamsize>(psz));
566 }
567 ofs.flush();
568}

Member Data Documentation

◆ m_coherence_window_ns_

long long eventbuilder::EventBuilder::m_coherence_window_ns_
private
Initial value:
{
5000000}

Definition at line 41 of file EventBuilder.h.

41 {
42 5000000}; // 5 ms window for collecting fragments

◆ m_event_buffer_

FragmentBuffer eventbuilder::EventBuilder::m_event_buffer_
private

Definition at line 37 of file EventBuilder.h.

◆ m_event_id_

unsigned int eventbuilder::EventBuilder::m_event_id_
private

Definition at line 36 of file EventBuilder.h.

◆ m_event_start_time_

std::chrono::steady_clock::time_point eventbuilder::EventBuilder::m_event_start_time_
private

Definition at line 49 of file EventBuilder.h.

◆ m_events_since_last_report_

unsigned long eventbuilder::EventBuilder::m_events_since_last_report_ {0}
private

Definition at line 52 of file EventBuilder.h.

52{0};

◆ m_input_file_

std::string eventbuilder::EventBuilder::m_input_file_
private

Definition at line 38 of file EventBuilder.h.

◆ m_min_subsystems_

int eventbuilder::EventBuilder::m_min_subsystems_ {2}
private

Definition at line 45 of file EventBuilder.h.

45{2};

◆ m_output_name_

std::string eventbuilder::EventBuilder::m_output_name_ {"BuilderOutput"}
private

Definition at line 40 of file EventBuilder.h.

40{"BuilderOutput"};

◆ m_reader_

packing::utility::Reader eventbuilder::EventBuilder::m_reader_
private

Definition at line 39 of file EventBuilder.h.

◆ m_start_time_

std::chrono::steady_clock::time_point eventbuilder::EventBuilder::m_start_time_
private

Definition at line 48 of file EventBuilder.h.

◆ m_total_bytes_read_

uint64_t eventbuilder::EventBuilder::m_total_bytes_read_ {0}
private

Definition at line 50 of file EventBuilder.h.

50{0};

◆ m_total_events_built_

uint64_t eventbuilder::EventBuilder::m_total_events_built_ {0}
private

Definition at line 51 of file EventBuilder.h.

51{0};

◆ m_verbose_parse_

bool eventbuilder::EventBuilder::m_verbose_parse_
private

Definition at line 35 of file EventBuilder.h.

◆ m_window_bytes_read_

uint64_t eventbuilder::EventBuilder::m_window_bytes_read_ {0}
private

Definition at line 59 of file EventBuilder.h.

59{0};

◆ m_window_events_count_

uint64_t eventbuilder::EventBuilder::m_window_events_count_ {0}
private

Definition at line 58 of file EventBuilder.h.

58{0};

◆ m_window_start_time_

std::chrono::steady_clock::time_point eventbuilder::EventBuilder::m_window_start_time_
private

Definition at line 57 of file EventBuilder.h.

◆ WINDOW_SIZE

unsigned int eventbuilder::EventBuilder::WINDOW_SIZE
staticconstexprprivate
Initial value:
=
100

Definition at line 55 of file EventBuilder.h.


The documentation for this class was generated from the following files: