15 using Timestamp =
long long;
18 std::lock_guard<std::mutex> lock(m_mutex_);
21 if (m_fragments_.empty()) {
22 m_event_reference_time_ = fragment.header_.timestamp_;
25 m_fragments_[fragment.header_.timestamp_].push_back(std::move(fragment));
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()) {
34 auto it_oldest = m_fragments_.begin();
35 return it_oldest->first < reference_time - coherence_window_ns;
38 Timestamp getReferenceTime()
const {
39 std::lock_guard<std::mutex> lock(m_mutex_);
40 return m_event_reference_time_;
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;
50 Timestamp window_ref_time = m_event_reference_time_;
53 m_fragments_.lower_bound(window_ref_time - coherence_window_ns);
55 m_fragments_.upper_bound(window_ref_time + coherence_window_ns);
57 if (it_begin == it_end)
return false;
59 std::set<uint64_t> subsystems_found;
60 std::vector<Timestamp> timestamps_in_window;
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_);
71 if (subsystems_found.size() <
static_cast<size_t>(min_subsystems)) {
76 for (Timestamp ts : timestamps_in_window) {
77 for (
auto& frag : m_fragments_[ts]) {
78 built_fragments.push_back(std::move(frag));
80 m_fragments_.erase(ts);
84 if (!m_fragments_.empty()) {
85 m_event_reference_time_ = m_fragments_.begin()->first;
92 std::map<Timestamp, std::vector<DataFragment>> m_fragments_;
93 Timestamp m_event_reference_time_ = 0;
94 mutable std::mutex m_mutex_;