LDMX Software
FragmentBuffer.h
1#ifndef EVENTBUILDER_FRAGMENTBUFFER_H
2#define EVENTBUILDER_FRAGMENTBUFFER_H
3
4#include <map>
5#include <mutex>
6#include <set>
7#include <vector>
8
9#include "Fragment.h"
10
11namespace eventbuilder {
12
14 public:
15 using Timestamp = long long;
16
17 void addFragment(DataFragment&& fragment) {
18 std::lock_guard<std::mutex> lock(m_mutex_);
19
20 // Set reference time on first fragment
21 if (m_fragments_.empty()) {
22 m_event_reference_time_ = fragment.header_.timestamp_;
23 }
24
25 m_fragments_[fragment.header_.timestamp_].push_back(std::move(fragment));
26 }
27
28 bool hasExpiredFragments(Timestamp reference_time,
29 long long coherence_window_ns) {
30 std::lock_guard<std::mutex> lock(m_mutex_);
31 if (m_fragments_.empty()) {
32 return false;
33 }
34 auto it_oldest = m_fragments_.begin();
35 return it_oldest->first < reference_time - coherence_window_ns;
36 }
37
38 Timestamp getReferenceTime() const {
39 std::lock_guard<std::mutex> lock(m_mutex_);
40 return m_event_reference_time_;
41 }
42
43 bool tryBuildEvent(long long coherence_window_ns, int min_subsystems,
44 std::vector<DataFragment>& built_fragments) {
45 std::lock_guard<std::mutex> lock(m_mutex_);
46 if (m_fragments_.empty()) return false;
47
48 // Use the stored reference time from the first fragment in current
49 // collection
50 Timestamp window_ref_time = m_event_reference_time_;
51
52 auto it_begin =
53 m_fragments_.lower_bound(window_ref_time - coherence_window_ns);
54 auto it_end =
55 m_fragments_.upper_bound(window_ref_time + coherence_window_ns);
56
57 if (it_begin == it_end) return false;
58
59 std::set<uint64_t> subsystems_found;
60 std::vector<Timestamp> timestamps_in_window;
61
62 for (auto it = it_begin; it != it_end; ++it) {
63 timestamps_in_window.push_back(it->first);
64 for (const auto& frag : it->second) {
65 subsystems_found.insert(frag.header_.subsystem_id_);
66 }
67 }
68
69 // Require at least min_subsystems distinct subsystems in the window before
70 // assembling, so we don't emit partial events. Configurable.
71 if (subsystems_found.size() < static_cast<size_t>(min_subsystems)) {
72 return false;
73 }
74
75 // Collect fragments and remove them from buffer
76 for (Timestamp ts : timestamps_in_window) {
77 for (auto& frag : m_fragments_[ts]) {
78 built_fragments.push_back(std::move(frag));
79 }
80 m_fragments_.erase(ts);
81 }
82
83 // Reset reference time for next event
84 if (!m_fragments_.empty()) {
85 m_event_reference_time_ = m_fragments_.begin()->first;
86 }
87
88 return true;
89 }
90
91 private:
92 std::map<Timestamp, std::vector<DataFragment>> m_fragments_;
93 Timestamp m_event_reference_time_ = 0;
94 mutable std::mutex m_mutex_;
95};
96
97} // namespace eventbuilder
98
99#endif // FRAGMENTBUFFER_H