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;
116 m_event_start_time_ = std::chrono::steady_clock::now();
119 uint32_t current_event_errors = 0;
121 while (m_reader_ && !m_reader_.
eof()) {
124 frame_header.
read(m_reader_);
127 const long long frame_end =
128 static_cast<long long>(m_reader_.
tell()) + frame_header.
size();
133 if (m_verbose_parse_)
134 ldmx_log(debug) <<
"skipping non-data frame (channel="
135 << frame_header.
channel() <<
")";
136 m_reader_.
seek(frame_end);
142 bool parsed_ok =
false;
146 int pos_before_ror = m_reader_.
tell();
150 ror_header.
read(m_reader_);
152 fragment.header_.subsystem_id_ =
153 static_cast<uint64_t
>(ror_header.
subsystem());
154 fragment.header_.contributor_id_ =
156 fragment.header_.timestamp_ = ror_header.
timestamp();
158 if (m_verbose_parse_) {
159 ldmx_log(debug) <<
"parsed RoR header subsys="
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()),
173 fragment.payload_ = std::move(payload_data);
177 if (m_verbose_parse_)
179 <<
"RoR header parse failed, trying packing subsystem format";
181 m_reader_.
seek(pos_before_ror);
189 fragment.header_.subsystem_id_ =
static_cast<uint64_t
>(pkt.id());
190 fragment.header_.contributor_id_ =
192 fragment.header_.timestamp_ =
static_cast<uint64_t
>(pkt.
header()[1]) *
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()),
202 if (m_verbose_parse_) {
203 ldmx_log(debug) <<
"parsed packing subsystem pkt subsys=" << pkt.id()
204 <<
" data_size=" << data.size();
209 if (m_verbose_parse_)
210 ldmx_log(debug) <<
"failed to parse as either format, skipping frame";
212 m_reader_.
seek(frame_end);
218 m_reader_.
seek(frame_end);
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();
231 long long fragment_ts = fragment.header_.timestamp_;
232 long long buffer_ref_time = m_event_buffer_.getReferenceTime();
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_)) {
239 std::vector<DataFragment> assembled_event_fragments;
240 if (m_event_buffer_.tryBuildEvent(m_coherence_window_ns_,
242 assembled_event_fragments) &&
243 !assembled_event_fragments.empty()) {
247 uint64_t event_timestamp = 0;
248 for (
const auto& frag : assembled_event_fragments) {
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_;
259 event.getEventHeader().setIntParameter(
260 "RoR Timestamp",
static_cast<long long>(event_timestamp));
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();
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();
280 if (unique_sys.size() != assembled_event_fragments.size()) {
283 summary.
setNSystems(
static_cast<uint32_t
>(unique_sys.size()));
285 std::vector<uint64_t>(unique_sys.begin(), unique_sys.end()));
288 event.add(
"EventSummary", summary);
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 -
303 double events_per_sec =
304 (total_elapsed_s.count() > 0)
305 ?
static_cast<double>(m_total_events_built_) /
306 total_elapsed_s.count()
308 double mb_per_sec = (total_elapsed_s.count() > 0)
309 ?
static_cast<double>(m_total_bytes_read_) /
311 total_elapsed_s.count()
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_);
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;
333 m_window_start_time_ = now;
334 m_window_events_count_ = 0;
335 m_window_bytes_read_ = 0;
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_);
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)
354 current_event_errors = 0;
360 m_event_buffer_.addFragment(std::move(fragment));
366 m_event_buffer_.addFragment(std::move(fragment));
368 if (m_verbose_parse_)
369 ldmx_log(debug) <<
"frame added to buffer, searching for more frames...";
373 if (m_verbose_parse_)
374 ldmx_log(debug) <<
"reached EOF, flushing remaining events";
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;
386 uint64_t event_timestamp = 0;
387 for (
const auto& frag : assembled_event_fragments) {
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_;
398 event.getEventHeader().setIntParameter(
399 "RoR Timestamp",
static_cast<long long>(event_timestamp));
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();
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();
418 if (unique_sys.size() != assembled_event_fragments.size()) {
421 summary.
setNSystems(
static_cast<uint32_t
>(unique_sys.size()));
423 std::vector<uint64_t>(unique_sys.begin(), unique_sys.end()));
426 event.add(
"EventSummary", summary);
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_);
440 double events_per_sec = (total_elapsed_s.count() > 0)
441 ?
static_cast<double>(m_total_events_built_) /
442 total_elapsed_s.count()
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()
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_);
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;
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_);
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)
484 current_event_errors = 0;
486 assembled_event_fragments.clear();
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 -
496 double final_events_per_sec =
497 (total_time_s.count() > 0)
498 ?
static_cast<double>(m_total_events_built_) / total_time_s.count()
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()
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)
515 ldmx_log(info) <<
"=============================\n";
516 ldmx_log(info) <<
"Event building complete";