LDMX Software
EventBuilder.cxx
1#include "EventBuilder/EventBuilder.h"
2
3#include <algorithm>
4#include <chrono>
5#include <cstdint>
6#include <cstring>
7#include <fstream>
8#include <iomanip>
9#include <iostream>
10#include <set>
11#include <vector>
12
13#include "EventBuilder/Event/GenericDataBlock.h"
14#include "EventBuilder/Event/PhysicsEventData.h"
15#include "EventBuilder/Fragment.h"
17#include "Framework/Logger.h"
18#include "Packing/LDMXRoRHeader.h"
19#include "Packing/RawDataFile/SubsystemPacket.h"
20#include "Packing/RogueFrameHeader.h"
21
22using namespace eventbuilder;
23
24// Global flag to track if performance CSV header has been written
25static bool g_perf_csv_header_written = false;
26
27// Helper function to write performance metrics to CSV
28void writePerformanceMetric(unsigned int event_id, double build_time_ms,
29 double cum_events_per_sec, double cum_mb_per_sec,
30 double window_events_per_sec,
31 double window_mb_per_sec, uint64_t total_events,
32 uint64_t total_bytes) {
33 std::ofstream perf_csv("event_performance.csv", std::ios::app);
34 if (!perf_csv) return;
35
36 // Write header on first call
37 if (!g_perf_csv_header_written) {
38 perf_csv
39 << "event_id,event_build_time_ms,cum_events_per_sec,cum_mb_per_sec,"
40 << "window_events_per_sec,window_mb_per_sec,total_events,total_bytes_"
41 "mb\n";
42 g_perf_csv_header_written = true;
43 }
44
45 // Write data with consistent formatting
46 double total_mb = total_bytes / (1024.0 * 1024.0);
47 perf_csv << event_id << "," << std::fixed << std::setprecision(3)
48 << build_time_ms << "," << std::fixed << std::setprecision(2)
49 << cum_events_per_sec << "," << std::fixed << std::setprecision(3)
50 << cum_mb_per_sec << "," << std::fixed << std::setprecision(2)
51 << window_events_per_sec << "," << std::fixed << std::setprecision(3)
52 << window_mb_per_sec << "," << total_events << "," << std::fixed
53 << std::setprecision(2) << total_mb << "\n";
54 perf_csv.flush();
55}
56
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}
105
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}
520
521PhysicsEventData EventBuilder::assemblePayload(
522 const std::vector<DataFragment>& fragments) {
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}
541
542void EventBuilder::writeEventBinary(const PhysicsEventData& ev,
543 const std::string& path) {
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}
569
570// Register producer with the framework factory
#define DECLARE_PRODUCER(CLASS)
Macro which allows the framework to construct a producer given its name during configuration.
Class that provides a summary of event assembly with metadata and error flags.
void produce(framework::Event &event) override
Process the event and put new data products into it.
void configure(framework::config::Parameters &ps) override
Callback for the EventProcessor to configure itself from the given set of parameters.
void abortEvent()
Abort the event immediately.
Implements an event buffer system for storing event data.
Definition Event.h:40
Class encapsulating parameters for configuring a processor.
Definition Parameters.h:26
const T & get(const std::string &name) const
Retrieve the parameter of the given name.
Definition Parameters.h:75
bool exists(const std::string &name) const
Check to see if a parameter exists.
Definition Parameters.h:60
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
void open(const std::string &file_name)
Open a file with this reader.
Definition Reader.h:36
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