24#include "GaudiKernel/ClassID.h"
25#include "GaudiKernel/FileIncident.h"
26#include "GaudiKernel/IIncidentSvc.h"
27#include "GaudiKernel/IIoComponentMgr.h"
28#include "GaudiKernel/GaudiException.h"
29#include "GaudiKernel/GenericAddress.h"
30#include "GaudiKernel/StatusCode.h"
45 base_class(name, pSvcLocator)
57 m_inputCollectionsChanged =
false;
61 if (this->FSMState() != Gaudi::StateMachine::OFFLINE) {
62 m_inputCollectionsChanged =
true;
75 m_autoRetrieveTools =
false;
76 m_checkToolDeps =
false;
79 ATH_MSG_DEBUG(
"Initializing secondary event selector " << name());
87 ATH_MSG_FATAL(
"Use the property: EventSelector.InputCollections = "
88 <<
"[ \"<collectionName>\" ] (list of collections)");
89 return StatusCode::FAILURE;
92 for(
const std::string&
r: ranges ) {
94 auto from_iter = fromto.begin();
95 long from = std::stol(*from_iter);
97 if( ++from_iter != fromto.end() ) {
98 to = std::stol(*from_iter);
100 m_skipEventRanges.emplace_back(from, to);
104 m_skipEventRanges.emplace_back(v, v);
106 std::sort(m_skipEventRanges.begin(), m_skipEventRanges.end());
107 if( msgLvl(MSG::DEBUG) ) {
108 std::string skip_ranges_str;
109 for(
const auto& [first, second] : m_skipEventRanges ) {
110 if( !skip_ranges_str.empty() ) skip_ranges_str +=
", ";
111 skip_ranges_str += std::to_string(first);
112 if( first != second) skip_ranges_str += std::format(
"-{}", second);
114 if( !skip_ranges_str.empty() )
130 std::vector<std::string> propVal;
132 bool foundCnvSvc =
false;
133 for (
const auto& property : propVal) {
138 if (!epSvc->setProperty(
"CnvServices", Gaudi::Utils::toString(propVal)).isSuccess()) {
139 ATH_MSG_FATAL(
"Cannot set EventPersistencySvc Property for CnvServices");
140 return StatusCode::FAILURE;
151 for (
const auto& inputCollection : incol) {
152 if (!iomgr->io_register(
this, IIoComponentMgr::IoMode::READ, inputCollection, inputCollection).isSuccess()) {
153 ATH_MSG_FATAL(
"could not register [" << inputCollection <<
"] for output !");
156 ATH_MSG_VERBOSE(
"io_register[" << this->name() <<
"](" << inputCollection <<
") [ok]");
160 return StatusCode::FAILURE;
166 return StatusCode::FAILURE;
170 if (!
reinit().isSuccess()) {
171 return StatusCode::FAILURE;
176 m_incidentSvc->addListener(
this, IncidentType::BeginProcessing, 0);
177 m_incidentSvc->addListener(
this, IncidentType::EndProcessing, 0);
178 return StatusCode::SUCCESS;
191 if (!m_firstEvt.empty()) {
194 m_inputCollectionsChanged =
false;
196 m_headerIterator = 0;
197 bool retError =
false;
198 for (
auto& tool : m_helperTools) {
199 if (!tool->postInitialize().isSuccess()) {
200 ATH_MSG_FATAL(
"Failed to postInitialize() " << tool->name());
206 return StatusCode::FAILURE;
211 if (!m_poolCollectionConverter) {
212 ATH_MSG_INFO(
"No Events found in any Input Collections");
218 FileIncident firstInputFileIncident(name(),
"FirstInputFile", *m_inputCollectionsIterator);
223 return StatusCode::SUCCESS;
227 m_headerIterator = m_poolCollectionConverter->selectAll();
228 }
catch (std::exception &e) {
229 ATH_MSG_FATAL(
"Cannot open input collection - check data/software version.");
231 return StatusCode::FAILURE;
233 while (m_headerIterator ==
nullptr || m_headerIterator->next() == 0) {
234 if (m_poolCollectionConverter) {
235 m_poolCollectionConverter->disconnectDb().ignore();
236 m_poolCollectionConverter.reset();
238 ++m_inputCollectionsIterator;
240 if (m_poolCollectionConverter) {
241 m_headerIterator = m_poolCollectionConverter->selectAll();
246 if (!m_poolCollectionConverter || m_headerIterator ==
nullptr) {
250 if (!m_poolCollectionConverter) {
251 return StatusCode::SUCCESS;
253 m_headerIterator = m_poolCollectionConverter->selectAll();
254 while (m_headerIterator ==
nullptr || m_headerIterator->next() == 0) {
255 if (m_poolCollectionConverter) {
256 m_poolCollectionConverter->disconnectDb().ignore();
257 m_poolCollectionConverter.reset();
259 ++m_inputCollectionsIterator;
261 if (m_poolCollectionConverter) {
262 m_headerIterator = m_poolCollectionConverter->selectAll();
268 if (!m_poolCollectionConverter || m_headerIterator ==
nullptr) {
269 return StatusCode::SUCCESS;
271 const Token& headRef = m_headerIterator->eventRef();
278 FileIncident firstInputFileIncident(name(),
"FirstInputFile",
"FID:" + fid, fid);
282 return StatusCode::SUCCESS;
286 if (m_poolCollectionConverter) {
288 m_poolCollectionConverter->disconnectDb().ignore();
289 m_poolCollectionConverter.reset();
294 if (!m_poolCollectionConverter) {
295 ATH_MSG_INFO(
"No Events found in any Input Collections");
298 --m_inputCollectionsIterator;
301 m_headerIterator = m_poolCollectionConverter->selectAll();
307 return StatusCode::SUCCESS;
313 m_inputFileGuard.reset();
315 IEvtSelector::Context* ctxt(
nullptr);
319 return StatusCode::SUCCESS;
327 for (
auto& tool : m_helperTools) {
328 if (!tool->preFinalize().isSuccess()) {
333 m_headerIterator =
nullptr;
334 if (m_poolCollectionConverter) {
335 m_poolCollectionConverter.reset();
338 return ::AthService::finalize();
344 return StatusCode::SUCCESS;
348 std::lock_guard<CallMutex> lockGuard(
m_callLock);
349 for (
const auto& tool : m_helperTools) {
350 if (!tool->preNext().isSuccess()) {
357 if (
sc.isRecoverable()) {
360 if (
sc.isFailure()) {
361 return StatusCode::FAILURE;
369 && (m_skipEventRanges.empty() ||
m_evtCount < m_skipEventRanges.front().first))
374 return StatusCode::FAILURE;
377 StatusCode status = StatusCode::SUCCESS;
378 for (
const auto& tool : m_helperTools) {
379 StatusCode toolStatus = tool->postNext();
380 if (toolStatus.isRecoverable()) {
381 ATH_MSG_INFO(
"Request skipping event from: " << tool->name());
382 if (status.isSuccess()) {
383 status = StatusCode::RECOVERABLE;
385 }
else if (toolStatus.isFailure()) {
387 status = StatusCode::FAILURE;
390 if (status.isRecoverable()) {
392 }
else if (status.isFailure()) {
401 while( !m_skipEventRanges.empty() &&
m_evtCount >= m_skipEventRanges.front().second ) {
402 m_skipEventRanges.erase(m_skipEventRanges.begin());
407 return StatusCode::SUCCESS;
412 for (
int i = 0; i < jump; i++) {
415 return StatusCode::SUCCESS;
417 return StatusCode::FAILURE;
422 if( m_inputCollectionsChanged ) {
424 if(
rc != StatusCode::SUCCESS )
return rc;
428 if (m_headerIterator ==
nullptr || m_headerIterator->next() == 0) {
429 m_headerIterator =
nullptr;
431 m_poolCollectionConverter.reset();
434 m_inputFileGuard.reset();
440 if( m_inputCollectionsChanged ) {
442 if(
rc != StatusCode::SUCCESS )
return rc;
445 ++m_inputCollectionsIterator;
448 if (!m_poolCollectionConverter) {
452 return StatusCode::FAILURE;
455 m_headerIterator = m_poolCollectionConverter->selectAll();
458 return StatusCode::RECOVERABLE;
462 const Token& headRef = m_headerIterator->eventRef();
466 if (guid != m_guid) {
476 m_activeEventsPerSource[guid.toString()] = 0;
477 if (!
m_athenaPoolCnvSvc->setInputAttributes(*m_inputCollectionsIterator).isSuccess()) {
479 return StatusCode::FAILURE;
483 *m_inputCollectionsIterator, m_guid.toString(),
487 return StatusCode::SUCCESS;
496 if (
sc.isRecoverable()) {
499 if (
sc.isFailure()) {
500 return StatusCode::FAILURE;
510 && (m_skipEventRanges.empty() ||
m_evtCount < m_skipEventRanges.front().first))
512 return StatusCode::SUCCESS;
514 while( !m_skipEventRanges.empty() &&
m_evtCount >= m_skipEventRanges.front().second ) {
515 m_skipEventRanges.erase(m_skipEventRanges.begin());
525 return StatusCode::SUCCESS;
530 return StatusCode::FAILURE;
535 for (
int i = 0; i < jump; i++) {
538 return StatusCode::SUCCESS;
540 return StatusCode::FAILURE;
544 if (ctxt.identifier() ==
m_endIter->identifier()) {
546 return StatusCode::SUCCESS;
548 return StatusCode::FAILURE;
554 return StatusCode::SUCCESS;
558 IOpaqueAddress*& iop)
const {
559 std::string tokenStr;
563 tokenStr = (*attrList)[
"eventRef"].data<std::string>();
564 ATH_MSG_DEBUG(
"found AthenaAttribute, name = eventRef = " << tokenStr);
565 }
catch (std::exception &e) {
567 return StatusCode::FAILURE;
571 tokenStr = m_headerIterator->eventRef().toString();
573 auto token = std::make_unique<Token>();
574 token->fromString(tokenStr);
575 m_incidentSvc->fireIncident(Incident(tokenStr,
"ProcessEventAttributes"));
577 return StatusCode::SUCCESS;
581 return StatusCode::SUCCESS;
585 IEvtSelector::Context& )
const {
586 return StatusCode::SUCCESS;
591 if( m_inputCollectionsChanged ) {
593 if(
rc != StatusCode::SUCCESS )
return rc;
601 m_headerIterator =
nullptr;
603 m_inputFileGuard.reset();
604 return StatusCode::RECOVERABLE;
608 m_poolCollectionConverter->disconnectDb().ignore();
610 m_poolCollectionConverter.reset();
615 <<
"\" from the collection list.");
619 m_poolCollectionConverter = std::make_unique<PoolCollectionConverter>(
623 if (!m_poolCollectionConverter || !m_poolCollectionConverter->initialize().isSuccess()) {
624 m_headerIterator =
nullptr;
625 ATH_MSG_ERROR(
"seek: Unable to initialize PoolCollectionConverter.");
626 return StatusCode::FAILURE;
629 m_headerIterator = m_poolCollectionConverter->selectAll();
632 next(*beginIter).ignore();
633 ATH_MSG_DEBUG(
"Token " << m_headerIterator->eventRef().toString());
634 }
catch (std::exception &e) {
635 m_headerIterator =
nullptr;
637 return StatusCode::FAILURE;
641 if (m_headerIterator->seek(evtNum - m_firstEvt[
m_curCollection]) == 0) {
642 m_headerIterator =
nullptr;
644 return StatusCode::FAILURE;
648 return StatusCode::SUCCESS;
660 for (std::size_t i = 0,
imax = m_numEvt.size(); i <
imax; i++) {
661 if (m_numEvt[i] == -1) {
669 int collection_size = 0;
671 std::unique_ptr<pool::ICollectionCursor> hi = pcc.
selectAll();
672 collection_size = hi->size();
678 m_firstEvt[i] = m_firstEvt[i - 1] + m_numEvt[i - 1];
682 m_numEvt[i] = collection_size;
684 if (evtNum >= m_firstEvt[i] && evtNum < m_firstEvt[i] + m_numEvt[i]) {
695 return std::accumulate(m_numEvt.begin(), m_numEvt.end(), 0);
698std::unique_ptr<PoolCollectionConverter>
706 ATH_MSG_DEBUG(
"Try item: \"" << *m_inputCollectionsIterator <<
"\" from the collection list.");
707 auto pCollCnv = std::make_unique<PoolCollectionConverter>(
708 *m_inputCollectionsIterator,
711 StatusCode status = pCollCnv->initialize();
712 if (!status.isSuccess()) {
715 if (!status.isRecoverable()) {
716 ATH_MSG_ERROR(
"Unable to initialize PoolCollectionConverter.");
717 throw GaudiException(
"Unable to read: " + *m_inputCollectionsIterator, name(), StatusCode::FAILURE);
719 ATH_MSG_ERROR(
"Unable to open: " << *m_inputCollectionsIterator);
720 throw GaudiException(
"Unable to open: " + *m_inputCollectionsIterator, name(), StatusCode::FAILURE);
723 if (!pCollCnv->isValid().isSuccess()) {
725 ATH_MSG_DEBUG(
"No events found in: " << *m_inputCollectionsIterator <<
" skipped!!!");
729 *m_inputCollectionsIterator, {},
730 "eventless " + *m_inputCollectionsIterator);
732 m_poolSvc->disconnectDb(*m_inputCollectionsIterator).ignore();
733 ++m_inputCollectionsIterator;
743 if (!
eventStore()->clearStore().isSuccess()) {
749 const coral::AttributeList& attrList = m_headerIterator->currentRow().attributeList();
756 ATH_CHECK(wh.record(std::move(athAttrList)));
757 return StatusCode::SUCCESS;
762 const auto& row = m_headerIterator->currentRow();
763 attrList->extend( row.tokenName() + suffix,
"string" );
764 (*attrList)[ row.tokenName() + suffix ].data<std::string>() = row.token().toString();
765 ATH_MSG_DEBUG(
"record AthenaAttribute, name = " << row.tokenName() + suffix <<
" = " << row.token().toString() <<
".");
767 std::string eventRef =
"eventRef";
769 eventRef.append(suffix);
771 attrList->extend(eventRef,
"string");
772 (*attrList)[eventRef].data<std::string>() = m_headerIterator->eventRef().toString();
773 ATH_MSG_DEBUG(
"record AthenaAttribute, name = " + eventRef +
" = " << m_headerIterator->eventRef().toString() <<
".");
776 const coral::AttributeList& sourceAttrList = m_headerIterator->currentRow().attributeList();
777 for (
const auto &attr : sourceAttrList) {
778 attrList->extend(attr.specification().name() + suffix, attr.specification().type());
779 (*attrList)[attr.specification().name() + suffix] = attr;
783 return StatusCode::SUCCESS;
788 if (m_poolCollectionConverter) {
789 m_poolCollectionConverter->disconnectDb().ignore();
790 m_poolCollectionConverter.reset();
792 m_headerIterator =
nullptr;
794 if (!iomgr.retrieve().isSuccess()) {
796 return StatusCode::FAILURE;
798 if (!iomgr->io_hasitem(
this)) {
799 ATH_MSG_FATAL(
"IoComponentMgr does not know about myself !");
800 return StatusCode::FAILURE;
803 std::set<std::size_t> updatedIndexes;
805 if (updatedIndexes.find(i) != updatedIndexes.end())
continue;
806 std::string savedName = inputCollections[i];
807 std::string &fname = inputCollections[i];
808 if (!iomgr->io_contains(
this, fname)) {
809 ATH_MSG_ERROR(
"IoComponentMgr does not know about [" << fname <<
"] !");
810 return StatusCode::FAILURE;
812 if (!iomgr->io_retrieve(
this, fname).isSuccess()) {
813 ATH_MSG_FATAL(
"Could not retrieve new value for [" << fname <<
"] !");
814 return StatusCode::FAILURE;
816 updatedIndexes.insert(i);
817 for (std::size_t j = i + 1; j <
imax; j++) {
818 if (inputCollections[j] == savedName) {
819 inputCollections[j] = fname;
820 updatedIndexes.insert(j);
833 m_inputFileGuard.reset();
834 if (m_poolCollectionConverter) {
835 m_poolCollectionConverter->disconnectDb().ignore();
836 m_poolCollectionConverter.reset();
838 return StatusCode::SUCCESS;
850 if (inc.type() == IncidentType::BeginProcessing) {
861 ATH_MSG_WARNING(
"could not read event source ID from incident event context");
864 if( m_activeEventsPerSource.find( fid ) == m_activeEventsPerSource.end()) {
865 ATH_MSG_DEBUG(
"Incident handler ignoring unknown input FID: " << fid );
868 ATH_MSG_DEBUG(
"** MN Incident handler " << inc.type() <<
" Event source ID=" << fid );
869 if( inc.type() == IncidentType::BeginProcessing ) {
871 m_activeEventsPerSource[fid]++;
872 }
else if( inc.type() == IncidentType::EndProcessing ) {
873 m_activeEventsPerSource[fid]--;
877 if( msgLvl(MSG::DEBUG) ) {
878 for(
auto& source: m_activeEventsPerSource )
879 msg(MSG::DEBUG) <<
"SourceID: " << source.first <<
" active events: " << source.second <<
endmsg;
890 if( m_activeEventsPerSource.find(fid) != m_activeEventsPerSource.end()
891 && m_activeEventsPerSource[fid] <= 0 && m_guid != fid ) {
897 m_activeEventsPerSource.erase( fid );
#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_VERBOSE(x,...)
#define ATH_MSG_INFO(x,...)
#define ATH_MSG_FATAL(x,...)
This file contains the class definition for the EventContextAthenaPool class.
This file contains the class definition for the EventSelectorAthenaPool class.
This file contains the class definition for the IPoolSvc interface class.
This file contains the class definition for the PoolCollectionConverter class.
Handle class for reading from StoreGate.
Handle class for recording to StoreGate.
This file contains the class definition for the TokenAddress class.
This file contains the class definition for the Token class (migrated from POOL).
An AttributeList represents a logical row of attributes in a metadata table.
This class provides the context to access an event from POOL persistent store.
EventSelectorAthenaPool(const std::string &name, ISvcLocator *pSvcLocator)
Standard Service Constructor.
Gaudi::CheckedProperty< uint64_t > m_firstEventNo
virtual StatusCode start() override
std::unique_ptr< PoolCollectionConverter > getCollectionCnv(bool throwIncidents=false) const
Return pointer to new PoolCollectionConverter.
Gaudi::Property< bool > m_isSecondary
IsSecondary, know if this is an instance of secondary event selector.
virtual StatusCode nextWithSkip(IEvtSelector::Context &ctxt) const override
Go to next event and skip if necessary.
virtual StatusCode initialize() override
Required of all Gaudi Services.
Gaudi::Property< int > m_skipEvents
SkipEvents, numbers of events to skip: default = 0.
Gaudi::Property< bool > m_processMetadata
ProcessMetadata, switch on firing of FileIncidents which will trigger processing of metadata: default...
virtual int curEvent(const Context &ctxt) const override
Return the current event number.
virtual StatusCode createContext(IEvtSelector::Context *&ctxt) const override
create context
virtual StatusCode io_finalize() override
Callback method to finalize the internal state of the component for I/O purposes (e....
Gaudi::Property< std::vector< long > > m_skipEventSequenceProp
Gaudi::CheckedProperty< uint32_t > m_initTimeStamp
void inputCollectionsHandler(Gaudi::Details::PropertyBase &)
virtual StatusCode io_reinit() override
Callback method to reinitialize the internal state of the component for I/O purposes (e....
virtual StatusCode resetCriteria(const std::string &criteria, IEvtSelector::Context &ctxt) const override
Set a selection criteria.
StatusCode reinit() const
Reinitialize the service when a fork() occurred/was-issued.
StoreGateSvc * eventStore() const
Return pointer to active event SG.
virtual StatusCode releaseContext(IEvtSelector::Context *&ctxt) const override
virtual StatusCode stop() override
ServiceHandle< IIncidentSvc > m_incidentSvc
virtual StatusCode next(IEvtSelector::Context &ctxt) const override
virtual ~EventSelectorAthenaPool()
Destructor.
ToolHandle< IAthenaSelectorTool > m_counterTool
ServiceHandle< IPoolSvc > m_poolSvc
virtual StatusCode fillAttributeList(coral::AttributeList *attrList, const std::string &suffix, bool copySource) const override
Fill AttributeList with specific items from the selector and a suffix.
std::string m_attrListKey
AttributeList SG key.
virtual StatusCode seek(Context &ctxt, int evtnum) const override
Seek to a given event number.
SG::SlotSpecificObj< SG::SourceID > m_sourceID
virtual StatusCode last(IEvtSelector::Context &ctxt) const override
Gaudi::CheckedProperty< uint64_t > m_eventsPerRun
Gaudi::CheckedProperty< uint32_t > m_eventsPerLB
virtual StatusCode createAddress(const IEvtSelector::Context &ctxt, IOpaqueAddress *&iop) const override
virtual void handle(const Incident &incident) override
Incident service handle listening for BeginProcessing and EndProcessing.
virtual StatusCode previous(IEvtSelector::Context &ctxt) const override
ServiceHandle< IAthenaPoolCnvSvc > m_athenaPoolCnvSvc
virtual StatusCode rewind(IEvtSelector::Context &ctxt) const override
Gaudi::Property< bool > m_keepInputFilesOpen
KeepInputFilesOpen, boolean flag to keep files open after PoolCollection reaches end: default = false...
Gaudi::CheckedProperty< uint32_t > m_oldRunNo
virtual StatusCode finalize() override
virtual bool disconnectIfFinished(const SG::SourceID &fid) const override
Disconnect DB if all events from the source FID were processed and the Selector moved to another file...
Gaudi::Property< std::vector< std::string > > m_inputCollectionsProp
InputCollections, vector with names of the input collections.
EventContextAthenaPool * m_endIter
std::atomic_int m_evtCount
std::atomic_bool m_firedIncident
std::atomic_long m_curCollection
virtual StatusCode nextHandleFileTransition(IEvtSelector::Context &ctxt) const override
Handle file transition at the next iteration.
virtual StatusCode recordAttributeList() const override
Record AttributeList in StoreGate.
Gaudi::CheckedProperty< uint32_t > m_runNo
The following are included for compatibility with McEventSelector and are not really used.
virtual int size(Context &ctxt) const override
Return the size of the collection.
int findEvent(int evtNum) const
Search for event with number evtNum.
Gaudi::Property< std::string > m_skipEventRangesProp
Skip Events - comma separated list of event to skip, ranges with '-': <start> - <end>.
Gaudi::CheckedProperty< uint32_t > m_firstLBNo
This class provides a encapsulation of a GUID/UUID/CLSID/IID data structure (128 bit number).
static const Guid & null() noexcept
NULL-Guid: static class method.
constexpr void toString(std::span< char, StrLen > buf, bool uppercase=true) const noexcept
Automatic conversion to string representation.
This class provides an interface to POOL collections.
std::unique_ptr< pool::ICollectionCursor > selectAll()
StatusCode isValid() const
Check whether has valid pool::ICollection*.
const std::string & lastError() const
StatusCode initialize()
Required by all Gaudi Services.
virtual bool isValid() override final
Can the handle be successfully dereferenced?
The Athena Transient Store API.
static StoreGateSvc * currentStoreGate()
get current StoreGate
This class provides a Generic Transient Address for POOL tokens.
This class provides a token that identifies in a unique way objects on the persistent storage.
const std::string toString() const
Retrieve the string representation of the token.
int technology() const
Access technology type.
const Guid & dbID() const
Access database identifier.
const ExtendedEventContext & getExtendedEventContext(const EventContext &ctx)
Retrieve an extended context from a context object.
bool hasExtendedEventContext(const EventContext &ctx)
Test whether a context object has an extended context installed.
std::vector< std::string > tokenize(std::string_view the_str, std::string_view delimiters)
Splits the string into smaller substrings.
StatusCode parse(std::tuple< Tup... > &tup, const Gaudi::Parsers::InputData &input)
static const DbType POOL_StorageType
void sort(typename DataModel_detail::iterator< DVL > beg, typename DataModel_detail::iterator< DVL > end)
Specialization of sort for DataVector/List.
static constexpr CLID ID()