ATLAS Offline Software
Loading...
Searching...
No Matches
OutputStreamSequencerSvc.cxx
Go to the documentation of this file.
1/*
2 Copyright (C) 2002-2026 CERN for the benefit of the ATLAS collaboration
3*/
4
9
11#include "MetaDataSvc.h"
12
13#include "GaudiKernel/IIncidentSvc.h"
14#include "GaudiKernel/FileIncident.h"
15#include "GaudiKernel/ConcurrencyFlags.h"
16
17#include <charconv>
18#include <format>
19#include <print>
20#include <string_view>
21#include <sstream>
22
23
24//________________________________________________________________________________
25OutputStreamSequencerSvc::OutputStreamSequencerSvc(const std::string& name, ISvcLocator* pSvcLocator)
26 : base_class(name, pSvcLocator),
27 m_metaDataSvc("MetaDataSvc", name),
29{
30}
31
32//__________________________________________________________________________
35//__________________________________________________________________________
37 ATH_MSG_DEBUG("Initializing {}", name());
38
39 // Set to be listener for end of event
40 ServiceHandle<IIncidentSvc> incsvc("IncidentSvc", this->name());
41 if (!incsvc.retrieve().isSuccess()) {
42 ATH_MSG_FATAL("Cannot get IncidentSvc.");
43 return(StatusCode::FAILURE);
44 }
45 if( !incidentName().empty() ) {
46 incsvc->addListener(this, incidentName(), 100);
47 incsvc->addListener(this, IncidentType::BeginProcessing, 100);
48 ATH_MSG_DEBUG("Listening to {} incidents", incidentName() );
49 ATH_MSG_DEBUG("Reporting is {}", (m_reportingOn.value()? "ON" : "OFF") );
50 // Retrieve MetaDataSvc
51 if( !m_metaDataSvc.isValid() and !m_metaDataSvc.retrieve().isSuccess() ) {
52 ATH_MSG_ERROR("Cannot get MetaDataSvc");
53 return StatusCode::FAILURE;
54 }
55 }
56
58 ATH_MSG_DEBUG("Concurrent events mode");
59 } else {
60 ATH_MSG_VERBOSE("Sequential events mode");
61 }
62
64
65 return(StatusCode::SUCCESS);
66}
67//__________________________________________________________________________
69 // Release MetaDataSvc
70 if (!m_metaDataSvc.release().isSuccess()) {
71 ATH_MSG_WARNING("Cannot release MetaDataSvc.");
72 }
73 return(StatusCode::SUCCESS);
74}
75
76//__________________________________________________________________________
78 return Gaudi::Concurrency::ConcurrencyFlags::numConcurrentEvents() > 1;
79}
80
81//__________________________________________________________________________
83 std::lock_guard lockg( m_mutex );
84 return m_fileSequenceNumber >= 0;
85}
86
87//__________________________________________________________________________
88void OutputStreamSequencerSvc::handle(const Incident& inc)
89{
90 const EventContext& ctx = inc.context();
91 m_lastIncident = inc.type();
92 ATH_MSG_INFO("Handling incident of type " << m_lastIncident << " for " << ctx);
93
94 if( inc.type() == incidentName() ) { // NextEventRange
95 std::string rangeID;
96 const FileIncident* fileInc = dynamic_cast<const FileIncident*>(&inc);
97 if (fileInc != nullptr) {
98 rangeID = fileInc->fileName();
99 // Handle BeginInputFile
100 if (inc.type() == IncidentType::BeginInputFile) {
101 rangeID = "INFILE";
102 }
104 "Requested (through incident) Next Event Range filename extension: {}",
105 rangeID);
106 }
107
108 if( rangeID == "dummy" ) {
109 if( not inConcurrentEventsMode() ) {
110 // finish the previous Range here only in SEQUENTIAL (threads<2) event processing
111 // Write metadata on the incident finishing a Range (filename=="dummy") in ES MP
112 ATH_MSG_DEBUG("MetaData transition");
113 // immediate write and disconnect for ES, otherwise do it after Event write is done
114 bool disconnect { true };
115 std::lock_guard lockg( m_mutex );
116 if( !m_metaDataSvc->transitionMetaDataFile( m_lastFileName, disconnect ).isSuccess() ) {
117 throw GaudiException("Cannot transition MetaData", name(), StatusCode::FAILURE);
118 }
119 }
120 // exit now, wait for the next (real) incident that will start the next range
121 return;
122 }
123 {
124 // start a new range
125 std::lock_guard lockg( m_mutex );
127 if( rangeID.empty() ) {
128 std::ostringstream n;
129 std::print (n, "_{:04}", m_fileSequenceNumber);
130 rangeID = n.str();
131 ATH_MSG_DEBUG("Default next event range filename extension: {}", rangeID);
132 }
133 else if (rangeID == "INFILE") {
134 rangeID = std::to_string(m_fileSequenceNumber);
135 }
136 // from now on new events will use the new rangeID
137 m_currentRangeID = rangeID;
138 // for ESMT these incidents are asynchronous, so wait for BeginProcessing to update the range map
139 if( not inConcurrentEventsMode() or ctx.valid() ) {
140 *m_rangeIDinSlot.get(ctx) = std::move(rangeID);
141 }
142 }
143 if( not inConcurrentEventsMode() and not fileInc ) {
144 // non-file incident case (filename=="") in regular SP LoopMgr
145 ATH_MSG_DEBUG("MetaData transition");
146 bool disconnect { false };
147 // MN: may not know the full filename yet, but that is only needed for disconnect==true
148 if( !m_metaDataSvc->transitionMetaDataFile( "" /*m_lastFileName*/, disconnect ).isSuccess() ) {
149 throw GaudiException("Cannot transition MetaData", name(), StatusCode::FAILURE);
150 }
151 }
152 }
153 else if( inc.type() == IncidentType::BeginProcessing ) {
154 // new event start - assing current rangeId to its slot
155 std::lock_guard lockg( m_mutex );
156 ATH_MSG_DEBUG("Assigning rangeID = {} to slot {}",
157 m_currentRangeID, ctx.slot());
159 }
160}
161
162//__________________________________________________________________________
163std::string OutputStreamSequencerSvc::buildSequenceFileName(const EventContext& ctx, const std::string& orgFileName)
164{
165 if( !inUse() ) {
166 // Event sequences not in use, just return the original filename
167 return orgFileName;
168 }
169 std::string rangeID = currentRangeID(ctx);
170 std::lock_guard lockg( m_mutex );
171 if (!m_replaceRangeMode) {
172 // build the full output file name for this event range
173 std::string fileNameCore = orgFileName, fileNameExt;
174 std::size_t sepPos = orgFileName.find('[');
175 if (sepPos != std::string::npos) {
176 fileNameCore = orgFileName.substr(0, sepPos);
177 fileNameExt = orgFileName.substr(sepPos);
178 }
179 m_lastFileName = fileNameCore + "." + rangeID + fileNameExt;
180 } else {
181 std::string_view origFileNameView = orgFileName;
182 std::size_t open = origFileNameView.find('[');
183 std::size_t close = origFileNameView.find(']');
184 // If we don't find a [ ] enclosed section, just append the rangeID to the
185 // end
186 if (open == std::string_view::npos || close == std::string_view::npos) {
187 m_lastFileName = std::format("{}.{}", origFileNameView, rangeID);
188 } else {
189 // build list of elems to substitute from
190 ATH_MSG_DEBUG("Building element list");
191 std::vector<std::string_view> elems{};
192 std::size_t pos = open + 1;
193 for (std::size_t comma = origFileNameView.find(',', pos);
194 comma < close;
195 comma = origFileNameView.find(',', pos)) {
196 std::string_view item = origFileNameView.substr(pos, comma - pos);
197 ATH_MSG_DEBUG("(start) pos = {}, (end) comma = comma, item = {}",
198 pos, comma, item);
199 elems.push_back(item);
200 pos = comma + 1;
201 }
202 std::string_view last_item = origFileNameView.substr(pos, close - pos);
203 ATH_MSG_DEBUG("(start) pos = {}, (end) close = {}, item = {}",
204 pos, close, last_item);
205 elems.push_back(last_item);
206 // substitute
207 std::size_t rangeIdx{};
208 auto rangeIdxParseRes = std::from_chars(
209 rangeID.data(), rangeID.data() + rangeID.size(), rangeIdx);
210 if (rangeIdxParseRes.ec != std::errc()) {
212 "Error parsing rangeID to integer. Replacing [] list with "
213 "rangeID.");
215 std::format("{}{}{}", origFileNameView.substr(0, open), rangeID,
216 origFileNameView.substr(close + 1));
217 } else if (rangeIdx >= elems.size()) {
219 "Number of elements in [] list <= rangeID. Replacing [] list with "
220 "rangeID.");
222 std::format("{}{}{}", origFileNameView.substr(0, open), rangeID,
223 origFileNameView.substr(close + 1));
224 } else {
225 m_lastFileName = std::format(
226 "{}{}{}", origFileNameView.substr(0, open), elems.at(rangeIdx),
227 origFileNameView.substr(close + 1));
228 ATH_MSG_DEBUG("Output file: {}", m_lastFileName);
229 }
230 }
231 }
232
233 if( m_reportingOn.value() ) {
234 m_fnToRangeId.insert( std::pair(m_lastFileName, rangeID) );
235 }
236
237 return m_lastFileName;
238}
239
240
241std::string OutputStreamSequencerSvc::currentRangeID(const EventContext& ctx) const
242{
243 if( !inUse() ) return "";
244 return *m_rangeIDinSlot.get(ctx);
245}
246
247
248std::string OutputStreamSequencerSvc::setRangeID(const EventContext& ctx, const std::string & rangeID)
249{
250 std::string* rangeid = m_rangeIDinSlot.get(ctx);
251 const std::string oldrange = *rangeid;
252 *rangeid = rangeID;
253 return oldrange;
254}
255
256
257void OutputStreamSequencerSvc::publishRangeReport(const std::string& outputFile)
258{
259 std::lock_guard lockg( m_mutex );
260 m_finishedRange = m_fnToRangeId.find(outputFile);
261}
262
264{
265 RangeReport_ptr report;
266 if( !m_reportingOn.value() ) {
267 ATH_MSG_WARNING("Reporting not turned on - set {} to True",
268 m_reportingOn.name());
269 } else {
270 std::lock_guard lockg( m_mutex );
271 if(m_finishedRange!=m_fnToRangeId.end()) {
272 report = std::make_unique<RangeReport_t>(m_finishedRange->second,m_finishedRange->first);
275 }
276 }
277 return report;
278}
#define ATH_MSG_DEBUG(x,...)
#define ATH_MSG_ERROR(x,...)
#define ATH_MSG_WARNING(x,...)
#define ATH_MSG_VERBOSE(x,...)
#define ATH_MSG_INFO(x,...)
#define ATH_MSG_FATAL(x,...)
This file contains the class definition for the MetaDataSvc class.
This file contains the class definition for the OutputStreamSequencerSvc class.
static const Attributes_t empty
void publishRangeReport(const std::string &outputFile)
SG::SlotSpecificObj< std::string, SG::InvalidSlot::Enabled > m_rangeIDinSlot
EventRange ID for all slots.
bool inUse() const
Is the service in active use? (true after the first range incident is handled).
OutputStreamSequencerSvc(const std::string &name, ISvcLocator *pSvcLocator)
Standard Service Constructor.
std::string buildSequenceFileName(const EventContext &ctx, const std::string &)
Returns sequenced file name for output stream.
virtual void handle(const Incident &) override final
Incident service handle.
std::string m_lastFileName
Recently constructed full file name (useful in single threaded processing).
std::string currentRangeID(const EventContext &ctx) const
The current Event Range ID (only one range is returned).
ServiceHandle< MetaDataSvc > m_metaDataSvc
int m_fileSequenceNumber
The event sequence number.
BooleanProperty m_replaceRangeMode
Flag to put in ReplaceRangeMode (i.e.
virtual StatusCode finalize() override final
Required of all Gaudi services:
std::string incidentName() const
The name of the incident that starts a new event sequence.
std::string m_currentRangeID
Current EventRange ID constructed on the last NextRange incident.
virtual ~OutputStreamSequencerSvc()
Destructor.
BooleanProperty m_reportingOn
Flag to switch on storage of reporting info in fnToRangeId.
static bool inConcurrentEventsMode()
Are there concurrent events? (threads>1).
std::string m_lastIncident
Last incident type that was handled.
std::string setRangeID(const EventContext &ctx, const std::string &rangeID)
set the RangeID (possibly temporarily) so the right Range Filename may be generated
std::unique_ptr< RangeReport_t > RangeReport_ptr
std::map< std::string, std::string >::iterator m_finishedRange
virtual StatusCode initialize() override final
Required of all Gaudi services:
std::map< std::string, std::string > m_fnToRangeId