Process the event and put new data products into it.
108 {
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;
115 }
116
117
118 m_event_start_time = std::chrono::steady_clock::now();
119
120
121 uint32_t current_event_errors = 0;
122
123 while (m_reader && !m_reader.
eof()) {
124
126 frame_header.
read(m_reader);
127
128
129 const long long frame_end =
130 static_cast<long long>(m_reader.
tell()) + frame_header.
size();
131
132
134
135 if (m_verbose_parse) ldmx_log(debug) <<
"skipping non-data frame (channel=" << 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 =
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();
155
156 if (m_verbose_parse) {
157 ldmx_log(debug) <<
"parsed RoR header subsys=" << (int)ror_header.
subsystem()
160 }
161
162
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);
168 }
169 fragment.payload = std::move(payload_data);
170 parsed_ok = true;
171 } catch (...) {
172
173 if (m_verbose_parse) ldmx_log(debug) << "RoR header parse failed, trying packing subsystem format";
175 m_reader.
seek(pos_before_ror);
176 }
177
178 if (!parsed_ok) {
179
181 try {
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;
186
187
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);
191
192 if (m_verbose_parse) {
193 ldmx_log(debug) << "parsed packing subsystem pkt subsys=" << pkt.id()
194 << " data_size=" << data.size();
195 }
196 parsed_ok = true;
197 } catch (...) {
198
199 if (m_verbose_parse) ldmx_log(debug) << "failed to parse as either format, skipping frame";
201 m_reader.
seek(frame_end);
202 continue;
203 }
204 }
205
206 if (!parsed_ok) {
207 m_reader.
seek(frame_end);
208 continue;
209 }
210
211 if (m_verbose_parse) ldmx_log(debug) << "adding fragment subsys=" << fragment.header.subsystem_id
212 << " ts=" << fragment.header.timestamp << " bytes=" << fragment.payload.size();
213
214
215
216 long long fragment_ts = fragment.header.timestamp;
217 long long buffer_ref_time = m_event_buffer.get_reference_time();
218
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)) {
221
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()) {
225
226 ++m_event_id;
227
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;
237 }
238 }
239 event.getEventHeader().setIntParameter("RoR Timestamp", static_cast<long long>(event_timestamp));
240
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();
245
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();
254 }
255
256 if (unique_sys.size() != assembled_event_fragments.size()) {
258 }
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);
264
265
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);
272
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;
275
276
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);
280
281
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;
288
289
290 m_window_start_time = now;
291 m_window_events_count = 0;
292 m_window_bytes_read = 0;
293 }
294
295
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);
301
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;
307 }
308
309 current_event_errors = 0;
310
311
312
313
314
315 m_event_buffer.add_fragment(std::move(fragment));
316 return;
317 }
318 }
319
320
321 m_event_buffer.add_fragment(std::move(fragment));
322
323 if (m_verbose_parse) ldmx_log(debug) << "frame added to buffer, searching for more frames...";
324 }
325
326
327 if (m_verbose_parse) ldmx_log(debug) << "reached EOF, flushing remaining events";
328
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;
332
333
335
336 ++m_event_id;
337
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;
347 }
348 }
349 event.getEventHeader().setIntParameter("RoR Timestamp", static_cast<long long>(event_timestamp));
350
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();
355
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();
364 }
365
366 if (unique_sys.size() != assembled_event_fragments.size()) {
368 }
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);
374
375
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);
382
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;
385
386
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);
390
391
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;
398 }
399
400
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);
406
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;
412 }
413
414 current_event_errors = 0;
415
416 assembled_event_fragments.clear();
417 return;
418 }
419
420
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;
425
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";
435
437}
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.