59 ATH_MSG_ERROR(
"Failed to create background event selector context");
60 return StatusCode::FAILURE;
65 SmartIF<IProxyProviderSvc> proxyProviderSvc{
66 serviceLocator()->service(std::format(
"ProxyProviderSvc/BkgPPSvc_{}", name()))
72 if (!addressProvider) {
74 "Could not cast background event selector to IAddressProvider");
76 proxyProviderSvc->addProvider(addressProvider);
79 SmartIF<IAddressProvider> athPoolAP{
80 serviceLocator()->service(std::format(
"AthenaPoolAddressProviderSvc/BkgAPAPSvc_{}", name()))
84 "Could not cast AthenaPoolAddressProviderSvc to IAddressProvider");
86 proxyProviderSvc->addProvider(athPoolAP);
89 SmartIF<IAddressProvider> addRemapAP{
90 serviceLocator()->service(std::format(
"AddressRemappingSvc/BkgARSvc_{}", name()))
93 ATH_MSG_WARNING(
"Could not cast AddressRemappingSvc to IAddressProvider");
95 proxyProviderSvc->addProvider(addRemapAP);
102 auto& sgs =
m_empty_caches.emplace_back(std::make_unique<SGHandleArray>());
103 sgs->reserve(mbBatchSize);
104 for (
int j = 0; j < mbBatchSize; ++j) {
106 auto& sg = sgs->emplace_back(
107 std::format(
"StoreGateSvc/StoreGate_{}_{}_{}", name(), i, j), name());
110 sg->setProxyProviderSvc(proxyProviderSvc);
120 auto skipEvent_callback = [
this, mbBatchSize](
123 using namespace std::chrono_literals;
124 auto evts = std::ranges::subrange(begin, end);
125 ATH_MSG_INFO(
"Skipping " << end - begin <<
" HS events.");
130 std::vector<std::tuple<int, int>> batches_with_counts{};
132 for (
int batch : batches_all) {
134 if (batches_with_counts.empty()) {
135 batches_with_counts.emplace_back(batch, 1);
139 auto& last_entry = batches_with_counts.back();
140 if (batch == std::get<0>(last_entry)) {
141 std::get<1>(last_entry) += 1;
144 batches_with_counts.emplace_back(batch, 1);
152 for (
const auto& [batch,
count] : batches_with_counts) {
153 if (
m_cache.count(batch) != 0) {
160 std::this_thread::sleep_for(50ms);
173 return StatusCode::FAILURE;
180 return StatusCode::SUCCESS;
185 return StatusCode::SUCCESS;
189 const EventContext& ctx) {
194 bool sf_updated_throwaway;
195 const float beam_lumi_sf =
197 ctx.eventID().lumi_block(),
198 sf_updated_throwaway)
200 std::vector<float> avg_num_mb_by_bunch(n_bunches,
210 avg_num_mb_by_bunch[idx] *=
m_beamInt->normFactor(bunch);
215 num_mb_by_bunch.clear();
216 num_mb_by_bunch.resize(n_bunches);
219 std::transform(avg_num_mb_by_bunch.begin(), avg_num_mb_by_bunch.end(),
220 num_mb_by_bunch.begin(), [&prng](
float avg) {
221 return std::poisson_distribution<std::uint64_t>(avg)(prng);
224 std::transform(avg_num_mb_by_bunch.begin(), avg_num_mb_by_bunch.end(),
225 num_mb_by_bunch.begin(), [](
float f) {
226 return static_cast<std::uint64_t>(std::round(f));
230 std::uint64_t num_mb = std::accumulate(num_mb_by_bunch.begin(), num_mb_by_bunch.end(), std::uint64_t{0});
231 std::vector<std::uint64_t>& index_array = *
m_idx_lists.get(ctx);
234 if (num_mb > mbBatchSize) {
237 rv::iota(0ULL, num_mb_by_bunch.size()) |
238 rv::filter([center_bunch, &num_mb_by_bunch](
int idx) {
239 bool good = idx != center_bunch;
241 good && num_mb_by_bunch[idx] > 0;
244 std::ranges::to<std::vector<std::size_t>>();
246 std::ranges::stable_sort(indices, std::greater{},
247 [center_bunch](std::size_t idx) {
248 return std::size_t(std::abs(
int(idx) - center_bunch));
251 for (
auto idx : indices) {
252 const std::uint64_t max_to_subtract = num_mb - mbBatchSize;
253 const std::uint64_t num_subtracted =
254 std::min(max_to_subtract, num_mb_by_bunch[idx]);
255 num_mb_by_bunch[idx] -= num_subtracted;
256 num_mb -= num_subtracted;
257 if (num_mb <= mbBatchSize) {
262 ATH_MSG_ERROR(
"We need " << num_mb <<
" events but the batch size is "
263 << mbBatchSize <<
". Restricting to "
264 << mbBatchSize <<
" events!");
267 rv::iota(std::uint64_t{0}, mbBatchSize) | std::ranges::to<std::vector<std::uint64_t>>();
269 index_array.reserve(num_mb);
270 std::sample(all_indices.begin(), all_indices.end(), std::back_inserter(index_array),
static_cast<std::size_t
>(num_mb), prng);
271 std::ranges::shuffle(index_array, prng);
272 ATH_MSG_DEBUG(
"HS ID " << hs_id <<
" uses " << num_mb <<
" events");
282 using namespace std::chrono_literals;
283 bool first_wait =
true;
284 std::chrono::steady_clock::time_point cache_wait_start{};
285 std::chrono::steady_clock::time_point order_wait_start{};
286 const std::int64_t hs_id =
get_hs_id(ctx);
291 if (
m_cache.count(batch) != 0) {
296 return StatusCode::SUCCESS;
300 ATH_MSG_INFO(
"Waiting to prevent out-of-order loading of batches");
301 order_wait_start = std::chrono::steady_clock::now();
303 std::this_thread::sleep_for(50ms);
305 auto wait_time = std::chrono::steady_clock::now() - order_wait_start;
307 "Waited {:%M:%S} to prevent out-of-order loading", wait_time));
313 if (empty_caches_lock.owns_lock()) {
316 empty_caches_lock.unlock();
320 cache_wait_start = std::chrono::steady_clock::now();
323 std::this_thread::sleep_for(100ms);
327 auto wait_time = std::chrono::steady_clock::now() - cache_wait_start;
329 std::format(
"Waited {:%M:%S} for a free cache", wait_time));
333 ATH_MSG_INFO(
"Reading next batch in event " << ctx.evt() <<
", slot "
334 << ctx.slot() <<
" (hs_id "
337 auto start_time = std::chrono::system_clock::now();
342 for (
auto&& sg : *
m_cache[batch]) {
345 SG::CurrentEventStore::Push reader_sg_ces(sg.get());
350 return StatusCode::FAILURE;
352 IOpaqueAddress* addr =
nullptr;
356 return StatusCode::FAILURE;
358 if (addr ==
nullptr) {
360 return StatusCode::FAILURE;
366 for (
const auto* proxy_ptr : sg->proxies()) {
367 if (!proxy_ptr->isValid()) {
372 sg->proxy_exact(proxy_ptr->sgkey())->accessData();
380 "Reading {} events took {:%OMm %OSs}",
m_cache[batch]->
size(),
381 std::chrono::system_clock::now() - start_time));
384 return StatusCode::SUCCESS;
387 return StatusCode::SUCCESS;