Skip to content
11 changes: 9 additions & 2 deletions artdaq/ArtModules/ArtdaqSharedMemoryService_service.cc
Original file line number Diff line number Diff line change
Expand Up @@ -93,6 +93,7 @@ class ArtdaqSharedMemoryService : public ArtdaqSharedMemoryServiceInterface
private:
std::unique_ptr<artdaq::SharedMemoryEventReceiver> incoming_events_;
std::list<std::shared_ptr<ArtdaqEvent>> event_ordering_;
std::set<artdaq::Fragment::sequence_id_t> released_broadcast_sequence_ids_;
size_t read_timeout_;
size_t subrun_closure_threshold_{1};
double safety_valve_timeout_s_{10.0};
Expand Down Expand Up @@ -309,9 +310,15 @@ std::shared_ptr<ArtdaqEvent> ArtdaqSharedMemoryService::ReceiveEvent(bool broadc
{
// First Fragment is broadcast (begin/end run/subrun), but there's more in event ordering!
TLOG(TLVL_RECEIVEEVENT) << "Returning Broadcast Fragment due to subrun closure";
output_event = event_ordering_.front();

if (released_broadcast_sequence_ids_.count(event_ordering_.front()->header->sequence_id) == 0)
{
output_event = event_ordering_.front();
released_broadcast_sequence_ids_.insert(event_ordering_.front()->header->sequence_id);
while (released_broadcast_sequence_ids_.size() > 1000) { released_broadcast_sequence_ids_.erase(released_broadcast_sequence_ids_.begin()); }
}
event_ordering_.pop_front();
break; // while(output_event == nullptr)
continue; // while(output_event == nullptr)
}
}
else if (current_subrun_ != 0 && first_sr > current_subrun_ + 1)
Expand Down