109 static int produce_call_count = 0;
110 produce_call_count++;
111 ldmx_log(debug) <<
"produce() called, count=" << produce_call_count;
112 ldmx_log(debug) <<
"verbose_parse=" << (m_verbose_parse ?
"true" :
"false");
113 if (m_verbose_parse || produce_call_count <= 3) {
114 ldmx_log(info) <<
"produce() call #" << produce_call_count;
118 m_event_start_time = std::chrono::steady_clock::now();
121 uint32_t current_event_errors = 0;
123 while (m_reader && !m_reader.
eof()) {
126 frame_header.
read(m_reader);
129 const long long frame_end =
130 static_cast<long long>(m_reader.
tell()) + frame_header.
size();
135 if (m_verbose_parse) ldmx_log(debug) <<
"skipping non-data frame (channel=" << 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 =
static_cast<uint64_t
>(ror_header.
subsystem());
153 fragment.header.contributor_id =
static_cast<uint64_t
>(ror_header.
contributor());
154 fragment.header.timestamp = ror_header.
timestamp();
156 if (m_verbose_parse) {
157 ldmx_log(debug) <<
"parsed RoR header subsys=" << (int)ror_header.
subsystem()
163 std::vector<uint8_t> payload_data;
164 long long payload_size = frame_end - m_reader.
tell();
165 if (payload_size > 0) {
166 payload_data.resize(payload_size);
167 m_reader.
read(
reinterpret_cast<char*
>(payload_data.data()), payload_size);
169 fragment.payload = std::move(payload_data);
173 if (m_verbose_parse) ldmx_log(debug) <<
"RoR header parse failed, trying packing subsystem format";
175 m_reader.
seek(pos_before_ror);
183 fragment.header.subsystem_id =
static_cast<uint64_t
>(pkt.id());
184 fragment.header.contributor_id = 0;
185 fragment.header.timestamp =
static_cast<uint64_t
>(pkt.
header()[1]) * 1000000000ULL;
188 const auto& data = pkt.
data();
189 fragment.payload.resize(data.size() * 4);
190 std::memcpy(fragment.payload.data(),
reinterpret_cast<const char*
>(data.data()), data.size() * 4);
192 if (m_verbose_parse) {
193 ldmx_log(debug) <<
"parsed packing subsystem pkt subsys=" << pkt.id()
194 <<
" data_size=" << data.size();
199 if (m_verbose_parse) ldmx_log(debug) <<
"failed to parse as either format, skipping frame";
201 m_reader.
seek(frame_end);
207 m_reader.
seek(frame_end);
211 if (m_verbose_parse) ldmx_log(debug) <<
"adding fragment subsys=" << fragment.header.subsystem_id
212 <<
" ts=" << fragment.header.timestamp <<
" bytes=" << fragment.payload.size();
216 long long fragment_ts = fragment.header.timestamp;
217 long long buffer_ref_time = m_event_buffer.get_reference_time();
219 if (buffer_ref_time > 0 && (fragment_ts < buffer_ref_time - m_coherence_window_ns ||
220 fragment_ts > buffer_ref_time + m_coherence_window_ns)) {
222 std::vector<DataFragment> assembled_event_fragments;
223 if (m_event_buffer.try_build_event(m_coherence_window_ns, m_min_subsystems, assembled_event_fragments) &&
224 !assembled_event_fragments.empty()) {
228 uint64_t event_timestamp = 0;
229 for (
const auto &frag : assembled_event_fragments) {
231 static_cast<uint8_t
>(frag.header.subsystem_id),
232 static_cast<uint8_t
>(frag.header.contributor_id));
233 std::vector<uint8_t> payload_bytes(frag.payload.begin(), frag.payload.end());
234 event.add(subsys_name, payload_bytes);
235 if (event_timestamp == 0) {
236 event_timestamp = frag.header.timestamp;
239 event.getEventHeader().setIntParameter(
"RoR Timestamp",
static_cast<long long>(event_timestamp));
242 write_event_binary(final_event,
"events.bin");
243 ldmx_log(info) <<
"assembled event id=" << m_event_id <<
" timestamp=" << final_event.timestamp
244 <<
" fragments=" << assembled_event_fragments.size() <<
" systems=" << final_event.systems_readout.size();
248 summary.
setTimestampNs(
static_cast<uint64_t
>(final_event.timestamp));
249 std::set<uint64_t> unique_sys;
250 uint64_t total_payload = 0;
251 for (
const auto &f : assembled_event_fragments) {
252 unique_sys.insert(f.header.subsystem_id);
253 total_payload += f.payload.size();
256 if (unique_sys.size() != assembled_event_fragments.size()) {
259 summary.
setNSystems(
static_cast<uint32_t
>(unique_sys.size()));
260 summary.
setSystemIds(std::vector<uint64_t>(unique_sys.begin(), unique_sys.end()));
263 event.add(
"EventSummary", summary);
266 m_total_bytes_read += total_payload;
267 m_total_events_built++;
268 m_events_since_last_report++;
269 std::chrono::steady_clock::time_point now = std::chrono::steady_clock::now();
270 std::chrono::milliseconds event_build_time_ms = std::chrono::duration_cast<std::chrono::milliseconds>(now - m_event_start_time);
271 std::chrono::seconds total_elapsed_s = std::chrono::duration_cast<std::chrono::seconds>(now - m_start_time);
273 double events_per_sec = (total_elapsed_s.count() > 0) ?
static_cast<double>(m_total_events_built) / total_elapsed_s.count() : 0.0;
274 double mb_per_sec = (total_elapsed_s.count() > 0) ?
static_cast<double>(m_total_bytes_read) / (1024.0 * 1024.0) / total_elapsed_s.count() : 0.0;
277 m_window_events_count++;
278 m_window_bytes_read += total_payload;
279 std::chrono::seconds window_elapsed_s = std::chrono::duration_cast<std::chrono::seconds>(now - m_window_start_time);
282 double window_events_per_sec = 0.0;
283 double window_mb_per_sec = 0.0;
284 if (m_window_events_count >= WINDOW_SIZE) {
285 double window_time_sec = std::max(1.0,
static_cast<double>(window_elapsed_s.count()));
286 window_events_per_sec =
static_cast<double>(m_window_events_count) / window_time_sec;
287 window_mb_per_sec =
static_cast<double>(m_window_bytes_read) / (1024.0 * 1024.0) / window_time_sec;
290 m_window_start_time = now;
291 m_window_events_count = 0;
292 m_window_bytes_read = 0;
296 long long event_build_time_ms_val = event_build_time_ms.count();
297 write_performance_metric(m_event_id,
static_cast<double>(event_build_time_ms_val),
298 events_per_sec, mb_per_sec,
299 window_events_per_sec, window_mb_per_sec,
300 m_total_events_built, m_total_bytes_read);
302 if (m_verbose_parse || m_events_since_last_report % 100 == 0) {
303 ldmx_log(info) <<
"Performance: "
304 <<
"total_events=" << m_total_events_built <<
", "
305 <<
"events_per_sec=" << std::fixed << std::setprecision(2) << events_per_sec <<
", "
306 <<
"mb_per_sec=" << std::fixed << std::setprecision(3) << mb_per_sec;
309 current_event_errors = 0;
315 m_event_buffer.add_fragment(std::move(fragment));
321 m_event_buffer.add_fragment(std::move(fragment));
323 if (m_verbose_parse) ldmx_log(debug) <<
"frame added to buffer, searching for more frames...";
327 if (m_verbose_parse) ldmx_log(debug) <<
"reached EOF, flushing remaining events";
329 std::vector<DataFragment> assembled_event_fragments;
330 while (m_event_buffer.try_build_event(m_coherence_window_ns, m_min_subsystems, assembled_event_fragments)) {
331 if (assembled_event_fragments.empty())
break;
338 uint64_t event_timestamp = 0;
339 for (
const auto &frag : assembled_event_fragments) {
341 static_cast<uint8_t
>(frag.header.subsystem_id),
342 static_cast<uint8_t
>(frag.header.contributor_id));
343 std::vector<uint8_t> payload_bytes(frag.payload.begin(), frag.payload.end());
344 event.add(subsys_name, payload_bytes);
345 if (event_timestamp == 0) {
346 event_timestamp = frag.header.timestamp;
349 event.getEventHeader().setIntParameter(
"RoR Timestamp",
static_cast<long long>(event_timestamp));
352 write_event_binary(final_event,
"events.bin");
353 ldmx_log(info) <<
"assembled event id=" << m_event_id <<
" timestamp=" << final_event.timestamp
354 <<
" fragments=" << assembled_event_fragments.size() <<
" systems=" << final_event.systems_readout.size();
358 summary.
setTimestampNs(
static_cast<uint64_t
>(final_event.timestamp));
359 std::set<uint64_t> unique_sys;
360 uint64_t total_payload = 0;
361 for (
const auto &f : assembled_event_fragments) {
362 unique_sys.insert(f.header.subsystem_id);
363 total_payload += f.payload.size();
366 if (unique_sys.size() != assembled_event_fragments.size()) {
369 summary.
setNSystems(
static_cast<uint32_t
>(unique_sys.size()));
370 summary.
setSystemIds(std::vector<uint64_t>(unique_sys.begin(), unique_sys.end()));
373 event.add(
"EventSummary", summary);
376 m_total_bytes_read += total_payload;
377 m_total_events_built++;
378 m_events_since_last_report++;
379 std::chrono::steady_clock::time_point now = std::chrono::steady_clock::now();
380 std::chrono::milliseconds event_build_time_ms = std::chrono::duration_cast<std::chrono::milliseconds>(now - m_event_start_time);
381 std::chrono::seconds total_elapsed_s = std::chrono::duration_cast<std::chrono::seconds>(now - m_start_time);
383 double events_per_sec = (total_elapsed_s.count() > 0) ?
static_cast<double>(m_total_events_built) / total_elapsed_s.count() : 0.0;
384 double mb_per_sec = (total_elapsed_s.count() > 0) ?
static_cast<double>(m_total_bytes_read) / (1024.0 * 1024.0) / total_elapsed_s.count() : 0.0;
387 m_window_events_count++;
388 m_window_bytes_read += total_payload;
389 std::chrono::seconds window_elapsed_s = std::chrono::duration_cast<std::chrono::seconds>(now - m_window_start_time);
392 double window_events_per_sec = 0.0;
393 double window_mb_per_sec = 0.0;
394 if (m_window_events_count >= WINDOW_SIZE || window_elapsed_s.count() > 0) {
395 double window_time_sec = std::max(1.0,
static_cast<double>(window_elapsed_s.count()));
396 window_events_per_sec =
static_cast<double>(m_window_events_count) / window_time_sec;
397 window_mb_per_sec =
static_cast<double>(m_window_bytes_read) / (1024.0 * 1024.0) / window_time_sec;
401 long long event_build_time_ms_val = event_build_time_ms.count();
402 write_performance_metric(m_event_id,
static_cast<double>(event_build_time_ms_val),
403 events_per_sec, mb_per_sec,
404 window_events_per_sec, window_mb_per_sec,
405 m_total_events_built, m_total_bytes_read);
407 if (m_verbose_parse || m_events_since_last_report % 100 == 0) {
408 ldmx_log(info) <<
"Performance: "
409 <<
"total_events=" << m_total_events_built <<
", "
410 <<
"events_per_sec=" << std::fixed << std::setprecision(2) << events_per_sec <<
", "
411 <<
"mb_per_sec=" << std::fixed << std::setprecision(3) << mb_per_sec;
414 current_event_errors = 0;
416 assembled_event_fragments.clear();
421 std::chrono::steady_clock::time_point final_time = std::chrono::steady_clock::now();
422 std::chrono::seconds total_time_s = std::chrono::duration_cast<std::chrono::seconds>(final_time - m_start_time);
423 double final_events_per_sec = (total_time_s.count() > 0) ?
static_cast<double>(m_total_events_built) / total_time_s.count() : 0.0;
424 double final_mb_per_sec = (total_time_s.count() > 0) ?
static_cast<double>(m_total_bytes_read) / (1024.0 * 1024.0) / total_time_s.count() : 0.0;
426 ldmx_log(info) <<
"\n===== FINAL STATISTICS =====";
427 ldmx_log(info) <<
"Total events built: " << m_total_events_built;
428 ldmx_log(info) <<
"Total bytes read: " << m_total_bytes_read / (1024.0 * 1024.0) <<
" MB";
429 ldmx_log(info) <<
"Total time: " << total_time_s.count() <<
" seconds";
430 ldmx_log(info) <<
"Average throughput: "
431 <<
"events_per_sec=" << std::fixed << std::setprecision(2) << final_events_per_sec <<
", "
432 <<
"mb_per_sec=" << std::fixed << std::setprecision(3) << final_mb_per_sec;
433 ldmx_log(info) <<
"=============================\n";
434 ldmx_log(info) <<
"Event building complete";