ATLAS Offline Software
Loading...
Searching...
No Matches
DelayedConditionsCleanerSvc.cxx
Go to the documentation of this file.
1/*
2 Copyright (C) 2002-2026 CERN for the benefit of the ATLAS collaboration
3*/
10
11
16#include "GaudiKernel/EventContext.h"
17#include "GaudiKernel/ServiceHandle.h"
18#include <algorithm>
19#include <unordered_set>
20
21
22// We wanted to try to allow cleaning to run as asynchronous tasks,
23// but there are some race conditions that appear to be difficult
24// to resolve. These stem from the fact that conditions are only
25// written from CondInputLoader. If we clean after that, and don't
26// have the full set of IOV keys for all executing events, then we can
27// clean an item which the slot won't be able to recover.
28// Leave the code commented-out for now while we think about this further.
29#define USE_ASYNC_TASK 0
30
31
32#if USE_ASYNC_TASK
33#include "tbb/task.h"
34#endif
35
36
37namespace Athena {
38
39
41 : public AthProperties<DelayedConditionsCleanerSvc>
42{
43public:
46
48 Gaudi::Property<size_t> m_ringSize
49 { parent(), "RingSize", 100,
50 "Number of previous events for which to remember IOV history." };
51
54 Gaudi::Property<size_t> m_cleanDelay
55 { parent(), "CleanDelay", 100,
56 "Number of events after adding a conditions object we try to clean its container." };
57
59 Gaudi::Property<size_t> m_lookAhead
60 { parent(), "LookAhead", 10,
61 "Maximum number of events to consolodate together when cleaning." };
62
63
65#if USE_ASYNC_TASK
66 Gaudi::Property<bool> m_async
67 { parent(), "Async", false,
68 "If true, run cleaning asynchronously in an MT job." };
69#else
70 bool m_async = false;
71#endif
72
75 { parent(), "RCUSvc", "Athena::RCUSvc",
76 "The RCU service." };
77};
78
79
80#if USE_ASYNC_TASK
84class DelayedConditionsCleanerTask
85 : public tbb::task
86{
87public:
95 DelayedConditionsCleanerTask (DelayedConditionsCleanerSvc& cleaner,
96 std::vector<DelayedConditionsCleanerSvc::CondContInfo*>&& cis,
97 DelayedConditionsCleanerSvc::twoKeys_t&& keys);
98
102 tbb::task* execute() override;
103
104
105private:
107 DelayedConditionsCleanerSvc& m_cleaner;
108
110 std::vector<DelayedConditionsCleanerSvc::CondContInfo*> m_cis;
111
113 DelayedConditionsCleanerSvc::twoKeys_t m_keys;
114};
115
116
123DelayedConditionsCleanerTask::DelayedConditionsCleanerTask
125 std::vector<DelayedConditionsCleanerSvc::CondContInfo*>&& cis,
127 : m_cleaner (cleaner),
128 m_cis (cis),
129 m_keys (keys)
130{
131}
132
133
137tbb::task* DelayedConditionsCleanerTask::execute()
138{
139 // Do the cleaning.
140 m_cleaner.cleanContainers (std::move (m_cis), std::move (m_keys));
141
142 // This task is terminating.
143 --m_cleaner.m_cleanTasks;
144 return nullptr;
145}
146#endif // USE_ASYNC_TASK
147
148
155 ISvcLocator* svc)
156 : base_class (name, svc),
157 m_props (std::make_unique<DelayedConditionsCleanerSvcProps> (this))
158{
159}
160
161
166{
167 // Set the ring buffer sizes.
168 m_runlbn.reset (m_props->m_ringSize);
169 m_timestamp.reset (m_props->m_ringSize);
170
171 ATH_CHECK( m_props->m_rcu.retrieve() );
172 size_t nslots = m_props->m_rcu->getNumSlots();
173 m_slotLBN.resize (nslots);
174 m_slotTimestamp.resize (nslots);
175
176 return StatusCode::SUCCESS;
177}
178
179
186StatusCode
187DelayedConditionsCleanerSvc::event (const EventContext& ctx, bool allowAsync)
188{
189 // Push the IOV key for the current event into the ring buffers.
190 // Also save in the per-slot arrays.
191 key_type key_lbn = CondContBase::keyFromRunLBN (ctx.eventID());
192 key_type key_ts = CondContBase::keyFromTimestamp (ctx.eventID());
193 m_runlbn.push (key_lbn);
194 m_timestamp.push (key_ts);
195 EventContext::ContextID_t slot = ctx.slot();
196 if (slot != EventContext::INVALID_CONTEXT_ID) {
197 m_slotLBN[slot] = key_lbn;
198 m_slotTimestamp[slot] = key_ts;
199 }
200
201 // Return now if an asynchronous cleaning task is still running ---
202 // we don't want to start a new one yet. We'll check pending work
203 // on the next call.
204 if (m_cleanTasks > 0) {
205 return StatusCode::SUCCESS;
206 }
207
208 // Collect conditions containers in need of cleaning.
209 std::vector<CondContInfo*> ci_vec;
210 {
212 // Is it time to clean the container at the top of the work queue?
213 if (!m_work.empty() && m_work.top().m_evt <= ctx.evt()) {
214 ++m_nEvents;
215 size_t sz = m_work.size();
216 m_queueSum += sz;
217 m_maxQueue = std::max (m_maxQueue, sz);
218
219 // Yes. Put it on the correct list. Also look ahead in the queue
220 // a bit; if there are other containers that we want to clean soon,
221 // go ahead and do them now.
222 do {
223 CondContInfo* ci = m_work.top().m_ci;
224 switch (ci->m_cc.keyType()) {
225 case KeyType::SINGLE:
226 break;
227 case KeyType::RUNLBN:
228 case KeyType::MIXED:
229 case KeyType::TIMESTAMP:
230 ci_vec.push_back (ci);
231 break;
232 default:
233 std::abort();
234 }
235 m_work.pop();
237 } while (!m_work.empty() && m_work.top().m_evt <= ctx.evt() + m_props->m_lookAhead);
238 }
239 }
240
241 // Clean the containers.
242 if (!ci_vec.empty()) {
243 scheduleClean (std::move (ci_vec), getKeys(m_runlbn,m_timestamp),
244 allowAsync);
245 }
246 return StatusCode::SUCCESS;
247}
248
249
255StatusCode DelayedConditionsCleanerSvc::condObjAdded (const EventContext& ctx,
256 CondContBase& cc)
257{
258 // Add this container to the priority queue.
260 CCInfoMap_t::iterator it = m_ccinfo.find (&cc);
261 if (it == m_ccinfo.end()) {
262 it = m_ccinfo.emplace (&cc, CondContInfo (cc)).first;
263 }
264
265 EventContext::ContextEvt_t evt = ctx.evt();
266 m_work.emplace (evt + m_props->m_cleanDelay, it->second);
267 return StatusCode::SUCCESS;
268}
269
270
277{
278 // Suppress output if we didn't actually do anything.
279 if (m_nEvents == 0) {
280 return StatusCode::SUCCESS;
281 }
282
283 ATH_MSG_INFO( "Conditions container statistics" );
284 ATH_MSG_INFO( " Work q: Max size: {} ({} queries) ", m_maxQueue, m_nEvents);
285 size_t den = std::max (m_nEvents, 1lu );
286 ATH_MSG_INFO( " Avg size: {:.2f} / Avg removed: {:.2f}",
287 static_cast<float>(m_queueSum)/den,
288 static_cast<float>(m_workRemoved)/den );
289
290 std::vector<const CondContInfo*> infos;
291 for (const auto& p : m_ccinfo) {
292 infos.push_back (&p.second);
293 }
294 std::sort (infos.begin(), infos.end(),
295 [](const CondContInfo* a, const CondContInfo* b)
296 { return a->m_cc.id().key() < b->m_cc.id().key(); });
297
298 for (const CondContInfo* ci : infos) {
299 ATH_MSG_INFO( " {:<20} nInserts {:6} maxSize {:3}",
300 ci->m_cc.id().key().c_str(),
301 ci->m_cc.nInserts(),
302 ci->m_cc.maxSize() );
303 den = std::max (ci->m_nClean, 1lu);
304 ATH_MSG_INFO( " nClean {} avgRemoved {:.2f} 0/1/2+ {}/{}/{}",
305 ci->m_nClean,
306 static_cast<float> (ci->m_nRemoved) / den,
307 ci->m_removed0,
308 ci->m_removed1,
309 ci->m_removed2plus );
310 }
311
312 return StatusCode::SUCCESS;
313}
314
315
321{
322 m_runlbn.reset (m_props->m_ringSize);
323 m_timestamp.reset (m_props->m_ringSize);
324
325 std::fill (m_slotLBN.begin(), m_slotLBN.end(), 0);
326 std::fill (m_slotTimestamp.begin(), m_slotTimestamp.end(), 0);
327
328 m_ccinfo.clear();
329 std::priority_queue<QueueItem> tmp;
330 m_work.swap (tmp);
331
332 m_nEvents = 0;
333 m_queueSum = 0;
334 m_workRemoved = 0;
335 m_maxQueue = 0;
336 m_cleanTasks = 0;
337
338 return StatusCode::SUCCESS;
339}
340
341
343DelayedConditionsCleanerSvc::getKeys(const Ring& runLBRing, const Ring& TSRing) const {
344
345 // Get a copy of the contents of the ring buffer holding runLumi and time-stamp keys
346 std::vector<key_type> runLBKeys=runLBRing.getKeysDedup();
347 std::vector<key_type> TSKeys=TSRing.getKeysDedup();
348
349 // Add in the keys for the currently-executing slots.
350 // These are very likely to already be in the ring, but that's
351 // not absolutely guaranteed.
352 // FIXME: This probably does another memory allocation, due to
353 // growing the buffer. Would be nice to avoid that.
354 runLBKeys.insert (runLBKeys.end(), m_slotLBN.begin(), m_slotLBN.end());
355 TSKeys.insert(TSKeys.end(), m_slotTimestamp.begin(), m_slotTimestamp.end());
356
357 twoKeys_t result{std::move(runLBKeys), std::move(TSKeys)};
358
363 for ( auto& keys : result ) {
364 std::sort (keys.begin(), keys.end());
365 auto end = std::unique (keys.begin(), keys.end());
366 keys.resize (end - keys.begin());
367 }
368
369
370 return result;
371}
372
373
383void
384DelayedConditionsCleanerSvc::scheduleClean (std::vector<CondContInfo*>&& cis,
385 twoKeys_t&& twoKeys,
386 bool allowAsync)
387{
388 // Remove any duplicates from the list of containers.
389 std::sort (cis.begin(), cis.end());
390 auto pos = std::unique (cis.begin(), cis.end());
391 cis.resize (pos - cis.begin());
392
393 if (allowAsync && m_props->m_async) {
394#if USE_ASYNC_TASK
395 // Queue cleaning as a TBB task.
396 // Count that we have another executing task.
397 ++m_cleanTasks;
398
399 // Create the TBB task and queue it.
400 // TBB will delete the task object after it completes.
401 tbb::task* t = new (tbb::task::allocate_root())
402 DelayedConditionsCleanerTask (*this, std::move (cis),
403 std::move (twoKeys));
404 tbb::task::enqueue (*t);
405#endif
406 }
407 else
408 {
409 // Call cleaning directly.
410 cleanContainers (std::move (cis), std::move (twoKeys));
411 }
412}
413
414
420void
421DelayedConditionsCleanerSvc::cleanContainers (std::vector<CondContInfo*>&& cis,
422 twoKeys_t&& twoKeys)
423{
424 // FIXME: Some conditions objects have pointers to parts of other
425 // conditions objects, which violates the lifetime guarantees of
426 // conditions containers (a pointer you get from a conditions container
427 // is guaranteed to be valid until the end of the current event, but
428 // not past that). In some cases, this can be dealt with by replacing
429 // the pointers with CondLink, but that is sometimes rather inconvenient.
430 // Try to work around this for now by ensuring that when we delete an object,
431 // we also try to clean the containers of other objects that may depend
432 // on it.
433 std::vector<CondContInfo*> toclean = std::move (cis);
434 std::unordered_set<CondContInfo*> cleaned (toclean.begin(), toclean.end());
435 while (!toclean.empty()) {
436 std::vector<CondContInfo*> newclean;
437 for (CondContInfo* ci : toclean) {
438 if (cleanContainer (ci, twoKeys)) {
440 for (CondContBase* dep : ci->m_cc.getDeps()) {
441 CCInfoMap_t::iterator it = m_ccinfo.find (dep);
442 // If we don't find it, then dep must have no conditions objects.
443 if (it != m_ccinfo.end()) {
444 CondContInfo* ci_dep = &it->second;
445 if (cleaned.insert (ci_dep).second) {
446 newclean.push_back (ci_dep);
447 }
448 }
449 }
450 }
451 }
452 toclean = std::move (newclean);
453 }
454}
455
456
465bool
467 const twoKeys_t& twoKeys) const
468{
469 size_t n = ci->m_cc.trim (twoKeys[0],twoKeys[1]);
470
471 ++ci->m_nClean;
472 ci->m_nRemoved += n;
473 switch (n) {
474 case 0:
475 ++ci->m_removed0;
476 break;
477 case 1:
478 ++ci->m_removed1;
479 break;
480 default:
481 ++ci->m_removed2plus;
482 }
483
484 return n > 0;
485}
486
490
492
497{
499 return StatusCode::SUCCESS;
500}
501
502
503} // namespace Athena
#define ATH_CHECK
Evaluate an expression and check for errors.
#define ATH_MSG_INFO(x,...)
pimpl-style holder for component properties.
Hold mappings of ranges to condition objects.
Clean conditions containers after a delay.
virtual void lock()=0
Interface to allow an object to lock itself when made const in SG.
read-copy-update (RCU) style synchronization for Athena.
static Double_t sz
static Double_t a
DelayedConditionsCleanerSvc * parent()
AthProperties(DelayedConditionsCleanerSvc *parent)
bool m_async
Property: If true, run cleaning asynchronously in an MT job.
Gaudi::Property< size_t > m_cleanDelay
Property: Number of events after adding a conditions object we try to clean its container.
ServiceHandle< Athena::IRCUSvc > m_rcu
Property: RCU Service.
DelayedConditionsCleanerSvcProps(DelayedConditionsCleanerSvc *parent)
Gaudi::Property< size_t > m_ringSize
Property: Number of previous events for which to remember IOV history.
Gaudi::Property< size_t > m_lookAhead
Property: Maximum number of events to consolodate together when cleaning.
Information that we maintain about each conditions container.
Clean conditions containers after a delay.
DelayedConditionsCleanerSvc(const std::string &name, ISvcLocator *svc)
Standard Gaudi constructor.
std::priority_queue< QueueItem > m_work
Priority queue of pending cleaning requests.
size_t m_nEvents
Priority queue statistics.
virtual StatusCode finalize() override
Standard Gaudi finalize method.
twoKeys_t getKeys(const Ring &runLBRing, const Ring &TSRing) const
virtual StatusCode condObjAdded(const EventContext &ctx, CondContBase &cc) override
Called after a conditions object has been added.
virtual StatusCode initialize() override
Standard Gaudi initialize method.
Ring m_runlbn
Two ring buffers for recent IOV keys, one for run+LBN and one for timestamp.
virtual StatusCode event(const EventContext &ctx, bool allowAsync) override
Called at the start of each event.
void scheduleClean(std::vector< CondContInfo * > &&cis, twoKeys_t &&twoKeys, bool allowAsync)
Do cleaning for a set of containers.
std::unique_ptr< DelayedConditionsCleanerSvcProps > m_props
Component properties.
CxxUtils::Ring< key_type > Ring
Ring buffer holding most recent IOV keys of a given type.
CondContBase::key_type key_type
Packed key type.
std::atomic< int > m_cleanTasks
Number of active asynchronous cleaning tasks.
virtual StatusCode reset() override
Clear the internal state of the service.
std::array< std::vector< key_type >, 2 > twoKeys_t
void cleanContainers(std::vector< CondContInfo * > &&cis, twoKeys_t &&twoKeys)
Clean a set of containers.
bool cleanContainer(CondContInfo *ci, const twoKeys_t &keys) const
Clean a single container.
std::vector< key_type > m_slotLBN
IOV keys currently in use for each slot.
virtual StatusCode printStats() const override
Print some statistics about the garbage collection.
~DelayedConditionsCleanerSvc()
Standard destructor.
Base class for all conditions containers.
Definition CondCont.h:140
static key_type keyFromTimestamp(const EventIDBase &b)
Make a timestamp key from an EventIDBase.
static key_type keyFromRunLBN(const EventIDBase &b)
Make a run+lbn key from an EventIDBase.
KeyType keyType() const
Return the key type for this container.
std::vector< T > getKeysDedup() const
Return a copy of keys in the buffer.
Some weak symbol referencing magic... These are declared in AthenaKernel/getMessageSvc....
Definition AthDsoUtils.h:10
::StatusCode StatusCode
StatusCode definition for legacy code.
STL namespace.
DataModel_detail::iterator< DVL > unique(typename DataModel_detail::iterator< DVL > beg, typename DataModel_detail::iterator< DVL > end)
Specialization of unique for DataVector/List.
void sort(typename DataModel_detail::iterator< DVL > beg, typename DataModel_detail::iterator< DVL > end)
Specialization of sort for DataVector/List.