Process the event and put new data products into it.
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
116 m_event_start_time_ = std::chrono::steady_clock::now();
117
118
119 uint32_t current_event_errors = 0;
120
121 while (m_reader_ && !m_reader_.
eof()) {
122
124 frame_header.
read(m_reader_);
125
126
127 const long long frame_end =
128 static_cast<long long>(m_reader_.
tell()) + frame_header.
size();
129
130
132
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
142 bool parsed_ok = false;
143
144
146 int pos_before_ror = m_reader_.
tell();
147
148
149 try {
150 ror_header.
read(m_reader_);
151
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();
157
158 if (m_verbose_parse_) {
159 ldmx_log(debug) << "parsed RoR header subsys="
163 }
164
165
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
177 if (m_verbose_parse_)
178 ldmx_log(debug)
179 << "RoR header parse failed, trying packing subsystem format";
181 m_reader_.
seek(pos_before_ror);
182 }
183
184 if (!parsed_ok) {
185
187 try {
189 fragment.header_.subsystem_id_ = static_cast<uint64_t>(pkt.id());
190 fragment.header_.contributor_id_ =
191 0;
192 fragment.header_.timestamp_ =
static_cast<uint64_t
>(pkt.
header()[1]) *
193 1000000000ULL;
194
195
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
209 if (m_verbose_parse_)
210 ldmx_log(debug) << "failed to parse as either format, skipping frame";
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
229
230
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
238
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
245 ++m_event_id_;
246
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_;
257 }
258 }
259 event.getEventHeader().setIntParameter(
260 "RoR Timestamp", static_cast<long long>(event_timestamp));
261
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
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
280 if (unique_sys.size() != assembled_event_fragments.size()) {
282 }
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);
289
290
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
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
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
333 m_window_start_time_ = now;
334 m_window_events_count_ = 0;
335 m_window_bytes_read_ = 0;
336 }
337
338
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;
355
356
357
358
359
360 m_event_buffer_.addFragment(std::move(fragment));
361 return;
362 }
363 }
364
365
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
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
383
384 ++m_event_id_;
385
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_;
396 }
397 }
398 event.getEventHeader().setIntParameter(
399 "RoR Timestamp", static_cast<long long>(event_timestamp));
400
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
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
418 if (unique_sys.size() != assembled_event_fragments.size()) {
420 }
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);
427
428
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
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
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
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;
485
486 assembled_event_fragments.clear();
487 return;
488 }
489
490
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
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)
SubsystemPacket structure.
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.
bool eof()
check if file is done
void seek(std::streampos off, std::ios_base::seekdir dir=std::ios::beg)
Go ("seek") a specific position in the file.
Reader & read(WordType *w, std::size_t count)
Read the next 'count' words into the input handle.