ATLAS Offline Software
Loading...
Searching...
No Matches
AthenaOutputStreamTool.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
12// Gaudi
13#include "GaudiKernel/IOpaqueAddress.h"
14#include "GaudiKernel/INamedInterface.h"
15#include "GaudiKernel/IClassIDSvc.h"
16#include "GaudiKernel/ThreadLocalContext.h"
17
18// Athena
22#include "SGTools/DataProxy.h"
23#include "SGTools/SGIFolder.h"
27
28
31 const std::string& name,
32 const IInterface* parent) : base_class(type, name, parent),
33 m_clidSvc("ClassIDSvc", name),
34 m_decSvc("DecisionSvc/DecisionSvc", name) {
35}
36//__________________________________________________________________________
38 ATH_MSG_INFO("Initializing {}", name());
39
40 ATH_CHECK( m_clidSvc.retrieve() );
41 ATH_CHECK( m_conversionSvc.retrieve() );
42
43 // Autoconfigure
44 if (m_dataHeaderKey.empty()) {
45 m_dataHeaderKey.setValue(name());
46 // Remove "ToolSvc." from m_dataHeaderKey.
47 if (m_dataHeaderKey.value().starts_with( "ToolSvc.")) {
48 m_dataHeaderKey.setValue(m_dataHeaderKey.value().substr(8));
49 // Remove "Tool" from m_dataHeaderKey.
50 if (m_dataHeaderKey.value().find("Tool") == m_dataHeaderKey.size() - 4) {
51 m_dataHeaderKey.setValue(m_dataHeaderKey.value().substr(0, m_dataHeaderKey.size() - 4));
52 }
53 } else {
54 const INamedInterface* parentAlg = dynamic_cast<const INamedInterface*>(parent());
55 if (parentAlg != 0) {
56 m_dataHeaderKey.setValue(parentAlg->name());
57 }
58 }
59 }
60 if (m_processTag.empty()) {
62 }
63
64 { // handle the AttrKey overwrite
65 const std::string keyword = "[AttributeListKey=";
66 std::string::size_type pos = m_outputName.value().find(keyword);
67 if( (pos != std::string::npos) ) {
68 ATH_MSG_INFO("The AttrListKey will be overwritten/set by the value from the OutputName: {}", m_outputName);
69 const std::string attrListKey = m_outputName.value().substr(pos + keyword.size(),
70 m_outputName.value().find(']', pos + keyword.size()) - pos - keyword.size());
71 m_attrListKey = attrListKey;
72 }
73 }
74 if ( ! m_attrListKey.key().empty() ) {
75 ATH_CHECK(m_attrListKey.initialize());
76 m_attrListWrite = m_attrListKey.key() + "Decisions";
77 //ATH_CHECK(m_attrListWrite.initialize());
78 }
79
80 if (!m_outputCollection.value().empty()) {
81 m_outputAttributes += "[OutputCollection=" + m_outputCollection.value() + "]";
82 }
83 if (!m_containerPrefix.value().empty()) {
84 m_outputAttributes += "[PoolContainerPrefix=" + m_containerPrefix.value() + "]";
85 }
86 if (m_containerNameHint.value() != "0") {
87 m_outputAttributes += "[TopLevelContainerName=" + m_containerNameHint.value() + "]";
88 }
89 if (m_branchNameHint.value() != "0") {
90 m_outputAttributes += "[SubLevelBranchName=" + m_branchNameHint.value() + "]";
91 }
92 if (!m_metaDataOutputCollection.value().empty()) {
93 m_metaDataOutputAttributes += "[OutputCollection=" + m_metaDataOutputCollection.value() + "]";
94 }
95 if (!m_metaDataContainerPrefix.value().empty()) {
96 m_metaDataOutputAttributes += "[PoolContainerPrefix=" + m_metaDataContainerPrefix.value() + "]";
97 }
98
99 if (m_extend) ATH_CHECK(m_decSvc.retrieve());
100 return(StatusCode::SUCCESS);
101}
102//__________________________________________________________________________
103StatusCode AthenaOutputStreamTool::connectServices(const std::string& dataStore,
104 const std::string& cnvSvc,
105 bool extendProvenenceRecord) {
106 // Release old data store
107 if (m_store.isValid()) {
108 if (m_store.release().isFailure()) {
109 ATH_MSG_ERROR("Could not release {} store", m_store.typeAndName());
110 }
111 }
112 m_store = ServiceHandle<StoreGateSvc>(dataStore, this->name());
113 if (cnvSvc != m_conversionSvc.type() && cnvSvc != "EventPersistencySvc") {
114 if (m_conversionSvc.release().isFailure()) {
115 ATH_MSG_ERROR("Could not release {}", m_conversionSvc.type());
116 }
117 m_conversionSvc = ServiceHandle<IConversionSvc>(cnvSvc, this->name());
118 if (m_conversionSvc.retrieve().isFailure() || m_conversionSvc == 0) {
119 ATH_MSG_ERROR("Could not locate {}", m_conversionSvc.type());
120 return(StatusCode::FAILURE);
121 }
122 }
123 m_extendProvenanceRecord = extendProvenenceRecord;
124 auto pprop = dynamic_cast<const IProperty*>(parent());
125 if (not pprop){
126 ATH_MSG_ERROR("'parent' could not be cast to IProperty");
127 return(StatusCode::FAILURE);
128 }
129 auto keep = dynamic_cast<const StringProperty&>( pprop->getProperty("KeepProvenanceTagsRegEx") );
130 m_keepProvenancesStr = keep.value();
131 // create RegEx pattern from the property value, specify extended grammar
132 m_keepProvenancesRE = std::regex(m_keepProvenancesStr, std::regex::extended);
133
134 return(connectServices());
135}
136//__________________________________________________________________________
138 // Find the data store
139 if (m_store.retrieve().isFailure() || m_store == 0) {
140 ATH_MSG_ERROR("Could not locate {} store", m_store.typeAndName());
141 return(StatusCode::FAILURE);
142 }
143 return(StatusCode::SUCCESS);
144}
145//__________________________________________________________________________
146StatusCode AthenaOutputStreamTool::connectOutput(const std::string& outputName) {
147 ATH_MSG_DEBUG("In connectOutput {}", outputName);
148
149 // Use arg if not empty, save the output name
150 if (!outputName.empty()) {
151 m_outputName.setValue(outputName);
152 }
153 if (m_outputName.value().empty()) {
154 ATH_MSG_ERROR("No OutputName provided");
155 return(StatusCode::FAILURE);
156 }
157 // Connect services if not already available
158 if (m_store == 0 || m_conversionSvc == 0) {
160 }
161 // Connect the output file to the service
162 if (m_conversionSvc->connectOutput(m_outputName.value()).isFailure()) {
163 ATH_MSG_ERROR("Unable to connect output {}", m_outputName.value());
164 return(StatusCode::FAILURE);
165 } else {
166 ATH_MSG_DEBUG("Connected to {}", m_outputName.value());
167 }
168
169 // Remove DataHeader with same key if it exists
170 if (m_store->contains<DataHeader>(m_dataHeaderKey)) {
171 const DataHeader* preDh = nullptr;
172 if (m_store->retrieve(preDh, m_dataHeaderKey).isSuccess()) {
173 if (m_store->removeDataAndProxy(preDh).isFailure()) {
174 ATH_MSG_ERROR("Unable to get proxy for the DataHeader with key " << m_dataHeaderKey);
175 return(StatusCode::FAILURE);
176 }
177 ATH_MSG_DEBUG("Released DataHeader with key " << m_dataHeaderKey);
178 }
179 }
180
181 // Create new DataHeader
182 m_dataHeader = new DataHeader();
183 m_dataHeader->setProcessTag(m_processTag);
184
185 // Retrieve all existing DataHeaders from StoreGate
186 const DataHeader* dh = nullptr;
187 std::vector<std::string> dhKeys;
188 m_store->keys<DataHeader>(dhKeys);
189 //construct string out of loop
190 const std::string boolTypeStr{"bool"};
191 for (const std::string& dhKey : dhKeys) {
192 bool primaryDH = false;
193 if (!m_store->transientContains<DataHeader>(dhKey)) {
194 if (dhKey == "EventSelector") primaryDH = true;
195 ATH_MSG_DEBUG("No transientContains DataHeader with key {}", dhKey);
196 }
197 if (m_store->retrieve(dh, dhKey).isFailure()) {
198 ATH_MSG_DEBUG("Unable to retrieve the DataHeader with key {}", dhKey);
199 }
200 // Propagate provenance from file inputs and the primary event-selector
201 // header. Headers produced earlier in this job are intentionally not
202 // treated as inputs; for example, an AOD written after an ESD in the
203 // same job will not automatically retain provenance back to that ESD.
204 // Revisit this policy if same-job output chaining needs support again.
205 if (dh->isInput() || primaryDH) {
206 propagateProvenance( *dh );
207 }
208 }
209
210 // Attach the attribute list to the DataHeader if requested
211 if (!m_attrListKey.key().empty() && m_store->storeID() == StoreID::EVENT_STORE) {
212 auto attrListHandle = SG::makeHandle(m_attrListKey);
213 if (!attrListHandle.isValid()) {
214 ATH_MSG_WARNING("Unable to retrieve AttributeList with key {}", m_attrListKey);
215 } else {
216 m_dataHeader->setAttributeList(attrListHandle.cptr());
217 if (m_extend) { // Add streaming decisions
218 ATH_MSG_DEBUG("Adding stream decisions to {}", m_attrListWrite);
219 // Look for attribute list created for mini-EventInfo
220 const AthenaAttributeList* attlist(attrListHandle.cptr());
221
222 // Build new attribute list for modification
223 AthenaAttributeList* newone = new AthenaAttributeList(attlist->specification());
224 newone->copyData(*attlist);
225
226 // Now loop over stream definitions and add decisions
227 auto streams = m_decSvc->getStreams();
228 for (auto it = streams.begin();
229 it != streams.end(); ++it) {
230 newone->extend(*it,boolTypeStr);
231 (*newone)[*it].data<bool>() = m_decSvc->isEventAccepted(*it,Gaudi::Hive::currentContext());
232 ATH_MSG_DEBUG("Added stream decision for {} to {}",
233 *it, m_attrListKey);
234 }
235 // record new attribute list with old key + suffix
236 const AthenaAttributeList* attrList2 = nullptr;
238 if (m_store->record(newone,m_attrListWrite).isFailure()) {
239 ATH_MSG_ERROR("Unable to record att list {}", m_attrListWrite);
240 }
241 } else {
242 ATH_MSG_DEBUG("Decisions already added by a different stream");
243 }
244 if (m_store->retrieve(attrList2,m_attrListWrite).isFailure()) {
245 ATH_MSG_ERROR("Unable to record att list {}", m_attrListWrite);
246 } else {
247 m_dataHeader->setAttributeList(attrList2);
248 }
249 } // list extend check
250 } // list retrieve check
251 } // list property check
252
253 // Record DataHeader in StoreGate
255 if (wh.record(std::unique_ptr<DataHeader>(m_dataHeader)).isFailure()) {
256 ATH_MSG_ERROR("Unable to record DataHeader with key " << m_dataHeaderKey);
257 return(StatusCode::FAILURE);
258 } else {
259 ATH_MSG_DEBUG("Recorded DataHeader with key " << m_dataHeaderKey);
260 }
262 // Set flag that connection is open
263 m_connectionOpen = true;
264 return(StatusCode::SUCCESS);
265}
266
267//__________________________________________________________________________
269{
270 // keep track of provenance entries inserted into the new DataHeader
271 std::set<std::string> insertedTags{};
272 // Add DataHeader token to the new DataHeader
274 std::string pTag;
275 std::unique_ptr<SG::TransientAddress> dhTransAddr;
276 for (const DataHeaderElement& dhe : src_dh) {
277 if (dhe.getPrimaryClassID() == ClassID_traits<DataHeader>::ID()) {
278 pTag = dhe.getKey();
279 dhTransAddr.reset( dhe.getAddress( m_conversionSvc->repSvcType() ) );
280 }
281 }
282 // Update dhTransAddr to handle fast merged files.
283 if( auto dhProxy=m_store->proxy(&src_dh); dhProxy && dhProxy->address() ) {
284 DataHeaderElement dhe(dhProxy, dhProxy->address(), pTag);
285 m_dataHeader->insertProvenance(dhe);
286 insertedTags.insert(std::move(pTag));
287 }
288 else if( dhTransAddr ) {
289 DataHeaderElement dhe(dhTransAddr.get(), dhTransAddr->address(), pTag);
290 m_dataHeader->insertProvenance(dhe);
291 insertedTags.insert(std::move(pTag));
292 }
293 }
294
295 // empty regexpr means do not keep any provenance
296 if( !m_keepProvenancesStr.empty() ) {
297 // Each stream tag is written only once in the provenance record
298 // In files where there are multiple entries per stream tag
299 // the record is in reverse, i.e., the latest appears first.
300 // Therefore, only keep the first entry if there are multiple
301 // matches so that we retain the latest one.
302 for(auto iter=src_dh.beginProvenance(), iEnd=src_dh.endProvenance(); iter != iEnd; ++iter) {
303 const auto & currentKey = (*iter).getKey();
304 if( insertedTags.insert(currentKey).second ) {
305 // first prov with that tag. Now check if we want to keep that tag
306 bool keep = false;
307 auto it = m_keepProvenanceMatch.find( currentKey );
308 if( it != m_keepProvenanceMatch.end() ) {
309 keep = it->second;
310 } else {
311 keep = std::regex_search(currentKey, m_keepProvenancesRE);
312 m_keepProvenanceMatch[currentKey] = keep;
313 }
314 if( keep ) {
315 m_dataHeader->insertProvenance(*iter);
316 }
317 }
318 }
319 }
320}
321
322//__________________________________________________________________________
323StatusCode AthenaOutputStreamTool::commitOutput(bool doCommit) {
324 ATH_MSG_DEBUG("In commitOutput");
325 // Connect the output file to the service
326 if (m_conversionSvc->commitOutput(m_outputName.value(), doCommit).isFailure()) {
327 ATH_MSG_ERROR("Unable to commit output {}", m_outputName.value());
328 return(StatusCode::FAILURE);
329 }
330 // Set flag that connection is closed
331 m_connectionOpen = false;
332 return(StatusCode::SUCCESS);
333}
334//__________________________________________________________________________
336 AthCnvSvc* athConversionSvc = dynamic_cast<AthCnvSvc*>(m_conversionSvc.get());
337 if (athConversionSvc != 0) {
338 if (athConversionSvc->disconnectOutput(m_outputName.value()).isFailure()) {
339 ATH_MSG_ERROR("Unable to finalize output {}", m_outputName.value());
340 return(StatusCode::FAILURE);
341 }
342 }
343 return(StatusCode::SUCCESS);
344}
345//__________________________________________________________________________
346StatusCode AthenaOutputStreamTool::streamObjects(const TypeKeyPairs& typeKeys, const std::string& outputName) {
347 ATH_MSG_DEBUG("In streamObjects");
348 // Check that a connection has been opened
349 if (!m_connectionOpen) {
350 ATH_MSG_ERROR("Connection NOT open. Please open a connection before streaming out objects.");
351 return(StatusCode::FAILURE);
352 }
353 // Use arg if not empty, save the output name
354 if (!outputName.empty()) {
355 m_outputName.setValue(outputName);
356 }
357 if (m_outputName.value().empty()) {
358 ATH_MSG_ERROR("No OutputName provided");
359 return(StatusCode::FAILURE);
360 }
361 // Now iterate over the type/key pairs and stream out each object
362 std::vector<DataObject*> dataObjects;
363 for (TypeKeyPairs::const_iterator first = typeKeys.begin(), last = typeKeys.end();
364 first != last; ++first) {
365 const std::string& type = (*first).first;
366 const std::string& key = (*first).second;
367 // Find the clid for type name from the CLIDSvc
368 CLID clid;
369 if (m_clidSvc->getIDOfTypeName(type, clid).isFailure()) {
370 ATH_MSG_ERROR("Could not get clid for typeName {}", type);
371 return(StatusCode::FAILURE);
372 }
373 DataObject* dObj = 0;
374 // Two options: no key or explicit key
375 if (key.empty()) {
376 ATH_MSG_DEBUG("Get data object with no key");
377 // Get DataObject without key
378 dObj = m_store->accessData(clid);
379 } else {
380 ATH_MSG_DEBUG("Get data object with key");
381 // Get DataObjects with key
382 dObj = m_store->accessData(clid, key);
383 }
384 if (dObj == 0) {
385 // No object - print warning and return
386 ATH_MSG_DEBUG("No object found for type {} key {}", type, key);
387 return(StatusCode::SUCCESS);
388 } else {
389 ATH_MSG_DEBUG("Found object for type {} key {}", type, key);
390 }
391 // Save the dObj
392 dataObjects.push_back(dObj);
393 }
394 // Stream out objects
395 if (dataObjects.size() == 0) {
396 ATH_MSG_DEBUG("No data objects found");
397 return(StatusCode::SUCCESS);
398 }
399 StatusCode status = streamObjects(dataObjects, m_outputName.value());
400 if (!status.isSuccess()) {
401 ATH_MSG_ERROR("Could not stream out objects");
402 return(status);
403 }
404 return(StatusCode::SUCCESS);
405}
406//__________________________________________________________________________
407StatusCode AthenaOutputStreamTool::streamObjects(const DataObjectVec& dataObjects, const std::string& outputName) {
408 // Check that a connection has been opened
409 if (!m_connectionOpen) {
410 ATH_MSG_ERROR("Connection NOT open. Please open a connection before streaming out objects.");
411 return(StatusCode::FAILURE);
412 }
413 // Connect the output file to the service
414 std::string outputConnectionString = outputName;
415 const std::string defaultMetaDataString = "[OutputCollection=MetaDataHdr][PoolContainerPrefix=MetaData]";
416 if (std::string::size_type mpos = outputConnectionString.find(defaultMetaDataString); mpos!=std::string::npos) {
417 // If we're in here we're writing MetaData
418 // Now let's see if we should be overwriting the MetaData attributes
419 // For the time-being this happens when we're writing MetaData in the augmentation mode
420 if (!m_metaDataOutputAttributes.empty()) {
421 // Note: This won't work quite right if only one attribute is set though!
422 outputConnectionString.replace(mpos, defaultMetaDataString.length(), m_metaDataOutputAttributes);
423 }
424 }
425 for (std::string::size_type pos = m_outputAttributes.find('['); pos != std::string::npos; pos = m_outputAttributes.find('[', ++pos)) {
426 if (outputConnectionString.find(m_outputAttributes.substr(pos, m_outputAttributes.find('=', pos) + 1 - pos)) == std::string::npos) {
427 outputConnectionString += m_outputAttributes.substr(pos, m_outputAttributes.find(']', pos) + 1 - pos);
428 }
429 }
430
431 // Check that the DataHeader is still valid
432 DataObject* dataHeaderObj = m_store->accessData(ClassID_traits<DataHeader>::ID(), m_dataHeaderKey);
433 std::map<DataObject*, IOpaqueAddress*> written;
434 for (DataObject* dobj : dataObjects) {
435 // Do not write the DataHeader via the explicit list
436 if (dobj->clID() == ClassID_traits<DataHeader>::ID()) {
437 ATH_MSG_DEBUG("Explicit request to write DataHeader: {} - skipping it.",
438 dobj->name());
439 // Do not stream out same object twice
440 } else if (written.find(dobj) != written.end()) {
441 // Print warning and skip
442 ATH_MSG_DEBUG("Trying to write DataObject twice (clid/key): {} {}",
443 dobj->clID(), dobj->name());
444 ATH_MSG_DEBUG(" Skipping this one.");
445 } else {
446 // Write object
447 IOpaqueAddress* addr = new TokenAddress(0, dobj->clID(), outputConnectionString);
448 addr->addRef();
449 if (m_conversionSvc->createRep(dobj, addr).isSuccess()) {
450 written.insert(std::pair<DataObject*, IOpaqueAddress*>(dobj, addr));
451 } else {
452 ATH_MSG_ERROR("Could not create Rep for DataObject (clid/key):{} {}",
453 dobj->clID(), dobj->name());
454 return(StatusCode::FAILURE);
455 }
456 }
457 }
458 // End of loop over DataObjects, write DataHeader
459 if ((m_conversionSvc.type() == "AthenaPoolCnvSvc" || m_conversionSvc.type() == "AthenaPoolSharedIOCnvSvc") && dataHeaderObj != nullptr) {
460 IOpaqueAddress* addr = new TokenAddress(0, dataHeaderObj->clID(), outputConnectionString);
461 addr->addRef();
462 if (m_conversionSvc->createRep(dataHeaderObj, addr).isSuccess()) {
463 written.insert(std::pair<DataObject*, IOpaqueAddress*>(dataHeaderObj, addr));
464 } else {
465 ATH_MSG_ERROR("Could not create Rep for DataHeader");
466 return(StatusCode::FAILURE);
467 }
468 }
469 for (DataObject* dobj : dataObjects) {
470 // call fillRepRefs of persistency service
471 SG::DataProxy* proxy = dynamic_cast<SG::DataProxy*>(dobj->registry());
472 if (proxy != nullptr && written.find(dobj) != written.end()) {
473 IOpaqueAddress* addr(written.find(dobj)->second);
474 if ((m_conversionSvc->fillRepRefs(addr, dobj)).isSuccess()) {
475 if (dobj->clID() != 1 || addr->par()[0] != "\n") {
476 if (dobj->clID() != ClassID_traits<DataHeader>::ID()) {
477 m_dataHeader->insert(proxy, addr);
478 } else {
479 m_dataHeader->insert(proxy, addr, m_processTag);
480 }
481 if (proxy->address() == nullptr) {
482 proxy->setAddress(addr);
483 }
484 addr->release();
485 }
486 } else {
487 ATH_MSG_ERROR("Could not fill Object Refs for DataObject (clid/key):{} {}",
488 dobj->clID(), dobj->name());
489 return(StatusCode::FAILURE);
490 }
491 } else {
492 ATH_MSG_WARNING("Could cast DataObject {} {}",
493 dobj->clID(), dobj->name());
494 }
495 }
496 m_dataHeader->addHash(&*m_store);
497 if ((m_conversionSvc.type() == "AthenaPoolCnvSvc" || m_conversionSvc.type() == "AthenaPoolSharedIOCnvSvc") && dataHeaderObj != nullptr) {
498 // End of DataObjects, fill refs for DataHeader
499 SG::DataProxy* proxy = dynamic_cast<SG::DataProxy*>(dataHeaderObj->registry());
500 if (proxy != nullptr && written.find(dataHeaderObj) != written.end()) {
501 IOpaqueAddress* addr(written.find(dataHeaderObj)->second);
502 if ((m_conversionSvc->fillRepRefs(addr, dataHeaderObj)).isSuccess()) {
503 if (dataHeaderObj->clID() != 1 || addr->par()[0] != "\n") {
504 if (dataHeaderObj->clID() != ClassID_traits<DataHeader>::ID()) {
505 m_dataHeader->insert(proxy, addr);
506 } else {
507 m_dataHeader->insert(proxy, addr, m_processTag);
508 }
509 addr->release();
510 }
511 } else {
512 ATH_MSG_ERROR("Could not fill Object Refs for DataHeader");
513 return(StatusCode::FAILURE);
514 }
515 } else {
516 ATH_MSG_ERROR("Could not cast DataHeader");
517 return(StatusCode::FAILURE);
518 }
519 }
520 return(StatusCode::SUCCESS);
521}
522//__________________________________________________________________________
524 const std::string hltKey = "HLTAutoKey";
527 if (m_store->retrieve(beg, ending).isFailure() || beg == ending) {
528 ATH_MSG_DEBUG("No DataHeaders present in StoreGate");
529 } else {
530 for ( ; beg != ending; ++beg) {
531 if (m_store->transientContains<DataHeader>(beg.key()) && beg->isInput()) {
532 for (std::vector<DataHeaderElement>::const_iterator it = beg->begin(), itLast = beg->end();
533 it != itLast; ++it) {
534 // Only insert the primary clid, not the ones for the symlinks!
535 CLID clid = it->getPrimaryClassID();
536 if (clid != ClassID_traits<DataHeader>::ID()) {
537 //check the typename is known ... we make an exception if the key contains 'Aux.' ... aux containers may not have their keys known yet in some cases
538 //see https://its.cern.ch/jira/browse/ATLASG-59 for the solution
539 std::string typeName;
540 if (m_clidSvc->getTypeNameOfID(clid, typeName).isFailure() && it->getKey().find("Aux.") == std::string::npos) {
541 if (m_skippedItems.find(it->getKey()) == m_skippedItems.end()) {
542 ATH_MSG_WARNING("Skipping {} with unknown clid {} . Further warnings for this item are suppressed",
543 it->getKey(), clid);
544 m_skippedItems.insert(it->getKey());
545 }
546 continue;
547 }
548 ATH_MSG_DEBUG("Adding {}#{} (clid {}) to itemlist",
549 typeName, it->getKey(), clid);
550 const std::string keyName = it->getKey();
551 if (keyName.size() > 10 && keyName.compare(0, 10,hltKey)==0) {
552 p2BWrittenFromTool->add(clid, hltKey + "*").ignore();
553 } else if (keyName.size() > 10 && keyName.compare(keyName.size() - 10, 10, hltKey)==0) {
554 p2BWrittenFromTool->add(clid, "*" + hltKey).ignore();
555 } else {
556 p2BWrittenFromTool->add(clid, keyName).ignore();
557 }
558 }
559 }
560 }
561 }
562 }
563 ATH_MSG_DEBUG("Adding DataHeader for stream {}", name());
564 return(StatusCode::SUCCESS);
565}
#define ATH_CHECK
Evaluate an expression and check for errors.
#define ATH_MSG_DEBUG(x,...)
#define ATH_MSG_ERROR(x,...)
#define ATH_MSG_WARNING(x,...)
#define ATH_MSG_INFO(x,...)
This is the implementation of IAthenaOutputStreamTool.
This file contains the class definition for the DataHeader and DataHeaderElement classes.
uint32_t CLID
The Class ID type.
An AttributeList represents a logical row of attributes in a metadata table.
This file contains the class definition for the TokenAddress class.
Base class for all conversion services.
Definition AthCnvSvc.h:66
virtual StatusCode disconnectOutput(const std::string &output)
Disconnect output files from the service.
An AttributeList represents a logical row of attributes in a metadata table.
StringProperty m_metaDataContainerPrefix
std::map< std::string, bool > m_keepProvenanceMatch
Cache provenance RegEx matching result in a map.
StatusCode finalizeOutput() override
Finalize the output stream after the last commit, e.g.
AthenaOutputStreamTool(const std::string &type, const std::string &name, const IInterface *parent)
Standard AlgTool Constructor.
std::set< std::string > m_skippedItems
set of skipped item keys, because of missing CLID
bool m_extendProvenanceRecord
Flag as to whether to extend provenance via the DataHeader.
DataHeader * m_dataHeader
Current DataHeader for streamed objects.
std::vector< DataObject * > DataObjectVec
Stream out a vector of objects Must convert to DataObject, e.g.
StatusCode connectOutput(const std::string &outputName="") override
Connect to the output stream Must connectOutput BEFORE streaming Only specify "outputName" if one wan...
ServiceHandle< IDecisionSvc > m_decSvc
Ref to DecisionSvc.
ServiceHandle< IClassIDSvc > m_clidSvc
Ref to ClassIDSvc to convert type name to clid.
std::vector< TypeKeyPair > TypeKeyPairs
virtual StatusCode streamObjects(const TypeKeyPairs &typeKeys, const std::string &outputName="") override
bool m_connectionOpen
Flag to tell whether connectOutput has been called.
ServiceHandle< IConversionSvc > m_conversionSvc
Keep reference to the data conversion service.
std::regex m_keepProvenancesRE
RegEx pattern created from m_keepProvenancesStr.
StatusCode connectServices()
Do the real connection to services.
void propagateProvenance(const DataHeader &src_dh)
copy provenance records when creating new DataHeaders
ServiceHandle< StoreGateSvc > m_store
virtual StatusCode initialize() override
AthAlgTool Interface method implementations:
virtual StatusCode getInputItemList(SG::IFolder *m_p2BWrittenFromTool) override
std::string m_keepProvenancesStr
RegEx string to match provenance tags to keep in the output DataHeader. Retrieved from an OutputStrea...
SG::ReadHandleKey< AthenaAttributeList > m_attrListKey
StatusCode commitOutput(bool doCommit=false) override
Commit the output stream after having streamed out objects Must commitOutput AFTER streaming.
StringProperty m_metaDataOutputCollection
This class provides a persistent form for the TransientAddress.
Definition DataHeader.h:37
This class provides the layout for summary information stored for data written to POOL.
Definition DataHeader.h:123
std::vector< DataHeaderElement >::const_iterator beginProvenance() const
std::vector< DataHeaderElement >::const_iterator endProvenance() const
bool isInput() const
Check whether StatusFlag is "Input".
a const_iterator facade to DataHandle.
Definition SGIterator.h:164
a run-time configurable list of data objects
Definition SGIFolder.h:21
virtual StatusCode add(const std::string &typeName, const std::string &skey)=0
add a data object identifier to the list
@ EVENT_STORE
Definition StoreID.h:26
This class provides a Generic Transient Address for POOL tokens.
SG::ReadCondHandle< T > makeHandle(const SG::ReadCondHandleKey< T > &key, const EventContext &ctx=Gaudi::Hive::currentContext())