ATLAS Offline Software
Loading...
Searching...
No Matches
EventSelectorAthenaPool.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
13
18#include "PoolSvc/IPoolSvc.h"
22
23// Framework
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"
32
33// Pool
36#include "StorageSvc/DbType.h"
37
39#include <algorithm>
40#include <format>
41#include <vector>
42
43//________________________________________________________________________________
44EventSelectorAthenaPool::EventSelectorAthenaPool(const std::string& name, ISvcLocator* pSvcLocator) :
45 base_class(name, pSvcLocator)
46{
47 // TODO: validate if those are even used
48 m_runNo.verifier().setLower(0);
49 m_oldRunNo.verifier().setLower(0);
50 m_eventsPerRun.verifier().setLower(0);
51 m_firstEventNo.verifier().setLower(1);
52 m_firstLBNo.verifier().setLower(0);
53 m_eventsPerLB.verifier().setLower(0);
54 m_initTimeStamp.verifier().setLower(0);
55
57 m_inputCollectionsChanged = false;
58}
59//________________________________________________________________________________
60void EventSelectorAthenaPool::inputCollectionsHandler(Gaudi::Details::PropertyBase&) {
61 if (this->FSMState() != Gaudi::StateMachine::OFFLINE) {
62 m_inputCollectionsChanged = true;
63 }
64}
65//________________________________________________________________________________
68//________________________________________________________________________________
72//________________________________________________________________________________
74
75 m_autoRetrieveTools = false;
76 m_checkToolDeps = false;
77
78 if (m_isSecondary.value()) {
79 ATH_MSG_DEBUG("Initializing secondary event selector " << name());
80 } else {
81 ATH_MSG_DEBUG("Initializing " << name());
82 }
83
84 ATH_CHECK(::AthService::initialize());
85 // Check for input collection
86 if (m_inputCollectionsProp.value().empty()) {
87 ATH_MSG_FATAL("Use the property: EventSelector.InputCollections = "
88 << "[ \"<collectionName>\" ] (list of collections)");
89 return StatusCode::FAILURE;
90 }
91 auto ranges = CxxUtils::tokenize(m_skipEventRangesProp.value(), ',');
92 for( const std::string& r: ranges ) {
93 auto fromto = CxxUtils::tokenize(r, '-');
94 auto from_iter = fromto.begin();
95 long from = std::stol(*from_iter);
96 long to = from;
97 if( ++from_iter != fromto.end() ) {
98 to = std::stol(*from_iter);
99 }
100 m_skipEventRanges.emplace_back(from, to);
101 }
102
103 for( auto v : m_skipEventSequenceProp.value() ) {
104 m_skipEventRanges.emplace_back(v, v);
105 }
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);
113 }
114 if( !skip_ranges_str.empty() )
115 ATH_MSG_DEBUG("Events to skip: " << skip_ranges_str);
116 }
117
118 // Get AthenaPoolCnvSvc
119 ATH_CHECK(m_athenaPoolCnvSvc.retrieve());
120 ATH_CHECK(m_poolSvc.retrieve());
121 // Get CounterTool (if configured)
122 if (!m_counterTool.empty()) {
123 ATH_CHECK(m_counterTool.retrieve());
124 }
125 // Get HelperTools
126 ATH_CHECK(m_helperTools.retrieve());
127
128 // Ensure the xAODCnvSvc is listed in the EventPersistencySvc
129 ServiceHandle<IProperty> epSvc("EventPersistencySvc", name());
130 std::vector<std::string> propVal;
131 ATH_CHECK(Gaudi::Parsers::parse(propVal , epSvc->getProperty("CnvServices").toString()));
132 bool foundCnvSvc = false;
133 for (const auto& property : propVal) {
134 if (property == m_athenaPoolCnvSvc.type()) { foundCnvSvc = true; }
135 }
136 if (!foundCnvSvc) {
137 propVal.push_back(m_athenaPoolCnvSvc.type());
138 if (!epSvc->setProperty("CnvServices", Gaudi::Utils::toString(propVal)).isSuccess()) {
139 ATH_MSG_FATAL("Cannot set EventPersistencySvc Property for CnvServices");
140 return StatusCode::FAILURE;
141 }
142 }
143
144 // Register this service for 'I/O' events
145 ServiceHandle<IIoComponentMgr> iomgr("IoComponentMgr", name());
146 ATH_CHECK(iomgr.retrieve());
147 ATH_CHECK(iomgr->io_register(this));
148 // Register input file's names with the I/O manager
149 const std::vector<std::string>& incol = m_inputCollectionsProp.value();
150 bool allGood = true;
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 !");
154 allGood = false;
155 } else {
156 ATH_MSG_VERBOSE("io_register[" << this->name() << "](" << inputCollection << ") [ok]");
157 }
158 }
159 if (!allGood) {
160 return StatusCode::FAILURE;
161 }
162
163 // Connect to PersistencySvc
164 if (!m_poolSvc->connect(Io::READ, IPoolSvc::kInputStream).isSuccess()) {
165 ATH_MSG_FATAL("Cannot connect to POOL PersistencySvc.");
166 return StatusCode::FAILURE;
167 }
168 // Jump to reinit() to execute common init/reinit actions
169 m_guid = Guid::null();
170 if (!reinit().isSuccess()) {
171 return StatusCode::FAILURE;
172 }
173 // Get IncidentSvc
174 ATH_CHECK(m_incidentSvc.retrieve());
175 // Listen to the Event Processing incidents
176 m_incidentSvc->addListener(this, IncidentType::BeginProcessing, 0);
177 m_incidentSvc->addListener(this, IncidentType::EndProcessing, 0);
178 return StatusCode::SUCCESS;
179}
180//________________________________________________________________________________
182 ATH_MSG_DEBUG("reinitialization...");
183
184 // reset markers
185 m_numEvt.resize(m_inputCollectionsProp.value().size(), -1);
186 m_firstEvt.resize(m_inputCollectionsProp.value().size(), -1);
187
188 // Initialize InputCollectionsIterator
189 m_inputCollectionsIterator = m_inputCollectionsProp.value().begin();
190 m_curCollection = 0;
191 if (!m_firstEvt.empty()) {
192 m_firstEvt[0] = 0;
193 }
194 m_inputCollectionsChanged = false;
195 m_evtCount = 0;
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());
201 retError = true;
202 }
203 }
204 if (retError) {
205 ATH_MSG_FATAL("Failed to postInitialize() helperTools");
206 return StatusCode::FAILURE;
207 }
208
209 // Create an m_poolCollectionConverter to read the objects in
210 m_poolCollectionConverter = getCollectionCnv();
211 if (!m_poolCollectionConverter) {
212 ATH_MSG_INFO("No Events found in any Input Collections");
213 if (m_processMetadata.value()) {
214 m_inputCollectionsIterator = m_inputCollectionsProp.value().end();
215 if (!m_inputCollectionsProp.value().empty()) --m_inputCollectionsIterator;
216 //NOTE (wb may 2016): this will make the FirstInputFile incident correspond to last file in the collection ... if want it to be first file then move iterator to begin and then move above two lines below this incident firing
217 if( !m_firedIncident && !m_inputCollectionsProp.value().empty() ) {
218 FileIncident firstInputFileIncident(name(), "FirstInputFile", *m_inputCollectionsIterator);
219 m_incidentSvc->fireIncident(firstInputFileIncident);
220 m_firedIncident = true;
221 }
222 }
223 return StatusCode::SUCCESS;
224 }
225 // Get DataHeader iterator
226 try {
227 m_headerIterator = m_poolCollectionConverter->selectAll();
228 } catch (std::exception &e) {
229 ATH_MSG_FATAL("Cannot open input collection - check data/software version.");
230 ATH_MSG_ERROR(e.what());
231 return StatusCode::FAILURE;
232 }
233 while (m_headerIterator == nullptr || m_headerIterator->next() == 0) { // no selected events
234 if (m_poolCollectionConverter) {
235 m_poolCollectionConverter->disconnectDb().ignore();
236 m_poolCollectionConverter.reset();
237 }
238 ++m_inputCollectionsIterator;
239 m_poolCollectionConverter = getCollectionCnv();
240 if (m_poolCollectionConverter) {
241 m_headerIterator = m_poolCollectionConverter->selectAll();
242 } else {
243 break;
244 }
245 }
246 if (!m_poolCollectionConverter || m_headerIterator == nullptr) { // no event selected in any collection
247 m_inputCollectionsIterator = m_inputCollectionsProp.value().begin();
248 m_curCollection = 0;
249 m_poolCollectionConverter = getCollectionCnv();
250 if (!m_poolCollectionConverter) {
251 return StatusCode::SUCCESS;
252 }
253 m_headerIterator = m_poolCollectionConverter->selectAll();
254 while (m_headerIterator == nullptr || m_headerIterator->next() == 0) { // empty collection
255 if (m_poolCollectionConverter) {
256 m_poolCollectionConverter->disconnectDb().ignore();
257 m_poolCollectionConverter.reset();
258 }
259 ++m_inputCollectionsIterator;
260 m_poolCollectionConverter = getCollectionCnv();
261 if (m_poolCollectionConverter) {
262 m_headerIterator = m_poolCollectionConverter->selectAll();
263 } else {
264 break;
265 }
266 }
267 }
268 if (!m_poolCollectionConverter || m_headerIterator == nullptr) {
269 return StatusCode::SUCCESS;
270 }
271 const Token& headRef = m_headerIterator->eventRef();
272 const std::string fid = headRef.dbID().toString();
273 const int tech = headRef.technology();
274 ATH_MSG_VERBOSE("reinit(): First DataHeder Token=" << headRef.toString() );
275
276 // Check if File is BS, for which Incident is thrown by SingleEventInputSvc
277 if (tech != 0x00001000 && m_processMetadata.value() && !m_firedIncident) {
278 FileIncident firstInputFileIncident(name(), "FirstInputFile", "FID:" + fid, fid);
279 m_incidentSvc->fireIncident(firstInputFileIncident);
280 m_firedIncident = true;
281 }
282 return StatusCode::SUCCESS;
283}
284//________________________________________________________________________________
286 if (m_poolCollectionConverter) {
287 // Reset iterators and apply new query
288 m_poolCollectionConverter->disconnectDb().ignore();
289 m_poolCollectionConverter.reset();
290 }
291 m_inputCollectionsIterator = m_inputCollectionsProp.value().begin();
292 m_curCollection = 0;
293 m_poolCollectionConverter = getCollectionCnv(true);
294 if (!m_poolCollectionConverter) {
295 ATH_MSG_INFO("No Events found in any Input Collections");
296 m_inputCollectionsIterator = m_inputCollectionsProp.value().end();
297 if (!m_inputCollectionsProp.value().empty()) {
298 --m_inputCollectionsIterator; //leave iterator in state of last input file
299 }
300 } else {
301 m_headerIterator = m_poolCollectionConverter->selectAll();
302 }
303 m_evtCount = 0;
304 delete m_endIter;
305 m_endIter = nullptr;
306 m_endIter = new EventContextAthenaPool(nullptr);
307 return StatusCode::SUCCESS;
308}
309//________________________________________________________________________________
311 // Fire EndInputFile for any file still open (the event loop may end
312 // before the file is fully read).
313 m_inputFileGuard.reset();
314
315 IEvtSelector::Context* ctxt(nullptr);
316 if (!releaseContext(ctxt).isSuccess()) {
317 ATH_MSG_WARNING("Cannot release context");
318 }
319 return StatusCode::SUCCESS;
320}
321
322//________________________________________________________________________________
324 if (!m_counterTool.empty() && !m_counterTool->preFinalize().isSuccess()) {
325 ATH_MSG_WARNING("Failed to preFinalize() CounterTool");
326 }
327 for (auto& tool : m_helperTools) {
328 if (!tool->preFinalize().isSuccess()) {
329 ATH_MSG_WARNING("Failed to preFinalize() " << tool->name());
330 }
331 }
332 delete m_endIter; m_endIter = nullptr;
333 m_headerIterator = nullptr;
334 if (m_poolCollectionConverter) {
335 m_poolCollectionConverter.reset();
336 }
337 // Finalize the Service base class.
338 return ::AthService::finalize();
339}
340
341//________________________________________________________________________________
342StatusCode EventSelectorAthenaPool::createContext(IEvtSelector::Context*& ctxt) const {
343 ctxt = new EventContextAthenaPool(this);
344 return StatusCode::SUCCESS;
345}
346//________________________________________________________________________________
347StatusCode EventSelectorAthenaPool::next(IEvtSelector::Context& ctxt) const {
348 std::lock_guard<CallMutex> lockGuard(m_callLock);
349 for (const auto& tool : m_helperTools) {
350 if (!tool->preNext().isSuccess()) {
351 ATH_MSG_WARNING("Failed to preNext() " << tool->name());
352 }
353 }
354 for (;;) {
355 // Handle possible file transition
356 StatusCode sc = nextHandleFileTransition(ctxt);
357 if (sc.isRecoverable()) {
358 continue; // handles empty files
359 }
360 if (sc.isFailure()) {
361 return StatusCode::FAILURE;
362 }
363 // Increase event count
364 ++m_evtCount;
365 if (!m_counterTool.empty() && !m_counterTool->preNext().isSuccess()) {
366 ATH_MSG_WARNING("Failed to preNext() CounterTool.");
367 }
369 && (m_skipEventRanges.empty() || m_evtCount < m_skipEventRanges.front().first))
370 {
371 if (!m_isSecondary.value()) {
372 if (!this->recordAttributeList().isSuccess()) {
373 ATH_MSG_ERROR("Failed to record AttributeList.");
374 return StatusCode::FAILURE;
375 }
376 }
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;
384 }
385 } else if (toolStatus.isFailure()) {
386 ATH_MSG_WARNING("Failed to postNext() " << tool->name());
387 status = StatusCode::FAILURE;
388 }
389 }
390 if (status.isRecoverable()) {
391 ATH_MSG_INFO("skipping event " << m_evtCount);
392 } else if (status.isFailure()) {
393 ATH_MSG_WARNING("Failed to postNext() HelperTool.");
394 } else {
395 if (!m_counterTool.empty() && !m_counterTool->postNext().isSuccess()) {
396 ATH_MSG_WARNING("Failed to postNext() CounterTool.");
397 }
398 break;
399 }
400 } else {
401 while( !m_skipEventRanges.empty() && m_evtCount >= m_skipEventRanges.front().second ) {
402 m_skipEventRanges.erase(m_skipEventRanges.begin());
403 }
404 ATH_MSG_INFO("skipping event " << m_evtCount);
405 }
406 }
407 return StatusCode::SUCCESS;
408}
409//________________________________________________________________________________
410StatusCode EventSelectorAthenaPool::next(IEvtSelector::Context& ctxt, int jump) const {
411 if (jump > 0) {
412 for (int i = 0; i < jump; i++) {
413 ATH_CHECK(next(ctxt));
414 }
415 return StatusCode::SUCCESS;
416 }
417 return StatusCode::FAILURE;
418}
419//________________________________________________________________________________
420StatusCode EventSelectorAthenaPool::nextHandleFileTransition(IEvtSelector::Context& ctxt) const
421{
422 if( m_inputCollectionsChanged ) {
423 StatusCode rc = reinit();
424 if( rc != StatusCode::SUCCESS ) return rc;
425 }
426 else { // advance to the next (not needed after reinit)
427 // Check if we're at the end of file
428 if (m_headerIterator == nullptr || m_headerIterator->next() == 0) {
429 m_headerIterator = nullptr;
430 // Close previous collection.
431 m_poolCollectionConverter.reset();
432
433 // Fire EndInputFile while data is still accessible, then disconnect
434 m_inputFileGuard.reset();
435 const SG::SourceID old_guid = m_guid.toString();
436 m_guid = Guid::null();
437 disconnectIfFinished( old_guid );
438
439 // check if somebody updated Inputs in the EOF incident (like VP1 does)
440 if( m_inputCollectionsChanged ) {
441 StatusCode rc = reinit();
442 if( rc != StatusCode::SUCCESS ) return rc;
443 } else {
444 // Open next file from inputCollections list.
445 ++m_inputCollectionsIterator;
446 // Create PoolCollectionConverter for input file
447 m_poolCollectionConverter = getCollectionCnv(true);
448 if (!m_poolCollectionConverter) {
449 // Return end iterator
450 ctxt = *m_endIter;
451 // This is not a real failure but a Gaudi way of handling "end of job"
452 return StatusCode::FAILURE;
453 }
454 // Get DataHeader iterator
455 m_headerIterator = m_poolCollectionConverter->selectAll();
456
457 // Return RECOVERABLE to mark we should still continue
458 return StatusCode::RECOVERABLE;
459 }
460 }
461 }
462 const Token& headRef = m_headerIterator->eventRef();
463 const Guid guid = headRef.dbID();
464 ATH_MSG_VERBOSE("next(): DataHeder Token=" << headRef.toString() );
465
466 if (guid != m_guid) {
467 // we are starting reading from a new DB. Check if the old one needs to be retired
468 if (m_guid != Guid::null()) {
469 // zero the current DB ID (m_guid) before trying disconnect() to indicate it is no longer in use
470 const SG::SourceID old_guid = m_guid.toString();
471 m_guid = Guid::null();
472 // EndInputFile is fired by the guard transition() below; just disconnect here
473 disconnectIfFinished( old_guid );
474 }
475 m_guid = guid;
476 m_activeEventsPerSource[guid.toString()] = 0;
477 if (!m_athenaPoolCnvSvc->setInputAttributes(*m_inputCollectionsIterator).isSuccess()) {
478 ATH_MSG_ERROR("Failed to set input attributes.");
479 return StatusCode::FAILURE;
480 }
481 if(m_processMetadata.value()) {
482 InputFileIncidentGuard::transition(m_inputFileGuard, *m_incidentSvc, name(),
483 *m_inputCollectionsIterator, m_guid.toString(),
484 /*endFileName=*/{});
485 }
486 } // end if (guid != m_guid)
487 return StatusCode::SUCCESS;
488}
489//________________________________________________________________________________
490StatusCode EventSelectorAthenaPool::nextWithSkip(IEvtSelector::Context& ctxt) const {
491 ATH_MSG_DEBUG("EventSelectorAthenaPool::nextWithSkip");
492
493 for (;;) {
494 // Check if we're at the end of file
495 StatusCode sc = nextHandleFileTransition(ctxt);
496 if (sc.isRecoverable()) {
497 continue; // handles empty files
498 }
499 if (sc.isFailure()) {
500 return StatusCode::FAILURE;
501 }
502
503 // Increase event count
504 ++m_evtCount;
505
506 if (!m_counterTool.empty() && !m_counterTool->preNext().isSuccess()) {
507 ATH_MSG_WARNING("Failed to preNext() CounterTool.");
508 }
510 && (m_skipEventRanges.empty() || m_evtCount < m_skipEventRanges.front().first))
511 {
512 return StatusCode::SUCCESS;
513 } else {
514 while( !m_skipEventRanges.empty() && m_evtCount >= m_skipEventRanges.front().second ) {
515 m_skipEventRanges.erase(m_skipEventRanges.begin());
516 }
517 if (m_isSecondary.value()) {
518 ATH_MSG_INFO("skipping secondary event " << m_evtCount);
519 } else {
520 ATH_MSG_INFO("skipping event " << m_evtCount);
521 }
522 }
523 }
524
525 return StatusCode::SUCCESS;
526}
527//________________________________________________________________________________
528StatusCode EventSelectorAthenaPool::previous(IEvtSelector::Context& /*ctxt*/) const {
529 ATH_MSG_ERROR("previous() not implemented");
530 return StatusCode::FAILURE;
531}
532//________________________________________________________________________________
533StatusCode EventSelectorAthenaPool::previous(IEvtSelector::Context& ctxt, int jump) const {
534 if (jump > 0) {
535 for (int i = 0; i < jump; i++) {
536 ATH_CHECK(previous(ctxt));
537 }
538 return StatusCode::SUCCESS;
539 }
540 return StatusCode::FAILURE;
541}
542//________________________________________________________________________________
543StatusCode EventSelectorAthenaPool::last(IEvtSelector::Context& ctxt) const {
544 if (ctxt.identifier() == m_endIter->identifier()) {
545 ATH_MSG_DEBUG("last(): Last event in InputStream.");
546 return StatusCode::SUCCESS;
547 }
548 return StatusCode::FAILURE;
549}
550//________________________________________________________________________________
551StatusCode EventSelectorAthenaPool::rewind(IEvtSelector::Context& ctxt) const {
552 ATH_CHECK(reinit());
553 ctxt = EventContextAthenaPool(this);
554 return StatusCode::SUCCESS;
555}
556//________________________________________________________________________________
557StatusCode EventSelectorAthenaPool::createAddress(const IEvtSelector::Context& /*ctxt*/,
558 IOpaqueAddress*& iop) const {
559 std::string tokenStr;
561 if (attrList.isValid()) {
562 try {
563 tokenStr = (*attrList)["eventRef"].data<std::string>();
564 ATH_MSG_DEBUG("found AthenaAttribute, name = eventRef = " << tokenStr);
565 } catch (std::exception &e) {
566 ATH_MSG_ERROR(e.what());
567 return StatusCode::FAILURE;
568 }
569 } else {
570 ATH_MSG_WARNING("Cannot find AthenaAttribute, key = " << m_attrListKey);
571 tokenStr = m_headerIterator->eventRef().toString();
572 }
573 auto token = std::make_unique<Token>();
574 token->fromString(tokenStr);
575 m_incidentSvc->fireIncident(Incident(tokenStr, "ProcessEventAttributes"));
576 iop = new TokenAddress(pool::POOL_StorageType.type(), ClassID_traits<DataHeader>::ID(), "", "EventSelector", IPoolSvc::kInputStream, std::move(token));
577 return StatusCode::SUCCESS;
578}
579//________________________________________________________________________________
580StatusCode EventSelectorAthenaPool::releaseContext(IEvtSelector::Context*& /*ctxt*/) const {
581 return StatusCode::SUCCESS;
582}
583//________________________________________________________________________________
584StatusCode EventSelectorAthenaPool::resetCriteria(const std::string& /*criteria*/,
585 IEvtSelector::Context& /*ctxt*/) const {
586 return StatusCode::SUCCESS;
587}
588//__________________________________________________________________________
589StatusCode EventSelectorAthenaPool::seek(Context& /*ctxt*/, int evtNum) const {
590
591 if( m_inputCollectionsChanged ) {
592 StatusCode rc = reinit();
593 if( rc != StatusCode::SUCCESS ) return rc;
594 }
595
596 long newColl = findEvent(evtNum);
597 if (newColl == -1 && evtNum >= m_firstEvt[m_curCollection] && evtNum < m_evtCount - 1) {
598 newColl = m_curCollection;
599 }
600 if (newColl == -1) {
601 m_headerIterator = nullptr;
602 ATH_MSG_INFO("seek: Reached end of Input.");
603 m_inputFileGuard.reset();
604 return StatusCode::RECOVERABLE;
605 }
606 if (newColl != m_curCollection) {
607 if (!m_keepInputFilesOpen.value() && m_poolCollectionConverter) {
608 m_poolCollectionConverter->disconnectDb().ignore();
609 }
610 m_poolCollectionConverter.reset();
611 m_curCollection = newColl;
612 try {
613 ATH_MSG_DEBUG("Seek to item: \""
615 << "\" from the collection list.");
616 // Reset input collection iterator to the right place
617 m_inputCollectionsIterator = m_inputCollectionsProp.value().begin();
618 m_inputCollectionsIterator += m_curCollection;
619 m_poolCollectionConverter = std::make_unique<PoolCollectionConverter>(
622 m_poolSvc.get());
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;
627 }
628 // Create DataHeader iterators
629 m_headerIterator = m_poolCollectionConverter->selectAll();
630 EventContextAthenaPool* beginIter = new EventContextAthenaPool(this);
631 m_evtCount = m_firstEvt[m_curCollection];
632 next(*beginIter).ignore();
633 ATH_MSG_DEBUG("Token " << m_headerIterator->eventRef().toString());
634 } catch (std::exception &e) {
635 m_headerIterator = nullptr;
636 ATH_MSG_ERROR(e.what());
637 return StatusCode::FAILURE;
638 }
639 }
640
641 if (m_headerIterator->seek(evtNum - m_firstEvt[m_curCollection]) == 0) {
642 m_headerIterator = nullptr;
643 ATH_MSG_ERROR("Did not find event, evtNum = " << evtNum);
644 return StatusCode::FAILURE;
645 } else {
646 m_evtCount = evtNum + 1;
647 }
648 return StatusCode::SUCCESS;
649}
650//__________________________________________________________________________
651int EventSelectorAthenaPool::curEvent (const Context& /*ctxt*/) const {
652 return(m_evtCount);
653}
654//__________________________________________________________________________
655// Search for event number evtNum.
656// Return the index of the collection containing it, or -1 if not found.
657// Note: passing -1 for evtNum will always yield failure,
658// but this can be used to force filling in the entire m_numEvt array.
660 for (std::size_t i = 0, imax = m_numEvt.size(); i < imax; i++) {
661 if (m_numEvt[i] == -1) {
663 m_inputCollectionsProp.value()[i],
665 m_poolSvc.get());
666 if (!pcc.initialize().isSuccess()) {
667 break;
668 }
669 int collection_size = 0;
670 if (pcc.isValid()) {
671 std::unique_ptr<pool::ICollectionCursor> hi = pcc.selectAll();
672 collection_size = hi->size();
673 }
674 else {
675 ATH_MSG_ERROR( pcc.lastError() );
676 }
677 if (i > 0) {
678 m_firstEvt[i] = m_firstEvt[i - 1] + m_numEvt[i - 1];
679 } else {
680 m_firstEvt[i] = 0;
681 }
682 m_numEvt[i] = collection_size;
683 }
684 if (evtNum >= m_firstEvt[i] && evtNum < m_firstEvt[i] + m_numEvt[i]) {
685 return(i);
686 }
687 }
688 return(-1);
689}
690
691//__________________________________________________________________________
692int EventSelectorAthenaPool::size(Context& /*ctxt*/) const {
693 // Fetch sizes of all collections.
694 findEvent(-1);
695 return std::accumulate(m_numEvt.begin(), m_numEvt.end(), 0);
696}
697//__________________________________________________________________________
698std::unique_ptr<PoolCollectionConverter>
700 while (m_inputCollectionsIterator != m_inputCollectionsProp.value().end()) {
701 if (m_curCollection != 0) {
702 m_numEvt[m_curCollection] = m_evtCount - m_firstEvt[m_curCollection];
704 m_firstEvt[m_curCollection] = m_evtCount;
705 }
706 ATH_MSG_DEBUG("Try item: \"" << *m_inputCollectionsIterator << "\" from the collection list.");
707 auto pCollCnv = std::make_unique<PoolCollectionConverter>(
708 *m_inputCollectionsIterator,
710 m_poolSvc.get());
711 StatusCode status = pCollCnv->initialize();
712 if (!status.isSuccess()) {
713 // Close previous collection.
714 pCollCnv.reset();
715 if (!status.isRecoverable()) {
716 ATH_MSG_ERROR("Unable to initialize PoolCollectionConverter.");
717 throw GaudiException("Unable to read: " + *m_inputCollectionsIterator, name(), StatusCode::FAILURE);
718 } else {
719 ATH_MSG_ERROR("Unable to open: " << *m_inputCollectionsIterator);
720 throw GaudiException("Unable to open: " + *m_inputCollectionsIterator, name(), StatusCode::FAILURE);
721 }
722 } else {
723 if (!pCollCnv->isValid().isSuccess()) {
724 pCollCnv.reset();
725 ATH_MSG_DEBUG("No events found in: " << *m_inputCollectionsIterator << " skipped!!!");
726 if (throwIncidents && m_processMetadata.value()) {
727 // Scoped guard: fires BeginInputFile now, EndInputFile at scope exit
728 auto guard = InputFileIncidentGuard::begin(*m_incidentSvc, name(),
729 *m_inputCollectionsIterator, {},
730 "eventless " + *m_inputCollectionsIterator);
731 }
732 m_poolSvc->disconnectDb(*m_inputCollectionsIterator).ignore();
733 ++m_inputCollectionsIterator;
734 } else {
735 return(pCollCnv);
736 }
737 }
738 }
739 return(nullptr);
740}
741//__________________________________________________________________________
743 if (!eventStore()->clearStore().isSuccess()) {
744 ATH_MSG_WARNING("Cannot clear Store");
745 }
746 // Get access to AttributeList
747 ATH_MSG_DEBUG("Get AttributeList from the collection");
748 // MN: accessing only attribute list, ignoring token list
749 const coral::AttributeList& attrList = m_headerIterator->currentRow().attributeList();
750 ATH_MSG_DEBUG("AttributeList size " << attrList.size());
751 std::unique_ptr<AthenaAttributeList> athAttrList(new AthenaAttributeList(attrList));
752 // Fill the new attribute list
753 ATH_CHECK(fillAttributeList(athAttrList.get(), "", false));
754 // Write the AttributeList
756 ATH_CHECK(wh.record(std::move(athAttrList)));
757 return StatusCode::SUCCESS;
758}
759//__________________________________________________________________________
760StatusCode EventSelectorAthenaPool::fillAttributeList(coral::AttributeList *attrList, const std::string &suffix, bool copySource) const
761{
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() << ".");
766
767 std::string eventRef = "eventRef";
768 if (m_isSecondary.value()) {
769 eventRef.append(suffix);
770 }
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() << ".");
774
775 if (copySource) {
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;
780 }
781 }
782
783 return StatusCode::SUCCESS;
784}
785//__________________________________________________________________________
787 ATH_MSG_INFO("I/O reinitialization...");
788 if (m_poolCollectionConverter) {
789 m_poolCollectionConverter->disconnectDb().ignore();
790 m_poolCollectionConverter.reset();
791 }
792 m_headerIterator = nullptr;
793 ServiceHandle<IIoComponentMgr> iomgr("IoComponentMgr", name());
794 if (!iomgr.retrieve().isSuccess()) {
795 ATH_MSG_FATAL("Could not retrieve IoComponentMgr !");
796 return StatusCode::FAILURE;
797 }
798 if (!iomgr->io_hasitem(this)) {
799 ATH_MSG_FATAL("IoComponentMgr does not know about myself !");
800 return StatusCode::FAILURE;
801 }
802 std::vector<std::string> inputCollections = m_inputCollectionsProp.value();
803 std::set<std::size_t> updatedIndexes;
804 for (std::size_t i = 0, imax = m_inputCollectionsProp.value().size(); i < imax; i++) {
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;
811 }
812 if (!iomgr->io_retrieve(this, fname).isSuccess()) {
813 ATH_MSG_FATAL("Could not retrieve new value for [" << fname << "] !");
814 return StatusCode::FAILURE;
815 }
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);
821 }
822 }
823 }
824 // all good... copy over.
825 m_inputCollectionsProp = inputCollections;
826 m_guid = Guid::null();
827 return reinit();
828}
829//__________________________________________________________________________
831 ATH_MSG_INFO("I/O finalization...");
832 // Fire EndInputFile before disconnecting — file data is still accessible here
833 m_inputFileGuard.reset();
834 if (m_poolCollectionConverter) {
835 m_poolCollectionConverter->disconnectDb().ignore();
836 m_poolCollectionConverter.reset();
837 }
838 return StatusCode::SUCCESS;
839}
840
841//__________________________________________________________________________
842/* Listen to IncidentType::BeginProcessing and EndProcessing
843 Maintain counters of how many events from a given file are being processed.
844 Files are identified by SG::SourceID (string GUID).
845 When there are no more events from a file, see if it can be closed.
846*/
847void EventSelectorAthenaPool::handle(const Incident& inc)
848{
849 SG::SourceID fid;
850 if (inc.type() == IncidentType::BeginProcessing) {
851 if ( Atlas::hasExtendedEventContext(inc.context()) ) {
852 fid = Atlas::getExtendedEventContext(inc.context()).proxy()->sourceID();
853 }
854 *m_sourceID.get(inc.context()) = fid;
855 }
856 else {
857 fid = *m_sourceID.get(inc.context());
858 }
859
860 if( fid.empty() ) {
861 ATH_MSG_WARNING("could not read event source ID from incident event context");
862 return;
863 }
864 if( m_activeEventsPerSource.find( fid ) == m_activeEventsPerSource.end()) {
865 ATH_MSG_DEBUG("Incident handler ignoring unknown input FID: " << fid );
866 return;
867 }
868 ATH_MSG_DEBUG("** MN Incident handler " << inc.type() << " Event source ID=" << fid );
869 if( inc.type() == IncidentType::BeginProcessing ) {
870 // increment the events-per-file counter for FID
871 m_activeEventsPerSource[fid]++;
872 } else if( inc.type() == IncidentType::EndProcessing ) {
873 m_activeEventsPerSource[fid]--;
875 *m_sourceID.get(inc.context()) = "";
876 }
877 if( msgLvl(MSG::DEBUG) ) {
878 for( auto& source: m_activeEventsPerSource )
879 msg(MSG::DEBUG) << "SourceID: " << source.first << " active events: " << source.second << endmsg;
880 }
881}
882
883//__________________________________________________________________________
884/* Disconnect Database identifieed by a SG::SourceID when it is no longer in use:
885 m_guid is not pointing to it and there are no events from it being processed
886 (if the EventLoopMgr was not firing Begin/End incidents, this will just close the DB)
887*/
889{
890 if( m_activeEventsPerSource.find(fid) != m_activeEventsPerSource.end()
891 && m_activeEventsPerSource[fid] <= 0 && m_guid != fid ) {
892 // Explicitly disconnect file corresponding to old FID to release memory.
893 // EndInputFile is handled by the InputFileIncidentGuard.
894 if( !m_keepInputFilesOpen.value() ) {
895 ATH_MSG_INFO("Disconnecting input sourceID: " << fid );
896 m_poolSvc->disconnectDb("FID:" + fid, IPoolSvc::kInputStream).ignore();
897 m_activeEventsPerSource.erase( fid );
898 return true;
899 }
900 }
901 return false;
902}
#define endmsg
#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 DataHeader and DataHeaderElement classes.
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.
static Double_t sc
static Double_t rc
This file contains the class definition for the PoolCollectionConverter class.
Handle class for reading from StoreGate.
Handle class for recording to StoreGate.
int imax(int i, int j)
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
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).
Definition Guid.h:25
static const Guid & null() noexcept
NULL-Guid: static class method.
Definition Guid.cxx:14
constexpr void toString(std::span< char, StrLen > buf, bool uppercase=true) const noexcept
Automatic conversion to string representation.
@ kInputStream
Definition IPoolSvc.h:39
static InputFileIncidentGuard begin(IIncidentSvc &incSvc, std::string_view source, std::string_view beginFileName, std::string_view guid, std::string_view endFileName={}, std::string_view beginType=IncidentType::BeginInputFile, std::string_view endType=IncidentType::EndInputFile)
Factory: fire the begin incident and return a guard whose destructor fires the matching end incident.
static void transition(std::optional< InputFileIncidentGuard > &guard, IIncidentSvc &incSvc, std::string_view source, std::string_view beginFileName, std::string_view guid, std::string_view endFileName={}, std::string_view beginType=IncidentType::BeginInputFile, std::string_view endType=IncidentType::EndInputFile)
Replace the guard in an optional, with strict End-before-Begin ordering.
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.
Definition Token.h:22
const std::string toString() const
Retrieve the string representation of the token.
Definition Token.cxx:135
int technology() const
Access technology type.
Definition Token.h:78
const Guid & dbID() const
Access database identifier.
Definition Token.h:65
int r
Definition globals.cxx:22
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
Definition DbType.h:84
void sort(typename DataModel_detail::iterator< DVL > beg, typename DataModel_detail::iterator< DVL > end)
Specialization of sort for DataVector/List.
MsgStream & msg
Definition testRead.cxx:32